From 6a3888c6523f44661257a4c03f6a250cd3449c7e Mon Sep 17 00:00:00 2001 From: Jhonathan Abreu Date: Tue, 11 Aug 2026 17:17:55 -0400 Subject: [PATCH] Improve consolidator registration ergonomics - Adapt trade bar consolidators and indicators to quote-only feeds (forex, cfd) by collapsing quote bars into mid-point trade bars with zero volume (new QuoteBarToTradeBarAdapter), instead of rejecting the registration with a type mismatch error - Validate fixed time span consolidator periods against the subscription period at registration time, mirroring the check PeriodCountConsolidatorBase performs when the first data point arrives, so invalid setups fail in Initialize instead of mid-backtest - Accept calendar periods (Calendar.WEEKLY/MONTHLY, custom functions) in RegisterIndicator and ResolveConsolidator, in C# overloads and in the Python register_indicator path, like Consolidate already does - Throw a clear ArgumentNullException for null symbols (e.g. Future.Mapped before the first mapping) instead of an NRE, and suggest Consolidate()/self.consolidate() in the type mismatch error --- ...idatorAutoAdaptationRegressionAlgorithm.cs | 158 ++++++++++++++++++ ...idatorAutoAdaptationRegressionAlgorithm.py | 65 +++++++ .../QuantConnect.Algorithm.Python.csproj | 1 + Algorithm/QCAlgorithm.Indicators.cs | 60 +++++++ Algorithm/QCAlgorithm.Python.cs | 14 +- Common/Data/Consolidators/ConsolidatorBase.cs | 9 + .../PeriodCountConsolidatorBase.cs | 6 + .../QuoteBarToTradeBarAdapter.cs | 112 +++++++++++++ Common/Data/SubscriptionManager.cs | 67 +++++++- Tests/Algorithm/AlgorithmIndicatorsTests.cs | 59 +++++++ Tests/Common/Data/SubscriptionManagerTests.cs | 144 ++++++++++++++++ Tests/Indicators/PythonIndicatorTests.cs | 10 +- 12 files changed, 694 insertions(+), 11 deletions(-) create mode 100644 Algorithm.CSharp/ConsolidatorAutoAdaptationRegressionAlgorithm.cs create mode 100644 Algorithm.Python/ConsolidatorAutoAdaptationRegressionAlgorithm.py create mode 100644 Common/Data/Consolidators/QuoteBarToTradeBarAdapter.cs diff --git a/Algorithm.CSharp/ConsolidatorAutoAdaptationRegressionAlgorithm.cs b/Algorithm.CSharp/ConsolidatorAutoAdaptationRegressionAlgorithm.cs new file mode 100644 index 000000000000..c593e63d6c44 --- /dev/null +++ b/Algorithm.CSharp/ConsolidatorAutoAdaptationRegressionAlgorithm.cs @@ -0,0 +1,158 @@ +/* + * 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.Generic; +using QuantConnect.Data; +using QuantConnect.Data.Consolidators; +using QuantConnect.Indicators; +using QuantConnect.Interfaces; + +namespace QuantConnect.Algorithm.CSharp +{ + /// + /// Regression algorithm asserting consolidator registration ergonomics on a quote-only feed (forex): + /// trade bar consolidators and indicators are fed collapsed quote bars instead of being rejected, + /// + /// accepts calendar periods, and consolidator periods smaller than the subscription period are + /// rejected at registration time instead of when the first data point arrives. + /// + public class ConsolidatorAutoAdaptationRegressionAlgorithm : QCAlgorithm, IRegressionAlgorithmDefinition + { + private TradeBarConsolidator _consolidator; + private OnBalanceVolume _obv; + private RelativeStrengthIndex _weeklyRsi; + private int _adaptedTradeBars; + + /// + /// Initialise the data and resolution required, as well as the cash and start-end dates for your algorithm. All algorithms must initialized. + /// + public override void Initialize() + { + SetStartDate(2014, 5, 5); + SetEndDate(2014, 5, 12); + + var eurusd = AddForex("EURUSD", Resolution.Minute, Market.Oanda).Symbol; + + // a trade bar consolidator on a quote-only feed: the engine feeds it quote bars collapsed + // into mid-point trade bars with zero volume + _consolidator = new TradeBarConsolidator(TimeSpan.FromHours(1)); + _consolidator.DataConsolidated += (_, bar) => + { + _adaptedTradeBars++; + if (bar.Volume != 0) + { + throw new RegressionTestException("Expected trade bars collapsed from quote bars to have zero volume"); + } + }; + SubscriptionManager.AddConsolidator(eurusd, _consolidator); + + // a trade bar indicator on a quote-only feed is fed collapsed trade bars as well + _obv = OBV(eurusd, Resolution.Hour); + + // RegisterIndicator accepts calendar periods, like Consolidate does + _weeklyRsi = new RelativeStrengthIndex(2); + RegisterIndicator(eurusd, _weeklyRsi, Calendar.Weekly); + + // a consolidator period smaller than the subscription period is rejected at registration time, + // instead of when the first data point arrives + try + { + SubscriptionManager.AddConsolidator(eurusd, new TradeBarConsolidator(TimeSpan.FromSeconds(10))); + throw new RegressionTestException($"Expected {nameof(ArgumentException)} for a consolidator period smaller than the subscription period"); + } + catch (ArgumentException) + { + // expected, all the required information is available at registration time + } + } + + public override void OnEndOfAlgorithm() + { + if (_adaptedTradeBars == 0) + { + throw new RegressionTestException("Expected the adapted trade bar consolidator to receive data"); + } + if (_obv.Samples == 0) + { + throw new RegressionTestException("Expected the OnBalanceVolume indicator to be updated with collapsed trade bars"); + } + if (_weeklyRsi.Samples == 0) + { + throw new RegressionTestException("Expected the weekly RSI to be updated when the week boundary was crossed"); + } + } + + /// + /// This is used by the regression test system to indicate if the open source Lean repository has the required data to run this algorithm. + /// + public bool CanRunLocally { get; } = true; + + /// + /// This is used by the regression test system to indicate which languages this algorithm is written in. + /// + public List Languages { get; } = new() { Language.CSharp, Language.Python }; + + /// + /// Data Points count of all timeslices of algorithm + /// + public long DataPoints => 8661; + + /// + /// Data Points count of the algorithm history + /// + public int AlgorithmHistoryDataPoints => 3; + + /// + /// Final status of the algorithm + /// + public AlgorithmStatus AlgorithmStatus => AlgorithmStatus.Completed; + + /// + /// This is used by the regression test system to indicate what the expected statistics are from running the algorithm + /// + public Dictionary ExpectedStatistics => new Dictionary + { + {"Total Orders", "0"}, + {"Average Win", "0%"}, + {"Average Loss", "0%"}, + {"Compounding Annual Return", "0%"}, + {"Drawdown", "0%"}, + {"Expectancy", "0"}, + {"Start Equity", "100000.00"}, + {"End Equity", "100000"}, + {"Net Profit", "0%"}, + {"Sharpe Ratio", "0"}, + {"Sortino Ratio", "0"}, + {"Probabilistic Sharpe Ratio", "0%"}, + {"Loss Rate", "0%"}, + {"Win Rate", "0%"}, + {"Profit-Loss Ratio", "0"}, + {"Alpha", "0"}, + {"Beta", "0"}, + {"Annual Standard Deviation", "0"}, + {"Annual Variance", "0"}, + {"Information Ratio", "-3.328"}, + {"Tracking Error", "0.091"}, + {"Treynor Ratio", "0"}, + {"Total Fees", "$0.00"}, + {"Estimated Strategy Capacity", "$0"}, + {"Lowest Capacity Asset", ""}, + {"Portfolio Turnover", "0%"}, + {"Drawdown Recovery", "0"}, + {"OrderListHash", "d41d8cd98f00b204e9800998ecf8427e"} + }; + } +} diff --git a/Algorithm.Python/ConsolidatorAutoAdaptationRegressionAlgorithm.py b/Algorithm.Python/ConsolidatorAutoAdaptationRegressionAlgorithm.py new file mode 100644 index 000000000000..1f49da826b84 --- /dev/null +++ b/Algorithm.Python/ConsolidatorAutoAdaptationRegressionAlgorithm.py @@ -0,0 +1,65 @@ +# 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. + +from AlgorithmImports import * + +### +### Regression algorithm asserting consolidator registration ergonomics on a quote-only feed (forex): +### trade bar consolidators and indicators are fed collapsed quote bars instead of being rejected, +### register_indicator accepts calendar periods, and consolidator periods smaller than the subscription +### period are rejected at registration time instead of when the first data point arrives. +### +class ConsolidatorAutoAdaptationRegressionAlgorithm(QCAlgorithm): + def initialize(self): + self.set_start_date(2014, 5, 5) + self.set_end_date(2014, 5, 12) + + eurusd = self.add_forex("EURUSD", Resolution.MINUTE, Market.OANDA).symbol + self._adapted_trade_bars = 0 + + # a trade bar consolidator on a quote-only feed: the engine feeds it quote bars collapsed + # into mid-point trade bars with zero volume + self._consolidator = TradeBarConsolidator(timedelta(hours=1)) + self._consolidator.data_consolidated += self._on_trade_bar + self.subscription_manager.add_consolidator(eurusd, self._consolidator) + + # a trade bar indicator on a quote-only feed is fed collapsed trade bars as well + self._obv = self.obv(eurusd, Resolution.HOUR) + + # register_indicator accepts calendar periods, like self.consolidate does + self._weekly_rsi = RelativeStrengthIndex(2) + self.register_indicator(eurusd, self._weekly_rsi, Calendar.WEEKLY) + + # a consolidator period smaller than the subscription period is rejected at registration time, + # instead of when the first data point arrives + rejected = False + try: + self.subscription_manager.add_consolidator(eurusd, TradeBarConsolidator(timedelta(seconds=10))) + except: + # expected, all the required information is available at registration time + rejected = True + if not rejected: + raise AssertionError("Expected an error for a consolidator period smaller than the subscription period") + + def _on_trade_bar(self, sender, bar): + self._adapted_trade_bars += 1 + if bar.volume != 0: + raise AssertionError("Expected trade bars collapsed from quote bars to have zero volume") + + def on_end_of_algorithm(self): + if self._adapted_trade_bars == 0: + raise AssertionError("Expected the adapted trade bar consolidator to receive data") + if self._obv.samples == 0: + raise AssertionError("Expected the OnBalanceVolume indicator to be updated with collapsed trade bars") + if self._weekly_rsi.samples == 0: + raise AssertionError("Expected the weekly RSI to be updated when the week boundary was crossed") diff --git a/Algorithm.Python/QuantConnect.Algorithm.Python.csproj b/Algorithm.Python/QuantConnect.Algorithm.Python.csproj index b5939b6a817d..3d45997b509b 100644 --- a/Algorithm.Python/QuantConnect.Algorithm.Python.csproj +++ b/Algorithm.Python/QuantConnect.Algorithm.Python.csproj @@ -53,6 +53,7 @@ + diff --git a/Algorithm/QCAlgorithm.Indicators.cs b/Algorithm/QCAlgorithm.Indicators.cs index e86ffd96b864..73febcb46bd2 100644 --- a/Algorithm/QCAlgorithm.Indicators.cs +++ b/Algorithm/QCAlgorithm.Indicators.cs @@ -3160,6 +3160,14 @@ public string CreateIndicatorName(Symbol symbol, string type, Resolution? resolu /// The SubscriptionDataConfig for the specified symbol private SubscriptionDataConfig GetSubscription(Symbol symbol, TickType? tickType = null) { + if (symbol == null) + { + // fail with a clear message instead of a null reference below, commonly a Future.Mapped + // or similar property that is still null when the consolidator/indicator is registered + throw new ArgumentNullException(nameof(symbol), "Cannot resolve a subscription because the given symbol is null. " + + "If the symbol comes from a property like Future.Mapped, note it can be null until the security receives data."); + } + if (!TryGetSubscription(symbol, tickType, out var subscription)) { // The symbol was not manually subscribed to. Mirror the behavior of order submission @@ -3231,6 +3239,21 @@ public void RegisterIndicator(Symbol symbol, IndicatorBase i RegisterIndicator(symbol, indicator, ResolveConsolidator(symbol, resolution), selector ?? (x => x.Value)); } + /// + /// Creates and registers a new consolidator to receive automatic updates on the given calendar as well as configures + /// the indicator to receive updates from the consolidator. + /// + /// The symbol to register against + /// The indicator to receive data from the consolidator + /// The consolidation calendar, for example + /// Selects a value from the BaseData to send into the indicator, if null defaults to the Value property of BaseData (x => x.Value) + [DocumentationAttribute(ConsolidatingData)] + [DocumentationAttribute(Indicators)] + public void RegisterIndicator(Symbol symbol, IndicatorBase indicator, Func calendar, Func selector = null) + { + RegisterIndicator(symbol, indicator, ResolveConsolidator(symbol, calendar), selector ?? (x => x.Value)); + } + /// /// Registers the consolidator to receive automatic updates as well as configures the indicator to receive updates /// from the consolidator. @@ -3303,6 +3326,22 @@ public void RegisterIndicator(Symbol symbol, IndicatorBase indicator, Time RegisterIndicator(symbol, indicator, ResolveConsolidator(symbol, resolution, typeof(T)), selector); } + /// + /// Creates and registers a new consolidator to receive automatic updates on the given calendar as well as configures + /// the indicator to receive updates from the consolidator. + /// + /// The symbol to register against + /// The indicator to receive data from the consolidator + /// The consolidation calendar, for example + /// Selects a value from the BaseData send into the indicator, if null defaults to a cast (x => (T)x) + [DocumentationAttribute(ConsolidatingData)] + [DocumentationAttribute(Indicators)] + public void RegisterIndicator(Symbol symbol, IndicatorBase indicator, Func calendar, Func selector = null) + where T : IBaseData + { + RegisterIndicator(symbol, indicator, ResolveConsolidator(symbol, calendar, typeof(T)), selector); + } + /// /// Registers the consolidator to receive automatic updates as well as configures the indicator to receive updates /// from the consolidator. @@ -3331,6 +3370,12 @@ public void RegisterIndicator(Symbol symbol, IndicatorBase indicator, IDat // so we use a smarter selector as in other API methods selectorToUse = consolidated => (T)(object)new IndicatorDataPoint(consolidated.Symbol, consolidated.EndTime, consolidated.Value); } + else if (type == typeof(TradeBar) && consolidator.OutputType == typeof(QuoteBar) && selector == null) + { + // trade bar indicator registered against a quote-only feed (forex, cfd): collapse each quote bar + // into a mid-point trade bar with zero volume instead of failing + selectorToUse = consolidated => (T)(object)((QuoteBar)consolidated).Collapse(); + } else { throw new ArgumentException($"Type mismatch found between consolidator and indicator for symbol: {symbol}." + @@ -3694,6 +3739,21 @@ public IDataConsolidator ResolveConsolidator(Symbol symbol, TimeSpan? timeSpan, return CreateConsolidator(symbol, null, tickType, timeSpan, null, null); } + /// + /// Gets the default consolidator for the specified symbol and consolidation calendar + /// + /// The symbol whose data is to be consolidated + /// The consolidation calendar, for example + /// The data type for this consolidator, if null, uses TradeBar over QuoteBar if present + /// The new default consolidator + [DocumentationAttribute(ConsolidatingData)] + [DocumentationAttribute(Indicators)] + public IDataConsolidator ResolveConsolidator(Symbol symbol, Func calendar, Type dataType = null) + { + var tickType = dataType != null ? LeanData.GetCommonTickTypeForCommonDataTypes(dataType, symbol.SecurityType) : (TickType?)null; + return CreateConsolidator(symbol, calendar, tickType, null, null, null); + } + /// /// Creates a new consolidator for the specified period, generating the requested output type. /// diff --git a/Algorithm/QCAlgorithm.Python.cs b/Algorithm/QCAlgorithm.Python.cs index 9859cbc4094b..e2aeb8a004e5 100644 --- a/Algorithm/QCAlgorithm.Python.cs +++ b/Algorithm/QCAlgorithm.Python.cs @@ -621,7 +621,7 @@ public void RegisterIndicator(Symbol symbol, PyObject indicator, TimeSpan? resol /// /// The symbol to register against /// The indicator to receive data from the consolidator - /// The python object that it is trying to register with, could be consolidator or a timespan + /// The python object that it is trying to register with, could be consolidator, timespan or a calendar /// Selects a value from the BaseData send into the indicator, if null defaults to a cast (x => (T)x) [DocumentationAttribute(Indicators)] [DocumentationAttribute(ConsolidatingData)] @@ -641,7 +641,7 @@ public void RegisterIndicator(Symbol symbol, PyObject indicator, PyObject pyObje } catch { - // Finally, since above didn't work, just try it as a timespan + // Since above didn't work, just try it as a timespan // Issue #4668 Fix using (Py.GIL()) { @@ -657,7 +657,15 @@ public void RegisterIndicator(Symbol symbol, PyObject indicator, PyObject pyObje } catch (Exception e) { - throw new ArgumentException("Invalid third argument, should be either a valid consolidator or timedelta object. The following exception was thrown: ", e); + // Finally, try it as a consolidation calendar, e.g. Calendar.WEEKLY, which Consolidate() also accepts + if (pyObject.TrySafeAs(out Func calendar)) + { + RegisterIndicator(symbol, indicator, ResolveConsolidator(symbol, calendar), selector); + return; + } + + throw new ArgumentException("Invalid third argument, should be either a valid consolidator, timedelta or " + + "calendar (e.g. Calendar.WEEKLY) object. The following exception was thrown: ", e); } } } diff --git a/Common/Data/Consolidators/ConsolidatorBase.cs b/Common/Data/Consolidators/ConsolidatorBase.cs index 0b853310efcc..aa651eab0e80 100644 --- a/Common/Data/Consolidators/ConsolidatorBase.cs +++ b/Common/Data/Consolidators/ConsolidatorBase.cs @@ -47,6 +47,15 @@ protected set /// public abstract IBaseData WorkingData { get; } + /// + /// The fixed time span between consolidated bars when the consolidator is driven by an explicit + /// time span period; null for count-based, mixed-mode, calendar and non-period consolidators. + /// Lets validate the period against the subscription resolution + /// at registration time instead of failing when the first data point arrives + /// (see ) + /// + internal virtual TimeSpan? FixedTimeSpanPeriod => null; + /// /// Gets the type consumed by this consolidator /// diff --git a/Common/Data/Consolidators/PeriodCountConsolidatorBase.cs b/Common/Data/Consolidators/PeriodCountConsolidatorBase.cs index 1f6f69e29c03..e8f572ab5e0c 100644 --- a/Common/Data/Consolidators/PeriodCountConsolidatorBase.cs +++ b/Common/Data/Consolidators/PeriodCountConsolidatorBase.cs @@ -293,6 +293,12 @@ public override void Reset() /// protected TimeSpan? Period => _period; + /// + /// The fixed time span period, only set when the user gave an explicit time span, matching the + /// data span validation performed in when the first data point arrives + /// + internal override TimeSpan? FixedTimeSpanPeriod => _periodSpecification is TimeSpanPeriodSpecification ? _period : null; + /// /// Determines whether or not the specified data should be processed /// diff --git a/Common/Data/Consolidators/QuoteBarToTradeBarAdapter.cs b/Common/Data/Consolidators/QuoteBarToTradeBarAdapter.cs new file mode 100644 index 000000000000..153bd659fbde --- /dev/null +++ b/Common/Data/Consolidators/QuoteBarToTradeBarAdapter.cs @@ -0,0 +1,112 @@ +/* + * 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 QuantConnect.Data.Market; + +namespace QuantConnect.Data.Consolidators +{ + /// + /// Adapts a consolidator with input so it can be fed from a quote-only + /// subscription (forex, cfd): each incoming is collapsed into a mid-point + /// with zero volume before being forwarded to the wrapped consolidator. + /// Used by + /// so trade bar consolidators and indicators work out of the box on quote-only security types. + /// + public class QuoteBarToTradeBarAdapter : IDataConsolidator + { + /// + /// The wrapped consolidator receiving the collapsed trade bars + /// + public IDataConsolidator Consolidator { get; } + + /// + /// Gets the most recently consolidated piece of data produced by the wrapped consolidator + /// + public IBaseData Consolidated => Consolidator.Consolidated; + + /// + /// Gets a clone of the data being currently consolidated by the wrapped consolidator + /// + public IBaseData WorkingData => Consolidator.WorkingData; + + /// + /// Gets the type consumed by this consolidator + /// + public Type InputType => typeof(QuoteBar); + + /// + /// Gets the type produced by the wrapped consolidator + /// + public Type OutputType => Consolidator.OutputType; + + /// + /// Event handler that fires when the wrapped consolidator produces a new piece of data + /// + public event DataConsolidatedHandler DataConsolidated + { + add { Consolidator.DataConsolidated += value; } + remove { Consolidator.DataConsolidated -= value; } + } + + /// + /// Creates a new adapter feeding collapsed quote bars into the given consolidator + /// + /// The consolidator to adapt, must consume input + public QuoteBarToTradeBarAdapter(IDataConsolidator consolidator) + { + if (!consolidator.InputType.IsAssignableFrom(typeof(TradeBar))) + { + throw new ArgumentException($"{nameof(QuoteBarToTradeBarAdapter)} requires a consolidator that accepts {nameof(TradeBar)} input " + + $"but was given one with input type {consolidator.InputType.Name}"); + } + Consolidator = consolidator; + } + + /// + /// Updates the wrapped consolidator with the collapsed version of the specified quote bar + /// + /// The new data for the consolidator + public void Update(IBaseData data) + { + Consolidator.Update(((QuoteBar)data).Collapse()); + } + + /// + /// Scans the wrapped consolidator to see if it should emit a bar due to time passing + /// + /// The current time in the local time zone (same as ) + public void Scan(DateTime currentLocalTime) + { + Consolidator.Scan(currentLocalTime); + } + + /// + /// Resets the wrapped consolidator + /// + public void Reset() + { + Consolidator.Reset(); + } + + /// + /// Disposes the wrapped consolidator + /// + public void Dispose() + { + Consolidator.Dispose(); + } + } +} diff --git a/Common/Data/SubscriptionManager.cs b/Common/Data/SubscriptionManager.cs index 33238d08b49c..9e56418abe96 100644 --- a/Common/Data/SubscriptionManager.cs +++ b/Common/Data/SubscriptionManager.cs @@ -164,6 +164,14 @@ public SubscriptionDataConfig Add( /// Desired tick type for the subscription public void AddConsolidator(Symbol symbol, IDataConsolidator consolidator, TickType? tickType = null) { + if (symbol == null) + { + // giving the raw NRE from the null symbol comparison below provides no clue of the actual issue, + // commonly a Future.Mapped or similar property that is still null when the consolidator is added + throw new ArgumentNullException(nameof(symbol), "Cannot add a consolidator because the given symbol is null. " + + "If the symbol comes from a property like Future.Mapped, note it can be null until the security receives data."); + } + // Find the right subscription and add the consolidator to it var subscriptions = Subscriptions.Where(x => x.Symbol == symbol).ToList(); @@ -184,6 +192,10 @@ public void AddConsolidator(Symbol symbol, IDataConsolidator consolidator, TickT // we need to be able to pipe data directly from the data feed into the consolidator if (IsSubscriptionValidForConsolidator(subscription, consolidator, tickType)) { + // fail here instead of months into the backtest when the first data point arrives, + // see PeriodCountConsolidatorBase.Update: all the information needed is already available + ValidateConsolidatorPeriod(symbol, consolidator, subscription); + subscription.Consolidators.Add(consolidator); var wrapper = _consolidators[consolidator] = @@ -198,6 +210,17 @@ public void AddConsolidator(Symbol symbol, IDataConsolidator consolidator, TickT } } + // a trade bar consolidator on a quote-only feed (forex, cfd) can still be fed by collapsing + // each quote bar into a mid-point trade bar instead of rejecting the registration + if (consolidator.InputType == typeof(TradeBar) && (tickType == null || tickType == TickType.Quote) && + subscriptions.Any(x => x.Type == typeof(QuoteBar))) + { + Logging.Log.Trace($"SubscriptionManager.AddConsolidator(): {symbol.Value} does not have a {nameof(TradeBar)} subscription: " + + $"feeding the {consolidator.GetType().Name} with quote bars collapsed into mid-point trade bars with zero volume."); + AddConsolidator(symbol, new QuoteBarToTradeBarAdapter(consolidator), TickType.Quote); + return; + } + string tickTypeException = null; if (tickType != null && !subscriptions.Where(x => x.TickType == tickType).Any()) { @@ -206,7 +229,25 @@ public void AddConsolidator(Symbol symbol, IDataConsolidator consolidator, TickT throw new ArgumentException(tickTypeException ?? ("Type mismatch found between consolidator and symbol. " + $"Symbol: {symbol.Value} does not support input type: {consolidator.InputType.Name}. " + - $"Supported types: {string.Join(",", subscriptions.Select(x => x.Type.Name))}.")); + $"Supported types: {string.Join(",", subscriptions.Select(x => x.Type.Name))}. " + + "Alternatively, consider using QCAlgorithm.Consolidate() ('self.consolidate()' in Python), which selects the correct consolidator type automatically.")); + } + + /// + /// Validates that a fixed time span period consolidator has a period of at least the subscription period, + /// mirroring the data span validation performs + /// when the first data point arrives, so invalid setups fail at registration time in Initialize instead + /// + private static void ValidateConsolidatorPeriod(Symbol symbol, IDataConsolidator consolidator, SubscriptionDataConfig subscription) + { + var period = ((consolidator as QuoteBarToTradeBarAdapter)?.Consolidator as ConsolidatorBase ?? consolidator as ConsolidatorBase)?.FixedTimeSpanPeriod; + if (period.HasValue && period.Value < subscription.Increment) + { + throw new ArgumentException($"Unable to add consolidator for symbol {symbol.Value}: the consolidator period ({period.Value}) " + + $"is smaller than the {subscription.Resolution.ResolutionToLower()} subscription period ({subscription.Increment}). " + + "Consolidators can only produce data of the same or lower resolution than their input data, " + + $"either request a higher resolution subscription for {symbol.Value} or use a consolidator period of at least {subscription.Increment}."); + } } /// @@ -237,7 +278,23 @@ public void RemoveConsolidator(Symbol symbol, IDataConsolidator consolidator) // remove consolidator from each subscription foreach (var subscription in _subscriptionManager.GetSubscriptionDataConfigs(symbol)) { - subscription.Consolidators.Remove(consolidator); + if (!subscription.Consolidators.Remove(consolidator)) + { + // the consolidator might have been registered wrapped in a quote to trade bar adapter, + // see AddConsolidator: users hold and remove the consolidator they created, not the adapter. + // Equals instead of reference equality so a fresh python wrapper for the same python object matches + var adapter = subscription.Consolidators + .OfType() + .FirstOrDefault(x => x.Consolidator.Equals(consolidator)); + if (adapter != null) + { + subscription.Consolidators.Remove(adapter); + if (_consolidators.Remove(adapter, out var adapterToScan)) + { + adapterToScan.Dispose(); + } + } + } if (_consolidators.Remove(consolidator, out var consolidatorsToScan)) { @@ -283,9 +340,11 @@ private IDataConsolidator FindPythonConsolidator(Symbol symbol, PyObject pyConso { foreach (var existing in subscription.Consolidators) { - if (existing is DataConsolidatorPythonWrapper && existing.Equals(pyConsolidator)) + // the python consolidator might have been registered wrapped in a quote to trade bar adapter + var candidate = (existing as QuoteBarToTradeBarAdapter)?.Consolidator ?? existing; + if (candidate is DataConsolidatorPythonWrapper && candidate.Equals(pyConsolidator)) { - return existing; + return candidate; } } } diff --git a/Tests/Algorithm/AlgorithmIndicatorsTests.cs b/Tests/Algorithm/AlgorithmIndicatorsTests.cs index dd0764671c4d..73108ef9ed23 100644 --- a/Tests/Algorithm/AlgorithmIndicatorsTests.cs +++ b/Tests/Algorithm/AlgorithmIndicatorsTests.cs @@ -21,6 +21,7 @@ using Python.Runtime; using QuantConnect.Algorithm; using QuantConnect.Data; +using QuantConnect.Data.Consolidators; using QuantConnect.Data.Market; using QuantConnect.Indicators; using QuantConnect.Lean.Engine.DataFeeds; @@ -845,5 +846,63 @@ def get_indicator(algo, symbol): var lastInput = testModule.GetAttr("LastInputTracker").GetAttr("last_input").GetAndDispose(); Assert.IsNotNull(lastInput); } + + [Test] + public void RegistersIndicatorWithCalendar() + { + var rsi = new RelativeStrengthIndex(2); + Assert.DoesNotThrow(() => _algorithm.RegisterIndicator(_equity, rsi, Calendar.Weekly)); + + // feed the registered consolidator trade bars crossing a week boundary so the weekly bar is emitted + var consolidator = _algorithm.SubscriptionManager.Subscriptions + .Single(x => x.Symbol == _equity && x.Consolidators.Count == 1).Consolidators.Single(); + var time = new DateTime(2013, 10, 7); + consolidator.Update(new TradeBar(time, _equity, 10, 20, 5, 15, 100, Time.OneDay)); + consolidator.Update(new TradeBar(time.AddDays(7), _equity, 10, 20, 5, 15, 100, Time.OneDay)); + + Assert.AreEqual(1, rsi.Samples); + } + + [Test] + public void RegistersIndicatorWithCalendarFromPython() + { + using var _ = Py.GIL(); + var testModule = PyModule.FromString("testModule", + @" +from AlgorithmImports import * + +def register(algo, symbol): + rsi = RelativeStrengthIndex(14) + algo.register_indicator(symbol, rsi, Calendar.WEEKLY) + return rsi +"); + + using var pyAlgo = _algorithm.ToPython(); + using var pySymbol = _equity.ToPython(); + Assert.DoesNotThrow(() => testModule.GetAttr("register").Invoke(pyAlgo, pySymbol).Dispose()); + + Assert.AreEqual(1, _algorithm.SubscriptionManager.Subscriptions.Sum(x => x.Consolidators.Count)); + } + + [Test] + public void TradeBarIndicatorOnQuoteOnlyFeedIsFedCollapsedTradeBars() + { + var eurusd = _algorithm.AddForex("EURUSD", market: Market.Oanda).Symbol; + + // OnBalanceVolume consumes trade bars but forex only has quote bars, used to throw + // 'Consolidator outputs type QuoteBar but indicator expects input type TradeBar' + var obv = new OnBalanceVolume(); + Assert.DoesNotThrow(() => _algorithm.RegisterIndicator(eurusd, obv, Resolution.Hour)); + + // feed the registered consolidator quote bars, the indicator receives collapsed trade bars + var consolidator = _algorithm.SubscriptionManager.Subscriptions + .Single(x => x.Symbol == eurusd && x.Consolidators.Count == 1).Consolidators.Single(); + var time = new DateTime(2013, 10, 7, 10, 0, 0); + var bar = new Bar(1m, 2m, 0.5m, 1.5m); + consolidator.Update(new QuoteBar(time, eurusd, bar, 0, bar, 0, Time.OneMinute)); + consolidator.Update(new QuoteBar(time.AddHours(1), eurusd, bar, 0, bar, 0, Time.OneMinute)); + + Assert.AreEqual(1, obv.Samples); + } } } diff --git a/Tests/Common/Data/SubscriptionManagerTests.cs b/Tests/Common/Data/SubscriptionManagerTests.cs index 33ec9a28b889..43c81f442ff4 100644 --- a/Tests/Common/Data/SubscriptionManagerTests.cs +++ b/Tests/Common/Data/SubscriptionManagerTests.cs @@ -709,6 +709,150 @@ def get_consolidator(): Assert.AreEqual(0, algorithm.SubscriptionManager.Subscriptions.Sum(x => x.Consolidators.Count)); } + [Test] + public void AddConsolidatorAdaptsTradeBarConsolidatorsForQuoteOnlySubscriptions() + { + var algorithm = new AlgorithmStub(); + var eurusd = algorithm.AddForex("EURUSD", market: QuantConnect.Market.Oanda).Symbol; + + using var consolidator = new TradeBarConsolidator(TimeSpan.FromHours(1)); + TradeBar consolidated = null; + consolidator.DataConsolidated += (_, bar) => consolidated = bar; + + algorithm.SubscriptionManager.AddConsolidator(eurusd, consolidator); + + // the registered consolidator is an adapter that accepts the quote bars of the subscription + var subscription = algorithm.SubscriptionManager.Subscriptions.Single(x => x.Symbol == eurusd); + var registered = subscription.Consolidators.Single(); + Assert.IsInstanceOf(registered); + Assert.AreEqual(typeof(QuoteBar), registered.InputType); + + // pump two quote bars an hour apart, the first hour bar is emitted when the second one arrives + var time = new DateTime(2014, 5, 5, 10, 0, 0); + var bid = new Bar(1m, 2m, 0.5m, 1.5m); + var ask = new Bar(1.1m, 2.1m, 0.6m, 1.7m); + registered.Update(new QuoteBar(time, eurusd, bid, 0, ask, 0, Time.OneMinute)); + registered.Update(new QuoteBar(time.AddHours(1), eurusd, bid, 0, ask, 0, Time.OneMinute)); + + Assert.IsNotNull(consolidated); + // trade bars collapsed from quote bars carry the mid-point prices and zero volume + Assert.AreEqual(1.6m, consolidated.Close); + Assert.AreEqual(0, consolidated.Volume); + } + + [Test] + public void RemoveConsolidatorRemovesTheQuoteToTradeBarAdapter() + { + var algorithm = new AlgorithmStub(); + var eurusd = algorithm.AddForex("EURUSD", market: QuantConnect.Market.Oanda).Symbol; + + var consolidator = new TradeBarConsolidator(TimeSpan.FromHours(1)); + algorithm.SubscriptionManager.AddConsolidator(eurusd, consolidator); + Assert.AreEqual(1, algorithm.SubscriptionManager.Subscriptions.Sum(x => x.Consolidators.Count)); + + // users hold and remove the consolidator they created, not the adapter wrapping it + algorithm.SubscriptionManager.RemoveConsolidator(eurusd, consolidator); + Assert.AreEqual(0, algorithm.SubscriptionManager.Subscriptions.Sum(x => x.Consolidators.Count)); + } + + [Test] + public void CanAddAndRemovePythonConsolidatorAdaptedToQuoteOnlyFeed() + { + using var _ = Py.GIL(); + var module = PyModule.FromString("testModule", + @" +from AlgorithmImports import * + +class CustomTradeBarConsolidator(PythonConsolidator): + + def __init__(self): + + #IDataConsolidator required vars for all consolidators + self.consolidated = None + self.working_data = None + self.input_type = TradeBar + self.output_type = TradeBar + + def update(self, data): + pass + + def scan(self, time): + pass + +def get_consolidator(): + return CustomTradeBarConsolidator() +"); + + var algorithm = new AlgorithmStub(); + var eurusd = algorithm.AddForex("EURUSD", market: QuantConnect.Market.Oanda).Symbol; + + var pyConsolidator = module.GetAttr("get_consolidator").Invoke(); + + // the trade bar python consolidator is adapted into the quote-only subscription + algorithm.SubscriptionManager.AddConsolidator(eurusd, pyConsolidator); + var subscription = algorithm.SubscriptionManager.Subscriptions.Single(x => x.Symbol == eurusd); + Assert.IsInstanceOf(subscription.Consolidators.Single()); + + algorithm.SubscriptionManager.RemoveConsolidator(eurusd, pyConsolidator); + Assert.AreEqual(0, algorithm.SubscriptionManager.Subscriptions.Sum(x => x.Consolidators.Count)); + } + + [Test] + public void AddConsolidatorRejectsPeriodsSmallerThanTheSubscriptionPeriodAtRegistration() + { + var algorithm = new AlgorithmStub(); + var spy = algorithm.AddEquity("SPY", Resolution.Daily).Symbol; + + using var consolidator = new TradeBarConsolidator(TimeSpan.FromMinutes(5)); + var exception = Assert.Throws(() => algorithm.SubscriptionManager.AddConsolidator(spy, consolidator)); + StringAssert.Contains("is smaller than the daily subscription period", exception.Message); + } + + [TestCase("period-equal")] + [TestCase("period-larger")] + [TestCase("count")] + [TestCase("calendar")] + public void AddConsolidatorAcceptsValidPeriodsAtRegistration(string testCase) + { + var algorithm = new AlgorithmStub(); + var spy = algorithm.AddEquity("SPY", Resolution.Daily).Symbol; + + using var consolidator = testCase switch + { + "period-equal" => new TradeBarConsolidator(TimeSpan.FromDays(1)), + "period-larger" => new TradeBarConsolidator(TimeSpan.FromDays(7)), + // count and calendar driven consolidators have no fixed time span period to validate + "count" => new TradeBarConsolidator(10), + _ => new TradeBarConsolidator(Calendar.Weekly) + }; + Assert.DoesNotThrow(() => algorithm.SubscriptionManager.AddConsolidator(spy, consolidator)); + } + + [Test] + public void AddConsolidatorThrowsClearMessageForNullSymbol() + { + var algorithm = new AlgorithmStub(); + algorithm.AddEquity("SPY"); + + // e.g. Future.Mapped can still be null when the consolidator is added, we used to throw a raw NRE + using var consolidator = new TradeBarConsolidator(TimeSpan.FromHours(1)); + var exception = Assert.Throws(() => algorithm.SubscriptionManager.AddConsolidator(null, consolidator)); + StringAssert.Contains("Future.Mapped", exception.Message); + } + + [Test] + public void AddConsolidatorTypeMismatchErrorSuggestsConsolidateApi() + { + var algorithm = new AlgorithmStub(); + // custom data is trade-only, and quote consolidators are not adapted into it + var symbol = algorithm.AddData("CUSTOM", Resolution.Daily).Symbol; + + using var consolidator = new QuoteBarConsolidator(TimeSpan.FromDays(1)); + var exception = Assert.Throws(() => algorithm.SubscriptionManager.AddConsolidator(symbol, consolidator)); + StringAssert.Contains("Type mismatch found between consolidator and symbol", exception.Message); + StringAssert.Contains("self.consolidate()", exception.Message); + } + [Test, Parallelizable(ParallelScope.None)] public void RunRemoveConsolidatorsRegressionAlgorithm() { diff --git a/Tests/Indicators/PythonIndicatorTests.cs b/Tests/Indicators/PythonIndicatorTests.cs index a2807ce13bd9..23ff9e5af099 100644 --- a/Tests/Indicators/PythonIndicatorTests.cs +++ b/Tests/Indicators/PythonIndicatorTests.cs @@ -278,11 +278,13 @@ public void AllPythonRegisterIndicatorCases() } //Test 1: Using a C# Consolidator; Should convert consolidator into IDataConsolidator and fail because of the InputType - [TestCase("consolidator", false, "Type mismatch found between consolidator and symbol. Symbol: SPY does not support input type: QuoteBar. Supported types: TradeBar.")] + [TestCase("consolidator", false, "Type mismatch found between consolidator and symbol. Symbol: SPY does not support input type: QuoteBar. Supported types: TradeBar. " + + "Alternatively, consider using QCAlgorithm.Consolidate() ('self.consolidate()' in Python), which selects the correct consolidator type automatically.")] //Test 2: Using a Python Consolidator; Should wrap consolidator and fail because of the InputType - [TestCase("CustomConsolidator", true, "Type mismatch found between consolidator and symbol. Symbol: SPY does not support input type: QuoteBar. Supported types: TradeBar.")] - //Test 3: Using an invalid consolidator; Should try to convert into C#, Python Consolidator and timedelta and fail as the type is invalid - [TestCase("InvalidConsolidator", true, "Invalid third argument, should be either a valid consolidator or timedelta object. The following exception was thrown: ")] + [TestCase("CustomConsolidator", true, "Type mismatch found between consolidator and symbol. Symbol: SPY does not support input type: QuoteBar. Supported types: TradeBar. " + + "Alternatively, consider using QCAlgorithm.Consolidate() ('self.consolidate()' in Python), which selects the correct consolidator type automatically.")] + //Test 3: Using an invalid consolidator; Should try to convert into C#, Python Consolidator, timedelta and calendar and fail as the type is invalid + [TestCase("InvalidConsolidator", true, "Invalid third argument, should be either a valid consolidator, timedelta or calendar (e.g. Calendar.WEEKLY) object. The following exception was thrown: ")] public void AllPythonRegisterIndicatorBadCases(string consolidatorName, bool needsInvoke, string expectedMessage) { //This test covers all three bad cases of registering a indicator through Python