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; + } } } }