using Backend.MarketDataRequest; using Npgsql; namespace Backend.DatabaseHandler { public sealed class PostgresHistoricalMarketDataCache : IHistoricalMarketDataCache { private readonly string _connectionString; public PostgresHistoricalMarketDataCache(IConfiguration configuration) { _connectionString = configuration.GetConnectionString("MarketDataDb") ?? throw new InvalidOperationException("Missing connection string: MarketDataDb"); } public async Task> GetBarsAsync( string symbol, HistoricalBarSize barSize, HistoricalWhatToShow whatToShow, bool useRth, DateTime startUtc, DateTime endUtc, CancellationToken cancellationToken = default) { var results = new List(); await using var conn = new NpgsqlConnection(_connectionString); await conn.OpenAsync(cancellationToken); const string sql = """ SELECT symbol, bar_size, what_to_show, use_rth, timestamp_utc, open, high, low, close, volume FROM historical_equity_bars WHERE symbol = @symbol AND bar_size = @bar_size AND what_to_show = @what_to_show AND use_rth = @use_rth AND timestamp_utc >= @start_utc AND timestamp_utc <= @end_utc ORDER BY timestamp_utc ASC; """; await using var cmd = new NpgsqlCommand(sql, conn); cmd.Parameters.AddWithValue("symbol", symbol); cmd.Parameters.AddWithValue("bar_size", barSize.ToIbString()); cmd.Parameters.AddWithValue("what_to_show", whatToShow.ToIbString()); cmd.Parameters.AddWithValue("use_rth", useRth); cmd.Parameters.AddWithValue("start_utc", startUtc); cmd.Parameters.AddWithValue("end_utc", endUtc); await using var reader = await cmd.ExecuteReaderAsync(cancellationToken); while (await reader.ReadAsync(cancellationToken)) { results.Add(new HistoricalBarRecord { Symbol = reader.GetString(0), BarSize = reader.GetString(1), WhatToShow = reader.GetString(2), UseRth = reader.GetBoolean(3), TimestampUtc = reader.GetDateTime(4), Open = reader.GetDouble(5), High = reader.GetDouble(6), Low = reader.GetDouble(7), Close = reader.GetDouble(8), Volume = reader.GetInt64(9) }); } return results; } public async Task UpsertBarsAsync( IEnumerable bars, CancellationToken cancellationToken = default) { var barList = bars.ToList(); if (barList.Count == 0) return; await using var conn = new NpgsqlConnection(_connectionString); await conn.OpenAsync(cancellationToken); await using var tx = await conn.BeginTransactionAsync(cancellationToken); const string sql = """ INSERT INTO historical_equity_bars ( symbol, bar_size, what_to_show, use_rth, timestamp_utc, open, high, low, close, volume ) VALUES ( @symbol, @bar_size, @what_to_show, @use_rth, @timestamp_utc, @open, @high, @low, @close, @volume ) ON CONFLICT (symbol, bar_size, what_to_show, use_rth, timestamp_utc) DO UPDATE SET open = EXCLUDED.open, high = EXCLUDED.high, low = EXCLUDED.low, close = EXCLUDED.close, volume = EXCLUDED.volume; """; foreach (var bar in barList) { await using var cmd = new NpgsqlCommand(sql, conn, tx); cmd.Parameters.AddWithValue("symbol", bar.Symbol); cmd.Parameters.AddWithValue("bar_size", bar.BarSize); cmd.Parameters.AddWithValue("what_to_show", bar.WhatToShow); cmd.Parameters.AddWithValue("use_rth", bar.UseRth); cmd.Parameters.AddWithValue("timestamp_utc", bar.TimestampUtc); cmd.Parameters.AddWithValue("open", bar.Open); cmd.Parameters.AddWithValue("high", bar.High); cmd.Parameters.AddWithValue("low", bar.Low); cmd.Parameters.AddWithValue("close", bar.Close); cmd.Parameters.AddWithValue("volume", bar.Volume); await cmd.ExecuteNonQueryAsync(cancellationToken); } await tx.CommitAsync(cancellationToken); } } }