126 lines
3.9 KiB
C#
126 lines
3.9 KiB
C#
using Backend.MarketDataRequest;
|
|
using Backend.Interface;
|
|
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);
|
|
}
|
|
}
|
|
} |