Skip to content

Commit 116a7f5

Browse files
Enable concurrent order processing (#29)
* Refactor `OandaBrokerage` to support concurrent message processing with `BrokerageConcurrentMessageHandler`. Simplify locks and enhance robustness. * replace some legacy methods * Fix `while` loop to check for infinite loop in `OandaRestApiBase` streaming task. Improve comment clarity.
1 parent 6257bea commit 116a7f5

4 files changed

Lines changed: 83 additions & 65 deletions

File tree

QuantConnect.OandaBrokerage/OandaBrokerage.cs

Lines changed: 8 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -89,10 +89,12 @@ public OandaBrokerage(IOrderProvider orderProvider, ISecurityProvider securityPr
8989
/// <summary>
9090
/// Returns true if we're currently connected to the broker
9191
/// </summary>
92-
public override bool IsConnected
93-
{
94-
get { return _api.IsConnected; }
95-
}
92+
public override bool IsConnected => _api.IsConnected;
93+
94+
/// <summary>
95+
/// Enables or disables concurrent processing of messages to and from the brokerage.
96+
/// </summary>
97+
public override bool ConcurrencyEnabled => true;
9698

9799
/// <summary>
98100
/// Returns the brokerage account's base currency
@@ -438,7 +440,7 @@ private void Initialize(IOrderProvider orderProvider, ISecurityProvider security
438440
if (environment != Environment.Trade && environment != Environment.Practice)
439441
throw new NotSupportedException("Oanda Environment not supported: " + environment);
440442

441-
_api = new OandaRestApiV20(_symbolMapper, orderProvider, securityProvider, aggregator, environment, accessToken, accountId, agent);
443+
_api = new OandaRestApiV20(_symbolMapper, orderProvider, securityProvider, aggregator, environment, accessToken, accountId, agent, ConcurrencyEnabled);
442444

443445
// forward events received from API
444446
_api.OrdersStatusChanged += (sender, orderEvents) => OnOrderEvents(orderEvents);
@@ -586,4 +588,4 @@ private static void ValidateSubscription()
586588
}
587589

588590
}
589-
}
591+
}

QuantConnect.OandaBrokerage/OandaRestApiBase.cs

Lines changed: 19 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -11,7 +11,7 @@
1111
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
1212
* See the License for the specific language governing permissions and
1313
* limitations under the License.
14-
*/
14+
*/
1515

1616
using System;
1717
using System.Collections.Concurrent;
@@ -37,19 +37,14 @@ namespace QuantConnect.Brokerages.Oanda
3737
public abstract class OandaRestApiBase : Brokerage, IDataQueueHandler
3838
{
3939
private static readonly TimeSpan SubscribeDelay = TimeSpan.FromMilliseconds(250);
40-
private readonly ManualResetEvent _refreshEvent = new ManualResetEvent(false);
41-
private readonly CancellationTokenSource _streamingCancellationTokenSource = new CancellationTokenSource();
42-
40+
private readonly ManualResetEventSlim _refreshEvent = new(false);
41+
private readonly CancellationTokenSource _streamingCancellationTokenSource = new();
4342
private bool _isConnected;
4443

45-
/// <summary>
46-
/// This lock is used to sync 'PlaceOrder' and callback 'OnTransactionDataReceived'
47-
/// </summary>
48-
protected readonly object Locker = new object();
4944
/// <summary>
5045
/// This container is used to keep pending to be filled market orders, so when the callback comes in we send the filled event
5146
/// </summary>
52-
protected readonly ConcurrentDictionary<int, OrderStatus> PendingFilledMarketOrders = new ConcurrentDictionary<int, OrderStatus>();
47+
protected readonly ConcurrentDictionary<int, OrderStatus> PendingFilledMarketOrders = new();
5348

5449
/// <summary>
5550
/// The connection handler for pricing
@@ -162,9 +157,9 @@ protected OandaRestApiBase(OandaSymbolMapper symbolMapper, IOrderProvider orderP
162157
Task.Factory.StartNew(
163158
() =>
164159
{
165-
do
160+
while (!_streamingCancellationTokenSource.IsCancellationRequested)
166161
{
167-
_refreshEvent.WaitOne();
162+
_refreshEvent.Wait(_streamingCancellationTokenSource.Token);
168163
Thread.Sleep(SubscribeDelay);
169164

170165
if (!_isConnected)
@@ -174,13 +169,20 @@ protected OandaRestApiBase(OandaSymbolMapper symbolMapper, IOrderProvider orderP
174169

175170
_refreshEvent.Reset();
176171

172+
// no need to subscribe again if cancellation called during timeout
173+
if (_streamingCancellationTokenSource.IsCancellationRequested)
174+
{
175+
break;
176+
}
177+
177178
var symbolsToSubscribe = SubscribedSymbols;
178179
// restart streaming session
179180
SubscribeSymbols(symbolsToSubscribe);
180-
181-
} while (!_streamingCancellationTokenSource.IsCancellationRequested);
181+
}
182182
},
183-
TaskCreationOptions.LongRunning
183+
_streamingCancellationTokenSource.Token,
184+
TaskCreationOptions.LongRunning,
185+
TaskScheduler.Default
184186
);
185187
}
186188

@@ -249,7 +251,7 @@ public override void Dispose()
249251
Aggregator.DisposeSafely();
250252
_refreshEvent.DisposeSafely();
251253

252-
_streamingCancellationTokenSource.Cancel();
254+
_streamingCancellationTokenSource.DisposeSafely();
253255

254256
PricingConnectionHandler.ConnectionLost -= OnPricingConnectionLost;
255257
PricingConnectionHandler.ConnectionRestored -= OnPricingConnectionRestored;
@@ -296,6 +298,7 @@ public override void Disconnect()
296298
StopPricingStream();
297299

298300
_isConnected = false;
301+
_streamingCancellationTokenSource.Cancel();
299302
}
300303

301304
/// <summary>
@@ -453,4 +456,4 @@ protected void EmitTick(Tick tick)
453456
Aggregator.Update(tick);
454457
}
455458
}
456-
}
459+
}

QuantConnect.OandaBrokerage/OandaRestApiV20.cs

Lines changed: 53 additions & 40 deletions
Original file line numberDiff line numberDiff line change
@@ -52,7 +52,8 @@ public class OandaRestApiV20 : OandaRestApiBase
5252

5353
private TransactionStreamSession _eventsSession;
5454
private PricingStreamSession _ratesSession;
55-
private readonly Dictionary<string, (Symbol Symbol, DateTimeZone ExchangeTimeZone)> _symbolsByInstrument = new Dictionary<string, (Symbol, DateTimeZone)>();
55+
private readonly Dictionary<string, (Symbol Symbol, DateTimeZone ExchangeTimeZone)> _symbolsByInstrument = new ();
56+
private BrokerageConcurrentMessageHandler<string> _messageHandler;
5657

5758
/// <summary>
5859
/// Initializes a new instance of the <see cref="OandaRestApiV20"/> class.
@@ -65,7 +66,8 @@ public class OandaRestApiV20 : OandaRestApiBase
6566
/// <param name="accessToken">The Oanda access token (can be the user's personal access token or the access token obtained with OAuth by QC on behalf of the user)</param>
6667
/// <param name="accountId">The account identifier.</param>
6768
/// <param name="agent">The Oanda agent string</param>
68-
public OandaRestApiV20(OandaSymbolMapper symbolMapper, IOrderProvider orderProvider, ISecurityProvider securityProvider, IDataAggregator aggregator, Environment environment, string accessToken, string accountId, string agent)
69+
/// <param name="concurrencyEnabled">Enables or disables concurrent processing of messages</param>
70+
public OandaRestApiV20(OandaSymbolMapper symbolMapper, IOrderProvider orderProvider, ISecurityProvider securityProvider, IDataAggregator aggregator, Environment environment, string accessToken, string accountId, string agent, bool concurrencyEnabled)
6971
: base(symbolMapper, orderProvider, securityProvider, aggregator, environment, accessToken, accountId, agent)
7072
{
7173
var basePathRest = environment == Environment.Trade ?
@@ -76,6 +78,7 @@ public OandaRestApiV20(OandaSymbolMapper symbolMapper, IOrderProvider orderProvi
7678
"https://stream-fxtrade.oanda.com/v3" :
7779
"https://stream-fxpractice.oanda.com/v3";
7880

81+
_messageHandler = new BrokerageConcurrentMessageHandler<string>(OnTransactionDataReceived, concurrencyEnabled);
7982
_apiRest = new DefaultApi(basePathRest);
8083
_apiRest.Configuration.AddDefaultHeader(OandaAgentKey, Agent);
8184

@@ -135,11 +138,11 @@ public override List<CashAmount> GetCashBalance()
135138
{
136139
var response = _apiRest.GetAccountSummary(Authorization, AccountId);
137140

138-
return new List<CashAmount>
139-
{
141+
return
142+
[
140143
new CashAmount(response.Account.Balance.ToDecimal(),
141144
response.Account.Currency)
142-
};
145+
];
143146
}
144147

145148
/// <summary>
@@ -172,8 +175,8 @@ public override bool PlaceOrder(Order order)
172175
var marketOrderStatus = OrderStatus.Filled;
173176
var request = GenerateOrderRequest(order);
174177

175-
lock (Locker)
176-
{
178+
var submitted = true;
179+
_messageHandler.WithLockedStream(() => {
177180
var response = _apiRest.CreateOrder(Authorization, AccountId, request);
178181
order.BrokerId.Add(response.Data.OrderCreateTransaction.Id);
179182

@@ -184,7 +187,8 @@ public override bool PlaceOrder(Order order)
184187
if (response.Data.OrderCancelTransaction != null && response.Data.OrderCancelTransaction.Type == OrderCancelTransaction.TypeEnum.ORDERCANCEL)
185188
{
186189
OnOrderEvent(new OrderEvent(order, DateTime.UtcNow, orderFee, response.Data.OrderCancelTransaction.Reason.ToString()) { Status = OrderStatus.Invalid });
187-
return false;
190+
submitted = false;
191+
return;
188192
}
189193

190194
var fill = response.Data.OrderFillTransaction;
@@ -213,7 +217,7 @@ public override bool PlaceOrder(Order order)
213217
PendingFilledMarketOrders[order.Id] = marketOrderStatus;
214218
}
215219
}
216-
}
220+
});
217221
OnOrderEvent(new OrderEvent(order, DateTime.UtcNow, orderFee) { Status = OrderStatus.Submitted });
218222

219223
// If 'marketOrderRemainingQuantity < order.AbsoluteQuantity' is false it means the order was not even PartiallyFilled, wait for callback
@@ -227,7 +231,7 @@ public override bool PlaceOrder(Order order)
227231
});
228232
}
229233

230-
return true;
234+
return submitted;
231235
}
232236

233237
/// <summary>
@@ -247,25 +251,27 @@ public override bool UpdateOrder(Order order)
247251
}
248252

249253
var request = GenerateOrderRequest(order);
254+
_messageHandler.WithLockedStream(() => {
250255

251-
var orderId = order.BrokerId.First();
252-
var response = _apiRest.ReplaceOrder(Authorization, AccountId, orderId, request);
256+
var orderId = order.BrokerId.First();
257+
var response = _apiRest.ReplaceOrder(Authorization, AccountId, orderId, request);
253258

254-
// replace the brokerage order id
255-
order.BrokerId[0] = response.Data.OrderCreateTransaction.Id;
259+
// replace the brokerage order id
260+
order.BrokerId[0] = response.Data.OrderCreateTransaction.Id;
256261

257-
OnOrderEvent(new OrderEvent(order, DateTime.UtcNow, OrderFee.Zero) { Status = OrderStatus.UpdateSubmitted });
262+
OnOrderEvent(new OrderEvent(order, DateTime.UtcNow, OrderFee.Zero) { Status = OrderStatus.UpdateSubmitted });
258263

259-
// check if the updated (marketable) order was filled
260-
if (response.Data.OrderFillTransaction != null)
261-
{
262-
OnOrderEvent(new OrderEvent(order, DateTime.UtcNow, OrderFee.Zero, "Oanda Fill Event")
264+
// check if the updated (marketable) order was filled
265+
if (response.Data.OrderFillTransaction != null)
263266
{
264-
Status = OrderStatus.Filled,
265-
FillPrice = response.Data.OrderFillTransaction.Price.ToDecimal(),
266-
FillQuantity = response.Data.OrderFillTransaction.Units.ConvertInvariant<decimal>()
267-
});
268-
}
267+
OnOrderEvent(new OrderEvent(order, DateTime.UtcNow, OrderFee.Zero, "Oanda Fill Event")
268+
{
269+
Status = OrderStatus.Filled,
270+
FillPrice = response.Data.OrderFillTransaction.Price.ToDecimal(),
271+
FillQuantity = response.Data.OrderFillTransaction.Units.ConvertInvariant<decimal>()
272+
});
273+
}
274+
});
269275

270276
return true;
271277
}
@@ -285,14 +291,16 @@ public override bool CancelOrder(Order order)
285291
return false;
286292
}
287293

288-
foreach (var orderId in order.BrokerId)
289-
{
290-
_apiRest.CancelOrder(Authorization, AccountId, orderId);
291-
OnOrderEvent(new OrderEvent(order,
292-
DateTime.UtcNow,
293-
OrderFee.Zero,
294-
"Oanda Cancel Order Event") { Status = OrderStatus.Canceled });
295-
}
294+
_messageHandler.WithLockedStream(() => {
295+
foreach (var orderId in order.BrokerId)
296+
{
297+
_apiRest.CancelOrder(Authorization, AccountId, orderId);
298+
OnOrderEvent(new OrderEvent(order,
299+
DateTime.UtcNow,
300+
OrderFee.Zero,
301+
"Oanda Cancel Order Event") { Status = OrderStatus.Canceled });
302+
}
303+
});
296304

297305
return true;
298306
}
@@ -303,7 +311,7 @@ public override bool CancelOrder(Order order)
303311
public override void StartTransactionStream()
304312
{
305313
_eventsSession = new TransactionStreamSession(this);
306-
_eventsSession.DataReceived += OnTransactionDataReceived;
314+
_eventsSession.DataReceived += OnDataReceived;
307315
_eventsSession.StartSession();
308316
}
309317

@@ -314,7 +322,7 @@ public override void StopTransactionStream()
314322
{
315323
if (_eventsSession != null)
316324
{
317-
_eventsSession.DataReceived -= OnTransactionDataReceived;
325+
_eventsSession.DataReceived -= OnDataReceived;
318326
_eventsSession.StopSession();
319327
}
320328
}
@@ -351,6 +359,15 @@ private static DateTime GetTickDateTimeFromString(string time)
351359
return OandaBrokerage.GetDateTimeFromString(time.Remove(25, 3));
352360
}
353361

362+
/// <summary>
363+
/// Handles the data received from the transaction stream session.
364+
/// </summary>
365+
/// <param name="data">The raw data string received from the transaction stream.</param>
366+
private void OnDataReceived(string data)
367+
{
368+
_messageHandler.HandleNewMessage(data);
369+
}
370+
354371
/// <summary>
355372
/// Event handler for streaming events
356373
/// </summary>
@@ -369,11 +386,7 @@ private void OnTransactionDataReceived(string json)
369386
case "ORDER_FILL":
370387
var transaction = obj.ToObject<OrderFillTransaction>();
371388

372-
Order order;
373-
lock (Locker)
374-
{
375-
order = OrderProvider.GetOrdersByBrokerageId(transaction.OrderID)?.SingleOrDefault();
376-
}
389+
Order order = OrderProvider.GetOrdersByBrokerageId(transaction.OrderID)?.SingleOrDefault();
377390
if (order != null)
378391
{
379392
OrderStatus status;
@@ -800,4 +813,4 @@ private string GenerateOrderRequest(Order order)
800813
return request;
801814
}
802815
}
803-
}
816+
}

QuantConnect.OandaBrokerage/RestV20/Client/ApiClient.cs

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -138,7 +138,7 @@ private RestRequest PrepareRequest(
138138

139139
if (postBody != null) // http body (model or byte[]) parameter
140140
{
141-
if (postBody.GetType() == typeof(String))
141+
if (postBody is string)
142142
{
143143
request.AddParameter("application/json", postBody, ParameterType.RequestBody);
144144
}
@@ -208,7 +208,7 @@ public async System.Threading.Tasks.Task<Object> CallApiAsync(
208208
path, method, queryParams, postBody, headerParams, formParams, fileParams,
209209
pathParams, contentType);
210210
InterceptRequest(request);
211-
var response = await RestClient.ExecuteTaskAsync(request);
211+
var response = await RestClient.ExecuteAsync(request);
212212
InterceptResponse(request, response);
213213
return (Object)response;
214214
}
@@ -480,4 +480,4 @@ public static string SanitizeFilename(string filename)
480480
}
481481
}
482482
}
483-
}
483+
}

0 commit comments

Comments
 (0)