SoftTraderBackend/MarketDataRequest/IBGatewayControl.cs

257 lines
6.3 KiB
C#

using IBApi;
using System.Collections.Concurrent;
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;
private readonly ConcurrentDictionary<int, List<Bar>> _barsByRequest = new();
private readonly ConcurrentDictionary<int, TaskCompletionSource<List<Bar>>> _historyRequests = 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}");
if (id > 0 && _historyRequests.TryRemove(id, out var tcs))
{
_barsByRequest.TryRemove(id, out _);
tcs.TrySetException(
new InvalidOperationException($"IB error for 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<Bar>();
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<Bar> 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<List<Bar>>(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));
}
}
}