using IBApi; using System.Collections.Concurrent; using System.Globalization; namespace Backend.MarketDataRequest { public sealed class IBGatewayControl : DefaultEWrapper, IDisposable { private readonly EReaderMonitorSignal _signal = new(); public EClientSocket Client { get; } private EReader? _reader; private Thread? _readerThread; private volatile bool _disposed; private volatile bool _readerRunning; private volatile bool _apiReady; public bool IsConnected => Client.IsConnected(); public bool IsApiReady => _apiReady; // Historical data private readonly ConcurrentDictionary> _BarsByRequest = new(); private readonly ConcurrentDictionary>> _historyRequests = new(); // Equity data private readonly ConcurrentDictionary _liveSnapshotsByRequest = new(); private readonly ConcurrentDictionary> _liveRequests = new(); // Fundamentals data private readonly ConcurrentDictionary _fundamentalsSnapshotsByRequest = new(); private readonly ConcurrentDictionary> _fundamentalsRequests = new(); public IBGatewayControl() { Client = new EClientSocket(this, _signal); } public void Start(string host = "127.0.0.1", int port = 4002, int clientId = 1) { ThrowIfDisposed(); if (Client.IsConnected()) return; Console.WriteLine($"Connecting to IB Gateway at {host}:{port} with clientId={clientId}..."); Client.eConnect(host, port, clientId); if (!Client.IsConnected()) throw new InvalidOperationException("Failed to connect to IB Gateway."); _reader = new EReader(Client, _signal); _reader.Start(); _readerRunning = true; _readerThread = new Thread(ReadLoop) { IsBackground = true, Name = "IBGatewayControl.EReader" }; _readerThread.Start(); Console.WriteLine("IB Gateway socket connected."); } public void Disconnect() { if (_disposed) return; _readerRunning = false; _apiReady = false; if (Client.IsConnected()) { Console.WriteLine("Disconnecting from IB Gateway..."); Client.eDisconnect(); } } private void ReadLoop() { try { while (_readerRunning && Client.IsConnected()) { _signal.waitForSignal(); _reader?.processMsgs(); } } catch (Exception ex) { Console.WriteLine($"[IBGatewayControl reader error] {ex}"); } } public override void connectAck() { Console.WriteLine("[connectAck]"); } public override void nextValidId(int orderId) { _apiReady = true; Console.WriteLine($"[nextValidId] {orderId}"); } public override void managedAccounts(string accountsList) { Console.WriteLine($"[managedAccounts] {accountsList}"); } public override void connectionClosed() { _apiReady = false; Console.WriteLine("[connectionClosed]"); } public override void error(Exception e) { Console.WriteLine($"[error exception] {e}"); } public override void error(string str) { Console.WriteLine($"[error string] {str}"); } public override void error(int id, long errorCode, int errorTime, string errorMsg, string advancedOrderRejectJson) { if (errorCode == 2104 || errorCode == 2106 || errorCode == 2158) return; Console.WriteLine($"[error] id={id}, code={errorCode}, time={errorTime}, msg={errorMsg}"); // Live requests if (id > 0 && _liveRequests.TryRemove(id, out var liveTcs)) { Client.cancelMktData(id); _liveSnapshotsByRequest.TryRemove(id, out _); liveTcs.TrySetException( new InvalidOperationException($"IB error for live reqId={id}: {errorCode} - {errorMsg}")); return; } // Historical requests 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; } // Fundamentals if (id > 0 && _fundamentalsRequests.TryRemove(id, out var fundamentalsTcs)) { Client.cancelMktData(id); _fundamentalsSnapshotsByRequest.TryRemove(id, out _); fundamentalsTcs.TrySetException( new InvalidOperationException($"IB error for fundamentals reqId={id}: {errorCode} - {errorMsg}")); } } public void RequestHistorical(HistoricalEquityRequest request) { ThrowIfDisposed(); request.Validate(); if (!Client.IsConnected()) throw new InvalidOperationException("IB Gateway is not connected."); if (!_apiReady) throw new InvalidOperationException("IB API is not ready yet."); var contract = new Contract { Symbol = request.Symbol, SecType = "STK", Exchange = "SMART", PrimaryExch = "NASDAQ", Currency = "USD" }; _BarsByRequest[request.RequestId] = new List(); Console.WriteLine( $"Requesting historical data: Symbol={request.Symbol}, End={request.EndDateTime}, Duration={request.Duration}, BarSize={request.BarSize.ToIbString()}"); Client.reqHistoricalData( request.RequestId, contract, request.EndDateTime, request.Duration, request.BarSize.ToIbString(), request.WhatToShow.ToIbString(), request.UseRth ? 1 : 0, request.FormatDate, request.KeepUpToDate, null ); } public List RequestHistoricalBlocking(HistoricalEquityRequest request, TimeSpan timeout) { ThrowIfDisposed(); request.Validate(); if (_historyRequests.ContainsKey(request.RequestId)) throw new InvalidOperationException($"RequestId {request.RequestId} is already in use."); var tcs = new TaskCompletionSource>(TaskCreationOptions.RunContinuationsAsynchronously); if (!_historyRequests.TryAdd(request.RequestId, tcs)) throw new InvalidOperationException($"Unable to register RequestId {request.RequestId}."); try { RequestHistorical(request); if (tcs.Task.Wait(timeout)) return tcs.Task.GetAwaiter().GetResult(); _historyRequests.TryRemove(request.RequestId, out _); _BarsByRequest.TryRemove(request.RequestId, out _); throw new TimeoutException($"Timed out waiting for historical data for reqId={request.RequestId}."); } catch { _historyRequests.TryRemove(request.RequestId, out _); _BarsByRequest.TryRemove(request.RequestId, out _); throw; } } public override void historicalData(int reqId, Bar Bar) { if (_BarsByRequest.TryGetValue(reqId, out var Bars)) { lock (Bars) { Bars.Add(Bar); } } Console.WriteLine( $"[historicalData] reqId={reqId}, Date={Bar.Time}, Open={Bar.Open}, High={Bar.High}, Low={Bar.Low}, Close={Bar.Close}, Volume={Bar.Volume}"); } public override void historicalDataEnd(int reqId, string start, string end) { Console.WriteLine($"[historicalDataEnd] reqId={reqId}, start={start}, end={end}"); if (_BarsByRequest.TryRemove(reqId, out var Bars)) { Console.WriteLine(); Console.WriteLine("=== Final Printed Bars ==="); lock (Bars) { foreach (var Bar in Bars) { Console.WriteLine( $"{Bar.Time} | O:{Bar.Open} H:{Bar.High} L:{Bar.Low} C:{Bar.Close} V:{Bar.Volume}"); } } if (_historyRequests.TryRemove(reqId, out var tcs)) { tcs.TrySetResult(Bars); } } } public void Dispose() { if (_disposed) return; Disconnect(); _disposed = true; GC.SuppressFinalize(this); } private void ThrowIfDisposed() { if (_disposed) throw new ObjectDisposedException(nameof(IBGatewayControl)); } private static bool IsValidNumber(double value) { return !double.IsNaN(value) && !double.IsInfinity(value) && value > -99999.0; } private static bool TryParseDouble(string raw, out double value) { return double.TryParse(raw, NumberStyles.Any, CultureInfo.InvariantCulture, out value); } private static Dictionary ParseFundamentalRatios(string raw) { var result = new Dictionary(StringComparer.OrdinalIgnoreCase); if (string.IsNullOrWhiteSpace(raw)) return result; string[] parts = raw.Split(';', StringSplitOptions.RemoveEmptyEntries); foreach (string part in parts) { int idx = part.IndexOf('='); if (idx <= 0 || idx >= part.Length - 1) continue; string key = part[..idx].Trim(); string value = part[(idx + 1)..].Trim(); result[key] = value; } return result; } private static void TryFillFundamentals(FundamentalsEquitySnapshot snapshot, string raw) { var ratios = ParseFundamentalRatios(raw); if (!snapshot.MarketCap.HasValue && ratios.TryGetValue("MKTCAP", out var marketCapRaw) && TryParseDouble(marketCapRaw, out var marketCap) && IsValidNumber(marketCap)) { snapshot.MarketCap = marketCap; } if (!snapshot.PERatio.HasValue) { if (ratios.TryGetValue("PEEXCLXOR", out var peRaw) && TryParseDouble(peRaw, out var pe) && IsValidNumber(pe)) { snapshot.PERatio = pe; } else if (ratios.TryGetValue("APENORM", out var peNormRaw) && TryParseDouble(peNormRaw, out var peNorm) && IsValidNumber(peNorm)) { snapshot.PERatio = peNorm; } } if (!snapshot.DividendYield.HasValue && ratios.TryGetValue("YIELD", out var yieldRaw) && TryParseDouble(yieldRaw, out var divYield) && IsValidNumber(divYield)) { snapshot.DividendYield = divYield; } } private void CompleteLiveRequestIfReady(int reqId) { if (!_liveSnapshotsByRequest.TryGetValue(reqId, out var snapshot)) return; if (!snapshot.IsComplete) return; Client.cancelMktData(reqId); _liveSnapshotsByRequest.TryRemove(reqId, out _); if (_liveRequests.TryRemove(reqId, out var tcs)) tcs.TrySetResult(snapshot); } private void CompleteFundamentalsRequestIfReady(int reqId) { if (!_fundamentalsSnapshotsByRequest.TryGetValue(reqId, out var snapshot)) return; if (!snapshot.IsComplete) return; Client.cancelMktData(reqId); _fundamentalsSnapshotsByRequest.TryRemove(reqId, out _); if (_fundamentalsRequests.TryRemove(reqId, out var tcs)) tcs.TrySetResult(snapshot); } public LiveEquitySnapshot RequestLiveEquityBlocking(LiveEquityRequest request, TimeSpan timeout) { ThrowIfDisposed(); request.Validate(); if (!Client.IsConnected()) throw new InvalidOperationException("IB Gateway is not connected."); if (!_apiReady) throw new InvalidOperationException("IB API is not ready yet."); if (_liveRequests.ContainsKey(request.RequestId)) throw new InvalidOperationException($"Live RequestId {request.RequestId} is already in use."); var tcs = new TaskCompletionSource( TaskCreationOptions.RunContinuationsAsynchronously); if (!_liveRequests.TryAdd(request.RequestId, tcs)) throw new InvalidOperationException($"Unable to register live request {request.RequestId}."); var snapshot = new LiveEquitySnapshot { Symbol = request.Symbol }; _liveSnapshotsByRequest[request.RequestId] = snapshot; try { var contract = new Contract { Symbol = request.Symbol, SecType = request.SecType, Exchange = request.Exchange, PrimaryExch = request.PrimaryExchange, Currency = request.Currency }; Console.WriteLine( $"Requesting live quote data: Symbol={request.Symbol}, Ticks=Last/Volume/Close/Open"); Client.reqMktData( request.RequestId, contract, "", false, request.RegulatorySnapshot, null ); if (tcs.Task.Wait(timeout)) return tcs.Task.GetAwaiter().GetResult(); Client.cancelMktData(request.RequestId); _liveRequests.TryRemove(request.RequestId, out _); _liveSnapshotsByRequest.TryRemove(request.RequestId, out _); throw new TimeoutException( $"Timed out waiting for live quote data for reqId={request.RequestId}."); } catch { Client.cancelMktData(request.RequestId); _liveRequests.TryRemove(request.RequestId, out _); _liveSnapshotsByRequest.TryRemove(request.RequestId, out _); throw; } } public FundamentalsEquitySnapshot RequestFundamentalsBlocking(FundamentalsEquityRequest request, TimeSpan timeout) { ThrowIfDisposed(); request.Validate(); if (!Client.IsConnected()) throw new InvalidOperationException("IB Gateway is not connected."); if (!_apiReady) throw new InvalidOperationException("IB API is not ready yet."); if (_fundamentalsRequests.ContainsKey(request.RequestId)) throw new InvalidOperationException($"Fundamentals RequestId {request.RequestId} is already in use."); var tcs = new TaskCompletionSource( TaskCreationOptions.RunContinuationsAsynchronously); if (!_fundamentalsRequests.TryAdd(request.RequestId, tcs)) throw new InvalidOperationException($"Unable to register fundamentals request {request.RequestId}."); var snapshot = new FundamentalsEquitySnapshot { Symbol = request.Symbol }; _fundamentalsSnapshotsByRequest[request.RequestId] = snapshot; try { var contract = new Contract { Symbol = request.Symbol, SecType = request.SecType, Exchange = request.Exchange, PrimaryExch = request.PrimaryExchange, Currency = request.Currency }; Console.WriteLine( $"Requesting fundamentals data: Symbol={request.Symbol}, GenericTicks=258"); Client.reqMktData( request.RequestId, contract, "258", false, request.RegulatorySnapshot, null ); if (tcs.Task.Wait(timeout)) return tcs.Task.GetAwaiter().GetResult(); Client.cancelMktData(request.RequestId); _fundamentalsRequests.TryRemove(request.RequestId, out _); _fundamentalsSnapshotsByRequest.TryRemove(request.RequestId, out _); throw new TimeoutException( $"Timed out waiting for fundamentals data for reqId={request.RequestId}."); } catch { Client.cancelMktData(request.RequestId); _fundamentalsRequests.TryRemove(request.RequestId, out _); _fundamentalsSnapshotsByRequest.TryRemove(request.RequestId, out _); throw; } } public override void tickPrice(int tickerId, int field, double price, TickAttrib attribs) { if (_liveSnapshotsByRequest.TryGetValue(tickerId, out var liveSnapshot)) { if (IsValidNumber(price)) { switch (field) { case 4: // Last Price liveSnapshot.LastPrice = price; break; case 9: // Close Price liveSnapshot.PreviousClose = price; break; case 14: // Open Tick liveSnapshot.Open = price; break; } } CompleteLiveRequestIfReady(tickerId); } Console.WriteLine($"[tickPrice] reqId={tickerId}, field={field}, price={price}"); } public override void tickSize(int tickerId, int field, decimal size) { if (_liveSnapshotsByRequest.TryGetValue(tickerId, out var liveSnapshot)) { switch (field) { case 8: // Volume liveSnapshot.Volume = (long)size; break; } CompleteLiveRequestIfReady(tickerId); } Console.WriteLine($"[tickSize] reqId={tickerId}, field={field}, size={size}"); } public override void tickString(int tickerId, int field, string value) { if (_fundamentalsSnapshotsByRequest.TryGetValue(tickerId, out var fundamentalsSnapshot)) { if (field == 47) { TryFillFundamentals(fundamentalsSnapshot, value); CompleteFundamentalsRequestIfReady(tickerId); } } Console.WriteLine($"[tickString] reqId={tickerId}, field={field}, value={value}"); } } }