1
0
mirror of https://github.com/JKorf/CryptoExchange.Net.git synced 2026-08-20 12:53:02 +00:00

Rework order book; added support for checksum

This commit is contained in:
Jan Korf
2020-06-20 20:34:48 +02:00
parent 5105e995e8
commit ddc4ebe638
4 changed files with 257 additions and 38 deletions
+85 -38
View File
@@ -40,7 +40,7 @@ namespace CryptoExchange.Net.OrderBook
private Task _processTask;
private AutoResetEvent _queueEvent;
private ConcurrentQueue<ProcessQueueItem> _processQueue;
private ConcurrentQueue<object> _processQueue;
/// <summary>
/// Order book implementation id
@@ -197,7 +197,7 @@ namespace CryptoExchange.Net.OrderBook
Id = options.OrderBookName;
processBuffer = new List<ProcessBufferRangeSequenceEntry>();
_processQueue = new ConcurrentQueue<ProcessQueueItem>();
_processQueue = new ConcurrentQueue<object>();
_queueEvent = new AutoResetEvent(false);
sequencesAreConsecutive = options.SequenceNumbersAreConsecutive;
@@ -243,7 +243,7 @@ namespace CryptoExchange.Net.OrderBook
private void Reset()
{
log.Write(LogVerbosity.Warning, $"{Id} order book {Symbol} connection lost");
Status = OrderBookStatus.Connecting;
Status = OrderBookStatus.Reconnecting;
_queueEvent.Set();
// Clear queue
while(_processQueue.TryDequeue(out _))
@@ -306,14 +306,58 @@ namespace CryptoExchange.Net.OrderBook
/// <returns></returns>
protected abstract Task<CallResult<bool>> DoResync();
/// <summary>
/// Validate a checksum with the current order book
/// </summary>
/// <param name="checksum"></param>
/// <returns></returns>
protected virtual bool DoChecksum(int checksum) => true;
private void ProcessQueue()
{
while(Status != OrderBookStatus.Disconnected)
{
_queueEvent.WaitOne();
while (_processQueue.TryDequeue(out var item))
ProcessQueueItem(item);
{
if (Status == OrderBookStatus.Disconnected)
break;
if (item is InitialOrderBookItem iobi)
ProcessInitialOrderBookItem(iobi);
if (item is ProcessQueueItem pqi)
ProcessQueueItem(pqi);
else if (item is ChecksumItem ci)
ProcessChecksum(ci);
}
}
}
private void ProcessInitialOrderBookItem(InitialOrderBookItem item)
{
lock (bookLock)
{
if (Status == OrderBookStatus.Connecting || Status == OrderBookStatus.Disconnected)
return;
asks.Clear();
foreach (var ask in item.Asks)
asks.Add(ask.Price, ask);
bids.Clear();
foreach (var bid in item.Bids)
bids.Add(bid.Price, bid);
LastSequenceNumber = item.EndUpdateId;
AskCount = asks.Count;
BidCount = bids.Count;
LastOrderBookUpdate = DateTime.UtcNow;
log.Write(LogVerbosity.Debug, $"{Id} order book {Symbol} data set: {BidCount} bids, {AskCount} asks. #{item.EndUpdateId}");
CheckProcessBuffer();
OnOrderBookUpdate?.Invoke(item.Asks, item.Bids);
OnBestOffersChanged?.Invoke(BestBid, BestAsk);
}
}
@@ -354,6 +398,20 @@ namespace CryptoExchange.Net.OrderBook
}
}
private void ProcessChecksum(ChecksumItem ci)
{
lock (bookLock)
{
var checksumResult = DoChecksum(ci.Checksum);
if(!checksumResult)
{
log.Write(LogVerbosity.Warning, $"{Id} order book {Symbol} out of sync. Resyncing");
_ = subscription?.Reconnect();
return;
}
}
}
/// <summary>
/// Set the initial data for the order book
/// </summary>
@@ -362,38 +420,10 @@ namespace CryptoExchange.Net.OrderBook
/// <param name="bidList">List of bids</param>
protected void SetInitialOrderBook(long orderBookSequenceNumber, IEnumerable<ISymbolOrderBookEntry> bidList, IEnumerable<ISymbolOrderBookEntry> askList)
{
lock (bookLock)
{
if (Status == OrderBookStatus.Connecting || Status == OrderBookStatus.Disconnected)
return;
bookSet = true;
asks.Clear();
foreach (var ask in askList)
asks.Add(ask.Price, ask);
bids.Clear();
foreach (var bid in bidList)
bids.Add(bid.Price, bid);
LastSequenceNumber = orderBookSequenceNumber;
AskCount = asks.Count;
BidCount = asks.Count;
bookSet = true;
LastOrderBookUpdate = DateTime.UtcNow;
log.Write(LogVerbosity.Debug, $"{Id} order book {Symbol} data set: {BidCount} bids, {AskCount} asks. #{orderBookSequenceNumber}");
CheckProcessBuffer();
OnOrderBookUpdate?.Invoke(bidList, askList);
OnBestOffersChanged?.Invoke(BestBid, BestAsk);
}
}
private void CheckBestOffersChanged(ISymbolOrderBookEntry prevBestBid, ISymbolOrderBookEntry prevBestAsk)
{
var (bestBid, bestAsk) = BestOffers;
if (bestBid.Price != prevBestBid.Price || bestBid.Quantity != prevBestBid.Quantity ||
bestAsk.Price != prevBestAsk.Price || bestAsk.Quantity != prevBestAsk.Quantity)
OnBestOffersChanged?.Invoke(bestBid, bestAsk);
_processQueue.Enqueue(new InitialOrderBookItem { StartUpdateId = orderBookSequenceNumber, EndUpdateId = orderBookSequenceNumber, Asks = askList, Bids = bidList });
_queueEvent.Set();
}
/// <summary>
@@ -408,6 +438,16 @@ namespace CryptoExchange.Net.OrderBook
_queueEvent.Set();
}
/// <summary>
/// Add a checksum to the process queue
/// </summary>
/// <param name="checksum"></param>
protected void AddChecksum(int checksum)
{
_processQueue.Enqueue(new ChecksumItem() { Checksum = checksum });
_queueEvent.Set();
}
/// <summary>
/// Update the order book using a first/last update id
/// </summary>
@@ -505,7 +545,6 @@ namespace CryptoExchange.Net.OrderBook
{
// Out of sync
log.Write(LogVerbosity.Warning, $"{Id} order book {Symbol} out of sync (expected { LastSequenceNumber + 1}, was {sequence}), reconnecting");
Status = OrderBookStatus.Connecting;
subscription?.Reconnect();
return false;
}
@@ -531,7 +570,7 @@ namespace CryptoExchange.Net.OrderBook
}
else
{
listToChange[entry.Price].Quantity = entry.Quantity;
listToChange[entry.Price] = entry;
}
}
@@ -557,6 +596,14 @@ namespace CryptoExchange.Net.OrderBook
return new CallResult<bool>(true, null);
}
private void CheckBestOffersChanged(ISymbolOrderBookEntry prevBestBid, ISymbolOrderBookEntry prevBestAsk)
{
var (bestBid, bestAsk) = BestOffers;
if (bestBid.Price != prevBestBid.Price || bestBid.Quantity != prevBestBid.Quantity ||
bestAsk.Price != prevBestAsk.Price || bestAsk.Quantity != prevBestAsk.Quantity)
OnBestOffersChanged?.Invoke(bestBid, bestAsk);
}
/// <summary>
/// Dispose the order book
/// </summary>