From 5e5b0b295919991eeb36034c37890d3b015e8dde Mon Sep 17 00:00:00 2001 From: Tom Longhurst <30480171+thomhurst@users.noreply.github.com> Date: Wed, 5 Aug 2026 01:57:56 +0100 Subject: [PATCH 1/5] Remove dead module awaiter cancellation --- .../Master/DistributedModuleExecutor.cs | 2 +- .../Worker/WorkerModuleScheduler.cs | 2 +- .../Engine/GeneratedModuleMetadata.cs | 7 --- .../Engine/IModuleScheduler.cs | 6 +- .../Engine/IModuleStateTracker.cs | 6 +- .../Engine/ModuleCompletionSourceCanceller.cs | 58 ------------------- src/ModularPipelines/Engine/ModuleExecutor.cs | 6 +- .../Engine/ModuleScheduler.cs | 7 +-- .../Engine/ModuleStateTracker.cs | 6 +- .../DistributedModuleExecutorTests.cs | 2 +- .../Engine/GeneratedModuleMetadataTests.cs | 15 ----- .../Engine/ModuleExecutorLoggingTests.cs | 10 ++-- .../Engine/TaskCompletionSourceTests.cs | 2 +- 13 files changed, 17 insertions(+), 112 deletions(-) delete mode 100644 src/ModularPipelines/Engine/ModuleCompletionSourceCanceller.cs diff --git a/src/ModularPipelines/Distributed/Master/DistributedModuleExecutor.cs b/src/ModularPipelines/Distributed/Master/DistributedModuleExecutor.cs index 0493b7d399b..283c7b1a9de 100644 --- a/src/ModularPipelines/Distributed/Master/DistributedModuleExecutor.cs +++ b/src/ModularPipelines/Distributed/Master/DistributedModuleExecutor.cs @@ -182,7 +182,7 @@ internal static void CompleteCancelledModules( IModuleResultRegistrar resultRegistrar, CancellationToken cancellationToken) { - var cancelledModules = scheduler.CancelPendingModules(cancelModuleResultAwaiters: false); + var cancelledModules = scheduler.CancelPendingModules(); resultRegistrar.RegisterTerminatedResultsForCancelledModules( cancelledModules, new OperationCanceledException(cancellationToken)); diff --git a/src/ModularPipelines/Distributed/Worker/WorkerModuleScheduler.cs b/src/ModularPipelines/Distributed/Worker/WorkerModuleScheduler.cs index 993b5957760..50e43bf453e 100644 --- a/src/ModularPipelines/Distributed/Worker/WorkerModuleScheduler.cs +++ b/src/ModularPipelines/Distributed/Worker/WorkerModuleScheduler.cs @@ -32,7 +32,7 @@ public void MarkModuleCompleted(Type moduleType, bool success, Exception? except public ModuleState? GetModuleState(Type moduleType) => null; - public IReadOnlyList CancelPendingModules(bool cancelModuleResultAwaiters = true) + public IReadOnlyList CancelPendingModules() { return []; } diff --git a/src/ModularPipelines/Engine/GeneratedModuleMetadata.cs b/src/ModularPipelines/Engine/GeneratedModuleMetadata.cs index efe066ef630..5ca2984510e 100644 --- a/src/ModularPipelines/Engine/GeneratedModuleMetadata.cs +++ b/src/ModularPipelines/Engine/GeneratedModuleMetadata.cs @@ -208,8 +208,6 @@ IModuleResult CreateFailure( ILogger GetOutputLogger(IServiceProvider serviceProvider); - void CancelCompletionSource(IModule module); - void SetCompletionSource(IModule module, IModuleResult result); Task ExecuteAsync( @@ -257,11 +255,6 @@ public ILogger GetOutputLogger(IServiceProvider serviceProvider) return serviceProvider.GetRequiredService>(); } - public void CancelCompletionSource(IModule module) - { - ((Module) module).CompletionSource.TrySetCanceled(); - } - public void SetCompletionSource(IModule module, IModuleResult result) { ((Module) module).CompletionSource.TrySetResult((ModuleResult) result); diff --git a/src/ModularPipelines/Engine/IModuleScheduler.cs b/src/ModularPipelines/Engine/IModuleScheduler.cs index dcf864ecf01..a012a9956df 100644 --- a/src/ModularPipelines/Engine/IModuleScheduler.cs +++ b/src/ModularPipelines/Engine/IModuleScheduler.cs @@ -52,10 +52,6 @@ internal interface IModuleScheduler : IDisposable /// Cancels all modules that are queued or pending (not yet executing) /// This is used when the pipeline is cancelled to ensure TaskCompletionSources are properly completed. /// - /// - /// Whether to cancel typed module result awaiters immediately. Set to - /// when terminated results will be registered after scheduler cancellation. - /// /// The modules transitioned to the completed state by cancellation. - IReadOnlyList CancelPendingModules(bool cancelModuleResultAwaiters = true); + IReadOnlyList CancelPendingModules(); } diff --git a/src/ModularPipelines/Engine/IModuleStateTracker.cs b/src/ModularPipelines/Engine/IModuleStateTracker.cs index b5a0c1fdf43..f6aff0b88fc 100644 --- a/src/ModularPipelines/Engine/IModuleStateTracker.cs +++ b/src/ModularPipelines/Engine/IModuleStateTracker.cs @@ -37,12 +37,8 @@ internal interface IModuleStateTracker /// This is used when the pipeline is cancelled to ensure TaskCompletionSources are properly completed. /// Note: AlwaysRun modules are not cancelled as they should be allowed to complete. /// - /// - /// Whether to cancel typed module result awaiters immediately. Set to - /// when terminated results will be registered after scheduler cancellation. - /// /// The modules transitioned to the completed state by cancellation. - IReadOnlyList CancelPendingModules(bool cancelModuleResultAwaiters = true); + IReadOnlyList CancelPendingModules(); /// /// Gets the state for a specific module. diff --git a/src/ModularPipelines/Engine/ModuleCompletionSourceCanceller.cs b/src/ModularPipelines/Engine/ModuleCompletionSourceCanceller.cs deleted file mode 100644 index 81b63d5943f..00000000000 --- a/src/ModularPipelines/Engine/ModuleCompletionSourceCanceller.cs +++ /dev/null @@ -1,58 +0,0 @@ -using System.Collections.Concurrent; -using System.Diagnostics.CodeAnalysis; -using System.Linq.Expressions; -using ModularPipelines.Modules; - -namespace ModularPipelines.Engine; - -/// -/// Cancels typed module completion sources when modules terminate before execution. -/// -internal static class ModuleCompletionSourceCanceller -{ - private static readonly ConcurrentDictionary> Cache = new(); - - /// - /// Cancels the module's typed completion source. - /// - /// The module instance to cancel. - /// The concrete module type. - public static void Cancel(IModule module, Type moduleType) - { - if (GeneratedModuleMetadata.TryGetRuntime(moduleType, out var runtime)) - { - runtime.CancelCompletionSource(module); - return; - } - - Cache.GetOrAdd(module.ResultType, CreateCanceller)(module); - } - - [UnconditionalSuppressMessage( - "AOT", - "IL3050", - Justification = "The dynamic completion-source canceller is a fallback for modules without generated runtime metadata.")] - [UnconditionalSuppressMessage( - "Trimming", - "IL2075", - Justification = "The dynamic completion-source canceller is a fallback for modules without generated runtime metadata.")] - private static Action CreateCanceller(Type resultType) - { - var moduleType = typeof(Module<>).MakeGenericType(resultType); - var moduleParameter = Expression.Parameter(typeof(IModule), "module"); - var typedModule = Expression.Convert(moduleParameter, moduleType); - var completionSourceProperty = moduleType.GetProperty( - "CompletionSource", - System.Reflection.BindingFlags.Instance | System.Reflection.BindingFlags.NonPublic) - ?? throw new InvalidOperationException($"CompletionSource property not found on {moduleType.Name}"); - var completionSource = Expression.Property(typedModule, completionSourceProperty); - var trySetCanceledMethod = completionSourceProperty.PropertyType.GetMethod( - "TrySetCanceled", - Type.EmptyTypes) - ?? throw new InvalidOperationException( - $"TrySetCanceled method not found on {completionSourceProperty.PropertyType.Name}"); - var cancel = Expression.Call(completionSource, trySetCanceledMethod); - - return Expression.Lambda>(cancel, moduleParameter).Compile(); - } -} diff --git a/src/ModularPipelines/Engine/ModuleExecutor.cs b/src/ModularPipelines/Engine/ModuleExecutor.cs index a72fc9c74b2..ea67ce4d857 100644 --- a/src/ModularPipelines/Engine/ModuleExecutor.cs +++ b/src/ModularPipelines/Engine/ModuleExecutor.cs @@ -189,7 +189,7 @@ private async Task> ExecuteWithSchedulerAsync( catch (Exception exception) { var cancelledModules = - scheduler.CancelPendingModules(cancelModuleResultAwaiters: false); + scheduler.CancelPendingModules(); _resultRegistrar.RegisterTerminatedResultsForCancelledModules( cancelledModules, exception); @@ -230,7 +230,7 @@ private async Task> ExecuteWithSchedulerAsync( private void RegisterCancellationCallback(CancellationTokenSource cancellationTokenSource, IModuleScheduler scheduler) { cancellationTokenSource.Token.Register( - () => scheduler.CancelPendingModules(cancelModuleResultAwaiters: false)); + () => scheduler.CancelPendingModules()); } private async Task ExecuteWorkerPoolAsync( @@ -277,7 +277,7 @@ await Parallel.ForEachAsync( if (isFirstFailure) { var cancelledModules = - scheduler.CancelPendingModules(cancelModuleResultAwaiters: false); + scheduler.CancelPendingModules(); _resultRegistrar.RegisterTerminatedResultsForCancelledModules( cancelledModules, ex); diff --git a/src/ModularPipelines/Engine/ModuleScheduler.cs b/src/ModularPipelines/Engine/ModuleScheduler.cs index 02d69951bd1..55fa163f78c 100644 --- a/src/ModularPipelines/Engine/ModuleScheduler.cs +++ b/src/ModularPipelines/Engine/ModuleScheduler.cs @@ -220,17 +220,14 @@ public void MarkModuleCompleted(Type moduleType, bool success, Exception? except /// This is used when the pipeline is cancelled to ensure TaskCompletionSources are properly completed /// Note: AlwaysRun modules are not cancelled as they should be allowed to complete. /// - /// - /// Whether to cancel typed module result awaiters immediately. - /// - public IReadOnlyList CancelPendingModules(bool cancelModuleResultAwaiters = true) + public IReadOnlyList CancelPendingModules() { if (IsDisposed) { return []; } - return _stateTracker.CancelPendingModules(cancelModuleResultAwaiters); + return _stateTracker.CancelPendingModules(); } public void Dispose() diff --git a/src/ModularPipelines/Engine/ModuleStateTracker.cs b/src/ModularPipelines/Engine/ModuleStateTracker.cs index 6491d97606b..5133af25d80 100644 --- a/src/ModularPipelines/Engine/ModuleStateTracker.cs +++ b/src/ModularPipelines/Engine/ModuleStateTracker.cs @@ -257,7 +257,7 @@ public void MarkModuleCompleted(Type moduleType, bool success, Exception? except } /// - public IReadOnlyList CancelPendingModules(bool cancelModuleResultAwaiters = true) + public IReadOnlyList CancelPendingModules() { List<(ModuleState State, ModuleExecutionState OriginalState)> cancelledModules; @@ -288,10 +288,6 @@ public IReadOnlyList CancelPendingModules(bool cancelModuleResultAwaite foreach (var (moduleState, _) in cancelledModules) { moduleState.CompletionSource.TrySetCanceled(); - if (cancelModuleResultAwaiters) - { - ModuleCompletionSourceCanceller.Cancel(moduleState.Module, moduleState.ModuleType); - } } // Logging outside lock diff --git a/test/ModularPipelines.UnitTests/Distributed/DistributedModuleExecutorTests.cs b/test/ModularPipelines.UnitTests/Distributed/DistributedModuleExecutorTests.cs index 72635f84d34..87c026603bc 100644 --- a/test/ModularPipelines.UnitTests/Distributed/DistributedModuleExecutorTests.cs +++ b/test/ModularPipelines.UnitTests/Distributed/DistributedModuleExecutorTests.cs @@ -30,7 +30,7 @@ public void CompleteCancelledModules_RegistersTerminatedResults() var scheduler = new Mock(); var resultRegistrar = new Mock(); scheduler - .Setup(x => x.CancelPendingModules(false)) + .Setup(x => x.CancelPendingModules()) .Returns(cancelledModules); DistributedModuleExecutor.CompleteCancelledModules( diff --git a/test/ModularPipelines.UnitTests/Engine/GeneratedModuleMetadataTests.cs b/test/ModularPipelines.UnitTests/Engine/GeneratedModuleMetadataTests.cs index 6a6084202d0..4e1535d5713 100644 --- a/test/ModularPipelines.UnitTests/Engine/GeneratedModuleMetadataTests.cs +++ b/test/ModularPipelines.UnitTests/Engine/GeneratedModuleMetadataTests.cs @@ -99,21 +99,6 @@ public async Task Cancelled_Result_Registration_Defers_AlwaysRun_Completion() } } - [Test] - public async Task Generated_Runtime_Cancels_Typed_Completion_Source() - { - var module = new GeneratedMetadataDependencyModule(); - - var found = GeneratedModuleMetadata.TryGetRuntime(module.GetType(), out var runtime); - runtime.CancelCompletionSource(module); - - using (Assert.Multiple()) - { - await Assert.That(found).IsTrue(); - await Assert.That(module.CompletionSource.Task.IsCanceled).IsTrue(); - } - } - [Test] public async Task Generated_Runtime_Resolves_Unbuffered_Output_Logger() { diff --git a/test/ModularPipelines.UnitTests/Engine/ModuleExecutorLoggingTests.cs b/test/ModularPipelines.UnitTests/Engine/ModuleExecutorLoggingTests.cs index 52e05e3c720..2158c2aeb52 100644 --- a/test/ModularPipelines.UnitTests/Engine/ModuleExecutorLoggingTests.cs +++ b/test/ModularPipelines.UnitTests/Engine/ModuleExecutorLoggingTests.cs @@ -103,7 +103,7 @@ public async Task SuccessfulCompletion_DoesNotLogCancellation() var logOutput = logs.ToString(); await Assert.That(logOutput).DoesNotContain("Cancellation triggered"); - scheduler.Verify(x => x.CancelPendingModules(false), Times.Once); + scheduler.Verify(x => x.CancelPendingModules(), Times.Once); } [Test] @@ -119,7 +119,7 @@ public async Task SchedulerAndAlwaysRunFaults_AreAggregatedAfterRegisteringTermi scheduler.SetupGet(x => x.ReadyModules).Returns(readyModules.Reader); scheduler.Setup(x => x.RunSchedulerAsync(It.IsAny())) .Returns(Task.CompletedTask); - scheduler.Setup(x => x.CancelPendingModules(false)) + scheduler.Setup(x => x.CancelPendingModules()) .Returns(cancelledModules); var schedulerFactory = new Mock(); schedulerFactory.Setup(x => x.Create()).Returns(scheduler.Object); @@ -668,7 +668,7 @@ public async Task CancelPendingModules_WithPendingModule_LogsCancellation() } [Test] - public async Task CancelPendingModules_CompletesPendingModuleAwaitable() + public async Task CancelPendingModules_LeavesModuleResultForRegistrarCompletion() { var module = new LaterModule(); var state = new ModuleState(module, module.GetType()); @@ -680,8 +680,8 @@ public async Task CancelPendingModules_CompletesPendingModuleAwaitable() using (Assert.Multiple()) { - await Assert.That(module.CompletionSource.Task.IsCanceled).IsTrue(); - await Assert.That(((IInternalModule) module).ResultTask.IsCompleted).IsTrue(); + await Assert.That(state.CompletionSource.Task.IsCanceled).IsTrue(); + await Assert.That(((IInternalModule) module).ResultTask.IsCompleted).IsFalse(); } } diff --git a/test/ModularPipelines.UnitTests/Engine/TaskCompletionSourceTests.cs b/test/ModularPipelines.UnitTests/Engine/TaskCompletionSourceTests.cs index cc7f8b139db..ebd1fc49bb4 100644 --- a/test/ModularPipelines.UnitTests/Engine/TaskCompletionSourceTests.cs +++ b/test/ModularPipelines.UnitTests/Engine/TaskCompletionSourceTests.cs @@ -292,7 +292,7 @@ private static async Task CreateAbandonedSchedulerFaultAsync(Exce scheduler.SetupGet(x => x.ReadyModules).Returns(readyModules.Reader); scheduler.Setup(x => x.RunSchedulerAsync(It.IsAny())) .Returns(schedulerTaskSource.Task); - scheduler.Setup(x => x.CancelPendingModules(false)).Returns([]); + scheduler.Setup(x => x.CancelPendingModules()).Returns([]); var schedulerFactory = new Mock(); schedulerFactory.Setup(x => x.Create()).Returns(scheduler.Object); var alwaysRunHandler = new Mock(); From 9ae309cca02a6a58206eb896c089f088a2372f50 Mon Sep 17 00:00:00 2001 From: Tom Longhurst <30480171+thomhurst@users.noreply.github.com> Date: Wed, 5 Aug 2026 02:05:29 +0100 Subject: [PATCH 2/5] test: update scheduler cancellation mocks --- .../Master/DistributedModuleExecutorTests.cs | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/test/ModularPipelines.Distributed.UnitTests/Master/DistributedModuleExecutorTests.cs b/test/ModularPipelines.Distributed.UnitTests/Master/DistributedModuleExecutorTests.cs index b7650ed53ad..33a1a54ed69 100644 --- a/test/ModularPipelines.Distributed.UnitTests/Master/DistributedModuleExecutorTests.cs +++ b/test/ModularPipelines.Distributed.UnitTests/Master/DistributedModuleExecutorTests.cs @@ -192,7 +192,7 @@ private static Mock CreateMockScheduler(params ModuleState[] m return tcs.Task; }); scheduler.Setup(s => s.MarkModuleStarted(It.IsAny())).Returns(true); - scheduler.Setup(s => s.CancelPendingModules(It.IsAny())).Returns([]); + scheduler.Setup(s => s.CancelPendingModules()).Returns([]); return scheduler; } @@ -898,7 +898,7 @@ public async Task Cancellation_After_Start_Claim_Completes_Module_Result( var module = new DistributedModule(); var moduleState = new ModuleState(module, typeof(DistributedModule)); var scheduler = CreateMockScheduler(moduleState); - scheduler.Setup(s => s.CancelPendingModules(false)) + scheduler.Setup(s => s.CancelPendingModules()) .Returns([]); scheduler.Setup(s => s.MarkModuleStarted(typeof(DistributedModule))) .Returns(() => @@ -954,7 +954,7 @@ public async Task Publish_Failure_Completes_Module_Result( var module = new DistributedModule(); var moduleState = new ModuleState(module, typeof(DistributedModule)); var scheduler = CreateMockScheduler(moduleState); - scheduler.Setup(s => s.CancelPendingModules(false)) + scheduler.Setup(s => s.CancelPendingModules()) .Returns([]); var publishException = new InvalidOperationException("Broker unavailable"); From 065815412e4d5630b5a5da40fdd1cc93c1b757d9 Mon Sep 17 00:00:00 2001 From: Tom Longhurst <30480171+thomhurst@users.noreply.github.com> Date: Mon, 10 Aug 2026 03:18:51 +0100 Subject: [PATCH 3/5] ci: retrigger cancelled checks Current head already fixes the stale scheduler-mock finding. Trigger fresh validation after the previous workflow was externally cancelled. Refs #3871 From 92920e8b96f404d6c05685eb4cbf11dc833caff6 Mon Sep 17 00:00:00 2001 From: Tom Longhurst <30480171+thomhurst@users.noreply.github.com> Date: Mon, 10 Aug 2026 04:01:11 +0100 Subject: [PATCH 4/5] docs(engine): clarify cancellation contract Document that pending-module cancellation only completes internal scheduler sources and requires result registration for public tasks. --- src/ModularPipelines/Engine/IModuleScheduler.cs | 5 +++-- src/ModularPipelines/Engine/IModuleStateTracker.cs | 3 ++- src/ModularPipelines/Engine/ModuleScheduler.cs | 5 +++-- 3 files changed, 8 insertions(+), 5 deletions(-) diff --git a/src/ModularPipelines/Engine/IModuleScheduler.cs b/src/ModularPipelines/Engine/IModuleScheduler.cs index a012a9956df..6b41c9c2184 100644 --- a/src/ModularPipelines/Engine/IModuleScheduler.cs +++ b/src/ModularPipelines/Engine/IModuleScheduler.cs @@ -49,8 +49,9 @@ internal interface IModuleScheduler : IDisposable ModuleState? GetModuleState(Type moduleType); /// - /// Cancels all modules that are queued or pending (not yet executing) - /// This is used when the pipeline is cancelled to ensure TaskCompletionSources are properly completed. + /// Cancels all modules that are queued or pending (not yet executing). + /// This cancels only the scheduler's internal completion sources. Call + /// RegisterTerminatedResultsForCancelledModules for the returned modules to complete their public result tasks. /// /// The modules transitioned to the completed state by cancellation. IReadOnlyList CancelPendingModules(); diff --git a/src/ModularPipelines/Engine/IModuleStateTracker.cs b/src/ModularPipelines/Engine/IModuleStateTracker.cs index f6aff0b88fc..c1243c5ac8f 100644 --- a/src/ModularPipelines/Engine/IModuleStateTracker.cs +++ b/src/ModularPipelines/Engine/IModuleStateTracker.cs @@ -34,7 +34,8 @@ internal interface IModuleStateTracker /// /// Cancels all modules that are queued or pending (not yet executing). - /// This is used when the pipeline is cancelled to ensure TaskCompletionSources are properly completed. + /// This cancels only the scheduler's internal completion sources. Call + /// RegisterTerminatedResultsForCancelledModules for the returned modules to complete their public result tasks. /// Note: AlwaysRun modules are not cancelled as they should be allowed to complete. /// /// The modules transitioned to the completed state by cancellation. diff --git a/src/ModularPipelines/Engine/ModuleScheduler.cs b/src/ModularPipelines/Engine/ModuleScheduler.cs index 55fa163f78c..226e4d0bb6a 100644 --- a/src/ModularPipelines/Engine/ModuleScheduler.cs +++ b/src/ModularPipelines/Engine/ModuleScheduler.cs @@ -216,8 +216,9 @@ public void MarkModuleCompleted(Type moduleType, bool success, Exception? except } /// - /// Cancels all modules that are queued or pending (not yet executing) - /// This is used when the pipeline is cancelled to ensure TaskCompletionSources are properly completed + /// Cancels all modules that are queued or pending (not yet executing). + /// This cancels only the scheduler's internal completion sources. Call + /// RegisterTerminatedResultsForCancelledModules for the returned modules to complete their public result tasks. /// Note: AlwaysRun modules are not cancelled as they should be allowed to complete. /// public IReadOnlyList CancelPendingModules() From b29f8a106bc51dbe0c4db6b5ab7eaaec3cdfa28c Mon Sep 17 00:00:00 2001 From: Tom Longhurst <30480171+thomhurst@users.noreply.github.com> Date: Mon, 10 Aug 2026 04:21:27 +0100 Subject: [PATCH 5/5] docs(engine): document AlwaysRun exclusion --- src/ModularPipelines/Engine/IModuleScheduler.cs | 1 + 1 file changed, 1 insertion(+) diff --git a/src/ModularPipelines/Engine/IModuleScheduler.cs b/src/ModularPipelines/Engine/IModuleScheduler.cs index 6b41c9c2184..ff9cfe913ac 100644 --- a/src/ModularPipelines/Engine/IModuleScheduler.cs +++ b/src/ModularPipelines/Engine/IModuleScheduler.cs @@ -50,6 +50,7 @@ internal interface IModuleScheduler : IDisposable /// /// Cancels all modules that are queued or pending (not yet executing). + /// AlwaysRun modules are excluded and are allowed to complete. /// This cancels only the scheduler's internal completion sources. Call /// RegisterTerminatedResultsForCancelledModules for the returned modules to complete their public result tasks. ///