diff --git a/CryptoExchange.Net/SharedApis/SharedUtils.cs b/CryptoExchange.Net/SharedApis/SharedUtils.cs index a4f9d757..58f62c4a 100644 --- a/CryptoExchange.Net/SharedApis/SharedUtils.cs +++ b/CryptoExchange.Net/SharedApis/SharedUtils.cs @@ -1,4 +1,5 @@ using CryptoExchange.Net.Objects; +using CryptoExchange.Net.Objects.Sockets; using Microsoft.Extensions.DependencyInjection; using System; using System.Collections.Generic; @@ -624,5 +625,187 @@ namespace CryptoExchange.Net.SharedApis .ParallelEnumerateAsync(); } + + /// + /// Subscribe to trade updates for all capabilities in parallel and return results as they arrive + /// + public static IAsyncEnumerable> SubscribeAllAsync( + this IEnumerable> capabilities, + SubscribeTradeRequest request, + Action> onData, + CancellationToken ct = default) + { + return capabilities + .Select(x => x.Capability.SubscribeToTradeUpdatesAsync(request, onData, ct)) + .ParallelEnumerateAsync(); + } + + /// + /// Subscribe to balance updates for all capabilities in parallel and return results as they arrive + /// + public static IAsyncEnumerable> SubscribeAllAsync( + this IEnumerable> capabilities, + SubscribeBalancesRequest request, + Action> onData, + CancellationToken ct = default) + { + return capabilities + .Select(x => x.Capability.SubscribeToBalanceUpdatesAsync(request, onData, ct)) + .ParallelEnumerateAsync(); + } + + /// + /// Subscribe to index price updates for all capabilities in parallel and return results as they arrive + /// + public static IAsyncEnumerable> SubscribeAllAsync( + this IEnumerable> capabilities, + SubscribeIndexPriceRequest request, + Action> onData, + CancellationToken ct = default) + { + return capabilities + .Select(x => x.Capability.SubscribeToIndexPriceUpdatesAsync(request, onData, ct)) + .ParallelEnumerateAsync(); + } + + /// + /// Subscribe to kline updates for all capabilities in parallel and return results as they arrive + /// + public static IAsyncEnumerable> SubscribeAllAsync( + this IEnumerable> capabilities, + SubscribeKlineRequest request, + Action> onData, + CancellationToken ct = default) + { + return capabilities + .Select(x => x.Capability.SubscribeToKlineUpdatesAsync(request, onData, ct)) + .ParallelEnumerateAsync(); + } + + /// + /// Subscribe to mark price updates for all capabilities in parallel and return results as they arrive + /// + public static IAsyncEnumerable> SubscribeAllAsync( + this IEnumerable> capabilities, + SubscribeMarkPriceRequest request, + Action> onData, + CancellationToken ct = default) + { + return capabilities + .Select(x => x.Capability.SubscribeToMarkPriceUpdatesAsync(request, onData, ct)) + .ParallelEnumerateAsync(); + } + + /// + /// Subscribe to order book updates for all capabilities in parallel and return results as they arrive + /// + public static IAsyncEnumerable> SubscribeAllAsync( + this IEnumerable> capabilities, + SubscribeOrderBookRequest request, + Action> onData, + CancellationToken ct = default) + { + return capabilities + .Select(x => x.Capability.SubscribeToOrderBookUpdatesAsync(request, onData, ct)) + .ParallelEnumerateAsync(); + } + + /// + /// Subscribe to incremental order book updates for all capabilities in parallel and return results as they arrive + /// + public static IAsyncEnumerable> SubscribeAllAsync( + this IEnumerable> capabilities, + SubscribeOrderBookRequest request, + Action> onData, + CancellationToken ct = default) + { + return capabilities + .Select(x => x.Capability.SubscribeToOrderBookUpdatesAsync(request, onData, ct)) + .ParallelEnumerateAsync(); + } + + /// + /// Subscribe to futures order updates for all capabilities in parallel and return results as they arrive + /// + public static IAsyncEnumerable> SubscribeAllAsync( + this IEnumerable> capabilities, + SubscribeFuturesOrderRequest request, + Action> onData, + CancellationToken ct = default) + { + return capabilities + .Select(x => x.Capability.SubscribeToFuturesOrderUpdatesAsync(request, onData, ct)) + .ParallelEnumerateAsync(); + } + + /// + /// Subscribe to spot order updates for all capabilities in parallel and return results as they arrive + /// + public static IAsyncEnumerable> SubscribeAllAsync( + this IEnumerable> capabilities, + SubscribeSpotOrderRequest request, + Action> onData, + CancellationToken ct = default) + { + return capabilities + .Select(x => x.Capability.SubscribeToSpotOrderUpdatesAsync(request, onData, ct)) + .ParallelEnumerateAsync(); + } + + /// + /// Subscribe to position updates for all capabilities in parallel and return results as they arrive + /// + public static IAsyncEnumerable> SubscribeAllAsync( + this IEnumerable> capabilities, + SubscribePositionRequest request, + Action> onData, + CancellationToken ct = default) + { + return capabilities + .Select(x => x.Capability.SubscribeToPositionUpdatesAsync(request, onData, ct)) + .ParallelEnumerateAsync(); + } + + /// + /// Subscribe to ticker updates for all capabilities in parallel and return results as they arrive + /// + public static IAsyncEnumerable> SubscribeAllAsync( + this IEnumerable> capabilities, + SubscribeTickerRequest request, + Action> onData, + CancellationToken ct = default) + { + return capabilities + .Select(x => x.Capability.SubscribeToTickerUpdatesAsync(request, onData, ct)) + .ParallelEnumerateAsync(); + } + + /// + /// Subscribe to all ticker updates for all capabilities in parallel and return results as they arrive + /// + public static IAsyncEnumerable> SubscribeAllAsync( + this IEnumerable> capabilities, + SubscribeAllTickersRequest request, + Action> onData, + CancellationToken ct = default) + { + return capabilities + .Select(x => x.Capability.SubscribeToAllTickersUpdatesAsync(request, onData, ct)) + .ParallelEnumerateAsync(); + } + + /// + /// Subscribe to user trade updates for all capabilities in parallel and return results as they arrive + /// + public static IAsyncEnumerable> SubscribeAllAsync( + this IEnumerable> capabilities, + SubscribeUserTradeRequest request, + Action> onData, + CancellationToken ct = default) + { + return capabilities + .Select(x => x.Capability.SubscribeToUserTradeUpdatesAsync(request, onData, ct)) + .ParallelEnumerateAsync(); + } } }