namespace QuantEngine.Infrastructure.Repositories; using Dapper; using QuantEngine.Core.Repositories; using System.Data; public class MarketDataRepository : IMarketDataRepository { private readonly IDbConnection _connection; private const string SchemaName = "quantengine"; public MarketDataRepository(IDbConnection connection) => _connection = connection; public async Task> GetByStockIdAsync( int stockId, int? sourceId = null, DateTime? start = null, DateTime? end = null) { var sql = $"SELECT * FROM {SchemaName}.market_data WHERE stock_id = @stock_id"; if (sourceId.HasValue) sql += " AND source_id = @source_id"; if (start.HasValue) sql += " AND recorded_at >= @start"; if (end.HasValue) sql += " AND recorded_at < @end"; sql += " ORDER BY recorded_at DESC"; var results = await _connection.QueryAsync(sql, new { stock_id = stockId, source_id = sourceId, start, end }); return results.ToList().AsReadOnly(); } public async Task GetLatestByTickerAsync(string ticker, int? sourceId = null) { var sql = $@" SELECT md.* FROM {SchemaName}.market_data md INNER JOIN {SchemaName}.stocks s ON md.stock_id = s.id WHERE s.ticker = @ticker"; if (sourceId.HasValue) sql += " AND md.source_id = @source_id"; sql += " ORDER BY md.recorded_at DESC LIMIT 1"; return await _connection.QueryFirstOrDefaultAsync(sql, new { ticker, source_id = sourceId }); } public async Task> GetLatestByStockIdsAsync( IEnumerable stockIds, DateTime? asOf = null) { if (!stockIds.Any()) return new(); var sql = $@" WITH latest AS ( SELECT stock_id, MAX(recorded_at) as latest_time FROM {SchemaName}.market_data WHERE stock_id = ANY(@stock_ids) GROUP BY stock_id ) SELECT md.* FROM {SchemaName}.market_data md INNER JOIN latest l ON md.stock_id = l.stock_id AND md.recorded_at = l.latest_time"; var results = await _connection.QueryAsync(sql, new { stock_ids = stockIds.ToList() }); return results.ToDictionary(r => r.StockId); } public async Task InsertAsync(MarketDataSnapshot snapshot) { var sql = $@" INSERT INTO {SchemaName}.market_data (stock_id, source_id, recorded_at, current_price, open_price, high_price, low_price, close_price, ask_price_1, ask_volume_1, bid_price_1, bid_volume_1, volume, trade_amount, individual_buy_volume, institutional_buy_volume, foreign_buy_volume, collected_at, notes) VALUES (@stock_id, @source_id, @recorded_at, @current_price, @open_price, @high_price, @low_price, @close_price, @ask_price_1, @ask_volume_1, @bid_price_1, @bid_volume_1, @volume, @trade_amount, @individual_buy_volume, @institutional_buy_volume, @foreign_buy_volume, @collected_at, @notes) RETURNING id"; return await _connection.QuerySingleAsync(sql, snapshot); } public async Task> InsertBatchAsync(IEnumerable snapshots) { var ids = new List(); foreach (var snapshot in snapshots) ids.Add(await InsertAsync(snapshot)); return ids.AsReadOnly(); } public async Task UpdateAsync(int marketDataId, MarketDataSnapshot updated) => true; public async Task> ValidateCompletenessAsync(int stockId, DateTime start, DateTime end) => new(); public async Task> DetectOutliersAsync( int stockId, DateTime start, DateTime end, double stdDevThreshold = 3.0) => new(); }