diff --git a/src/Daqifi.Core.Tests/Device/Internal/StreamingSessionSnapshotTests.cs b/src/Daqifi.Core.Tests/Device/Internal/StreamingSessionSnapshotTests.cs new file mode 100644 index 0000000..d9a5818 --- /dev/null +++ b/src/Daqifi.Core.Tests/Device/Internal/StreamingSessionSnapshotTests.cs @@ -0,0 +1,186 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using Daqifi.Core.Channel; +using Daqifi.Core.Device; +using Daqifi.Core.Device.Internal; +using Xunit; + +namespace Daqifi.Core.Tests.Device.Internal; + +/// +/// Unit tests for , which records what a streaming session +/// looked like at a drop and decides what putting it back implies (issue #379). These pin the +/// decision itself; DeviceReconnectTests pins the effects the device applies from it. +/// +public class StreamingSessionSnapshotTests +{ + private static AnalogChannel Analog(int number, bool enabled) + { + var channel = new AnalogChannel(number); + channel.IsEnabled = enabled; + return channel; + } + + private static DigitalChannel Digital(int number, bool enabled) + { + var channel = new DigitalChannel(number); + channel.IsEnabled = enabled; + return channel; + } + + private static ReconnectOptions Policy(bool resumeStreaming = true) => + new() { Enabled = true, ResumeStreaming = resumeStreaming }; + + [Fact] + public void Capture_RecordsOnlyTheEnabledChannels() + { + var channels = new IChannel[] + { + Analog(0, enabled: true), + Analog(1, enabled: false), + Digital(2, enabled: true), + }; + + var snapshot = StreamingSessionSnapshot.Capture(channels, isStreaming: true); + + Assert.Equal(2, snapshot.EnabledChannelCount); + Assert.True(snapshot.WasStreaming); + } + + [Fact] + public void Capture_TakesACopy_SoLaterChannelMutationCannotRewriteTheSession() + { + // This is the whole reason a snapshot exists: re-initialization after a reconnect goes on + // mutating these same channel objects (analog IsEnabled is resynced from the device's own + // reported mask on every status frame, #409), so a snapshot that held them by reference + // would describe the reconnected device rather than the session that was lost. + var enabled = Analog(0, enabled: true); + var disabled = Analog(1, enabled: false); + + var snapshot = StreamingSessionSnapshot.Capture(new IChannel[] { enabled, disabled }, isStreaming: true); + + enabled.IsEnabled = false; + disabled.IsEnabled = true; + + var plan = snapshot.PlanRestore(new IChannel[] { enabled, disabled }, Policy()); + + Assert.Equal(new[] { 0 }, plan.ChannelsToEnable.Select(c => c.ChannelNumber)); + } + + [Fact] + public void Capture_WithNoChannelsEnabled_PlansNothingToEnable() + { + var snapshot = StreamingSessionSnapshot.Capture( + new IChannel[] { Analog(0, enabled: false) }, isStreaming: false); + + var plan = snapshot.PlanRestore(new IChannel[] { Analog(0, enabled: false) }, Policy()); + + Assert.Equal(0, snapshot.EnabledChannelCount); + Assert.Empty(plan.ChannelsToEnable); + Assert.False(plan.ResumeStreaming); + } + + [Fact] + public void PlanRestore_ReturnsTheReconnectedDevicesOwnChannelObjects_NotTheCapturedOnes() + { + // A reconnect can replace the channel objects wholesale, so matching has to be by identity + // (type + number) and the plan has to hand back the objects the device has now — enabling + // the old ones would set a flag on garbage. + var beforeDrop = Analog(3, enabled: true); + var snapshot = StreamingSessionSnapshot.Capture(new IChannel[] { beforeDrop }, isStreaming: true); + + var afterReconnect = Analog(3, enabled: false); + var plan = snapshot.PlanRestore(new IChannel[] { afterReconnect }, Policy()); + + Assert.Same(afterReconnect, Assert.Single(plan.ChannelsToEnable)); + } + + [Fact] + public void PlanRestore_MatchesOnChannelTypeAsWellAsNumber() + { + // Analog 0 and digital 0 share a number and are different channels; a number-only match + // would enable the wrong one. + var snapshot = StreamingSessionSnapshot.Capture( + new IChannel[] { Analog(0, enabled: true), Digital(0, enabled: false) }, isStreaming: true); + + var plan = snapshot.PlanRestore( + new IChannel[] { Analog(0, enabled: false), Digital(0, enabled: false) }, Policy()); + + var restored = Assert.Single(plan.ChannelsToEnable); + Assert.Equal(ChannelType.Analog, restored.Type); + } + + [Fact] + public void PlanRestore_WhenTheDeviceComesBackSmaller_RestoresTheIntersection() + { + var snapshot = StreamingSessionSnapshot.Capture( + new IChannel[] { Analog(0, enabled: true), Analog(1, enabled: true) }, isStreaming: true); + + var plan = snapshot.PlanRestore(new IChannel[] { Analog(0, enabled: false) }, Policy()); + + Assert.Equal(new[] { 0 }, plan.ChannelsToEnable.Select(c => c.ChannelNumber)); + } + + [Fact] + public void PlanRestore_IgnoresChannelsTheSessionNeverHad() + { + var snapshot = StreamingSessionSnapshot.Capture( + new IChannel[] { Analog(0, enabled: true) }, isStreaming: false); + + var plan = snapshot.PlanRestore( + new IChannel[] { Analog(0, enabled: false), Analog(7, enabled: true) }, Policy()); + + Assert.Equal(new[] { 0 }, plan.ChannelsToEnable.Select(c => c.ChannelNumber)); + } + + [Theory] + [InlineData(true, true, true)] + [InlineData(true, false, false)] + [InlineData(false, true, false)] + [InlineData(false, false, false)] + public void PlanRestore_ResumesOnlyWhenTheSessionWasStreamingAndThePolicyAllowsIt( + bool wasStreaming, bool resumeStreaming, bool expected) + { + var snapshot = StreamingSessionSnapshot.Capture( + new IChannel[] { Analog(0, enabled: true) }, wasStreaming); + + var plan = snapshot.PlanRestore( + new IChannel[] { Analog(0, enabled: false) }, Policy(resumeStreaming)); + + Assert.Equal(expected, plan.ResumeStreaming); + } + + [Fact] + public void PlanRestore_CanBeCalledMoreThanOnce_WithoutConsumingTheSnapshot() + { + // The device keeps one snapshot field across reconnect attempts, so a failed attempt must + // not leave a snapshot that has already given up its contents. + var snapshot = StreamingSessionSnapshot.Capture( + new IChannel[] { Analog(0, enabled: true) }, isStreaming: true); + + var first = snapshot.PlanRestore(new IChannel[] { Analog(0, enabled: false) }, Policy()); + var second = snapshot.PlanRestore(new IChannel[] { Analog(0, enabled: false) }, Policy()); + + Assert.Single(first.ChannelsToEnable); + Assert.Single(second.ChannelsToEnable); + Assert.True(second.ResumeStreaming); + } + + [Fact] + public void Capture_WithNullChannels_Throws() + { + Assert.Throws( + () => StreamingSessionSnapshot.Capture(null!, isStreaming: false)); + } + + [Fact] + public void PlanRestore_WithNullArguments_Throws() + { + var snapshot = StreamingSessionSnapshot.Capture(Array.Empty(), isStreaming: false); + + Assert.Throws(() => snapshot.PlanRestore(null!, Policy())); + Assert.Throws( + () => snapshot.PlanRestore(new List(), null!)); + } +} diff --git a/src/Daqifi.Core/Device/DaqifiStreamingDevice.cs b/src/Daqifi.Core/Device/DaqifiStreamingDevice.cs index 2796188..e267c5c 100644 --- a/src/Daqifi.Core/Device/DaqifiStreamingDevice.cs +++ b/src/Daqifi.Core/Device/DaqifiStreamingDevice.cs @@ -514,42 +514,9 @@ private void ApplyAdcEnableMask(uint mask) /// private volatile StreamingSessionSnapshot? _sessionSnapshot; - /// - /// The subset of a streaming session that Core owns and can therefore put back: which - /// channels were enabled, and whether data was flowing. - /// - private sealed class StreamingSessionSnapshot - { - public StreamingSessionSnapshot(HashSet<(ChannelType Type, int Number)> enabledChannels, bool wasStreaming) - { - EnabledChannels = enabledChannels; - WasStreaming = wasStreaming; - } - - /// - /// The enabled channels, held by identity rather than by reference: a reconnect can - /// replace the channel objects, and a device that came back with a different channel - /// count should restore the intersection rather than fail. - /// - public HashSet<(ChannelType Type, int Number)> EnabledChannels { get; } - - public bool WasStreaming { get; } - } - /// - protected override void CaptureSessionSnapshot() - { - var enabled = new HashSet<(ChannelType, int)>(); - foreach (var channel in GetChannelsSnapshot()) - { - if (channel.IsEnabled) - { - enabled.Add((channel.Type, channel.ChannelNumber)); - } - } - - _sessionSnapshot = new StreamingSessionSnapshot(enabled, IsStreaming); - } + protected override void CaptureSessionSnapshot() => + _sessionSnapshot = StreamingSessionSnapshot.Capture(GetChannelsSnapshot(), IsStreaming); /// /// Re-applies the enabled-channel set recorded at the drop and, if the policy says so, @@ -557,16 +524,9 @@ protected override void CaptureSessionSnapshot() /// /// /// - /// The enable set has to be replayed from the snapshot rather than read back off the - /// channel objects: resyncs analog - /// IsEnabled from the device's own enabled mask on every status message (#409), so by - /// the time re-initialization is done the in-memory view reflects the freshly reconnected - /// device, not the session that was lost. - /// - /// - /// The streaming frequency needs no replay — it is a host-side setting that the drop never - /// touched — but it does have to reach the device again, which is what the resumed - /// does. + /// What to restore is 's decision; this + /// method owns the effects, and their order is the part that matters. See that method for + /// why the enable set is replayed from the snapshot rather than read back off the channels. /// /// /// A resumed stream is a genuinely new session: timestamp reconstruction re-anchors and the @@ -597,32 +557,24 @@ protected override Task RestoreSessionSnapshotAsync( cancellationToken.ThrowIfCancellationRequested(); // Normalize to a known state before re-applying: whatever the device came back with is - // not necessarily what it had, and the enable commands are set-replace anyway. + // not necessarily what it had, and the enable commands are set-replace anyway. The + // channel list is read afterwards so the plan is built against the post-reset objects. DisableAllChannels(); - var toEnable = new List(); - foreach (var channel in GetChannelsSnapshot()) - { - if (snapshot.EnabledChannels.Contains((channel.Type, channel.ChannelNumber))) - { - toEnable.Add(channel); - } - } - - if (toEnable.Count > 0) + var plan = snapshot.PlanRestore(GetChannelsSnapshot(), options); + if (plan.ChannelsToEnable.Count > 0) { - EnableChannels(toEnable); + EnableChannels(plan.ChannelsToEnable); } cancellationToken.ThrowIfCancellationRequested(); - var resumeStreaming = snapshot.WasStreaming && options.ResumeStreaming; - if (resumeStreaming) + if (plan.ResumeStreaming) { StartStreaming(); } - return Task.FromResult(resumeStreaming); + return Task.FromResult(plan.ResumeStreaming); } #endregion diff --git a/src/Daqifi.Core/Device/Internal/StreamingSessionSnapshot.cs b/src/Daqifi.Core/Device/Internal/StreamingSessionSnapshot.cs new file mode 100644 index 0000000..91a7d21 --- /dev/null +++ b/src/Daqifi.Core/Device/Internal/StreamingSessionSnapshot.cs @@ -0,0 +1,137 @@ +using System; +using System.Collections.Generic; +using Daqifi.Core.Channel; + +#nullable enable + +namespace Daqifi.Core.Device.Internal +{ + /// + /// What a captured session says should happen on the way back: the channels to re-enable, and + /// whether the stream itself should be restarted. + /// + /// + /// A plan is a decision, not an action. Nothing here has touched the device — the caller applies + /// it, in the order that matters, through the device's own members. + /// + internal readonly struct SessionRestorePlan + { + internal SessionRestorePlan(IReadOnlyList channelsToEnable, bool resumeStreaming) + { + ChannelsToEnable = channelsToEnable; + ResumeStreaming = resumeStreaming; + } + + /// + /// Gets the channel objects, drawn from the ones the device has now, that were + /// enabled when the connection dropped. Empty when none of them survived the reconnect. + /// + public IReadOnlyList ChannelsToEnable { get; } + + /// + /// Gets a value indicating whether the interrupted stream should be restarted — true only + /// when data really was flowing at the drop and the policy allows resuming. + /// + public bool ResumeStreaming { get; } + } + + /// + /// The subset of a streaming session that Core owns and can therefore put back: which channels + /// were enabled, and whether data was flowing (issue #379). + /// + /// + /// + /// Deciding is separated from applying, the same split + /// uses. Capturing the session and working out what restoring it implies are pure operations + /// over a channel list, pinned directly by StreamingSessionSnapshotTests; the effects + /// they imply — disabling everything first, sending the enable mask, restarting the stream — + /// stay on the device, which is the only thing that owns them. + /// + /// + /// The enabled set is held by channel identity (type and number) rather than by + /// reference, because a reconnect can replace the channel objects wholesale, and a device that + /// came back reporting a different channel count should restore the intersection rather than + /// fail. + /// + /// + internal sealed class StreamingSessionSnapshot + { + private readonly HashSet<(ChannelType Type, int Number)> _enabledChannels; + + private StreamingSessionSnapshot(HashSet<(ChannelType Type, int Number)> enabledChannels, bool wasStreaming) + { + _enabledChannels = enabledChannels; + WasStreaming = wasStreaming; + } + + /// Gets a value indicating whether data was flowing when the connection dropped. + public bool WasStreaming { get; } + + /// Gets how many channels were enabled at the drop. + public int EnabledChannelCount => _enabledChannels.Count; + + /// + /// Records the session as it stands: the identities of every enabled channel in + /// , plus . + /// + /// + /// The channel state is copied out immediately rather than held by reference, so a snapshot + /// keeps describing the instant it was taken even though the caller goes on mutating those + /// same channel objects — which is exactly what re-initialization after a reconnect does. + /// + /// The device's channels at the moment of the drop. + /// Whether the device was streaming at the moment of the drop. + public static StreamingSessionSnapshot Capture(IEnumerable channels, bool isStreaming) + { + ArgumentNullException.ThrowIfNull(channels); + + var enabled = new HashSet<(ChannelType, int)>(); + foreach (var channel in channels) + { + if (channel.IsEnabled) + { + enabled.Add((channel.Type, channel.ChannelNumber)); + } + } + + return new StreamingSessionSnapshot(enabled, isStreaming); + } + + /// + /// Works out what putting this session back means for the device as it is now. + /// + /// + /// + /// The enable set is replayed from the snapshot rather than read back off the channel + /// objects: resyncs analog + /// IsEnabled from the device's own enabled mask on every status message (#409), so by + /// the time re-initialization is done the in-memory view reflects the freshly reconnected + /// device, not the session that was lost. + /// + /// + /// The streaming frequency needs no replay — it is a host-side setting that the drop never + /// touched — but it does have to reach the device again, which is what the caller's resumed + /// does. + /// + /// + /// The channels the reconnected device is reporting. + /// The reconnect policy; gates the restart. + /// The channels to re-enable and whether to restart the stream. + public SessionRestorePlan PlanRestore(IEnumerable currentChannels, ReconnectOptions options) + { + ArgumentNullException.ThrowIfNull(currentChannels); + ArgumentNullException.ThrowIfNull(options); + + var toEnable = new List(); + foreach (var channel in currentChannels) + { + if (_enabledChannels.Contains((channel.Type, channel.ChannelNumber))) + { + toEnable.Add(channel); + } + } + + return new SessionRestorePlan(toEnable, WasStreaming && options.ResumeStreaming); + } + } +}