diff --git a/CLAUDE.md b/CLAUDE.md index 2b3aad3f..447c7c27 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -63,7 +63,7 @@ Rule of thumb: anything named `*.NetDaemon` is a thin adapter; keep logic in the - **`UseLightGroup`** (unrelated to `ForLights`) is the Zigbee-group optimisation: a `GroupNode` is appended last to each pipeline and `GroupNodeContext` buffers transitions for a window (default 20 ms); if every member gets an identical transition, one `ApplyTransition` goes to the group entity instead of N unicasts. - **Reactive nodes** (`ReactiveNode/`): a `ReactiveNode` swaps its active child based on an observable of **node factories** (`Func>`), each activation getting its own DI scope; `null` = deactivate/pass through. `On`, `TurnOffWhen`, `AddToggle`, `AddCycle`, `AddInteractionNode` and `AddNotifications` are all sugar over this. `InstantiationScope.Shared` builds nested nodes once for all lights via the composite factory maps (lazy, locked, `ScopedPipelineNode` disposal clears the map); `PerChild` builds per light per activation. Dimmers (`AddReactiveDimmer`) wrap the reactive node in `ReactiveDimmerPipeline`; for a dimmer shared across lights the first configurator's options win. - **Interaction node**: `AddInteractionNode` watches `light.StateChanges()` and, when brightness hits 0 while `LightPipelineContext` says the pipeline didn't output that off, inserts a `TurnOffThenPassThroughNode` — physical switch wins until the next upstream change. -- **Threading**: `ReactiveNode` serialises all state mutation through a `Subject` + `Synchronize()` queue (replaced a `Lock` that deadlocked — see commit `0b4571e`). Don't reintroduce locks held across callbacks into nodes. `Lock` remains in `GroupNodeContext`, `RegistrationManager` and the composite factory maps. +- **Threading**: `ReactiveNode` and `LightTransitionNode` serialise all state mutation through a `SerializedActionQueue` (`Utils/`; callers on other threads enqueue and return, the draining thread runs the action; replaced a `Lock` that deadlocked — see commit `0b4571e`). Internal `LightTransitionNode` subclasses that react to timers or observables wrap their callbacks in `RunSerialized` (see `ResettableTimeoutNode`). Don't reintroduce locks held across callbacks into nodes. `Lock` remains in `GroupNodeContext`, `RegistrationManager` and the composite factory maps. - **Hierarchy/logging**: all configurators implement internal `IPipelineHierarchyContext`; `ActionExtensions.ApplyHierarchySettings` propagates names/logging into nested configurators before user code runs. `PipelineLogger` logs at Trace, prefixed `[light.id] Path->To->Node`. ### Light model conventions (`CodeCasa.Lights`, `Lights.NetDaemon`) diff --git a/README.md b/README.md index aee1d11f..490dfec7 100644 --- a/README.md +++ b/README.md @@ -390,7 +390,7 @@ lightPipelineFactory.SetupLightPipeline(lightEntities.HallwayLight, pipeline => .When(motionDetected, LightParameters.Dimmed) .AddReactiveNode(node => node .On(brightButtonPressed, sp => sp.CreateAutoPassThroughLightNode( - LightParameters.Bright, TimeSpan.FromMinutes(10), motionDetected))); + LightParameters.Bright, TimeSpan.FromMinutes(10)))); }); ``` diff --git a/src/CodeCasa.AutomationPipelines.Lights/Extensions/LightTransitionNodeExtensions.cs b/src/CodeCasa.AutomationPipelines.Lights/Extensions/LightTransitionNodeExtensions.cs index 17cf4a60..1e97f00f 100644 --- a/src/CodeCasa.AutomationPipelines.Lights/Extensions/LightTransitionNodeExtensions.cs +++ b/src/CodeCasa.AutomationPipelines.Lights/Extensions/LightTransitionNodeExtensions.cs @@ -21,7 +21,7 @@ public static class LightTransitionNodeExtensions public static IPipelineNode TurnOffAfter(this IPipelineNode node, TimeSpan timeSpan, IScheduler scheduler) { - return new ResettableTimeoutNode(node, timeSpan, Observable.Empty(), scheduler); + return node.TurnOffAfter(timeSpan, Observable.Empty(), scheduler); } /// @@ -56,7 +56,7 @@ public static IPipelineNode TurnOffAfter(this IPipelineNode PassThroughAfter(this IPipelineNode node, TimeSpan timeSpan, IScheduler scheduler) { - return new ResettableTimeoutNode(node, timeSpan, Observable.Empty(), scheduler, TimeoutBehaviour.PassThrough); + return node.PassThroughAfter(timeSpan, Observable.Empty(), scheduler); } /// diff --git a/src/CodeCasa.AutomationPipelines.Lights/Extensions/ServiceProviderExtensions.cs b/src/CodeCasa.AutomationPipelines.Lights/Extensions/ServiceProviderExtensions.cs index d45fb10c..dbf058b9 100644 --- a/src/CodeCasa.AutomationPipelines.Lights/Extensions/ServiceProviderExtensions.cs +++ b/src/CodeCasa.AutomationPipelines.Lights/Extensions/ServiceProviderExtensions.cs @@ -3,6 +3,7 @@ using CodeCasa.Lights; using Microsoft.Extensions.DependencyInjection; using System.Reactive.Concurrency; +using System.Reactive.Linq; namespace CodeCasa.AutomationPipelines.Lights.Extensions; @@ -70,9 +71,7 @@ public static IPipelineNode CreateAutoOffLightNode(this IServic LightParameters lightParameters, TimeSpan timeSpan) { - var scheduler = serviceProvider.GetRequiredService(); - var innerNode = new StaticLightTransitionNode(lightParameters.AsTransition(), scheduler); - return innerNode.TurnOffAfter(timeSpan, scheduler); + return serviceProvider.CreateAutoOffLightNode(lightParameters, timeSpan, Observable.Empty()); } /// @@ -89,9 +88,7 @@ public static IPipelineNode CreateAutoOffLightNode(this IServic LightParameters lightParameters, TimeSpan timeSpan, IObservable persistObservable) { - var scheduler = serviceProvider.GetRequiredService(); - var innerNode = new StaticLightTransitionNode(lightParameters.AsTransition(), scheduler); - return innerNode.TurnOffAfter(timeSpan, persistObservable, scheduler); + return serviceProvider.CreateTimeoutLightNode(lightParameters, timeSpan, persistObservable, TimeoutBehaviour.TurnOff); } /// @@ -114,9 +111,7 @@ public static IPipelineNode CreateAutoPassThroughLightNode(this LightParameters lightParameters, TimeSpan timeSpan) { - var scheduler = serviceProvider.GetRequiredService(); - var innerNode = new StaticLightTransitionNode(lightParameters.AsTransition(), scheduler); - return innerNode.PassThroughAfter(timeSpan, scheduler); + return serviceProvider.CreateAutoPassThroughLightNode(lightParameters, timeSpan, Observable.Empty()); } /// @@ -140,9 +135,16 @@ public static IPipelineNode CreateAutoPassThroughLightNode(this public static IPipelineNode CreateAutoPassThroughLightNode(this IServiceProvider serviceProvider, LightParameters lightParameters, TimeSpan timeSpan, IObservable persistObservable) + { + return serviceProvider.CreateTimeoutLightNode(lightParameters, timeSpan, persistObservable, TimeoutBehaviour.PassThrough); + } + + private static IPipelineNode CreateTimeoutLightNode(this IServiceProvider serviceProvider, + LightParameters lightParameters, TimeSpan timeSpan, IObservable persistObservable, + TimeoutBehaviour timeoutBehaviour) { var scheduler = serviceProvider.GetRequiredService(); var innerNode = new StaticLightTransitionNode(lightParameters.AsTransition(), scheduler); - return innerNode.PassThroughAfter(timeSpan, persistObservable, scheduler); + return new ResettableTimeoutNode(innerNode, timeSpan, persistObservable, scheduler, timeoutBehaviour); } } \ No newline at end of file diff --git a/src/CodeCasa.AutomationPipelines.Lights/Nodes/LightTransitionNode.cs b/src/CodeCasa.AutomationPipelines.Lights/Nodes/LightTransitionNode.cs index 9378f7dc..83003bee 100644 --- a/src/CodeCasa.AutomationPipelines.Lights/Nodes/LightTransitionNode.cs +++ b/src/CodeCasa.AutomationPipelines.Lights/Nodes/LightTransitionNode.cs @@ -2,6 +2,7 @@ using System.Reactive.Linq; using System.Reactive.Subjects; using CodeCasa.AutomationPipelines.Lights.Extensions; +using CodeCasa.AutomationPipelines.Lights.Utils; using CodeCasa.Lights; namespace CodeCasa.AutomationPipelines.Lights.Nodes; @@ -13,13 +14,16 @@ namespace CodeCasa.AutomationPipelines.Lights.Nodes; public abstract class LightTransitionNode(IScheduler scheduler) : IPipelineNode { private readonly Subject _newOutputSubject = new(); + // Inputs, outputs and scheduled continuations arrive on different threads; see SerializedActionQueue for why this is not a lock. + private readonly SerializedActionQueue _stateQueue = new(); private LightParameters? _inputLightDestinationParameters; private DateTime? _inputStartOfTransition; private DateTime? _inputEndOfTransition; private LightTransition? _output; private bool _passThroughNextInput; private IDisposable? _scheduledAction; - private bool _isDisposed; + private int _scheduleGeneration; + private volatile bool _isDisposed; /// /// Gets the source light parameters from the previous input, useful for interpolating transitions. @@ -38,14 +42,9 @@ public abstract class LightTransitionNode(IScheduler scheduler) : IPipelineNode< public LightTransition? Input { get; - set + set => RunSerialized(() => { - if (_isDisposed) - { - return; - } - - _scheduledAction?.Dispose(); // Always cancel scheduled actions when the input changes. + CancelScheduledAction(); // Always cancel scheduled actions when the input changes. // We save additional information on the light transition that we can later use to continue the transition if it would be interrupted. InputLightSourceParameters = _inputLightDestinationParameters; field = value; @@ -69,7 +68,7 @@ public LightTransition? Input } InputReceived(field); - } + }); } /// @@ -105,30 +104,33 @@ protected void PassInputThrough() public LightTransition? Output { get => _output; - protected set + protected set => RunSerialized(() => { - if (_isDisposed) - { - return; - } - - _scheduledAction?.Dispose(); // Always cancel scheduled actions when the output is changed directly. + CancelScheduledAction(); // Always cancel scheduled actions when the output is changed directly. PassThrough = false; SetOutputInternal(value); - } + }); } /// /// Schedules an interpolated light transition that will animate from source to desired parameters using the input's transition time. /// + /// + /// The remainder of the transition is cancelled by the next input, so a node that keeps overriding the input has to + /// schedule again from . + /// /// The source light parameters to transition from. /// The desired light parameters to transition to. protected void ScheduleInterpolatedLightTransitionUsingInputTransitionTime(LightParameters? sourceLightParameters, LightParameters? desiredLightParameters) { - PassThrough = false; - _scheduledAction = scheduler.ScheduleInterpolatedLightTransition(sourceLightParameters, - desiredLightParameters, _inputStartOfTransition, _inputEndOfTransition, SetOutputInternal); + RunSerialized(() => + { + // The PassThrough setter only cancels when the value changes, which it does not when scheduling twice in a row. + CancelScheduledAction(); + PassThrough = false; + ScheduleInterpolated(sourceLightParameters, desiredLightParameters); + }); } /// @@ -138,13 +140,8 @@ protected void ScheduleInterpolatedLightTransitionUsingInputTransitionTime(Light public bool PassThrough { get; - set + set => RunSerialized(() => { - if (_isDisposed) - { - return; - } - // Always reset _passThroughNextInput when PassThrough is explicitly called. _passThroughNextInput = false; @@ -153,16 +150,14 @@ public bool PassThrough return; } - _scheduledAction?.Dispose(); // Always cancel scheduled actions when the pass through value changes. + CancelScheduledAction(); // Always cancel scheduled actions when the pass through value changes. field = value; if (field) { - _scheduledAction = scheduler.ScheduleInterpolatedLightTransition(InputLightSourceParameters, - _inputLightDestinationParameters, _inputStartOfTransition, _inputEndOfTransition, - SetOutputInternal); + ScheduleInterpolated(InputLightSourceParameters, _inputLightDestinationParameters); } - } + }); } /// @@ -172,8 +167,11 @@ public bool PassThrough /// The output light transition to set. protected void ChangeOutputAndTurnOnPassThroughOnNextInput(LightTransition? output) { - Output = output; - TurnOnPassThroughOnNextInput(); + RunSerialized(() => + { + Output = output; + TurnOnPassThroughOnNextInput(); + }); } /// @@ -182,17 +180,66 @@ protected void ChangeOutputAndTurnOnPassThroughOnNextInput(LightTransition? outp /// protected void TurnOnPassThroughOnNextInput() { - if (PassThrough) + RunSerialized(() => + { + if (PassThrough) + { + return; + } + + _passThroughNextInput = true; + }); + } + + /// + /// Runs serialised with every other state change of this node, so a node that reacts to + /// timers or observables can change several things without an input or a continuation running in between. The action + /// is skipped once the node is disposed. + /// + private protected void RunSerialized(Action action) + { + _stateQueue.Run(() => { + if (!_isDisposed) + { + action(); + } + }); + } + + private void CancelScheduledAction() + { + // Disposing does not stop a continuation that already started on the scheduler thread; the generation makes it drop its output. + _scheduleGeneration++; + _scheduledAction?.Dispose(); + _scheduledAction = null; + } + + private void ScheduleInterpolated(LightParameters? sourceLightParameters, LightParameters? desiredLightParameters) + { + var generation = _scheduleGeneration; + var scheduledAction = scheduler.ScheduleInterpolatedLightTransition(sourceLightParameters, + desiredLightParameters, _inputStartOfTransition, _inputEndOfTransition, output => RunSerialized(() => + { + // Checked inside the queue: a cancellation on another thread is either fully applied by now or runs after this output. + if (generation == _scheduleGeneration) + { + SetOutputInternal(output); + } + })); + + // The first output is emitted synchronously, so a downstream reaction may already have cancelled or replaced this schedule. + if (generation != _scheduleGeneration) + { + scheduledAction?.Dispose(); return; } - - _passThroughNextInput = true; + _scheduledAction = scheduledAction; } private void SetOutputInternal(LightTransition? output) { - // Scheduled continuations may still fire after disposal; a disposed subject would throw on the scheduler thread. + // Scheduled continuations may still fire after disposal and must not emit anymore. if (_isDisposed) { return; @@ -203,7 +250,7 @@ private void SetOutputInternal(LightTransition? output) } /// - public override string ToString() => GetType().Name; + public override string ToString() => Name ?? GetType().Name; /// public virtual ValueTask DisposeAsync() @@ -213,10 +260,13 @@ public virtual ValueTask DisposeAsync() return ValueTask.CompletedTask; } _isDisposed = true; - _scheduledAction?.Dispose(); - _scheduledAction = null; - _newOutputSubject.OnCompleted(); - _newOutputSubject.Dispose(); + // Queued behind a state change that is still running on another thread, so completion is always the last thing the subject sees. + _stateQueue.Run(() => + { + CancelScheduledAction(); + // The subject is completed but not disposed, so subscribing to a disposed node completes instead of throwing. + _newOutputSubject.OnCompleted(); + }); return ValueTask.CompletedTask; } } \ No newline at end of file diff --git a/src/CodeCasa.AutomationPipelines.Lights/Nodes/ResettableTimeoutNode.cs b/src/CodeCasa.AutomationPipelines.Lights/Nodes/ResettableTimeoutNode.cs index a943cc75..460e85e4 100644 --- a/src/CodeCasa.AutomationPipelines.Lights/Nodes/ResettableTimeoutNode.cs +++ b/src/CodeCasa.AutomationPipelines.Lights/Nodes/ResettableTimeoutNode.cs @@ -1,3 +1,4 @@ +using CodeCasa.AutomationPipelines.Lights.Utils; using CodeCasa.Lights; using System.Reactive.Concurrency; using System.Reactive.Disposables; @@ -12,16 +13,20 @@ internal class ResettableTimeoutNode : LightTransitionNode private readonly CompositeDisposable _disposables = new(); private readonly SerialDisposable _timerSubscription = new(); private bool _isPersisting; + private bool _hasHandedOver; + private int _timerGeneration; + private int _isChildDisposed; private bool _isDisposed; - public ResettableTimeoutNode(IPipelineNode childNode, TimeSpan turnOffTime, + public ResettableTimeoutNode(IPipelineNode childNode, TimeSpan timeout, IObservable persistObservable, IScheduler scheduler, TimeoutBehaviour timeoutBehaviour = TimeoutBehaviour.TurnOff) : base(scheduler) { _childNode = childNode; + var childName = childNode.Name ?? childNode.ToString(); Name = timeoutBehaviour == TimeoutBehaviour.PassThrough - ? $"{childNode.Name} (passes through after timeout)" - : $"{childNode.Name} (resets after timeout)"; + ? $"{childName} (passes through after timeout)" + : $"{childName} (resets after timeout)"; _timerSubscription.DisposeWith(_disposables); // The initial output is set synchronously: a reactive node reads Output right after activating this node, and would @@ -29,29 +34,45 @@ public ResettableTimeoutNode(IPipelineNode childNode, TimeSpan Output = childNode.Output; RestartTimer(); + /* + * The scheduler may run the callbacks below and the timer on different threads. Each one runs serialised with the + * node's input, and checks _hasHandedOver and the timer generation there: disposing a subscription does not stop + * a callback that already started, and such a callback must not claim the light again after a hand-over. + */ childNode.OnNewOutput .ObserveOn(scheduler) - .Subscribe(output => + .Subscribe(output => RunSerialized(() => { + if (_hasHandedOver) + { + return; + } + Output = output; RestartTimer(); - }).DisposeWith(_disposables); + })).DisposeWith(_disposables); persistObservable .ObserveOn(scheduler) .DistinctUntilChanged() - .Subscribe(persist => + .Subscribe(persist => RunSerialized(() => { + if (_hasHandedOver) + { + return; + } + _isPersisting = persist; if (persist) { + _timerGeneration++; _timerSubscription.Disposable = null; } else { RestartTimer(); } - }).DisposeWith(_disposables); + })).DisposeWith(_disposables); void RestartTimer() { @@ -60,8 +81,15 @@ void RestartTimer() return; } - _timerSubscription.Disposable = Observable.Timer(turnOffTime, scheduler) - .Subscribe(_ => OnTimeout()); + var generation = ++_timerGeneration; + _timerSubscription.Disposable = Observable.Timer(timeout, scheduler) + .Subscribe(_ => RunSerialized(() => + { + if (generation == _timerGeneration) + { + OnTimeout(); + } + })); } void OnTimeout() @@ -73,16 +101,31 @@ void OnTimeout() } // Passing through ends the override for good: a later child output or persist change must not claim the light again. + _hasHandedOver = true; _disposables.Dispose(); PassInputThrough(); + DisposeChildNode().GetAwaiter().GetResult(); } } protected override void OnInputChanged(LightTransition? input) { + if (_hasHandedOver) + { + return; + } + _childNode.Input = input; } + private Task DisposeChildNode() + { + // A hand-over on the scheduler thread and DisposeAsync can get here at the same time. + return Interlocked.Exchange(ref _isChildDisposed, 1) == 0 + ? _childNode.DisposeOrDisposeAsync() + : Task.CompletedTask; + } + public override async ValueTask DisposeAsync() { if (_isDisposed) @@ -92,7 +135,7 @@ public override async ValueTask DisposeAsync() _isDisposed = true; _disposables.Dispose(); - await _childNode.DisposeAsync(); + await DisposeChildNode(); await base.DisposeAsync(); } } diff --git a/src/CodeCasa.AutomationPipelines.Lights/ReactiveNode/ReactiveNode.cs b/src/CodeCasa.AutomationPipelines.Lights/ReactiveNode/ReactiveNode.cs index fbde08c1..2bd38358 100644 --- a/src/CodeCasa.AutomationPipelines.Lights/ReactiveNode/ReactiveNode.cs +++ b/src/CodeCasa.AutomationPipelines.Lights/ReactiveNode/ReactiveNode.cs @@ -1,5 +1,4 @@ -using System.Collections.Concurrent; -using System.Reactive; +using System.Reactive; using System.Reactive.Linq; using System.Reactive.Subjects; using CodeCasa.AutomationPipelines.Lights.Utils; @@ -18,8 +17,7 @@ public class ReactiveNode : PipelineNode private readonly ILogger? _logger; private readonly IEqualityComparer? _equalityComparer; private readonly Subject _nodeChangedSubject = new(); - private readonly ConcurrentQueue _stateQueue = new(); - private int _drainingThreadId; + private readonly SerializedActionQueue _stateQueue = new(); private volatile bool _isDisposed; private IDisposable? _nodeObservableSubscription; private IDisposable? _activeNodeSubscription; @@ -103,13 +101,7 @@ protected override void InputReceived(LightTransition? input) }); } - /* - * All state mutation is serialised through this queue. Unlike a lock, a caller on another thread never waits: - * it enqueues and returns, and the thread that is already draining runs the action. Holding a lock while calling - * into child nodes deadlocked nested reactive nodes (outer input vs. inner trigger, see 0b4571e for the same - * problem in Pipeline). Actions enqueued from within a running action execute inline, which keeps synchronous - * emissions during activation ordered exactly as before. - */ + // All state mutation is serialised through the queue; see SerializedActionQueue for why this is not a lock. private void EnqueueStateChange(Action action) { if (_isDisposed) @@ -117,28 +109,7 @@ private void EnqueueStateChange(Action action) return; } - var currentThreadId = Environment.CurrentManagedThreadId; - if (Volatile.Read(ref _drainingThreadId) == currentThreadId) - { - RunStateChange(action); - return; - } - - _stateQueue.Enqueue(action); - while (!_stateQueue.IsEmpty && Interlocked.CompareExchange(ref _drainingThreadId, currentThreadId, 0) == 0) - { - try - { - while (_stateQueue.TryDequeue(out var next)) - { - RunStateChange(next); - } - } - finally - { - Volatile.Write(ref _drainingThreadId, 0); - } - } + _stateQueue.Run(() => RunStateChange(action)); } private void RunStateChange(Action action) @@ -221,7 +192,8 @@ public override async ValueTask DisposeAsync() } _stateQueue.Clear(); - _nodeChangedSubject.Dispose(); + // The subject is completed but not disposed: a state change still draining on another thread would otherwise throw. + _nodeChangedSubject.OnCompleted(); await base.DisposeAsync(); } diff --git a/src/CodeCasa.AutomationPipelines.Lights/Utils/SerializedActionQueue.cs b/src/CodeCasa.AutomationPipelines.Lights/Utils/SerializedActionQueue.cs new file mode 100644 index 00000000..789cd2aa --- /dev/null +++ b/src/CodeCasa.AutomationPipelines.Lights/Utils/SerializedActionQueue.cs @@ -0,0 +1,56 @@ +using System.Collections.Concurrent; +using System.Runtime.ExceptionServices; + +namespace CodeCasa.AutomationPipelines.Lights.Utils; + +/* + * Runs actions one at a time. Unlike a lock, a caller on another thread never waits: it enqueues and returns, and + * the thread that is already draining runs the action. Holding a lock while calling into child nodes deadlocked + * nested reactive nodes (outer input vs. inner trigger, see 0b4571e for the same problem in Pipeline). Actions + * enqueued from within a running action execute inline, which keeps synchronous emissions ordered as they would be + * without the queue. + */ +internal sealed class SerializedActionQueue +{ + private readonly ConcurrentQueue _queue = new(); + private int _drainingThreadId; + + public void Run(Action action) + { + var currentThreadId = Environment.CurrentManagedThreadId; + if (Volatile.Read(ref _drainingThreadId) == currentThreadId) + { + action(); + return; + } + + ExceptionDispatchInfo? failure = null; + _queue.Enqueue(action); + while (!_queue.IsEmpty && Interlocked.CompareExchange(ref _drainingThreadId, currentThreadId, 0) == 0) + { + try + { + while (_queue.TryDequeue(out var next)) + { + try + { + next(); + } + catch (Exception e) + { + // A throwing action must not strand the actions other threads queued behind it. + failure ??= ExceptionDispatchInfo.Capture(e); + } + } + } + finally + { + Volatile.Write(ref _drainingThreadId, 0); + } + } + + failure?.Throw(); + } + + public void Clear() => _queue.Clear(); +} diff --git a/src/CodeCasa.AutomationPipelines/PipelineNode.cs b/src/CodeCasa.AutomationPipelines/PipelineNode.cs index 3ee7cf25..f4dcaa16 100644 --- a/src/CodeCasa.AutomationPipelines/PipelineNode.cs +++ b/src/CodeCasa.AutomationPipelines/PipelineNode.cs @@ -140,7 +140,7 @@ protected void TurnOnPassThroughOnNextInput() private void SetOutputInternal(TState? output) { - // Scheduler-driven nodes may still fire after the pipeline was torn down; a disposed subject would throw. + // Scheduler-driven nodes may still fire after the pipeline was torn down and must not emit anymore. if (_isDisposed) { return; @@ -161,8 +161,8 @@ public virtual ValueTask DisposeAsync() return ValueTask.CompletedTask; } _isDisposed = true; + // The subject is completed but not disposed: an output that already passed the disposed check on another thread would otherwise throw. _newOutputSubject.OnCompleted(); - _newOutputSubject.Dispose(); return ValueTask.CompletedTask; } } diff --git a/tests/CodeCasa.AutomationPipelines.Lights.Tests/LightTransitionNodeTests.cs b/tests/CodeCasa.AutomationPipelines.Lights.Tests/LightTransitionNodeTests.cs new file mode 100644 index 00000000..9e4db863 --- /dev/null +++ b/tests/CodeCasa.AutomationPipelines.Lights.Tests/LightTransitionNodeTests.cs @@ -0,0 +1,128 @@ +using System.Reactive.Concurrency; +using CodeCasa.AutomationPipelines.Lights.Nodes; +using CodeCasa.Lights; +using Microsoft.Reactive.Testing; + +namespace CodeCasa.AutomationPipelines.Lights.Tests; + +[TestClass] +public sealed class LightTransitionNodeTests +{ + private static readonly TimeSpan LongTransition = TimeSpan.FromSeconds(10); + + [TestMethod] + public void ScheduleInterpolated_Twice_OnlyLastContinuationEmits() + { + var scheduler = new TestScheduler(); + var node = CreateNodeWithLongInputTransition(scheduler); + var outputs = new List(); + node.OnNewOutput.Subscribe(outputs.Add); + + node.Schedule(Parameters(50), Parameters(150)); + scheduler.AdvanceBy(TimeSpan.FromMilliseconds(200).Ticks); + node.Schedule(Parameters(40), Parameters(140)); + outputs.Clear(); + + scheduler.AdvanceBy(TimeSpan.FromSeconds(1).Ticks); + + Assert.AreEqual(1, outputs.Count); + Assert.AreEqual(140d, outputs[0]!.LightParameters.Brightness); + } + + [TestMethod] + public void ScheduleInterpolated_TwiceThenNewInput_NoStaleContinuationEmits() + { + var scheduler = new TestScheduler(); + var node = CreateNodeWithLongInputTransition(scheduler); + + node.Schedule(Parameters(50), Parameters(150)); + scheduler.AdvanceBy(TimeSpan.FromMilliseconds(200).Ticks); + node.Schedule(Parameters(40), Parameters(140)); + + var outputs = new List(); + node.OnNewOutput.Subscribe(outputs.Add); + node.Input = Parameters(10).AsTransition(); + scheduler.AdvanceBy(TimeSpan.FromSeconds(1).Ticks); + + Assert.AreEqual(0, outputs.Count); + } + + [TestMethod] + public void ScheduleInterpolated_TwiceThenPassThrough_NoStaleContinuationOverridesInput() + { + var scheduler = new TestScheduler(); + var node = CreateNodeWithLongInputTransition(scheduler); + + node.Schedule(Parameters(50), Parameters(150)); + scheduler.AdvanceBy(TimeSpan.FromMilliseconds(200).Ticks); + node.Schedule(Parameters(40), Parameters(140)); + + var outputs = new List(); + node.OnNewOutput.Subscribe(outputs.Add); + node.PassThrough = true; + scheduler.AdvanceBy(TimeSpan.FromSeconds(1).Ticks); + + Assert.IsFalse(outputs.Any(o => o?.LightParameters.Brightness is 150d or 140d)); + Assert.AreEqual(200d, node.Output!.LightParameters.Brightness); + } + + [TestMethod] + public async Task ScheduleInterpolated_AfterDispose_DoesNotEmit() + { + var scheduler = new TestScheduler(); + var node = CreateNodeWithLongInputTransition(scheduler); + node.Schedule(Parameters(50), Parameters(150)); + var outputAtDisposal = node.Output; + + await node.DisposeAsync(); + node.Schedule(Parameters(40), Parameters(140)); + scheduler.AdvanceBy(TimeSpan.FromSeconds(1).Ticks); + + Assert.AreEqual(outputAtDisposal, node.Output); + } + + [TestMethod] + public async Task OnNewOutput_SubscribeAfterDispose_Completes() + { + var node = new TestNode(new TestScheduler()); + await node.DisposeAsync(); + + var completed = false; + node.OnNewOutput.Subscribe(_ => { }, () => completed = true); + + Assert.IsTrue(completed); + } + + [TestMethod] + public void ToString_NameSet_ReturnsName() + { + var node = new TestNode(new TestScheduler()) { Name = "My Node" }; + + Assert.AreEqual("My Node", node.ToString()); + } + + [TestMethod] + public void ToString_NameNotSet_ReturnsTypeName() + { + var node = new TestNode(new TestScheduler()); + + Assert.AreEqual(nameof(TestNode), node.ToString()); + } + + private static TestNode CreateNodeWithLongInputTransition(TestScheduler scheduler) + { + var node = new TestNode(scheduler); + // Two inputs are needed so the node knows both ends of the in-flight transition it has to continue. + node.Input = Parameters(100).AsTransition(); + node.Input = Parameters(200).AsTransition(LongTransition); + return node; + } + + private static LightParameters Parameters(double brightness) => new() { Brightness = brightness, ColorTempKelvin = 3000 }; + + private sealed class TestNode(IScheduler scheduler) : LightTransitionNode(scheduler) + { + public void Schedule(LightParameters? source, LightParameters? desired) => + ScheduleInterpolatedLightTransitionUsingInputTransitionTime(source, desired); + } +} diff --git a/tests/CodeCasa.AutomationPipelines.Lights.Tests/ResettableTimeoutNodePassThroughTests.cs b/tests/CodeCasa.AutomationPipelines.Lights.Tests/ResettableTimeoutNodePassThroughTests.cs index 63bda7ed..c43e6c56 100644 --- a/tests/CodeCasa.AutomationPipelines.Lights.Tests/ResettableTimeoutNodePassThroughTests.cs +++ b/tests/CodeCasa.AutomationPipelines.Lights.Tests/ResettableTimeoutNodePassThroughTests.cs @@ -178,6 +178,40 @@ public async Task AfterTimeout_PersistAndChildOutputAreIgnored() Assert.IsFalse(_childOutputSubject.HasObservers); } + [TestMethod] + public async Task TimeoutElapsed_DisposesChildNode() + { + await using var node = CreateNode(); + + _scheduler.AdvanceBy(DefaultTimeout.Ticks + 1); + + _childNodeMock.Verify(x => x.DisposeAsync(), Times.Once); + } + + [TestMethod] + public async Task InputAfterTimeout_IsNotForwardedToChildNode() + { + await using var node = CreateNode(); + node.Input = UpstreamOutput; + _scheduler.AdvanceBy(DefaultTimeout.Ticks + 1); + var newInput = new LightParameters { Brightness = 80 }.AsTransition(); + + node.Input = newInput; + + _childNodeMock.VerifySet(x => x.Input = UpstreamOutput, Times.Once); + _childNodeMock.VerifySet(x => x.Input = newInput, Times.Never); + } + + [TestMethod] + public async Task Name_ChildNodeWithoutToStringOverride_UsesChildNodeName() + { + _childNodeMock.Setup(x => x.Name).Returns("Movie scene"); + + await using var node = CreateNode(); + + Assert.AreEqual("Movie scene (passes through after timeout)", node.Name); + } + [TestMethod] public async Task DisposeAsync_BeforeTimeout_TimerNoLongerFiresAndChildIsDisposed() { diff --git a/tests/CodeCasa.AutomationPipelines.Lights.Tests/SerializedActionQueueTests.cs b/tests/CodeCasa.AutomationPipelines.Lights.Tests/SerializedActionQueueTests.cs new file mode 100644 index 00000000..8cc92e07 --- /dev/null +++ b/tests/CodeCasa.AutomationPipelines.Lights.Tests/SerializedActionQueueTests.cs @@ -0,0 +1,82 @@ +using CodeCasa.AutomationPipelines.Lights.Utils; + +namespace CodeCasa.AutomationPipelines.Lights.Tests; + +[TestClass] +public sealed class SerializedActionQueueTests +{ + [TestMethod] + public void Run_FromWithinRunningAction_ExecutesInline() + { + var queue = new SerializedActionQueue(); + var order = new List(); + + queue.Run(() => + { + order.Add("outer start"); + queue.Run(() => order.Add("inner")); + order.Add("outer end"); + }); + + CollectionAssert.AreEqual(new[] { "outer start", "inner", "outer end" }, order); + } + + [TestMethod] + public void Run_FromOtherThreadWhileDraining_DoesNotWaitAndRunsAfterTheRunningAction() + { + var queue = new SerializedActionQueue(); + var order = new List(); + using var firstActionStarted = new ManualResetEventSlim(); + using var secondActionEnqueued = new ManualResetEventSlim(); + + var draining = Task.Run(() => queue.Run(() => + { + order.Add("first start"); + firstActionStarted.Set(); + Assert.IsTrue(secondActionEnqueued.Wait(TimeSpan.FromSeconds(30))); + order.Add("first end"); + })); + Assert.IsTrue(firstActionStarted.Wait(TimeSpan.FromSeconds(30))); + + queue.Run(() => order.Add("second")); + secondActionEnqueued.Set(); + draining.Wait(TimeSpan.FromSeconds(30)); + + CollectionAssert.AreEqual(new[] { "first start", "first end", "second" }, order); + } + + [TestMethod] + public void Run_ConcurrentCallers_ActionsNeverOverlapAndAllRun() + { + var queue = new SerializedActionQueue(); + const int iterations = 50_000; + var running = 0; + var overlapped = false; + var count = 0; + + Parallel.For(0, iterations, _ => queue.Run(() => + { + if (Interlocked.Increment(ref running) != 1) + { + overlapped = true; + } + count++; + Interlocked.Decrement(ref running); + })); + + Assert.IsFalse(overlapped); + Assert.AreEqual(iterations, count); + } + + [TestMethod] + public void Run_ActionThrows_RethrowsAndKeepsProcessingLaterActions() + { + var queue = new SerializedActionQueue(); + var laterActionRan = false; + + Assert.ThrowsExactly(() => queue.Run(() => throw new InvalidOperationException())); + queue.Run(() => laterActionRan = true); + + Assert.IsTrue(laterActionRan); + } +}