Files
quantconnect--lean/Engine/DataFeeds/Enumerators/BaseDataCollectionAggregatorEnumerator.cs
Jhonathan Abreu e29bb2c5e0
API Tests / build (push) Has been cancelled
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
File-based options universe (#8212)
* Initial options universe with greeks implementation

* Options universe improvements

* Address peer review

* File based options universe fixes and improvements.

- Adjust OptionUniverse start-end times and period.
- Adapt unit tests and some algorithms to pass with new options universe selection.

* Updated options regression algorithms stats for new universe data

* Updated options regression algorithms stats for new universe data

* Updated options regression algorithms stats for new universe data

* Updated options regression algorithms stats for new universe data

* Updated options regression algorithms stats for new universe data

* Option chain provider with new options universe

* Allow canonical option history requests

* Address peer review

* Address peer review

* Fix symbols parsing in OptionUniverse

* Fix universe selection subscriptions start time to not include extended market hours

* Minor changes

* Minor changes

* Peer recommended changes and fixes

* Update regression algorithm stats

* Update regression algorithms stats and minor fixes

* Fix option chain provider history request

* Round option indicators values

* Added option universe csv header property

* Update regression algorithms stats

* Update regression algorithms stats

* Data fixes and regression algos stats update

* Unit test fixes

* Minor changes

* Option chain handling in live trading data feed

* Minor changes

* Added processed data provider

* Fix thread-safety violation in Slice class

* Minor change

* Update options filter universe API to use OptionUniverse data

Add new filter methods for greeks, IV and open interest

* Option filter universe api updates

* Add OptionUniverse history regression algorithms

* Add regression algorithms for new options filter universe api methods

* Added options greeks data and updated regression algorithms

* Address peer review

* Address peer review

* Add more assertions to new options filter api regression algorithms

* Minor performance improvement.

Reduce greeks binomial model steps to 140

* Minor tests updates

* Greeks numerical models performance improvements

* Greeks numerical models performance improvements

* Revert array pool change for option pricing numerical models

* Update default dividend yield provider depending on option type

* [TEST]

* Add helper method con calculate time till expiration

* Use double in price option numerical models

* Implied volatility calculation improvements

- Adjust root finding method accuracy as a factor of the option price
- Use BSM to get a first guess

* Cleanup

* Some regression algorithms and unit tests cleanup

* Regression tests updates after rebasing from master

* Add universe files

* Self review and cleanup

* Minor regression tests updates after rebase

* Fix: set data time zone to same as exchange tz for options universes

* Minor change

* Minor change

* Fix for live trading options universe selection

* Keep underlying when aggregating collections in BaseDataCollectionAggregatorEnumerator

* Update index options regression algorithms stats

* Minor change

* Address peer review

* Memory usage improvements

* Minor build fix

* Minor changes and test fixes

* Cache symbols in OptionUniverse

* Cleanup

* Fix index option creation in OptionUniverse

* Use cached underlying SID when parsing from string

* Abstract symbols cache to BaseDataCollection

* Return actual underlying symbol when mapping decomposing ICO ticker

* Address peer review

* Minor performance improvements reduce garbage

* Limit Symbols and SIDs cache size to help with memory usage

* Minor fix in symbols and sid cache cleanup

* Build fix

* Lazily parse greeks on individual access

* Cleanup and tests

* Address peer review

* Minor greeks fix

---------

Co-authored-by: Martin Molinero <martin.molinero1@gmail.com>
2024-09-09 12:39:31 -03:00

235 lines
8.9 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 QuantConnect.Data;
using System.Collections;
using System.Collections.Generic;
using QuantConnect.Data.UniverseSelection;
namespace QuantConnect.Lean.Engine.DataFeeds.Enumerators
{
/// <summary>
/// Provides an implementation of <see cref="IEnumerator{BaseDataCollection}"/>
/// that aggregates an underlying <see cref="IEnumerator{BaseData}"/> into a single
/// data packet
/// </summary>
public class BaseDataCollectionAggregatorEnumerator : IEnumerator<BaseDataCollection>
{
private bool _endOfStream;
private bool _needsMoveNext;
private bool _liveMode;
private readonly Symbol _symbol;
private readonly IEnumerator<BaseData> _enumerator;
/// <summary>
/// Initializes a new instance of the <see cref="BaseDataCollectionAggregatorEnumerator"/> class
/// This will aggregate instances emitted from the underlying enumerator and tag them with the
/// specified symbol
/// </summary>
/// <param name="enumerator">The underlying enumerator to aggregate</param>
/// <param name="symbol">The symbol to place on the aggregated collection</param>
/// <param name="liveMode">True if running in live mode</param>
public BaseDataCollectionAggregatorEnumerator(IEnumerator<BaseData> enumerator, Symbol symbol, bool liveMode = false)
{
_symbol = symbol;
_enumerator = enumerator;
_liveMode = liveMode;
_needsMoveNext = true;
}
/// <summary>
/// Advances the enumerator to the next element of the collection.
/// </summary>
/// <returns>
/// true if the enumerator was successfully advanced to the next element; false if the enumerator has passed the end of the collection.
/// </returns>
/// <exception cref="T:System.InvalidOperationException">The collection was modified after the enumerator was created. </exception><filterpriority>2</filterpriority>
public bool MoveNext()
{
if (_endOfStream)
{
return false;
}
BaseDataCollection collection = null;
while (true)
{
if (_needsMoveNext)
{
// move next if we dequeued the last item last time we were invoked
if (!_enumerator.MoveNext())
{
_endOfStream = true;
if (!IsValid(collection))
{
// we don't emit
collection = null;
}
break;
}
}
if (_enumerator.Current == null)
{
// the underlying returned null, stop here and start again on the next call
_needsMoveNext = true;
break;
}
if (collection == null)
{
// we have new data, set the collection's symbol/times
var current = _enumerator.Current;
collection = CreateCollection(_symbol, current.Time, current.EndTime);
}
if (collection.EndTime != _enumerator.Current.EndTime)
{
// the data from the underlying is at a different time, stop here
_needsMoveNext = false;
if (IsValid(collection))
{
// we emit
break;
}
// we try again
collection = null;
continue;
}
// this data belongs in this collection, keep going until null or bad time
Add(collection, _enumerator.Current);
_needsMoveNext = true;
}
Current = collection;
return _liveMode || collection != null;
}
/// <summary>
/// Sets the enumerator to its initial position, which is before the first element in the collection.
/// </summary>
/// <exception cref="T:System.InvalidOperationException">The collection was modified after the enumerator was created. </exception><filterpriority>2</filterpriority>
public void Reset()
{
_enumerator.Reset();
}
/// <summary>
/// Gets the element in the collection at the current position of the enumerator.
/// </summary>
/// <returns>
/// The element in the collection at the current position of the enumerator.
/// </returns>
public BaseDataCollection Current
{
get; private set;
}
/// <summary>
/// Gets the current element in the collection.
/// </summary>
/// <returns>
/// The current element in the collection.
/// </returns>
/// <filterpriority>2</filterpriority>
object IEnumerator.Current
{
get { return Current; }
}
/// <summary>
/// Performs application-defined tasks associated with freeing, releasing, or resetting unmanaged resources.
/// </summary>
/// <filterpriority>2</filterpriority>
public void Dispose()
{
_enumerator.Dispose();
}
/// <summary>
/// Creates a new, empty <see cref="BaseDataCollection"/>.
/// </summary>
/// <param name="symbol">The base data collection symbol</param>
/// <param name="time">The start time of the collection</param>
/// <param name="endTime">The end time of the collection</param>
/// <returns>A new, empty <see cref="BaseDataCollection"/></returns>
private BaseDataCollection CreateCollection(Symbol symbol, DateTime time, DateTime endTime)
{
return new BaseDataCollection
{
Symbol = symbol,
Time = time,
EndTime = endTime
};
}
/// <summary>
/// Adds the specified instance of <see cref="BaseData"/> to the current collection
/// </summary>
/// <param name="collection">The collection to be added to</param>
/// <param name="current">The data to be added</param>
private void Add(BaseDataCollection collection, BaseData current)
{
var baseDataCollection = current as BaseDataCollection;
if (_symbol.HasUnderlying && _symbol.Underlying == current.Symbol)
{
// if the underlying has been aggregated, even if it shouldn't need to be, let's handle it nicely
if (baseDataCollection != null)
{
collection.Underlying = baseDataCollection.Data[0];
}
else
{
collection.Underlying = current;
}
}
else
{
if (baseDataCollection != null)
{
// datapoint is already aggregated, let's see if it's a single point or a collection we can use already
if(baseDataCollection.Data.Count > 1)
{
collection.Data = baseDataCollection.Data;
}
else
{
collection.Data.Add(baseDataCollection.Data[0]);
}
// Let's keep the underlying in case it's already there
collection.Underlying ??= baseDataCollection.Underlying;
}
else
{
collection.Data.Add(current);
}
}
}
/// <summary>
/// Determines if a given data point is valid and can be emitted
/// </summary>
/// <param name="collection">The collection to be emitted</param>
/// <returns>True if its a valid data point</returns>
private static bool IsValid(BaseDataCollection collection)
{
return collection != null && collection.Data?.Count > 0;
}
}
}