SoftTraderBackend/DatabaseHandler/PostgresHistoricalMarketDataStore.cs
2026-03-20 10:42:48 -04:00

125 lines
4.2 KiB
C#

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<List<HistoricalBarRecord>> GetBarsAsync(
string symbol,
HistoricalBarSize barSize,
HistoricalWhatToShow whatToShow,
bool useRth,
DateTime startUtc,
DateTime endUtc,
CancellationToken cancellationToken = default)
{
var results = new List<HistoricalBarRecord>();
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<HistoricalBarRecord> 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);
}
}
}