-
Notifications
You must be signed in to change notification settings - Fork 1k
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
* Start pool * Prepare for neo-node * Fix double call * Prepare pool * Fix * Join oracle pool with oracle service * dotnet-format * Try to fix UT * UT pass * Fix UT * Fix p2p message * Unify collections * Relay p2p message * Clean code and fixes * Send tx to OracleService * RequestTx will wait for ResponseTx * Rename * Remove supervisor * Allow to put request, and response in the same block * Check the sender of OracleResponses * Remove OracleResponse message in OracleService * Organize code * Fix typo * Remove task count TODO * Remove Thread-safe TODO * Remove TODO * ResponseItem changes * Read the oracle contract on my response Receive oraclePayload * Clean code * Add Stop message * Save only one response for PublicKey/RequestTx Improve sorting response pool * Improve sort * Group sort methods in the same region * Ask for the request TX if I don't have it * Reorder code * Rename * First oracle TX
- Loading branch information
Showing
18 changed files
with
1,097 additions
and
167 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,84 @@ | ||
using System.Collections.Concurrent; | ||
using System.Collections.Generic; | ||
using System.Threading; | ||
|
||
namespace Neo.Ledger | ||
{ | ||
public class SortedBlockingCollection<TKey, TValue> | ||
{ | ||
/// <summary> | ||
/// _oracleTasks will consume from this pool | ||
/// </summary> | ||
private readonly BlockingCollection<TValue> _asyncPool = new BlockingCollection<TValue>(); | ||
|
||
/// <summary> | ||
/// Queue | ||
/// </summary> | ||
private readonly SortedConcurrentDictionary<TKey, TValue> _queue; | ||
|
||
/// <summary> | ||
/// Constructor | ||
/// </summary> | ||
/// <param name="comparer">Comparer</param> | ||
/// <param name="capacity">Capacity</param> | ||
public SortedBlockingCollection(IComparer<KeyValuePair<TKey, TValue>> comparer, int capacity) | ||
{ | ||
_queue = new SortedConcurrentDictionary<TKey, TValue>(comparer, capacity); | ||
} | ||
|
||
/// <summary> | ||
/// Add entry | ||
/// </summary> | ||
/// <param name="key">Key</param> | ||
/// <param name="value">Value</param> | ||
public void Add(TKey key, TValue value) | ||
{ | ||
if (_queue.TryAdd(key, value) && _asyncPool.Count <= 0) | ||
{ | ||
Pop(); | ||
} | ||
} | ||
|
||
/// <summary> | ||
/// Clear | ||
/// </summary> | ||
public void Clear() | ||
{ | ||
_queue.Clear(); | ||
|
||
while (_asyncPool.Count > 0) | ||
{ | ||
_asyncPool.TryTake(out _); | ||
} | ||
} | ||
|
||
/// <summary> | ||
/// Get consuming enumerable | ||
/// </summary> | ||
/// <param name="token">Token</param> | ||
public IEnumerable<TValue> GetConsumingEnumerable(CancellationToken token) | ||
{ | ||
foreach (var entry in _asyncPool.GetConsumingEnumerable(token)) | ||
{ | ||
// Prepare other item in _asyncPool | ||
|
||
Pop(); | ||
|
||
// Iterate items | ||
|
||
yield return entry; | ||
} | ||
} | ||
|
||
/// <summary> | ||
/// Move one item from the sorted queue to _asyncPool, this will ensure that the threads process the entries according to the priority | ||
/// </summary> | ||
private void Pop() | ||
{ | ||
if (_queue.TryPop(out var entry)) | ||
{ | ||
_asyncPool.Add(entry); | ||
} | ||
} | ||
} | ||
} |
Oops, something went wrong.