diff --git a/src/dotnet/QuantEngine.Infrastructure/Repositories/CollectionRepository.cs b/src/dotnet/QuantEngine.Infrastructure/Repositories/CollectionRepository.cs index 3c5254d5..e23cc7cf 100644 --- a/src/dotnet/QuantEngine.Infrastructure/Repositories/CollectionRepository.cs +++ b/src/dotnet/QuantEngine.Infrastructure/Repositories/CollectionRepository.cs @@ -17,11 +17,28 @@ namespace QuantEngine.Infrastructure.Repositories _connectionFactory = connectionFactory; } + private async Task ExecuteAsync(string sql, object? param = null) + { + using var conn = _connectionFactory.CreateConnection(); + await conn.ExecuteAsync(sql, param); + } + + private async Task> QueryListAsync(string sql, object? param = null) + { + using var conn = _connectionFactory.CreateConnection(); + return (await conn.QueryAsync(sql, param)).ToList(); + } + + private async Task QuerySingleOrDefaultAsync(string sql, object? param = null) + { + using var conn = _connectionFactory.CreateConnection(); + return await conn.QueryFirstOrDefaultAsync(sql, param); + } + public async Task SaveRunAsync(CollectionRunRecord run) { await EnsureTablesAsync(); - using var conn = _connectionFactory.CreateConnection(); - await conn.ExecuteAsync(@" + await ExecuteAsync(@" INSERT INTO quantengine.kis_collection_runs (run_id, status, started_at, finished_at, total_snapshots, total_errors, updated_at) VALUES (@RunId, @Status, @StartedAt, @FinishedAt, @TotalSnapshots, @TotalErrors, @UpdatedAt) ON CONFLICT (run_id) DO UPDATE SET @@ -36,8 +53,7 @@ namespace QuantEngine.Infrastructure.Repositories public async Task UpdateRunStatusAsync(string runId, string status, string? finishedAt = null, int? totalSnapshots = null, int? totalErrors = null) { - using var conn = _connectionFactory.CreateConnection(); - await conn.ExecuteAsync(@" + await ExecuteAsync(@" UPDATE quantengine.kis_collection_runs SET status = @Status, finished_at = @FinishedAt, total_snapshots = @TotalSnapshots, total_errors = @TotalErrors, updated_at = @UpdatedAt WHERE run_id = @RunId", @@ -47,8 +63,7 @@ namespace QuantEngine.Infrastructure.Repositories public async Task SaveSnapshotAsync(CollectionSnapshotRecord snapshot) { - using var conn = _connectionFactory.CreateConnection(); - await conn.ExecuteAsync(@" + await ExecuteAsync(@" INSERT INTO quantengine.kis_collection_snapshots (run_id, dataset_name, ticker, source_name, payload_json, captured_at, created_at) VALUES (@RunId, @DatasetName, @Ticker, @SourceName, @PayloadJson, @CapturedAt, @CreatedAt) ON CONFLICT (run_id, ticker, source_name) DO UPDATE SET @@ -60,8 +75,7 @@ namespace QuantEngine.Infrastructure.Repositories public async Task SaveErrorAsync(CollectionErrorRecord error) { - using var conn = _connectionFactory.CreateConnection(); - await conn.ExecuteAsync(@" + await ExecuteAsync(@" INSERT INTO quantengine.kis_collection_errors (run_id, source_name, error_kind, error_message, ticker, created_at) VALUES (@RunId, @SourceName, @ErrorKind, @ErrorMessage, @Ticker, @CreatedAt)", error @@ -70,34 +84,31 @@ namespace QuantEngine.Infrastructure.Repositories public async Task> GetRecentRunsAsync(int limit = 20) { - using var conn = _connectionFactory.CreateConnection(); - return (await conn.QueryAsync(@" + return await QueryListAsync(@" SELECT run_id as RunId, status, started_at as StartedAt, finished_at as FinishedAt, total_snapshots as TotalSnapshots, total_errors as TotalErrors, updated_at as UpdatedAt FROM quantengine.kis_collection_runs ORDER BY started_at DESC LIMIT @Limit", new { Limit = limit } - )).ToList(); + ); } public async Task> GetRunSnapshotsAsync(string runId) { - using var conn = _connectionFactory.CreateConnection(); - return (await conn.QueryAsync(@" + return await QueryListAsync(@" SELECT run_id as RunId, dataset_name as DatasetName, ticker, source_name as SourceName, payload_json as PayloadJson, captured_at as CapturedAt, created_at as CreatedAt FROM quantengine.kis_collection_snapshots WHERE run_id = @RunId ORDER BY captured_at DESC", new { RunId = runId } - )).ToList(); + ); } public async Task> GetRunErrorsAsync(string runId, int limit = 50) { - using var conn = _connectionFactory.CreateConnection(); - return (await conn.QueryAsync(@" + return await QueryListAsync(@" SELECT run_id as RunId, source_name as SourceName, error_kind as ErrorKind, error_message as ErrorMessage, ticker as Ticker, created_at as CreatedAt FROM quantengine.kis_collection_errors @@ -105,32 +116,30 @@ namespace QuantEngine.Infrastructure.Repositories ORDER BY created_at DESC LIMIT @Limit", new { RunId = runId, Limit = limit } - )).ToList(); + ); } public async Task GetDashboardStateAsync() { - using var conn = _connectionFactory.CreateConnection(); - - var lastRun = await conn.QueryFirstOrDefaultAsync(@" + var lastRun = await QuerySingleOrDefaultAsync(@" SELECT run_id as RunId, status, started_at as StartedAt, finished_at as FinishedAt, total_snapshots as TotalSnapshots, total_errors as TotalErrors, updated_at as UpdatedAt FROM quantengine.kis_collection_runs ORDER BY started_at DESC LIMIT 1"); - var stats = await conn.QueryFirstOrDefaultAsync(@" + var stats = await QuerySingleOrDefaultAsync(@" SELECT COALESCE(SUM(total_snapshots), 0) as TotalSnapshots, COALESCE(SUM(total_errors), 0) as TotalErrors FROM quantengine.kis_collection_runs"); - var recentErrors = (await conn.QueryAsync(@" + var recentErrors = await QueryListAsync(@" SELECT run_id as RunId, source_name as SourceName, error_kind as ErrorKind, error_message as ErrorMessage, ticker as Ticker, created_at as CreatedAt FROM quantengine.kis_collection_errors ORDER BY created_at DESC - LIMIT 5")).ToList(); + LIMIT 5"); return new CollectionDashboardStateRecord( LastRunId: lastRun?.RunId, @@ -144,8 +153,7 @@ namespace QuantEngine.Infrastructure.Repositories public async Task> GetLatestSnapshotsForTickerAsync(string ticker, int limit = 10) { - using var conn = _connectionFactory.CreateConnection(); - return (await conn.QueryAsync(@" + return await QueryListAsync(@" SELECT run_id as RunId, dataset_name as DatasetName, ticker, source_name as SourceName, payload_json as PayloadJson, captured_at as CapturedAt, created_at as CreatedAt FROM quantengine.kis_collection_snapshots @@ -153,13 +161,12 @@ namespace QuantEngine.Infrastructure.Repositories ORDER BY captured_at DESC LIMIT @Limit", new { Ticker = ticker, Limit = limit } - )).ToList(); + ); } public async Task SavePriceHistoryDailyAsync(PriceHistoryDailyRecord record) { - using var conn = _connectionFactory.CreateConnection(); - await conn.ExecuteAsync(@" + await ExecuteAsync(@" INSERT INTO quantengine.price_history_daily (ticker, trade_date, open, high, low, close, volume, source, provenance) VALUES (@Ticker, @TradeDate, @Open, @High, @Low, @Close, @Volume, @Source, @Provenance::jsonb) ON CONFLICT (ticker, trade_date) DO NOTHING", @@ -182,20 +189,18 @@ namespace QuantEngine.Infrastructure.Repositories public async Task> GetPriceHistorySummaryAsync() { - using var conn = _connectionFactory.CreateConnection(); - return (await conn.QueryAsync(@" + return await QueryListAsync(@" SELECT ticker AS Ticker, count(*)::int AS RowCount, min(trade_date) AS FirstDate, max(trade_date) AS LastDate FROM quantengine.price_history_daily GROUP BY ticker ORDER BY ticker", new { } - )).ToList(); + ); } private async Task EnsureTablesAsync() { - using var conn = _connectionFactory.CreateConnection(); - await conn.ExecuteAsync(@" + await ExecuteAsync(@" CREATE TABLE IF NOT EXISTS quantengine.kis_collection_runs ( run_id TEXT PRIMARY KEY, status TEXT NOT NULL,