diff --git a/Backend.cs b/Backend.cs index 77e8a42..1887402 100644 --- a/Backend.cs +++ b/Backend.cs @@ -1,3 +1,4 @@ +using Backend.DatabaseHandler; using Backend.ServiceHander; using Microsoft.AspNetCore.HttpOverrides; @@ -7,6 +8,7 @@ builder.Services.AddControllers(); builder.Services.AddEndpointsApiExplorer(); builder.Services.AddSwaggerGen(); +builder.Services.AddSingleton(); builder.Services.AddSingleton(); builder.Services.AddCors(options => diff --git a/Backend.csproj b/Backend.csproj index 309f8bb..fa2acbb 100644 --- a/Backend.csproj +++ b/Backend.csproj @@ -13,6 +13,7 @@ + diff --git a/Controller/MarketDataController.cs b/Controller/MarketDataController.cs index 7cff173..d299ff5 100644 --- a/Controller/MarketDataController.cs +++ b/Controller/MarketDataController.cs @@ -1,7 +1,6 @@ using Backend.ServiceHander; using Backend.MarketDataRequest; using Microsoft.AspNetCore.Mvc; -using IBApi; namespace BackendApi.Controller { @@ -31,7 +30,7 @@ namespace BackendApi.Controller } [HttpGet("historical")] - public IActionResult GetHistorical( + public async Task GetHistorical( [FromQuery] string symbol, [FromQuery] string endDateTime, [FromQuery] string duration, @@ -41,7 +40,7 @@ namespace BackendApi.Controller { try { - var data = _marketDataService.GetHistorical( + var data = await _marketDataService.GetHistoricalAsync( symbol.ToUpperInvariant(), endDateTime, duration, @@ -59,17 +58,17 @@ namespace BackendApi.Controller } [HttpGet("historical-line")] - public IActionResult GetHistoricalLine( - [FromQuery] string symbol, - [FromQuery] string endDateTime, - [FromQuery] string duration, - [FromQuery] HistoricalBarSize barSize, - [FromQuery] HistoricalWhatToShow whatToShow = HistoricalWhatToShow.Trades, - [FromQuery] bool useRth = true) + public async Task GetHistoricalLine( + [FromQuery] string symbol, + [FromQuery] string endDateTime, + [FromQuery] string duration, + [FromQuery] HistoricalBarSize barSize, + [FromQuery] HistoricalWhatToShow whatToShow = HistoricalWhatToShow.Trades, + [FromQuery] bool useRth = true) { try { - var data = _marketDataService.GetHistoricalLine( + var data = await _marketDataService.GetHistoricalLineAsync( symbol.ToUpperInvariant(), endDateTime, duration, diff --git a/DatabaseHandler/HistoricalBarRecord.cs b/DatabaseHandler/HistoricalBarRecord.cs new file mode 100644 index 0000000..a16e89e --- /dev/null +++ b/DatabaseHandler/HistoricalBarRecord.cs @@ -0,0 +1,18 @@ +namespace Backend.DatabaseHandler +{ + public sealed class HistoricalBarRecord + { + public string Symbol { get; init; } = string.Empty; + public string BarSize { get; init; } = string.Empty; + public string WhatToShow { get; init; } = string.Empty; + public bool UseRth { get; init; } + + public DateTime TimestampUtc { get; init; } + + public double Open { get; init; } + public double High { get; init; } + public double Low { get; init; } + public double Close { get; init; } + public long Volume { get; init; } + } +} \ No newline at end of file diff --git a/DatabaseHandler/HistoricalMarketDataCache.cs b/DatabaseHandler/HistoricalMarketDataCache.cs new file mode 100644 index 0000000..fc44641 --- /dev/null +++ b/DatabaseHandler/HistoricalMarketDataCache.cs @@ -0,0 +1,20 @@ +using Backend.MarketDataRequest; + +namespace Backend.DatabaseHandler +{ + public interface IHistoricalMarketDataCache + { + Task> GetBarsAsync( + string symbol, + HistoricalBarSize barSize, + HistoricalWhatToShow whatToShow, + bool useRth, + DateTime startUtc, + DateTime endUtc, + CancellationToken cancellationToken = default); + + Task UpsertBarsAsync( + IEnumerable bars, + CancellationToken cancellationToken = default); + } +} \ No newline at end of file diff --git a/DatabaseHandler/HistoricalRequestParser.cs b/DatabaseHandler/HistoricalRequestParser.cs new file mode 100644 index 0000000..21da50f --- /dev/null +++ b/DatabaseHandler/HistoricalRequestParser.cs @@ -0,0 +1,96 @@ +using Backend.MarketDataRequest; +using System.Globalization; + +namespace Backend.DatabaseHandler +{ + public static class HistoricalRequestParser + { + public static HistoricalRequestWindow BuildWindow( + string endDateTimeIsoUtc, + string duration, + HistoricalBarSize barSize) + { + DateTime endUtc = ParseRequestEndDateTimeUtc(endDateTimeIsoUtc); + TimeSpan durationSpan = ParseDuration(duration); + DateTime startUtc = endUtc - durationSpan; + + int expectedBarsMin = EstimateExpectedBars(durationSpan, barSize); + + return new HistoricalRequestWindow + { + StartUtc = startUtc, + EndUtc = endUtc, + ExpectedBarCountMin = expectedBarsMin + }; + } + + public static DateTime ParseRequestEndDateTimeUtc(string raw) + { + if (string.IsNullOrWhiteSpace(raw)) + throw new FormatException("endDateTime is empty."); + + raw = raw.Trim(); + +Console.WriteLine(raw); + if (!DateTimeOffset.TryParseExact( + raw, + new[] + { + "yyyy-MM-dd'T'HH:mm:ss'Z'", + "yyyy-MM-dd'T'HH:mm:ss.FFF'Z'", + "O" + }, + CultureInfo.InvariantCulture, + DateTimeStyles.AssumeUniversal | DateTimeStyles.AdjustToUniversal, + out var dto)) + { + throw new FormatException( + $"endDateTime must be ISO 8601 UTC, e.g. 2026-03-17T20:00:00Z. Received: '{raw}'"); + } + + return dto.UtcDateTime; + } + + private static TimeSpan ParseDuration(string duration) + { + if (string.IsNullOrWhiteSpace(duration)) + throw new ArgumentException("Duration is empty."); + + string normalized = duration.Trim().Replace("+", " "); + var parts = normalized.Split(' ', StringSplitOptions.RemoveEmptyEntries); + + if (parts.Length != 2) + throw new ArgumentException($"Invalid duration format: '{duration}'"); + + int value = int.Parse(parts[0], CultureInfo.InvariantCulture); + string unit = parts[1].ToUpperInvariant(); + + return unit switch + { + "S" => TimeSpan.FromSeconds(value), + "D" => TimeSpan.FromDays(value), + "W" => TimeSpan.FromDays(value * 7), + "M" => TimeSpan.FromDays(value * 30), + "Y" => TimeSpan.FromDays(value * 365), + _ => throw new ArgumentException($"Unsupported duration unit: '{unit}'") + }; + } + + private static int EstimateExpectedBars(TimeSpan duration, HistoricalBarSize barSize) + { + double secondsPerBar = barSize switch + { + HistoricalBarSize.OneMin => 60, + HistoricalBarSize.FiveMins => 300, + HistoricalBarSize.FifteenMins => 900, + HistoricalBarSize.OneHour => 3600, + HistoricalBarSize.OneDay => 86400, + HistoricalBarSize.FiveDays => 432000, + HistoricalBarSize.OneMonth => 2592000, + _ => throw new ArgumentOutOfRangeException(nameof(barSize)) + }; + + return Math.Max(1, (int)Math.Floor(duration.TotalSeconds / secondsPerBar)); + } + } +} \ No newline at end of file diff --git a/DatabaseHandler/HistoricalRequestWindow.cs b/DatabaseHandler/HistoricalRequestWindow.cs new file mode 100644 index 0000000..324300c --- /dev/null +++ b/DatabaseHandler/HistoricalRequestWindow.cs @@ -0,0 +1,9 @@ +namespace Backend.DatabaseHandler +{ + public sealed class HistoricalRequestWindow + { + public DateTime StartUtc { get; init; } + public DateTime EndUtc { get; init; } + public int ExpectedBarCountMin { get; init; } + } +} \ No newline at end of file diff --git a/DatabaseHandler/PostgresHistoricalMarketDataCache.cs b/DatabaseHandler/PostgresHistoricalMarketDataCache.cs new file mode 100644 index 0000000..ac449f7 --- /dev/null +++ b/DatabaseHandler/PostgresHistoricalMarketDataCache.cs @@ -0,0 +1,125 @@ +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); + } + } +} \ No newline at end of file diff --git a/MarketDataRequest/IBGatewayControl.cs b/MarketDataRequest/IBGatewayControl.cs index de250ac..8778d11 100644 --- a/MarketDataRequest/IBGatewayControl.cs +++ b/MarketDataRequest/IBGatewayControl.cs @@ -141,6 +141,15 @@ namespace Backend.MarketDataRequest return; } + if (id > 0 && _historyRequests.TryRemove(id, out var historyTcs)) + { + _BarsByRequest.TryRemove(id, out _); + + historyTcs.TrySetException( + new InvalidOperationException($"IB error for historical reqId={id}: {errorCode} - {errorMsg}")); + return; + } + if (id > 0 && _fundamentalsRequests.TryRemove(id, out var fundamentalsTcs)) { Client.cancelMktData(id); diff --git a/ServiceHandler/MarketDataHander.cs b/ServiceHandler/MarketDataHander.cs index ef99d15..3f85d8c 100644 --- a/ServiceHandler/MarketDataHander.cs +++ b/ServiceHandler/MarketDataHander.cs @@ -1,16 +1,21 @@ +using Backend.DatabaseHandler; using Backend.MarketDataRequest; using IBApi; +using System.Globalization; namespace Backend.ServiceHander { public sealed class MarketDataService { private readonly IBGatewayControl Gateway; + private readonly IHistoricalMarketDataCache Cache; private int NextRequestId = 1; private readonly object Lock = new(); - public MarketDataService() + public MarketDataService(IHistoricalMarketDataCache cache) { + Cache = cache; + Gateway = new IBGatewayControl(); try { @@ -31,12 +36,12 @@ namespace Backend.ServiceHander } } - public LiveEquitySnapshot GetLiveQuote(string Symbol) // + public LiveEquitySnapshot GetLiveQuote(string symbol) { - var Request = new LiveEquityRequest + var request = new LiveEquityRequest { RequestId = GetNextRequestId(), - Symbol = Symbol, + Symbol = symbol, SecType = "STK", Exchange = "SMART", PrimaryExchange = "NASDAQ", @@ -44,29 +49,28 @@ namespace Backend.ServiceHander RegulatorySnapshot = false }; - return Gateway.RequestLiveEquityBlocking(Request, TimeSpan.FromSeconds(10)); + return Gateway.RequestLiveEquityBlocking(request, TimeSpan.FromSeconds(10)); } - public List GetHistorical(string symbol, string endDateTime, string duration, HistoricalBarSize barSize, HistoricalWhatToShow whatToShow = HistoricalWhatToShow.Trades, bool useRth = true) + public async Task> GetHistoricalAsync( + string symbol, + string endDateTime, + string duration, + HistoricalBarSize barSize, + HistoricalWhatToShow whatToShow = HistoricalWhatToShow.Trades, + bool useRth = true) { - var Request = new HistoricalEquityRequest - { - RequestId = GetNextRequestId(), - Symbol = symbol, - EndDateTime = endDateTime, - Duration = duration, - BarSize = barSize, - WhatToShow = whatToShow, - UseRth = useRth, - KeepUpToDate = false, - FormatDate = 1 - }; + var bars = await GetHistoricalBarsAsync( + symbol, + endDateTime, + duration, + barSize, + whatToShow, + useRth); - var Bars = Gateway.RequestHistoricalBlocking(Request, TimeSpan.FromSeconds(20)); - - return Bars.Select(b => new + return bars.Select(b => new { - time = b.Time, + timestamp = ToIsoUtcString(b.TimestampUtc), open = b.Open, high = b.High, low = b.Low, @@ -75,14 +79,59 @@ namespace Backend.ServiceHander }).Cast().ToList(); } - public List GetHistoricalLine(string symbol, string endDateTime, string duration, HistoricalBarSize barSize, HistoricalWhatToShow whatToShow = HistoricalWhatToShow.Trades, bool useRth = true) + public async Task> GetHistoricalLineAsync( + string symbol, + string endDateTime, + string duration, + HistoricalBarSize barSize, + HistoricalWhatToShow whatToShow = HistoricalWhatToShow.Trades, + bool useRth = true) { - var Request = new HistoricalEquityRequest + var bars = await GetHistoricalBarsAsync( + symbol, + endDateTime, + duration, + barSize, + whatToShow, + useRth); + + return bars.Select(b => new + { + timestamp = ToIsoUtcString(b.TimestampUtc), + price = b.Close, + volume = b.Volume + }).Cast().ToList(); + } + + private async Task> GetHistoricalBarsAsync( + string symbol, + string endDateTimeIsoUtc, + string duration, + HistoricalBarSize barSize, + HistoricalWhatToShow whatToShow, + bool useRth) + { + var window = HistoricalRequestParser.BuildWindow(endDateTimeIsoUtc, duration, barSize); + + var cachedBars = await Cache.GetBarsAsync( + symbol, + barSize, + whatToShow, + useRth, + window.StartUtc, + window.EndUtc); + + if (cachedBars.Count >= window.ExpectedBarCountMin) + { + return cachedBars; + } + + var request = new HistoricalEquityRequest { RequestId = GetNextRequestId(), Symbol = symbol, - EndDateTime = endDateTime, - Duration = duration, + EndDateTime = ToIbHistoricalEndDateTime(window.EndUtc), + Duration = NormalizeDuration(duration), BarSize = barSize, WhatToShow = whatToShow, UseRth = useRth, @@ -90,14 +139,102 @@ namespace Backend.ServiceHander FormatDate = 1 }; - var Bars = Gateway.RequestHistoricalBlocking(Request, TimeSpan.FromSeconds(20)); + List ibBars = Gateway.RequestHistoricalBlocking(request, TimeSpan.FromSeconds(20)); - return Bars.Select(b => new + var mappedBars = ibBars.Select(b => new HistoricalBarRecord { - timestamp = b.Time, - price = b.Close, - volume = b.Volume - }).Cast().ToList(); + Symbol = symbol, + BarSize = barSize.ToIbString(), + WhatToShow = whatToShow.ToIbString(), + UseRth = useRth, + TimestampUtc = ParseIbBarTimeToUtc(b.Time), + Open = b.Open, + High = b.High, + Low = b.Low, + Close = b.Close, + Volume = (long)b.Volume + }) + .OrderBy(b => b.TimestampUtc) + .ToList(); + + await Cache.UpsertBarsAsync(mappedBars); + + return mappedBars + .Where(b => b.TimestampUtc >= window.StartUtc && b.TimestampUtc <= window.EndUtc) + .OrderBy(b => b.TimestampUtc) + .ToList(); + } + + private static string NormalizeDuration(string duration) + { + return duration.Trim().Replace("+", " "); + } + + private static string ToIsoUtcString(DateTime utc) + { + return utc.ToUniversalTime().ToString("yyyy-MM-ddTHH:mm:ssZ", CultureInfo.InvariantCulture); + } + + private static string ToIbHistoricalEndDateTime(DateTime utc) + { + var eastern = GetEasternTimeZone(); + var easternTime = TimeZoneInfo.ConvertTimeFromUtc(utc.ToUniversalTime(), eastern); + return easternTime.ToString("yyyyMMdd HH:mm:ss", CultureInfo.InvariantCulture) + " US/Eastern"; + } + + private static DateTime ParseIbBarTimeToUtc(string raw) + { + if (string.IsNullOrWhiteSpace(raw)) + throw new FormatException("IB bar time is empty."); + + raw = raw.Trim(); + + const string easternSuffix = " US/Eastern"; + if (raw.EndsWith(easternSuffix, StringComparison.Ordinal)) + { + string localPart = raw[..^easternSuffix.Length]; + + if (!DateTime.TryParseExact( + localPart, + "yyyyMMdd HH:mm:ss", + CultureInfo.InvariantCulture, + DateTimeStyles.None, + out var easternLocal)) + { + throw new FormatException($"Unable to parse IB datetime portion: {localPart}"); + } + + return TimeZoneInfo.ConvertTimeToUtc(easternLocal, GetEasternTimeZone()); + } + + if (DateTime.TryParseExact( + raw, + "yyyyMMdd", + CultureInfo.InvariantCulture, + DateTimeStyles.None, + out var dayOnly)) + { + // Daily bars: treat as midnight UTC for storage consistency. + return DateTime.SpecifyKind(dayOnly, DateTimeKind.Utc); + } + + if (DateTimeOffset.TryParseExact( + raw, + new[] { "yyyy-MM-ddTHH:mm:ss'Z'", "O" }, + CultureInfo.InvariantCulture, + DateTimeStyles.AssumeUniversal | DateTimeStyles.AdjustToUniversal, + out var iso)) + { + return iso.UtcDateTime; + } + + throw new FormatException($"Unable to parse IB bar time: {raw}"); + } + + private static TimeZoneInfo GetEasternTimeZone() + { + return TimeZoneInfo.FindSystemTimeZoneById( + OperatingSystem.IsWindows() ? "Eastern Standard Time" : "America/New_York"); } public void Dispose()