Files
quantconnect--lean/Common/Data/DataQueueHandlerSubscriptionManager.cs
T
Ricardo Andrés Marino Rojas 5fd021996a
Benchmarks / build (push) Has been cancelled
Build & Test Lean / build (push) Has been cancelled
Regression Tests / build (push) Has been cancelled
Report Generator Tests / build (push) Has been cancelled
Research Regression Tests / build (push) Has been cancelled
Python Virtual Environments / build (push) Has been cancelled
API Tests / build (push) Has been cancelled
Fix half of the CA1051 warnings (#8137)
* Fix half of the CA1051 warnings

This warning is about not declaring visible instance fields. There are
something about 500 warnings in the solution, mostly in the QuantConnect and QuantConnect.Algorithm.CSharp projects. I aim to fix one of them in this PR and the other half of them in a second one. To fix it, I'm changing the visible instancce fields for properties.

* fix bugs

* Addressing minor reviews

* More minor fixes

---------

Co-authored-by: Martin Molinero <martin.molinero1@gmail.com>
2024-07-03 15:43:17 -03:00

158 lines
5.8 KiB
C#

/*
* QUANTCONNECT.COM - Democratizing Finance, Empowering Individuals.
* Lean Algorithmic Trading Engine v2.0. Copyright 2014 QuantConnect Corporation.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
using QuantConnect.Interfaces;
using QuantConnect.Logging;
using System;
using System.Collections.Concurrent;
using System.Collections.Generic;
using System.Linq;
namespace QuantConnect.Data
{
/// <summary>
/// Count number of subscribers for each channel (Symbol, Socket) pair
/// </summary>
public abstract class DataQueueHandlerSubscriptionManager : IDisposable
{
/// <summary>
/// Counter
/// </summary>
protected ConcurrentDictionary<Channel, int> SubscribersByChannel { get; init; } = new ConcurrentDictionary<Channel, int>();
/// <summary>
/// Increment number of subscribers for current <see cref="TickType"/>
/// </summary>
/// <param name="dataConfig">defines the subscription configuration data.</param>
public void Subscribe(SubscriptionDataConfig dataConfig)
{
try
{
var channel = GetChannel(dataConfig);
int count;
if (SubscribersByChannel.TryGetValue(channel, out count))
{
SubscribersByChannel.TryUpdate(channel, count + 1, count);
return;
}
if (Subscribe(new[] { dataConfig.Symbol }, dataConfig.TickType))
{
SubscribersByChannel.AddOrUpdate(channel, 1);
}
}
catch (Exception exception)
{
Log.Error(exception);
throw;
}
}
/// <summary>
/// Decrement number of subscribers for current <see cref="TickType"/>
/// </summary>
/// <param name="dataConfig">defines the subscription configuration data.</param>
public void Unsubscribe(SubscriptionDataConfig dataConfig)
{
try
{
var channel = GetChannel(dataConfig);
int count;
if (SubscribersByChannel.TryGetValue(channel, out count))
{
if (count > 1)
{
SubscribersByChannel.TryUpdate(channel, count - 1, count);
return;
}
if (Unsubscribe(new[] { dataConfig.Symbol }, dataConfig.TickType))
{
SubscribersByChannel.TryRemove(channel, out count);
}
}
}
catch (Exception exception)
{
Log.Error(exception);
throw;
}
}
/// <summary>
/// Returns subscribed symbols
/// </summary>
/// <returns>list of <see cref="Symbol"/> currently subscribed</returns>
public IEnumerable<Symbol> GetSubscribedSymbols()
{
return SubscribersByChannel.Keys
.Select(c => c.Symbol)
.Distinct();
}
/// <summary>
/// Checks if there is existing subscriber for current channel
/// </summary>
/// <param name="symbol">Symbol</param>
/// <param name="tickType">Type of tick data</param>
/// <returns>return true if there is one subscriber at least; otherwise false</returns>
public bool IsSubscribed(Symbol symbol, TickType tickType)
{
return SubscribersByChannel.ContainsKey(GetChannel(
symbol,
tickType));
}
/// <summary>
/// Performs application-defined tasks associated with freeing, releasing, or resetting unmanaged resources.
/// </summary>
public virtual void Dispose()
{
}
/// <summary>
/// Describes the way <see cref="IDataQueueHandler"/> implements subscription
/// </summary>
/// <param name="symbols">Symbols to subscribe</param>
/// <param name="tickType">Type of tick data</param>
/// <returns>Returns true if subsribed; otherwise false</returns>
protected abstract bool Subscribe(IEnumerable<Symbol> symbols, TickType tickType);
/// <summary>
/// Describes the way <see cref="IDataQueueHandler"/> implements unsubscription
/// </summary>
/// <param name="symbols">Symbols to unsubscribe</param>
/// <param name="tickType">Type of tick data</param>
/// <returns>Returns true if unsubsribed; otherwise false</returns>
protected abstract bool Unsubscribe(IEnumerable<Symbol> symbols, TickType tickType);
/// <summary>
/// Brokerage maps <see cref="TickType"/> to real socket/api channel
/// </summary>
/// <param name="tickType">Type of tick data</param>
/// <returns></returns>
protected abstract string ChannelNameFromTickType(TickType tickType);
private Channel GetChannel(SubscriptionDataConfig dataConfig) => GetChannel(dataConfig.Symbol, dataConfig.TickType);
private Channel GetChannel(Symbol symbol, TickType tickType)
{
return new Channel(
ChannelNameFromTickType(tickType),
symbol);
}
}
}