eb1181f5f7
* Adds preliminary universe selection for Future Options
* Fixes scaling issues with Future Options
* Fixes scaling multiplying by 10000x instead of using _scaleFactor
* Fixes scaling for Tick
* Revert changes to Tick since it divides the scaling factor
* Changes stale method name to new method name after rebase
* Fixes selection bugs, adds new methods, and adds unit tests
* Fixes bug where Equity Symbol was created for an underlying
non-equity Symbol, resulting in equity data trying to be loaded
* Adds unit tests covering changes to Tick, QuoteBar, TradeBar and
LeanData
* Adds regression test for AddUniverseOption filter contract selection
for Future Options
* Addresses review - modifies the AddFutureOption signature
* Adds new AddUniverseOptions method overload
* Removes and adds a new unit test
* Misc. modifications to account for new changes
* Fixes bug where futures were loaded using default SID Date
* Refactors and removes unnecessary work
* Fixes regression algorithm, which previously made no trades
* Adds future option data
* Adds the corresponding underlying data, in this case, futures data
to enable usage of future options data
* Replaces data with new data (ES18Z20)
* Improves Future chain filtering and updates regression stats
* Add AddFutureOptionContract API
* Expands regression and unit tests to test in finer detail
* Adds Python regression algorithms for AddFutureOption[Contract] methods
* Adds new unit test for BacktestingOptionChainProvider
* Fixes bug with BacktesingOptionChainProvider where we
attempted to load the Trades option chain first, resulting
in breakage of backwards compatibility and limitation of the
option chain.
* Adds new regression algorithms (Py) to Algorithm.Python project
* Adds FutureOptionMarginBuyingPowerModel
* Modifies code paths used to select margin model
* Adds related unit tests for margin model
* Fixes issue with unit test and MHDB/SPDB lookup for Future Options
* Preliminary regression algorithm testing ITM call/put option buying
* Fixes bug where fee model used did not find non-US market
options fee model. We now use the futures fee model for future
options because IB charges the same commissions per contract
between futures and futures options
* Adds proper regression algorithm for ITM future options expiration
* Pushing broken algorithm for review
* Currently, algorithm does not fill forward, causing
a single future option to not get exercised when it is delisted.
* Adds FutureOptionPutITMExpiryRegressionAlgorithm
* Improves existing regression algorithm for call side
* Fixes bug in existing regression algorithm
* Adds AAPL daily data to advance enumerator for ^^^ fix
* Adds additional future option regression algorithms
* Adds Buy OTM expiration regression algorithms
* Adds Sell ITM/OTM expiration regression algorithms
* Adds missing Python regression algorithms
* Adds remaining Python regression algorithms and fixes issues
* Fixes naming issues and statistics
* Adds short option OTM regression algorithms (Py)
* Add license header and class comments to python algorithms
* Cleans up comments and docstrings
* Create Buy/Sell call intraday regression algo
* Redirects future options symbol properties to futures symbol properties
* Asserts exercise/assignment price and updates stats in regression algos
* Adds new unit test covering changes to SecurityService
* Adds comments and fixes failing test
* Partially fixes future option mis-calculated profit/loss
* Adjusts portfolio model to calculate FOP as a no upfront pay asset class
* Updates regression algorithm statistics
* Begin IB FOP support
* Initial support for FOP IB data streaming, live í¾
* Adds additional functionality to LiveOptionChainProvider
- Allows querying CME API to retrieve option chains for CME products
- Ultimately, it's also the groundwork for the CME
LiveFutureChainProvider
* Edits IDataQueueUniverseProvider interface to provide greater
control to implementors of it
* Misc. bug fixes required to get FOP data streaming through IB
* Adds comments, adds missing rategate call, and cleans up code
* Force exchange for FOP and Futures when no exchange is provided
* Fixes bug with Portfolio modeling across all asset classes
* Adds LiveOptionChainProvider tests for Future Options
* IB brokerage option symbol bug fixes and improvements
* Fixes contract multiplier lookup bug
* Fixes issue where we attempted to subscribe to IB data feed with canonical security
* Adds ES MHDB entry
* Reverts portfolio modeling changes for Futures Options
* Since IB eats into our account's cash balance when
a new FOP contract is purchased, we must model by applying funds
to our cash whenever a new purchase/sell occurs.
If we choose to model FOPs exactly as we do with futures, we
will end up with an invalid TotalPortfolioValue on algorithm
restart. By all means and purposes, FOPs are modeled exactly
the same as equity options with respect to the portfolio.
* Adds comments clarifying portfolio modeling and clarifies
existing portfolio modeling comments with additional context.
* Fixes IB symbol lookup for future options
* Fixes LiveOptionChainProvider looping 5 times per option chain
request, even on success
* Sets OptionChainedUniverseSelectionModel to produce a canonical
future/future option/option Symbol to avoid creating two Symbols
* Adds GLOBEX future option symbol mapping from future -> fop
* Fixes LiveOptionChainProvider loading wrong contract option chains
* Fixes loading of futures options ZIP files when backtesting
* Adds a string -> decimal JSON converter
* Additional fixes/refactoring to the LiveOptionChainProvider
* Adds tests for changes to Symbol and LeanData
* Reverts changes to IB-symbol-map
* Fixes Value for mapped future options tickers
* Fixes Symbol test
* Changes path of future options to future's expiry date
* Extra changes made to remove scaling from writing CSV
* Added method to map from FOP Globex -> FUT Globex
* Fixes MOO and MOC orders for future options
* Note: this order type might not be supported by IB or CME.
* Bug fixes and updates unit tests
* Update regression tests and data format
* Rebase changes
* 1. Multiple bug fixes for LiveOptionChainProvider, reverts IQFeed changes
2. Address review (partial): Code reuse and cleanup
1.
* Modifies check in
`AddFutureOptionShort(Call|Put)ITMExpiryRegressionAlgorithm`
to ensure no buys have negative quantity
* Code reuse changes in IB brokerage
* Bug fix in IB brokerage where we assigned the FOP expiry
as the futures expiry (requires verification)
* Doc changes and adds missing summaries/license banners
* Disposes of HTTP client resources in LiveOptionChainProvider
* Renames classes and adds FutureOption folder in Common/Securities
2.
* We revert back to the quotes API for the option chain,
since the settlement API sometimes had missing strikes.
* Fixes future option expiry being set as future's expiry
in LiveOptionChainProvider
* Fixes bug where wrong option chain was selected because of bad
expiry lookup in the futures expiries returned from CME
* Fixes multiple looping bug in LiveOptionChainProvider
* Adds strike price scaling for LiveOptionChainProvider
* Reverts IQFeed changes and simplifies interface upgrade changes
Some additional challenges we'll have to solve as part of FOPs:
- The `OptionSymbol.IsStandard` method makes the assumption that
weeklies contracts follow the pattern equities follows, which
does not apply to Futures Options
- The Subscription created in:
`OptionChainUniverseSubscriptionEnumeratorFactory`
...adds a Trade config. For illiquid contracts, this
will delay universe selection for the option symbol
until we get a trade. However, if we add a quote config,
the data would instead be loaded based on the first quote
we received from the brokerage.
But since we're currently using a trade config, illiquid
contracts won't start streaming data until it receives a trade.
NOTE: this commit is a WIP to addressing the reviews received in the PR,
but has been committed early for efficiency in the review process
* Fixes regression algorithms and misc. bugs
* Fixes map file lookup for non-equity options
* Adds extra assertion at end of algorithm to ensure no holdings are
left when the algorithm ends.
* Adds FutureOptionSymbol, allowing all contracts through as standard
* Changes SPDB to allow defaulting to underlying future symbol
properties if no entry is found for the given FOP
* Fixes calls to SPDB in SecurityService, IBBrokerage
* Reverts AAPL daily ZIP file to fix majority of regression algorithms
* Adds FOPs symbol properties
* Fixes existing symbol properties for a few futures
* Adds tests for changes to Symbol Properties Database
* Removes string SPDB lookup method
* Updates tests and misc callees of previous method
* Updates all regression tests to use data of already expired contracts
* Adds Futures Options Expiry Functions tests
* Adds required futures data for 2020-01-05
* Address review (partial): Expands test coverage and fixes tests
* Set option chain tests parallelism to fixture only
* Fixes broken test for contract month delta for FuturesOptionsExpiryFunctions
* Changes delisting date logic for Futures Options
* Address review: removes duplicate code, misc code fixes
* Bug fix in MarketHoursDatabase.GetDatabaseSymbolKey() where
we would use the underlying's Symbol for lookup in the MHDB
* Adds missing license banner
* Removes Futures Options entries from MHDB
* Adds new tests
* Adds SecurityType.FutureOption
* Converts any underlying comparisons and uses SecurityType directly
instead for FOP specific behavior
* Extra code modifications to acommodate new SecurityType
* Addresses review: fixes order fee bug on exercise
* Additional bug fixes and adding of SecurityType.FutureOption
* Updates regression algorithms OrderListHash
* Fixes various bugs in IB live implementation
* Fixes bug setting the right contract expiration date for FOP
generated by LiveOptionChainProvider
* Adds new function to FuturesOptionsExpiryFunctions
* Clarifies parameter names better in some functions/methods
* Fixes bugs in IB brokerage for FOPs
* Address review - code cleanup and refactor
* Remove MappingEventProvider, SplitEventProvider, and
DividendEventProvider for Futures Options in
CorporateEventEnumeratorFactory
* Address review: Use MHDB key resolver in SPDB
* Makes regression tests pass and adds comment for expiry issue
* Fixes MHDB lookup on string symbol method
* Adds Futures Options greeks regression algorithm (C# only)
* Adds explanitory comment on MHDB FOP lookup
* Remove python from FutureOptionCallITMGreeksExpiryRegressionAlgorithm
576 lines
26 KiB
C#
576 lines
26 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 System;
|
|
using System.Collections.Concurrent;
|
|
using System.Collections.Generic;
|
|
using System.Collections.Specialized;
|
|
using System.Linq;
|
|
using QuantConnect.Data;
|
|
using QuantConnect.Data.Auxiliary;
|
|
using QuantConnect.Data.UniverseSelection;
|
|
using QuantConnect.Interfaces;
|
|
using QuantConnect.Logging;
|
|
using QuantConnect.Securities;
|
|
using QuantConnect.Util;
|
|
|
|
namespace QuantConnect.Lean.Engine.DataFeeds
|
|
{
|
|
/// <summary>
|
|
/// DataManager will manage the subscriptions for both the DataFeeds and the SubscriptionManager
|
|
/// </summary>
|
|
public class DataManager : IAlgorithmSubscriptionManager, IDataFeedSubscriptionManager, IDataManager
|
|
{
|
|
private readonly IAlgorithmSettings _algorithmSettings;
|
|
private readonly IDataFeed _dataFeed;
|
|
private readonly MarketHoursDatabase _marketHoursDatabase;
|
|
private readonly ITimeKeeper _timeKeeper;
|
|
private readonly bool _liveMode;
|
|
private readonly IRegisteredSecurityDataTypesProvider _registeredTypesProvider;
|
|
private readonly IDataPermissionManager _dataPermissionManager;
|
|
|
|
/// There is no ConcurrentHashSet collection in .NET,
|
|
/// so we use ConcurrentDictionary with byte value to minimize memory usage
|
|
private readonly ConcurrentDictionary<SubscriptionDataConfig, SubscriptionDataConfig> _subscriptionManagerSubscriptions
|
|
= new ConcurrentDictionary<SubscriptionDataConfig, SubscriptionDataConfig>();
|
|
|
|
/// <summary>
|
|
/// Event fired when a new subscription is added
|
|
/// </summary>
|
|
public event EventHandler<Subscription> SubscriptionAdded;
|
|
|
|
/// <summary>
|
|
/// Event fired when an existing subscription is removed
|
|
/// </summary>
|
|
public event EventHandler<Subscription> SubscriptionRemoved;
|
|
|
|
/// <summary>
|
|
/// Creates a new instance of the DataManager
|
|
/// </summary>
|
|
public DataManager(
|
|
IDataFeed dataFeed,
|
|
UniverseSelection universeSelection,
|
|
IAlgorithm algorithm,
|
|
ITimeKeeper timeKeeper,
|
|
MarketHoursDatabase marketHoursDatabase,
|
|
bool liveMode,
|
|
IRegisteredSecurityDataTypesProvider registeredTypesProvider,
|
|
IDataPermissionManager dataPermissionManager)
|
|
{
|
|
_dataFeed = dataFeed;
|
|
UniverseSelection = universeSelection;
|
|
UniverseSelection.SetDataManager(this);
|
|
_algorithmSettings = algorithm.Settings;
|
|
AvailableDataTypes = SubscriptionManager.DefaultDataTypes();
|
|
_timeKeeper = timeKeeper;
|
|
_marketHoursDatabase = marketHoursDatabase;
|
|
_liveMode = liveMode;
|
|
_registeredTypesProvider = registeredTypesProvider;
|
|
_dataPermissionManager = dataPermissionManager;
|
|
|
|
// wire ourselves up to receive notifications when universes are added/removed
|
|
algorithm.UniverseManager.CollectionChanged += (sender, args) =>
|
|
{
|
|
switch (args.Action)
|
|
{
|
|
case NotifyCollectionChangedAction.Add:
|
|
foreach (var universe in args.NewItems.OfType<Universe>())
|
|
{
|
|
var config = universe.Configuration;
|
|
var start = algorithm.UtcTime;
|
|
|
|
var end = algorithm.LiveMode ? Time.EndOfTime
|
|
: algorithm.EndDate.ConvertToUtc(algorithm.TimeZone);
|
|
|
|
Security security;
|
|
if (!algorithm.Securities.TryGetValue(config.Symbol, out security))
|
|
{
|
|
// create a canonical security object if it doesn't exist
|
|
security = new Security(
|
|
_marketHoursDatabase.GetExchangeHours(config),
|
|
config,
|
|
algorithm.Portfolio.CashBook[algorithm.AccountCurrency],
|
|
SymbolProperties.GetDefault(algorithm.AccountCurrency),
|
|
algorithm.Portfolio.CashBook,
|
|
RegisteredSecurityDataTypesProvider.Null,
|
|
new SecurityCache()
|
|
);
|
|
}
|
|
AddSubscription(
|
|
new SubscriptionRequest(true,
|
|
universe,
|
|
security,
|
|
config,
|
|
start,
|
|
end));
|
|
}
|
|
break;
|
|
|
|
case NotifyCollectionChangedAction.Remove:
|
|
foreach (var universe in args.OldItems.OfType<Universe>())
|
|
{
|
|
// removing the subscription will be handled by the SubscriptionSynchronizer
|
|
// in the next loop as well as executing a UniverseSelection one last time.
|
|
if (!universe.DisposeRequested)
|
|
{
|
|
universe.Dispose();
|
|
}
|
|
}
|
|
break;
|
|
|
|
default:
|
|
throw new NotImplementedException("The specified action is not implemented: " + args.Action);
|
|
}
|
|
};
|
|
}
|
|
|
|
#region IDataFeedSubscriptionManager
|
|
|
|
/// <summary>
|
|
/// Gets the data feed subscription collection
|
|
/// </summary>
|
|
public SubscriptionCollection DataFeedSubscriptions { get; } = new SubscriptionCollection();
|
|
|
|
/// <summary>
|
|
/// Will remove all current <see cref="Subscription"/>
|
|
/// </summary>
|
|
public void RemoveAllSubscriptions()
|
|
{
|
|
// remove each subscription from our collection
|
|
foreach (var subscription in DataFeedSubscriptions)
|
|
{
|
|
try
|
|
{
|
|
RemoveSubscription(subscription.Configuration);
|
|
}
|
|
catch (Exception err)
|
|
{
|
|
Log.Error(err, "DataManager.RemoveAllSubscriptions():" +
|
|
$"Error removing: {subscription.Configuration}");
|
|
}
|
|
}
|
|
}
|
|
|
|
/// <summary>
|
|
/// Adds a new <see cref="Subscription"/> to provide data for the specified security.
|
|
/// </summary>
|
|
/// <param name="request">Defines the <see cref="SubscriptionRequest"/> to be added</param>
|
|
/// <returns>True if the subscription was created and added successfully, false otherwise</returns>
|
|
public bool AddSubscription(SubscriptionRequest request)
|
|
{
|
|
// guarantee the configuration is present in our config collection
|
|
// this is related to GH issue 3877: where we added a configuration which we also removed
|
|
_subscriptionManagerSubscriptions.TryAdd(request.Configuration, request.Configuration);
|
|
|
|
Subscription subscription;
|
|
if (DataFeedSubscriptions.TryGetValue(request.Configuration, out subscription))
|
|
{
|
|
// duplicate subscription request
|
|
subscription.AddSubscriptionRequest(request);
|
|
// only result true if the existing subscription is internal, we actually added something from the users perspective
|
|
return subscription.Configuration.IsInternalFeed;
|
|
}
|
|
|
|
// before adding the configuration to the data feed let's assert it's valid
|
|
_dataPermissionManager.AssertConfiguration(request.Configuration);
|
|
|
|
subscription = _dataFeed.CreateSubscription(request);
|
|
|
|
if (subscription == null)
|
|
{
|
|
Log.Trace($"DataManager.AddSubscription(): Unable to add subscription for: {request.Configuration}");
|
|
// subscription will be null when there's no tradeable dates for the security between the requested times, so
|
|
// don't even try to load the data
|
|
return false;
|
|
}
|
|
|
|
if (_liveMode)
|
|
{
|
|
OnSubscriptionAdded(subscription);
|
|
Log.Trace($"DataManager.AddSubscription(): Added {request.Configuration}." +
|
|
$" Start: {request.StartTimeUtc}. End: {request.EndTimeUtc}");
|
|
}
|
|
else if(Log.DebuggingEnabled)
|
|
{
|
|
// for performance lets not create the message string if debugging is not enabled
|
|
// this can be executed many times and its in the algorithm thread
|
|
Log.Debug($"DataManager.AddSubscription(): Added {request.Configuration}." +
|
|
$" Start: {request.StartTimeUtc}. End: {request.EndTimeUtc}");
|
|
}
|
|
|
|
return DataFeedSubscriptions.TryAdd(subscription);
|
|
}
|
|
|
|
/// <summary>
|
|
/// Removes the <see cref="Subscription"/>, if it exists
|
|
/// </summary>
|
|
/// <param name="configuration">The <see cref="SubscriptionDataConfig"/> of the subscription to remove</param>
|
|
/// <param name="universe">Universe requesting to remove <see cref="Subscription"/>.
|
|
/// Default value, null, will remove all universes</param>
|
|
/// <returns>True if the subscription was successfully removed, false otherwise</returns>
|
|
public bool RemoveSubscription(SubscriptionDataConfig configuration, Universe universe = null)
|
|
{
|
|
// remove the subscription from our collection, if it exists
|
|
Subscription subscription;
|
|
|
|
if (DataFeedSubscriptions.TryGetValue(configuration, out subscription))
|
|
{
|
|
// we remove the subscription when there are no other requests left
|
|
if (subscription.RemoveSubscriptionRequest(universe))
|
|
{
|
|
if (!DataFeedSubscriptions.TryRemove(configuration, out subscription))
|
|
{
|
|
Log.Error($"DataManager.RemoveSubscription(): Unable to remove {configuration}");
|
|
return false;
|
|
}
|
|
|
|
_dataFeed.RemoveSubscription(subscription);
|
|
|
|
if (_liveMode)
|
|
{
|
|
OnSubscriptionRemoved(subscription);
|
|
}
|
|
|
|
subscription.Dispose();
|
|
|
|
RemoveSubscriptionDataConfig(subscription);
|
|
|
|
if (_liveMode)
|
|
{
|
|
Log.Trace($"DataManager.RemoveSubscription(): Removed {configuration}");
|
|
}
|
|
else if(Log.DebuggingEnabled)
|
|
{
|
|
// for performance lets not create the message string if debugging is not enabled
|
|
// this can be executed many times and its in the algorithm thread
|
|
Log.Debug($"DataManager.RemoveSubscription(): Removed {configuration}");
|
|
}
|
|
return true;
|
|
}
|
|
}
|
|
else if (universe != null)
|
|
{
|
|
// a universe requested removal of a subscription which wasn't present anymore, this can happen when a subscription ends
|
|
// it will get removed from the data feed subscription list, but the configuration will remain until the universe removes it
|
|
// why? the effect I found is that the fill models are using these subscriptions to determine which data they could use
|
|
SubscriptionDataConfig config;
|
|
_subscriptionManagerSubscriptions.TryRemove(configuration, out config);
|
|
}
|
|
return false;
|
|
}
|
|
|
|
/// <summary>
|
|
/// Event invocator for the <see cref="SubscriptionAdded"/> event
|
|
/// </summary>
|
|
/// <param name="subscription">The added subscription</param>
|
|
private void OnSubscriptionAdded(Subscription subscription)
|
|
{
|
|
SubscriptionAdded?.Invoke(this, subscription);
|
|
}
|
|
|
|
/// <summary>
|
|
/// Event invocator for the <see cref="SubscriptionRemoved"/> event
|
|
/// </summary>
|
|
/// <param name="subscription">The removed subscription</param>
|
|
private void OnSubscriptionRemoved(Subscription subscription)
|
|
{
|
|
SubscriptionRemoved?.Invoke(this, subscription);
|
|
}
|
|
|
|
#endregion
|
|
|
|
#region IAlgorithmSubscriptionManager
|
|
|
|
/// <summary>
|
|
/// Gets all the current data config subscriptions that are being processed for the SubscriptionManager
|
|
/// </summary>
|
|
public IEnumerable<SubscriptionDataConfig> SubscriptionManagerSubscriptions =>
|
|
_subscriptionManagerSubscriptions.Select(x => x.Key);
|
|
|
|
/// <summary>
|
|
/// Gets existing or adds new <see cref="SubscriptionDataConfig" />
|
|
/// </summary>
|
|
/// <returns>Returns the SubscriptionDataConfig instance used</returns>
|
|
public SubscriptionDataConfig SubscriptionManagerGetOrAdd(SubscriptionDataConfig newConfig)
|
|
{
|
|
var config = _subscriptionManagerSubscriptions.GetOrAdd(newConfig, newConfig);
|
|
|
|
// if the reference is not the same, means it was already there and we did not add anything new
|
|
if (!ReferenceEquals(config, newConfig))
|
|
{
|
|
// for performance lets not create the message string if debugging is not enabled
|
|
// this can be executed many times and its in the algorithm thread
|
|
if (Log.DebuggingEnabled)
|
|
{
|
|
Log.Debug("DataManager.SubscriptionManagerGetOrAdd(): subscription already added: " + config);
|
|
}
|
|
}
|
|
else
|
|
{
|
|
// for performance, only count if we are above the limit
|
|
if (SubscriptionManagerCount() > _algorithmSettings.DataSubscriptionLimit)
|
|
{
|
|
// count data subscriptions by symbol, ignoring multiple data types.
|
|
// this limit was added due to the limits IB places on number of subscriptions
|
|
var uniqueCount = SubscriptionManagerSubscriptions
|
|
.Where(x => !x.Symbol.IsCanonical())
|
|
.DistinctBy(x => x.Symbol.Value)
|
|
.Count();
|
|
|
|
if (uniqueCount > _algorithmSettings.DataSubscriptionLimit)
|
|
{
|
|
throw new Exception(
|
|
$"The maximum number of concurrent market data subscriptions was exceeded ({_algorithmSettings.DataSubscriptionLimit})." +
|
|
"Please reduce the number of symbols requested or increase the limit using Settings.DataSubscriptionLimit.");
|
|
}
|
|
}
|
|
|
|
// add the time zone to our time keeper
|
|
_timeKeeper.AddTimeZone(newConfig.ExchangeTimeZone);
|
|
}
|
|
|
|
return config;
|
|
}
|
|
|
|
/// <summary>
|
|
/// Will try to remove a <see cref="SubscriptionDataConfig"/> and update the corresponding
|
|
/// consumers accordingly
|
|
/// </summary>
|
|
/// <param name="subscription">The <see cref="Subscription"/> owning the configuration to remove</param>
|
|
private void RemoveSubscriptionDataConfig(Subscription subscription)
|
|
{
|
|
// the subscription could of ended but might still be part of the universe
|
|
if (subscription.RemovedFromUniverse.Value)
|
|
{
|
|
SubscriptionDataConfig config;
|
|
_subscriptionManagerSubscriptions.TryRemove(subscription.Configuration, out config);
|
|
}
|
|
}
|
|
|
|
/// <summary>
|
|
/// Returns the amount of data config subscriptions processed for the SubscriptionManager
|
|
/// </summary>
|
|
public int SubscriptionManagerCount()
|
|
{
|
|
return _subscriptionManagerSubscriptions.Skip(0).Count();
|
|
}
|
|
|
|
#region ISubscriptionDataConfigService
|
|
|
|
/// <summary>
|
|
/// The different <see cref="TickType" /> each <see cref="SecurityType" /> supports
|
|
/// </summary>
|
|
public Dictionary<SecurityType, List<TickType>> AvailableDataTypes { get; }
|
|
|
|
/// <summary>
|
|
/// Creates and adds a list of <see cref="SubscriptionDataConfig" /> for a given symbol and configuration.
|
|
/// Can optionally pass in desired subscription data type to use.
|
|
/// If the config already existed will return existing instance instead
|
|
/// </summary>
|
|
public SubscriptionDataConfig Add(
|
|
Type dataType,
|
|
Symbol symbol,
|
|
Resolution? resolution = null,
|
|
bool fillForward = true,
|
|
bool extendedMarketHours = false,
|
|
bool isFilteredSubscription = true,
|
|
bool isInternalFeed = false,
|
|
bool isCustomData = false,
|
|
DataNormalizationMode dataNormalizationMode = DataNormalizationMode.Adjusted
|
|
)
|
|
{
|
|
return Add(symbol, resolution, fillForward, extendedMarketHours, isFilteredSubscription, isInternalFeed, isCustomData,
|
|
new List<Tuple<Type, TickType>> { new Tuple<Type, TickType>(dataType, LeanData.GetCommonTickTypeForCommonDataTypes(dataType, symbol.SecurityType))}, dataNormalizationMode)
|
|
.First();
|
|
}
|
|
|
|
/// <summary>
|
|
/// Creates and adds a list of <see cref="SubscriptionDataConfig" /> for a given symbol and configuration.
|
|
/// Can optionally pass in desired subscription data types to use.
|
|
/// If the config already existed will return existing instance instead
|
|
/// </summary>
|
|
public List<SubscriptionDataConfig> Add(
|
|
Symbol symbol,
|
|
Resolution? resolution = null,
|
|
bool fillForward = true,
|
|
bool extendedMarketHours = false,
|
|
bool isFilteredSubscription = true,
|
|
bool isInternalFeed = false,
|
|
bool isCustomData = false,
|
|
List<Tuple<Type, TickType>> subscriptionDataTypes = null,
|
|
DataNormalizationMode dataNormalizationMode = DataNormalizationMode.Adjusted
|
|
)
|
|
{
|
|
var dataTypes = subscriptionDataTypes ??
|
|
LookupSubscriptionConfigDataTypes(symbol.SecurityType, resolution ?? Resolution.Minute, symbol.IsCanonical());
|
|
|
|
if (!dataTypes.Any())
|
|
{
|
|
throw new ArgumentNullException(nameof(dataTypes), "At least one type needed to create new subscriptions");
|
|
}
|
|
|
|
var resolutionWasProvided = resolution.HasValue;
|
|
foreach (var typeTuple in dataTypes)
|
|
{
|
|
var baseInstance = typeTuple.Item1.GetBaseDataInstance();
|
|
baseInstance.Symbol = symbol;
|
|
if (!resolutionWasProvided)
|
|
{
|
|
var defaultResolution = baseInstance.DefaultResolution();
|
|
if (resolution.HasValue && resolution != defaultResolution)
|
|
{
|
|
// we are here because there are multiple 'dataTypes'.
|
|
// if we get different default resolutions lets throw, this shouldn't happen
|
|
throw new InvalidOperationException(
|
|
$"Different data types ({string.Join(",", dataTypes.Select(tuple => tuple.Item1))})" +
|
|
$" provided different default resolutions {defaultResolution} and {resolution}, this is an unexpected invalid operation.");
|
|
}
|
|
resolution = defaultResolution;
|
|
}
|
|
else
|
|
{
|
|
// only assert resolution in backtesting, live can use other data source
|
|
// for example daily data for options
|
|
if (!_liveMode)
|
|
{
|
|
var supportedResolutions = baseInstance.SupportedResolutions();
|
|
if (supportedResolutions.Contains(resolution.Value))
|
|
{
|
|
continue;
|
|
}
|
|
|
|
throw new ArgumentException($"Sorry {resolution.ToStringInvariant()} is not a supported resolution for {typeTuple.Item1.Name}" +
|
|
$" and SecurityType.{symbol.SecurityType.ToStringInvariant()}." +
|
|
$" Please change your AddData to use one of the supported resolutions ({string.Join(",", supportedResolutions)}).");
|
|
}
|
|
}
|
|
}
|
|
|
|
MarketHoursDatabase.Entry marketHoursDbEntry;
|
|
if (!_marketHoursDatabase.TryGetEntry(symbol.ID.Market, symbol, symbol.ID.SecurityType, out marketHoursDbEntry))
|
|
{
|
|
if (symbol.SecurityType == SecurityType.Base)
|
|
{
|
|
var baseInstance = dataTypes.Single().Item1.GetBaseDataInstance();
|
|
baseInstance.Symbol = symbol;
|
|
_marketHoursDatabase.SetEntryAlwaysOpen(symbol.ID.Market, null, SecurityType.Base, baseInstance.DataTimeZone());
|
|
}
|
|
|
|
marketHoursDbEntry = _marketHoursDatabase.GetEntry(symbol.ID.Market, symbol, symbol.ID.SecurityType);
|
|
}
|
|
|
|
var exchangeHours = marketHoursDbEntry.ExchangeHours;
|
|
if (symbol.ID.SecurityType == SecurityType.Option ||
|
|
symbol.ID.SecurityType == SecurityType.FutureOption ||
|
|
symbol.ID.SecurityType == SecurityType.Future)
|
|
{
|
|
dataNormalizationMode = DataNormalizationMode.Raw;
|
|
}
|
|
|
|
if (marketHoursDbEntry.DataTimeZone == null)
|
|
{
|
|
throw new ArgumentNullException(nameof(marketHoursDbEntry.DataTimeZone),
|
|
"DataTimeZone is a required parameter for new subscriptions. Set to the time zone the raw data is time stamped in.");
|
|
}
|
|
|
|
if (exchangeHours.TimeZone == null)
|
|
{
|
|
throw new ArgumentNullException(nameof(exchangeHours.TimeZone),
|
|
"ExchangeTimeZone is a required parameter for new subscriptions. Set to the time zone the security exchange resides in.");
|
|
}
|
|
|
|
var result = (from subscriptionDataType in dataTypes
|
|
let dataType = subscriptionDataType.Item1
|
|
let tickType = subscriptionDataType.Item2
|
|
select new SubscriptionDataConfig(
|
|
dataType,
|
|
symbol,
|
|
resolution.Value,
|
|
marketHoursDbEntry.DataTimeZone,
|
|
exchangeHours.TimeZone,
|
|
fillForward,
|
|
extendedMarketHours,
|
|
isInternalFeed,
|
|
isCustomData,
|
|
isFilteredSubscription: isFilteredSubscription,
|
|
tickType: tickType,
|
|
dataNormalizationMode: dataNormalizationMode)).ToList();
|
|
|
|
for (int i = 0; i < result.Count; i++)
|
|
{
|
|
result[i] = SubscriptionManagerGetOrAdd(result[i]);
|
|
|
|
// track all registered data types
|
|
_registeredTypesProvider.RegisterType(result[i].Type);
|
|
}
|
|
return result;
|
|
}
|
|
|
|
/// <summary>
|
|
/// Get the data feed types for a given <see cref="SecurityType" /> <see cref="Resolution" />
|
|
/// </summary>
|
|
/// <param name="symbolSecurityType">The <see cref="SecurityType" /> used to determine the types</param>
|
|
/// <param name="resolution">The resolution of the data requested</param>
|
|
/// <param name="isCanonical">Indicates whether the security is Canonical (future and options)</param>
|
|
/// <returns>Types that should be added to the <see cref="SubscriptionDataConfig" /></returns>
|
|
public List<Tuple<Type, TickType>> LookupSubscriptionConfigDataTypes(
|
|
SecurityType symbolSecurityType,
|
|
Resolution resolution,
|
|
bool isCanonical
|
|
)
|
|
{
|
|
if (isCanonical)
|
|
{
|
|
return new List<Tuple<Type, TickType>> { new Tuple<Type, TickType>(typeof(ZipEntryName), TickType.Quote) };
|
|
}
|
|
|
|
IEnumerable<TickType> availableDataType = AvailableDataTypes[symbolSecurityType];
|
|
// Equities will only look for trades in case of low resolutions.
|
|
if (symbolSecurityType == SecurityType.Equity && (resolution == Resolution.Daily || resolution == Resolution.Hour))
|
|
{
|
|
// we filter out quote tick type
|
|
availableDataType = availableDataType.Where(t => t != TickType.Quote);
|
|
}
|
|
|
|
return availableDataType
|
|
.Select(tickType => new Tuple<Type, TickType>(LeanData.GetDataType(resolution, tickType), tickType)).ToList();
|
|
}
|
|
|
|
/// <summary>
|
|
/// Gets a list of all registered <see cref="SubscriptionDataConfig"/> for a given <see cref="Symbol"/>
|
|
/// </summary>
|
|
/// <remarks>Will not return internal subscriptions by default</remarks>
|
|
public List<SubscriptionDataConfig> GetSubscriptionDataConfigs(Symbol symbol, bool includeInternalConfigs = false)
|
|
{
|
|
return SubscriptionManagerSubscriptions.Where(x => x.Symbol == symbol
|
|
&& (includeInternalConfigs || !x.IsInternalFeed)).ToList();
|
|
}
|
|
|
|
#endregion
|
|
|
|
#endregion
|
|
|
|
#region IDataManager
|
|
|
|
/// <summary>
|
|
/// Get the universe selection instance
|
|
/// </summary>
|
|
public UniverseSelection UniverseSelection { get; }
|
|
|
|
#endregion
|
|
}
|
|
}
|