diff --git a/Engine/Results/Analysis/ResultsAnalyzer.cs b/Engine/Results/Analysis/ResultsAnalyzer.cs
index cae7a0e2c96a..ccf5edfc9a14 100644
--- a/Engine/Results/Analysis/ResultsAnalyzer.cs
+++ b/Engine/Results/Analysis/ResultsAnalyzer.cs
@@ -267,8 +267,8 @@ public static ResultsAnalyzer CreateForInRunAnalysis(QCAlgorithm algorithm, Lang
/// 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.
+ /// 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.
/// The full list of log lines produced so far; the analyzer analyzes the
/// lines past the ones consumed by previous runs.
/// The current total algorithm performance, for analyses that read
diff --git a/Engine/Results/BacktestingResultHandler.cs b/Engine/Results/BacktestingResultHandler.cs
index fff50479125e..137fc806f2a7 100644
--- a/Engine/Results/BacktestingResultHandler.cs
+++ b/Engine/Results/BacktestingResultHandler.cs
@@ -74,6 +74,11 @@ public class BacktestingResultHandler : BaseResultsHandler, IResultHandler
///
protected bool RunResultsAnalysis { get; set; } = true;
+ ///
+ /// The delay after the handler starts before the first result is stored, which also runs the first in-run analysis
+ ///
+ protected virtual TimeSpan InitialResultStoreDelay { get; } = TimeSpan.FromSeconds(5);
+
///
/// A dictionary containing summary statistics
///
@@ -91,7 +96,7 @@ public BacktestingResultHandler()
_chartSeriesCount = new();
// Delay uploading first packet
- _nextS3Update = StartTime.AddSeconds(5);
+ _nextS3Update = StartTime.Add(InitialResultStoreDelay);
}
///
@@ -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(),
@@ -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);
@@ -384,7 +389,8 @@ protected void SendFinalResult()
if (Algorithm != null)
{
//Convert local dictionary:
- var charts = new Dictionary(Charts);
+ // The algorithm thread can still be sampling when it was stopped for exceeding a limit
+ var charts = CloneCharts();
var orders = new Dictionary(TransactionHandler.Orders);
var profitLoss = new SortedDictionary(Algorithm.Transactions.TransactionRecord);
var statisticsResults = GenerateStatisticsResults(charts, profitLoss, _capacityEstimate);
@@ -519,6 +525,17 @@ private List CloneLogs()
}
}
+ ///
+ /// Takes a snapshot of the charts under the chart lock.
+ ///
+ private Dictionary CloneCharts()
+ {
+ lock (ChartLock)
+ {
+ return Charts.ToDictionary(x => x.Key, x => x.Value.Clone());
+ }
+ }
+
///
/// Sends the in-run analysis findings to the browser in their own packet.
///
diff --git a/Tests/Engine/Results/BacktestingResultHandlerTests.cs b/Tests/Engine/Results/BacktestingResultHandlerTests.cs
index e62dd1551574..fe2474de238c 100644
--- a/Tests/Engine/Results/BacktestingResultHandlerTests.cs
+++ b/Tests/Engine/Results/BacktestingResultHandlerTests.cs
@@ -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
{
@@ -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();
+
+ ///
+ /// The findings of the first in-run analysis run, null when the analysis failed
+ ///
+ public IReadOnlyList InRunAnalysisFindings { get; private set; }
+
+ protected override IReadOnlyList RunInRunResultsAnalysis(BacktestResult completeResult,
+ AlgorithmPerformance totalPerformance)
+ {
+ InRunAnalysisFindings = base.RunInRunResultsAnalysis(completeResult, totalPerformance);
+ InRunAnalysisRan.Set();
+ return InRunAnalysisFindings;
+ }
}
}
}