using System; using System.Collections.Generic; using System.Linq; using System.Threading.Tasks; using Dapper; using QuantEngine.Core.Interfaces; using QuantEngine.Infrastructure.Data; namespace QuantEngine.Infrastructure.Repositories { public class CollectionRepository : ICollectionRepository, ICollectionReadRepository, ICollectionWriteRepository { private readonly IDbConnectionFactory _connectionFactory; public CollectionRepository(IDbConnectionFactory connectionFactory) { _connectionFactory = connectionFactory; } public async Task SaveRunAsync(CollectionRunRecord run) { await EnsureTablesAsync(); using var conn = _connectionFactory.CreateConnection(); await conn.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 status = EXCLUDED.status, finished_at = EXCLUDED.finished_at, total_snapshots = EXCLUDED.total_snapshots, total_errors = EXCLUDED.total_errors, updated_at = EXCLUDED.updated_at", run ); } 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(@" UPDATE quantengine.kis_collection_runs SET status = @Status, finished_at = @FinishedAt, total_snapshots = @TotalSnapshots, total_errors = @TotalErrors, updated_at = @UpdatedAt WHERE run_id = @RunId", new { RunId = runId, Status = status, FinishedAt = finishedAt, TotalSnapshots = totalSnapshots, TotalErrors = totalErrors, UpdatedAt = DateTime.UtcNow.ToString("o") } ); } public async Task SaveSnapshotAsync(CollectionSnapshotRecord snapshot) { using var conn = _connectionFactory.CreateConnection(); await conn.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 payload_json = EXCLUDED.payload_json, captured_at = EXCLUDED.captured_at", snapshot ); } public async Task SaveErrorAsync(CollectionErrorRecord error) { using var conn = _connectionFactory.CreateConnection(); await conn.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 ); } public async Task> GetRecentRunsAsync(int limit = 20) { using var conn = _connectionFactory.CreateConnection(); return (await conn.QueryAsync(@" 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(@" 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(@" 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 WHERE run_id = @RunId 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(@" 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(@" 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(@" 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(); return new CollectionDashboardStateRecord( LastRunId: lastRun?.RunId, LastRunStatus: lastRun?.Status, LastFinishedAt: lastRun?.FinishedAt, TotalSnapshots: stats?.TotalSnapshots ?? 0, TotalErrors: stats?.TotalErrors ?? 0, RecentErrors: recentErrors ); } public async Task> GetLatestSnapshotsForTickerAsync(string ticker, int limit = 10) { using var conn = _connectionFactory.CreateConnection(); return (await conn.QueryAsync(@" 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 ticker = @Ticker 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(@" 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", new { record.Ticker, // Dapper has no built-in type handler for System.DateOnly (throws // NotSupportedException) — pass as DateTime; the DATE column truncates the time part. TradeDate = record.TradeDate.ToDateTime(TimeOnly.MinValue), record.Open, record.High, record.Low, record.Close, record.Volume, record.Source, Provenance = record.ProvenanceJson ?? "{}" } ); } public async Task> GetPriceHistorySummaryAsync() { using var conn = _connectionFactory.CreateConnection(); return (await conn.QueryAsync(@" 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(@" CREATE TABLE IF NOT EXISTS quantengine.kis_collection_runs ( run_id TEXT PRIMARY KEY, status TEXT NOT NULL, started_at TEXT NOT NULL, finished_at TEXT, total_snapshots INTEGER, total_errors INTEGER, updated_at TEXT NOT NULL ); CREATE TABLE IF NOT EXISTS quantengine.kis_collection_snapshots ( run_id TEXT NOT NULL, dataset_name TEXT, ticker TEXT NOT NULL, source_name TEXT NOT NULL, payload_json TEXT NOT NULL, captured_at TEXT NOT NULL, created_at TEXT NOT NULL, PRIMARY KEY (run_id, ticker, source_name) ); CREATE TABLE IF NOT EXISTS quantengine.kis_collection_errors ( id SERIAL PRIMARY KEY, run_id TEXT NOT NULL, source_name TEXT NOT NULL, error_kind TEXT NOT NULL, error_message TEXT, ticker TEXT, created_at TEXT NOT NULL ); CREATE INDEX IF NOT EXISTS idx_kis_runs_started_at ON quantengine.kis_collection_runs(started_at DESC); CREATE INDEX IF NOT EXISTS idx_kis_snapshots_ticker ON quantengine.kis_collection_snapshots(ticker); CREATE INDEX IF NOT EXISTS idx_kis_snapshots_captured_at ON quantengine.kis_collection_snapshots(captured_at DESC); CREATE INDEX IF NOT EXISTS idx_kis_errors_run_id ON quantengine.kis_collection_errors(run_id); "); } } }