Skip to content
Merged
Original file line number Diff line number Diff line change
Expand Up @@ -16,11 +16,15 @@

using System;
using System.Collections.Generic;
using System.IO;
using System.Linq;
using Python.Runtime;
using QuantConnect.Configuration;
using QuantConnect.Data;
using QuantConnect.Data.UniverseSelection;
using QuantConnect.Interfaces;
using QuantConnect.Logging;
using QuantConnect.Securities;
using QuantConnect.Util;

namespace QuantConnect.Lean.Engine.DataFeeds.Enumerators.Factories
Expand All @@ -30,9 +34,15 @@ namespace QuantConnect.Lean.Engine.DataFeeds.Enumerators.Factories
/// </summary>
public class LiveCustomDataSubscriptionEnumeratorFactory : ISubscriptionEnumeratorFactory
{
// when the expected universe file is not available yet, we fall back to the backup universe file ("*.backup"),
// if any, as a last resort, when the market is open or within this time span before the next market open
private static readonly TimeSpan UniverseFileBackupFallbackWindow =
TimeSpan.FromMinutes(Config.GetInt("universe-file-backup-fallback-minutes", 30));

private readonly TimeSpan _minimumIntervalCheck;
private readonly ITimeProvider _timeProvider;
private readonly Func<DateTime, DateTime> _dateAdjustment;
private readonly BackupUniverseFileDataProvider _backupUniverseFileDataProvider;
private readonly IObjectStore _objectStore;

/// <summary>
Expand All @@ -42,13 +52,21 @@ public class LiveCustomDataSubscriptionEnumeratorFactory : ISubscriptionEnumerat
/// <param name="objectStore">The object store to use</param>
/// <param name="dateAdjustment">Func that allows adjusting the datetime to use</param>
/// <param name="minimumIntervalCheck">Allows specifying the minimum interval between each enumerator refresh and data check, default is 30 minutes</param>
/// <param name="fallBackToBackupUniverseFiles">Whether to fall back to the backup universe file ("*.backup"), if any, as a last resort
/// when the expected universe file is not available and the market is open or close to opening.
/// Only meaningful for universe subscriptions backed by local files</param>
public LiveCustomDataSubscriptionEnumeratorFactory(ITimeProvider timeProvider, IObjectStore objectStore,
Func<DateTime, DateTime> dateAdjustment = null, TimeSpan? minimumIntervalCheck = null)
Func<DateTime, DateTime> dateAdjustment = null, TimeSpan? minimumIntervalCheck = null,
bool fallBackToBackupUniverseFiles = false)
{
_timeProvider = timeProvider;
_dateAdjustment = dateAdjustment;
_minimumIntervalCheck = minimumIntervalCheck ?? TimeSpan.FromMinutes(30);
_objectStore = objectStore;
if (fallBackToBackupUniverseFiles)
{
_backupUniverseFileDataProvider = new BackupUniverseFileDataProvider();
}
}

/// <summary>
Expand All @@ -66,6 +84,7 @@ public IEnumerator<BaseData> CreateEnumerator(SubscriptionRequest request, IData
var frontier = Ref.Create(_dateAdjustment?.Invoke(request.StartTimeLocal) ?? request.StartTimeLocal);
var lastSourceRefreshTime = DateTime.MinValue;
var sourceFactory = config.GetBaseDataInstance();
var exchangeHours = request.Security.Exchange.Hours;

// this is refreshing the enumerator stack for each new source
var refresher = new RefreshEnumerator<BaseData>(() =>
Expand All @@ -79,16 +98,27 @@ public IEnumerator<BaseData> CreateEnumerator(SubscriptionRequest request, IData
}

lastSourceRefreshTime = utcNow;
var localDate = _dateAdjustment?.Invoke(utcNow.ConvertFromUtc(config.ExchangeTimeZone).Date) ?? utcNow.ConvertFromUtc(config.ExchangeTimeZone).Date;
var localTime = utcNow.ConvertFromUtc(config.ExchangeTimeZone);
var localDate = _dateAdjustment?.Invoke(localTime.Date) ?? localTime.Date;
var source = sourceFactory.GetSource(config, localDate, true);
if (source == null)
{
// a null source is equivalent to an unreachable source: no data this cycle, retry on the next refresh
return Enumerable.Empty<BaseData>().GetEnumerator();
}

var sourceDataProvider = dataProvider;
if (_backupUniverseFileDataProvider != null
&& source.TransportMedium == SubscriptionTransportMedium.LocalFile
&& IsUniverseFileBackupFallbackActive(exchangeHours, localTime))
{
// if the expected universe file is not available, the backup universe file, if any, will be read as a last resort
_backupUniverseFileDataProvider.SetDataProvider(dataProvider);
sourceDataProvider = _backupUniverseFileDataProvider;
}

// fetch the new source and enumerate the data source reader
var enumerator = EnumerateDataSourceReader(config, dataProvider, frontier, source, localDate, sourceFactory);
var enumerator = EnumerateDataSourceReader(config, sourceDataProvider, frontier, source, localDate, sourceFactory);

if (SourceRequiresFastForward(source))
{
Expand Down Expand Up @@ -202,6 +232,20 @@ IDataProvider dataProvider
return SubscriptionDataSourceReader.ForSource(source, dataCacheProvider, config, date, true, baseDataInstance, dataProvider, _objectStore);
}

/// <summary>
/// Determines whether the backup universe file fallback is active, which is only when the market is open or close to opening
/// (within <see cref="UniverseFileBackupFallbackWindow"/> of the next market open), when the expected universe file should already be available.
/// It is evaluated at the same cadence as the enumerator refreshes
/// </summary>
/// <param name="exchangeHours">The exchange hours of the security</param>
/// <param name="localTime">The current time in the exchange time zone</param>
private static bool IsUniverseFileBackupFallbackActive(SecurityExchangeHours exchangeHours, DateTime localTime)
{
return exchangeHours.IsOpen(localTime, extendedMarketHours: false)
// if the market is closed, GetNextMarketOpen returns the next day open
|| exchangeHours.GetNextMarketOpen(localTime, extendedMarketHours: false) - localTime <= UniverseFileBackupFallbackWindow;
}

private bool SourceRequiresFastForward(SubscriptionDataSource source)
{
return source.TransportMedium == SubscriptionTransportMedium.LocalFile
Expand All @@ -217,5 +261,50 @@ private static TimeSpan GetMaximumDataAge(TimeSpan increment)
{
return TimeSpan.FromTicks(Math.Max(increment.Ticks, TimeSpan.FromSeconds(5).Ticks));
}

/// <summary>
/// Data provider wrapper that falls back to the backup universe file ("*.backup"), if any,
/// when the expected universe file can't be fetched, as a last resort
/// </summary>
private sealed class BackupUniverseFileDataProvider : IDataProvider
{
private IDataProvider _dataProvider;

/// <summary>
/// Event raised each time data fetch is finished (successfully or not)
/// </summary>
public event EventHandler<DataProviderNewDataRequestEventArgs> NewDataRequest
{
add => _dataProvider?.NewDataRequest += value;
remove => _dataProvider?.NewDataRequest -= value;
}

/// <summary>
/// Sets the data provider to wrap, forwarding its <see cref="IDataProvider.NewDataRequest"/> events
/// </summary>
public void SetDataProvider(IDataProvider dataProvider)
{
_dataProvider = dataProvider;
}

public Stream Fetch(string key)
{
var stream = _dataProvider.Fetch(key);
if (stream != null)
{
return stream;
}

var backupKey = key + ".backup";
stream = _dataProvider.Fetch(backupKey);
if (stream != null)
{
Log.Trace($"LiveCustomDataSubscriptionEnumeratorFactory.BackupUniverseFileDataProvider.Fetch(): universe file '{key}' is not available, " +
$"falling back to backup universe file '{backupKey}'");
}

return stream;
}
}
}
}
4 changes: 3 additions & 1 deletion Engine/DataFeeds/LiveTradingDataFeed.cs
Original file line number Diff line number Diff line change
Expand Up @@ -356,7 +356,9 @@ request.Universe is OptionChainUniverse ||
_algorithm.ObjectStore,
// we adjust time to the previous tradable date
time => Time.GetStartTimeForTradeBars(request.Security.Exchange.Hours, time, Time.OneDay, 1, false, config.DataTimeZone, _algorithm.Settings.DailyPreciseEndTime),
TimeSpan.FromMinutes(10)
TimeSpan.FromMinutes(10),
// when the expected universe file is not available yet, fall back to the backup universe file as a last resort
fallBackToBackupUniverseFiles: true
);
var enumeratorStack = factory.CreateEnumerator(request, _dataProvider);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@

using System;
using System.Collections.Generic;
using System.IO;
using System.Linq;
using Moq;
using NUnit.Framework;
Expand Down Expand Up @@ -549,6 +550,129 @@ public void ToleratesNullSource()
enumerator.DisposeSafely();
}

[Test]
public void FallsBackToBackupUniverseFileWhenExpectedSourceIsNotAvailable()
{
// 10 am, the market is open, so the backup fallback is active
var referenceLocal = new DateTime(2017, 10, 12, 10, 0, 0);
var referenceUtc = referenceLocal.ConvertToUtc(TimeZones.NewYork);

var timeProvider = new ManualTimeProvider(referenceUtc);

var expectedSourceAvailable = false;
var dataProvider = new Mock<IDataProvider>();
dataProvider.Setup(dp => dp.Fetch("local.file.source")).Returns(() => expectedSourceAvailable ? new MemoryStream() : null);
dataProvider.Setup(dp => dp.Fetch("local.file.source.backup")).Returns(() => new MemoryStream());

var dataSourceReader = new Mock<ISubscriptionDataSourceReader>();
var factory = new TestableLiveCustomDataSubscriptionEnumeratorFactory(timeProvider, dataSourceReader.Object,
fallBackToBackupUniverseFiles: true);
SetUpDataSourceReader(dataSourceReader, factory, () => new LocalFileData { EndTime = timeProvider.GetUtcNow().ConvertFromUtc(TimeZones.NewYork).AddSeconds(1) });

var config = new SubscriptionDataConfig(typeof(LocalFileData), Symbols.SPY, Resolution.Daily, TimeZones.NewYork, TimeZones.NewYork, false, false, false);
var request = GetSubscriptionRequest(config, referenceUtc.AddSeconds(-1), referenceUtc.AddDays(1));

using var enumerator = factory.CreateEnumerator(request, dataProvider.Object);

// the expected source is not available, so the backup file is the one that gets read, in a single read of the source
Assert.IsTrue(enumerator.MoveNext());
Assert.IsNotNull(enumerator.Current);
VerifyGetSourceInvocationCount(dataSourceReader, 1, "local.file.source", SubscriptionTransportMedium.LocalFile, FileFormat.Csv);
dataProvider.Verify(dp => dp.Fetch("local.file.source"), Times.Once);
dataProvider.Verify(dp => dp.Fetch("local.file.source.backup"), Times.Once);

// the fallback is rate limited like the source refreshes
Assert.IsTrue(enumerator.MoveNext());
Assert.IsNull(enumerator.Current);
dataProvider.Verify(dp => dp.Fetch(It.IsAny<string>()), Times.Exactly(2));

// the expected source is preferred on the next refresh once it becomes available, without touching the backup file
expectedSourceAvailable = true;
timeProvider.Advance(TimeSpan.FromMinutes(30));
Assert.IsTrue(enumerator.MoveNext());
Assert.IsNotNull(enumerator.Current);
VerifyGetSourceInvocationCount(dataSourceReader, 2, "local.file.source", SubscriptionTransportMedium.LocalFile, FileFormat.Csv);
dataProvider.Verify(dp => dp.Fetch("local.file.source"), Times.Exactly(2));
dataProvider.Verify(dp => dp.Fetch("local.file.source.backup"), Times.Once);
}

[Test]
public void DoesNotFallBackToBackupUniverseFileFarFromMarketOpen()
{
// midnight, more than the fallback window away from the next market open, so the backup file is never tried
var referenceLocal = new DateTime(2017, 10, 12);
var referenceUtc = referenceLocal.ConvertToUtc(TimeZones.NewYork);

var timeProvider = new ManualTimeProvider(referenceUtc);

// the expected source is not available
var dataProvider = new Mock<IDataProvider>();

var dataSourceReader = new Mock<ISubscriptionDataSourceReader>();
var factory = new TestableLiveCustomDataSubscriptionEnumeratorFactory(timeProvider, dataSourceReader.Object,
fallBackToBackupUniverseFiles: true);
SetUpDataSourceReader(dataSourceReader, factory, () => new LocalFileData { EndTime = timeProvider.GetUtcNow().ConvertFromUtc(TimeZones.NewYork).AddSeconds(1) });

var config = new SubscriptionDataConfig(typeof(LocalFileData), Symbols.SPY, Resolution.Daily, TimeZones.NewYork, TimeZones.NewYork, false, false, false);
var request = GetSubscriptionRequest(config, referenceUtc.AddSeconds(-1), referenceUtc.AddDays(1));

using var enumerator = factory.CreateEnumerator(request, dataProvider.Object);

Assert.IsTrue(enumerator.MoveNext());
Assert.IsNull(enumerator.Current);

// only the expected source is tried
VerifyGetSourceInvocationCount(dataSourceReader, 1, "local.file.source", SubscriptionTransportMedium.LocalFile, FileFormat.Csv);
dataProvider.Verify(dp => dp.Fetch("local.file.source"), Times.Once);
dataProvider.Verify(dp => dp.Fetch("local.file.source.backup"), Times.Never);
}

[Test]
public void DoesNotFallBackToBackupUniverseFileWhenNotConfigured()
{
// 10 am, the market is open, but the factory is not configured to fall back to backup universe files
var referenceLocal = new DateTime(2017, 10, 12, 10, 0, 0);
var referenceUtc = referenceLocal.ConvertToUtc(TimeZones.NewYork);

var timeProvider = new ManualTimeProvider(referenceUtc);

// the expected source is not available
var dataProvider = new Mock<IDataProvider>();

var dataSourceReader = new Mock<ISubscriptionDataSourceReader>();
var factory = new TestableLiveCustomDataSubscriptionEnumeratorFactory(timeProvider, dataSourceReader.Object);
SetUpDataSourceReader(dataSourceReader, factory, () => new LocalFileData { EndTime = timeProvider.GetUtcNow().ConvertFromUtc(TimeZones.NewYork).AddSeconds(1) });

var config = new SubscriptionDataConfig(typeof(LocalFileData), Symbols.SPY, Resolution.Daily, TimeZones.NewYork, TimeZones.NewYork, false, false, false);
var request = GetSubscriptionRequest(config, referenceUtc.AddSeconds(-1), referenceUtc.AddDays(1));

using var enumerator = factory.CreateEnumerator(request, dataProvider.Object);

Assert.IsTrue(enumerator.MoveNext());
Assert.IsNull(enumerator.Current);

// only the expected source is tried
VerifyGetSourceInvocationCount(dataSourceReader, 1, "local.file.source", SubscriptionTransportMedium.LocalFile, FileFormat.Csv);
dataProvider.Verify(dp => dp.Fetch("local.file.source"), Times.Once);
dataProvider.Verify(dp => dp.Fetch("local.file.source.backup"), Times.Never);
}

/// <summary>
/// Sets up the mocked data source reader to fetch the source through the data cache provider the factory gave it,
/// like the real readers do, yielding a data point only when the source could be fetched
/// </summary>
private static void SetUpDataSourceReader(Mock<ISubscriptionDataSourceReader> dataSourceReader,
TestableLiveCustomDataSubscriptionEnumeratorFactory factory, Func<BaseData> dataFactory)
{
dataSourceReader.Setup(dsr => dsr.Read(It.IsAny<SubscriptionDataSource>()))
.Returns((SubscriptionDataSource source) =>
{
using var stream = factory.DataCacheProvider.Fetch(source.Source);
return stream == null ? Enumerable.Empty<BaseData>() : new[] { dataFactory() };
})
.Verifiable();
}

private static void VerifyGetSourceInvocationCount(Mock<ISubscriptionDataSourceReader> dataSourceReader, int count, string source, SubscriptionTransportMedium medium, FileFormat fileFormat)
{
dataSourceReader.Verify(dsr => dsr.Read(It.Is<SubscriptionDataSource>(sds =>
Expand Down Expand Up @@ -630,8 +754,14 @@ class TestableLiveCustomDataSubscriptionEnumeratorFactory : LiveCustomDataSubscr
{
private readonly ISubscriptionDataSourceReader _dataSourceReader;

public TestableLiveCustomDataSubscriptionEnumeratorFactory(ITimeProvider timeProvider, ISubscriptionDataSourceReader dataSourceReader, TimeSpan? minimumIntervalCheck = null)
: base(timeProvider, null, minimumIntervalCheck: minimumIntervalCheck)
/// <summary>
/// The data cache provider the last data source reader was created with
/// </summary>
public IDataCacheProvider DataCacheProvider { get; private set; }

public TestableLiveCustomDataSubscriptionEnumeratorFactory(ITimeProvider timeProvider, ISubscriptionDataSourceReader dataSourceReader,
TimeSpan? minimumIntervalCheck = null, bool fallBackToBackupUniverseFiles = false)
: base(timeProvider, null, minimumIntervalCheck: minimumIntervalCheck, fallBackToBackupUniverseFiles: fallBackToBackupUniverseFiles)
{
_dataSourceReader = dataSourceReader;
}
Expand All @@ -643,6 +773,7 @@ protected override ISubscriptionDataSourceReader GetSubscriptionDataSourceReader
BaseData baseData,
IDataProvider dataProvider)
{
DataCacheProvider = dataCacheProvider;
return _dataSourceReader;
}
}
Expand Down
Loading
Loading