Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions Engine/Results/Analysis/ResultsAnalyzer.cs
Original file line number Diff line number Diff line change
Expand Up @@ -267,8 +267,8 @@ public static ResultsAnalyzer CreateForInRunAnalysis(QCAlgorithm algorithm, Lang
/// <param name="result">The current intermediate backtest result. Its orders and order events
/// are windows truncated to the most recent ones, so the in-run analyses can miss orders and
/// events already evicted from them; the final analysis re-scans the complete data. Its charts
/// are the handler's live ones, read without synchronization: a torn read while the algorithm
/// thread updates them can fail a run, which the handler catches, and the next run retries.</param>
/// are a copy the handler takes under its chart lock, so the analyses can enumerate them while
/// the algorithm thread keeps sampling the live ones.</param>
/// <param name="logs">The full list of log lines produced so far; the analyzer analyzes the
/// lines past the ones consumed by previous runs.</param>
/// <param name="totalPerformance">The current total algorithm performance, for analyses that read
Expand Down
63 changes: 40 additions & 23 deletions Engine/Results/BacktestingResultHandler.cs
Original file line number Diff line number Diff line change
Expand Up @@ -74,6 +74,11 @@ public class BacktestingResultHandler : BaseResultsHandler, IResultHandler
/// </summary>
protected bool RunResultsAnalysis { get; set; } = true;

/// <summary>
/// The delay after the handler starts before the first result is stored, which also runs the first in-run analysis
/// </summary>
protected virtual TimeSpan InitialResultStoreDelay { get; } = TimeSpan.FromSeconds(5);

/// <summary>
/// A dictionary containing summary statistics
/// </summary>
Expand All @@ -91,7 +96,7 @@ public BacktestingResultHandler()
_chartSeriesCount = new();

// Delay uploading first packet
_nextS3Update = StartTime.AddSeconds(5);
_nextS3Update = StartTime.Add(InitialResultStoreDelay);
}

/// <summary>
Expand Down Expand Up @@ -225,8 +230,10 @@ private void Update()
const int maxOrders = 100;
var orderCount = TransactionHandler.Orders.Count;

// The in-run analyses enumerate the charts while the algorithm thread keeps sampling them,
// which fails the enumeration, so hand them a copy taken under the chart lock
var completeResult = new BacktestResult(new BacktestResultParameters(
Charts,
CloneCharts(),
orderCount > maxOrders ? TransactionHandler.Orders.Skip(orderCount - maxOrders).ToDictionary() : TransactionHandler.Orders.ToDictionary(),
Algorithm.Transactions.TransactionRecord,
new Dictionary<string, string>(),
Expand Down Expand Up @@ -332,28 +339,26 @@ protected override void StoreResult(Packet packet)
// Get Storage Location:
var key = $"{AlgorithmId}.json";

BacktestResult results;
lock (ChartLock)
// The charts are a snapshot taken under the chart lock, so they are cleaned up and stored without another copy
if (result.Results.Charts.TryGetValue(PortfolioMarginKey, out var marginChart))
{
results = new BacktestResult(new BacktestResultParameters(
result.Results.Charts.ToDictionary(x => x.Key, x => x.Value.Clone()),
result.Results.Orders,
result.Results.ProfitLoss,
result.Results.Statistics,
result.Results.RuntimeStatistics,
result.Results.RollingWindow,
null, // null order events, we store them separately
result.Results.TotalPerformance,
result.Results.AlgorithmConfiguration,
result.Results.State,
result.Results.Analysis,
result.Results.ServerStatistics));

if (result.Results.Charts.TryGetValue(PortfolioMarginKey, out var marginChart))
{
PortfolioMarginChart.RemoveSinglePointSeries(marginChart);
}
PortfolioMarginChart.RemoveSinglePointSeries(marginChart);
}

var results = new BacktestResult(new BacktestResultParameters(
result.Results.Charts,
result.Results.Orders,
result.Results.ProfitLoss,
result.Results.Statistics,
result.Results.RuntimeStatistics,
result.Results.RollingWindow,
null, // null order events, we store them separately
result.Results.TotalPerformance,
result.Results.AlgorithmConfiguration,
result.Results.State,
result.Results.Analysis,
result.Results.ServerStatistics));

// Save results
SaveResults(key, results);

Expand Down Expand Up @@ -384,7 +389,8 @@ protected void SendFinalResult()
if (Algorithm != null)
{
//Convert local dictionary:
var charts = new Dictionary<string, Chart>(Charts);
// The algorithm thread can still be sampling when it was stopped for exceeding a limit
var charts = CloneCharts();
var orders = new Dictionary<int, Order>(TransactionHandler.Orders);
var profitLoss = new SortedDictionary<DateTime, decimal>(Algorithm.Transactions.TransactionRecord);
var statisticsResults = GenerateStatisticsResults(charts, profitLoss, _capacityEstimate);
Expand Down Expand Up @@ -519,6 +525,17 @@ private List<string> CloneLogs()
}
}

/// <summary>
/// Takes a snapshot of the charts under the chart lock.
/// </summary>
private Dictionary<string, Chart> CloneCharts()
{
lock (ChartLock)
{
return Charts.ToDictionary(x => x.Key, x => x.Value.Clone());
}
}

/// <summary>
/// Sends the in-run analysis findings to the browser in their own packet.
/// </summary>
Expand Down
76 changes: 76 additions & 0 deletions Tests/Engine/Results/BacktestingResultHandlerTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -22,10 +22,13 @@
using QuantConnect.Logging;
using QuantConnect.Packets;
using QuantConnect.Report;
using QuantConnect.Statistics;
using QuantConnect.Tests.Engine.DataFeeds;
using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;

namespace QuantConnect.Tests.Engine.Results
{
Expand Down Expand Up @@ -582,11 +585,84 @@ public void RecapturesStartingPortfolioValueAfterWarmup()
}
}

[Test]
public void InRunAnalysisRunsWhileTheAlgorithmSamplesTheCharts()
{
using var api = new Api.Api();
using var messaging = new QuantConnect.Messaging.Messaging();
var resultHandler = new TestableBacktestingResultHandler();
resultHandler.Initialize(new(new BacktestNodePacket(), messaging, api, new BacktestingTransactionHandler(), null));

using var sampling = new CancellationTokenSource();
try
{
var algorithm = new AlgorithmStub();
resultHandler.SetAlgorithm(algorithm, 100000);

// A margin chart large enough for the analysis to still be reading it when the next sample lands
var series = new Series("SPY", SeriesType.StackedArea, "%");
var time = new DateTime(2024, 1, 1);
for (var i = 0; i < 100000; i++)
{
series.AddPoint(new ChartPoint(time.AddMinutes(i), 50));
}
var chart = new Chart(BaseResultsHandler.PortfolioMarginKey);
chart.AddSeries(series);
resultHandler.Charts[chart.Name] = chart;

// The algorithm thread samples under the chart lock. Sampling the last point again overwrites it,
// which keeps the chart bounded but still invalidates any enumeration of it in progress
var lastPoint = series.Values[^1];
var sampler = Task.Run(() =>
{
while (!sampling.IsCancellationRequested)
{
lock (resultHandler.ExposedChartLock)
{
series.AddPoint(new ChartPoint(lastPoint.Time, 50));
}
}
});

// Locking the algorithm enables the handler updates
algorithm.SetLocked();
Assert.IsTrue(resultHandler.InRunAnalysisRan.Wait(TimeSpan.FromSeconds(30)), "The in-run analysis did not run");

sampling.Cancel();
sampler.Wait();
Assert.IsNotNull(resultHandler.InRunAnalysisFindings, "The in-run analysis failed, see the logged error");
}
finally
{
sampling.Cancel();
resultHandler.Exit();
}
}

private class TestableBacktestingResultHandler : BacktestingResultHandler
{
public decimal ExposedStartingPortfolioValue => StartingPortfolioValue;
public decimal ExposedDailyPortfolioValue => DailyPortfolioValue;
public decimal ExposedCumulativeMaxPortfolioValue => CumulativeMaxPortfolioValue;
public object ExposedChartLock => ChartLock;

// Store the first result, which runs the first in-run analysis, on the first update instead of 5 seconds in
protected override TimeSpan InitialResultStoreDelay => TimeSpan.Zero;

public ManualResetEventSlim InRunAnalysisRan { get; } = new();

/// <summary>
/// The findings of the first in-run analysis run, null when the analysis failed
/// </summary>
public IReadOnlyList<QuantConnect.Analysis> InRunAnalysisFindings { get; private set; }

protected override IReadOnlyList<QuantConnect.Analysis> RunInRunResultsAnalysis(BacktestResult completeResult,
AlgorithmPerformance totalPerformance)
{
InRunAnalysisFindings = base.RunInRunResultsAnalysis(completeResult, totalPerformance);
InRunAnalysisRan.Set();
return InRunAnalysisFindings;
}
}
}
}
Loading