From e25311961894acef0a482c234ebdfbb272882560 Mon Sep 17 00:00:00 2001 From: Justin Puah Date: Thu, 23 Jul 2026 23:05:08 +1000 Subject: [PATCH 1/2] =?UTF-8?q?feat(tui):=20H3=20=E2=80=94=20SingleFlightC?= =?UTF-8?q?ache=20and=20Published,=20typed=20ADR=200008=20invariants=20(#4?= =?UTF-8?q?6)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- src/Cobalt.Tui/Tasks/Published.cs | 21 +++ src/Cobalt.Tui/Tasks/SingleFlightCache.cs | 104 +++++++++++ .../Cobalt.Tui.Tests/Tasks/PublishedTests.cs | 22 +++ .../Tasks/SingleFlightCacheTests.cs | 176 ++++++++++++++++++ 4 files changed, 323 insertions(+) create mode 100644 src/Cobalt.Tui/Tasks/Published.cs create mode 100644 src/Cobalt.Tui/Tasks/SingleFlightCache.cs create mode 100644 tests/Cobalt.Tui.Tests/Tasks/PublishedTests.cs create mode 100644 tests/Cobalt.Tui.Tests/Tasks/SingleFlightCacheTests.cs diff --git a/src/Cobalt.Tui/Tasks/Published.cs b/src/Cobalt.Tui/Tasks/Published.cs new file mode 100644 index 0000000..2c9b436 --- /dev/null +++ b/src/Cobalt.Tui/Tasks/Published.cs @@ -0,0 +1,21 @@ +namespace Cobalt.Tui.Tasks; + +/// +/// An atomic publish/read pair for view-model state: one volatile reference holding one +/// immutable . Generalizes the PrDiffViewModel.DiffState +/// pattern (ADR 0008): related values that must never tear live together inside a single +/// sealed record, published by a single reference write and read back by a single +/// reference read. Consumers snapshot once and destructure — reading +/// two properties through two separate reads would re-introduce the tear on the read side +/// (the ADR 0008 amendment behind CurrentDiffSnapshot). +/// +public sealed class Published where T : class +{ + private T? _current; + + /// The last published value, from a single volatile read; null before the first publish. + public T? Current => Volatile.Read(ref _current); + + /// Replaces the published value in a single volatile write; null clears it. + public void Publish(T? value) => Volatile.Write(ref _current, value); +} diff --git a/src/Cobalt.Tui/Tasks/SingleFlightCache.cs b/src/Cobalt.Tui/Tasks/SingleFlightCache.cs new file mode 100644 index 0000000..11dbac2 --- /dev/null +++ b/src/Cobalt.Tui/Tasks/SingleFlightCache.cs @@ -0,0 +1,104 @@ +namespace Cobalt.Tui.Tasks; + +/// +/// The typed SUPERSEDE primitive of ADR 0008: at most one logical in-flight fetch, where +/// scheduling a new key cancels and abandons the previous one and only the newest may publish. +/// Generalizes the _loadSeq stamp guard proven in PrListViewModel.LoadTabAsync — +/// cancellation is cooperative, so a fetch may ignore its token and still complete, and the +/// stamp comparison (not the cancel) is what keeps its stale result from landing. This is NOT +/// the diff dialog's join/dedup cache (the ConcurrentDictionary+Lazy single-flight +/// in ADR 0008 §"Single-flight diff fetches") — that shares one fetch between converging +/// callers; this one abandons the old fetch entirely. +/// +public sealed class SingleFlightCache : IDisposable where TKey : notnull +{ + private readonly object _gate = new(); + private readonly CancellationToken _lifetime; + private CancellationTokenSource? _cts; + private long _stamp; + private bool _disposed; + + /// Every per-schedule token source is created linked to , + /// so the owner's shutdown cancels whatever fetch is in flight. + public SingleFlightCache(CancellationToken lifetime) => _lifetime = lifetime; + + /// The monotonic supersede stamp, incremented once per schedule. + public long Stamp + { + get { lock (_gate) { return _stamp; } } + } + + /// + /// Supersedes any in-flight fetch (cancelling its token) and runs ; + /// is invoked only if this schedule is still the newest when the + /// fetch completes, checked atomically with any concurrent schedule's stamp increment. The + /// returned task completes successfully when the fetch published, was cancelled, or was + /// superseded — cancellation is swallowed ( + /// semantics). A fault from the fetch that is still the newest propagates; a fault from a + /// superseded fetch is observed and swallowed, so an abandoned fetch nobody awaits can never + /// reach the crash-log hook (ADR 0013; + /// precedent: PrDiffViewModel.ObserveFault). + /// + public async Task ScheduleAsync( + TKey key, Func> fetch, Action publish) + { + long stamp; + CancellationToken token; + lock (_gate) + { + ObjectDisposedException.ThrowIf(_disposed, this); + stamp = ++_stamp; + _cts?.Cancel(); + _cts?.Dispose(); + _cts = CancellationTokenSource.CreateLinkedTokenSource(_lifetime); + token = _cts.Token; + } + + TValue value; + try + { + value = await fetch(key, token).ConfigureAwait(false); + } + catch (OperationCanceledException) + { + return; // a superseded or shut-down fetch is never an error (IgnoreCancellationAsync semantics) + } + catch + { + lock (_gate) + { + if (stamp == _stamp) + { + throw; // still the newest — the caller owns surfacing this fault + } + } + return; // superseded — the fault is observed here (caught) and goes no further + } + + lock (_gate) + { + if (stamp != _stamp) + { + return; // superseded while mid-flight — drop the result, even if the fetch ignored its token + } + publish(key, value); + } + } + + /// Cancels and disposes the current token source; idempotent. A schedule after + /// dispose throws . + public void Dispose() + { + lock (_gate) + { + if (_disposed) + { + return; + } + _disposed = true; + _cts?.Cancel(); + _cts?.Dispose(); + _cts = null; + } + } +} diff --git a/tests/Cobalt.Tui.Tests/Tasks/PublishedTests.cs b/tests/Cobalt.Tui.Tests/Tasks/PublishedTests.cs new file mode 100644 index 0000000..ef7a4ec --- /dev/null +++ b/tests/Cobalt.Tui.Tests/Tasks/PublishedTests.cs @@ -0,0 +1,22 @@ +using Cobalt.Tui.Tasks; + +namespace Cobalt.Tui.Tests.Tasks; + +public class PublishedTests +{ + private sealed record Snapshot(string Key, int Value); + + [Fact] + public void Current_Is_Null_Initially_Publish_Sets_And_Null_Clears() + { + var published = new Published(); + Assert.Null(published.Current); + + var snapshot = new Snapshot("a", 1); + published.Publish(snapshot); + Assert.Same(snapshot, published.Current); + + published.Publish(null); + Assert.Null(published.Current); + } +} diff --git a/tests/Cobalt.Tui.Tests/Tasks/SingleFlightCacheTests.cs b/tests/Cobalt.Tui.Tests/Tasks/SingleFlightCacheTests.cs new file mode 100644 index 0000000..2a3dc37 --- /dev/null +++ b/tests/Cobalt.Tui.Tests/Tasks/SingleFlightCacheTests.cs @@ -0,0 +1,176 @@ +using Cobalt.Tui.Tasks; + +namespace Cobalt.Tui.Tests.Tasks; + +public class SingleFlightCacheTests +{ + [Fact] + public async Task Single_Schedule_Publishes_Once_With_The_Scheduled_Key() + { + using var cache = new SingleFlightCache(TestContext.Current.CancellationToken); + var published = new List<(string Key, int Value)>(); + + await cache.ScheduleAsync( + "a", + (key, _) => Task.FromResult(42), + (key, value) => published.Add((key, value))); + + Assert.Equal([("a", 42)], published); + } + + [Fact] + public async Task Rapid_Reschedule_Cancels_Previous_Fetches_And_Only_The_Newest_Publishes() + { + using var cache = new SingleFlightCache(TestContext.Current.CancellationToken); + var published = new List<(string Key, int Value)>(); + var tokens = new Dictionary(); + + // A and B are held open on TCSes that honour their token; C completes immediately. + var taskA = cache.ScheduleAsync("a", (key, ct) => Hold(key, ct, tokens), Publish(published)); + var taskB = cache.ScheduleAsync("b", (key, ct) => Hold(key, ct, tokens), Publish(published)); + var taskC = cache.ScheduleAsync("c", (key, ct) => Task.FromResult(3), Publish(published)); + + await Task.WhenAll(taskA, taskB, taskC); + + Assert.Equal([("c", 3)], published); + Assert.True(tokens["a"].IsCancellationRequested); + Assert.True(tokens["b"].IsCancellationRequested); + } + + [Fact] + public async Task Stamp_Guard_Drops_A_Superseded_Fetch_That_Ignores_Its_Token_And_Completes() + { + using var cache = new SingleFlightCache(TestContext.Current.CancellationToken); + var published = new List<(string Key, int Value)>(); + var b = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + + // B's fake IGNORES its token: cancellation alone cannot stop it from completing. + var taskB = cache.ScheduleAsync("b", (_, _) => b.Task, Publish(published)); + var taskC = cache.ScheduleAsync("c", (_, _) => Task.FromResult(3), Publish(published)); + await taskC; + + // B completes successfully AFTER C was scheduled — only the stamp guard can drop it. + b.SetResult(2); + await taskB; + + Assert.Equal([("c", 3)], published); + } + + [Fact] + public async Task Superseded_Fetch_That_Faults_Is_Consumed_Without_Throwing() + { + using var cache = new SingleFlightCache(TestContext.Current.CancellationToken); + var published = new List<(string Key, int Value)>(); + var b = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + + var taskB = cache.ScheduleAsync("b", (_, _) => b.Task, Publish(published)); + var taskC = cache.ScheduleAsync("c", (_, _) => Task.FromResult(3), Publish(published)); + await taskC; + + // B faults AFTER being superseded: nobody is waiting on the abandoned fetch, so its + // fault must be observed and swallowed rather than escaping (ADR 0013). + b.SetException(new InvalidOperationException("boom")); + await taskB; + + Assert.Equal([("c", 3)], published); + } + + [Fact] + public async Task Newest_Fetch_That_Faults_Propagates_Out_Of_Its_Returned_Task() + { + using var cache = new SingleFlightCache(TestContext.Current.CancellationToken); + + var task = cache.ScheduleAsync( + "a", + (_, _) => Task.FromException(new InvalidOperationException("boom")), + (_, _) => { }); + + var ex = await Assert.ThrowsAsync(() => task); + Assert.Equal("boom", ex.Message); + } + + [Fact] + public async Task Lifetime_Cancellation_Completes_The_Newest_Fetch_Without_Throwing_And_Without_Publishing() + { + using var lifetime = CancellationTokenSource.CreateLinkedTokenSource(TestContext.Current.CancellationToken); + using var cache = new SingleFlightCache(lifetime.Token); + var published = new List<(string Key, int Value)>(); + var tokens = new Dictionary(); + + var task = cache.ScheduleAsync("a", (key, ct) => Hold(key, ct, tokens), Publish(published)); + await lifetime.CancelAsync(); + await task; + + Assert.Empty(published); + } + + [Fact] + public async Task Stamp_Increments_Once_Per_Schedule() + { + using var cache = new SingleFlightCache(TestContext.Current.CancellationToken); + var before = cache.Stamp; + + await cache.ScheduleAsync("a", (_, _) => Task.FromResult(1), (_, _) => { }); + Assert.Equal(before + 1, cache.Stamp); + + await cache.ScheduleAsync("b", (_, _) => Task.FromResult(2), (_, _) => { }); + Assert.Equal(before + 2, cache.Stamp); + } + + [Fact] + public async Task Dispose_Cancels_The_InFlight_Fetch_And_A_Later_Schedule_Throws() + { + var cache = new SingleFlightCache(TestContext.Current.CancellationToken); + var published = new List<(string Key, int Value)>(); + var tokens = new Dictionary(); + + var task = cache.ScheduleAsync("a", (key, ct) => Hold(key, ct, tokens), Publish(published)); + cache.Dispose(); + await task; + + Assert.True(tokens["a"].IsCancellationRequested); + Assert.Empty(published); + await Assert.ThrowsAsync( + () => cache.ScheduleAsync("b", (_, _) => Task.FromResult(2), Publish(published))); + + cache.Dispose(); // double-dispose is a no-op + } + + private sealed record Snapshot(string Key, int Value); + + [Fact] + public async Task Composition_Rapid_A_B_C_Leaves_Only_C_In_The_Published_State() + { + // The ticket's acceptance (#46): a fake source of controllable latency, moving A->B->C + // fast, publishing snapshots into a Published — only C lands, A/B are cancelled, and + // no fault goes unobserved (every returned task completes cleanly). + using var cache = new SingleFlightCache(TestContext.Current.CancellationToken); + var state = new Published(); + var tokens = new Dictionary(); + void PublishSnapshot(string key, int value) => state.Publish(new Snapshot(key, value)); + + var taskA = cache.ScheduleAsync("a", (key, ct) => Hold(key, ct, tokens), PublishSnapshot); + var taskB = cache.ScheduleAsync("b", (key, ct) => Hold(key, ct, tokens), PublishSnapshot); + var taskC = cache.ScheduleAsync("c", (_, _) => Task.FromResult(3), PublishSnapshot); + + await Task.WhenAll(taskA, taskB, taskC); + + var current = state.Current; + Assert.NotNull(current); + Assert.Equal(new Snapshot("c", 3), current); + Assert.True(tokens["a"].IsCancellationRequested); + Assert.True(tokens["b"].IsCancellationRequested); + } + + /// A fetch held open forever that honours its token — completes only by cancellation. + private static Task Hold(string key, CancellationToken ct, Dictionary tokens) + { + tokens[key] = ct; + var tcs = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + ct.Register(() => tcs.TrySetCanceled(ct)); + return tcs.Task; + } + + private static Action Publish(List<(string Key, int Value)> published) => + (key, value) => published.Add((key, value)); +} From fcf03e20c7ee147c771ff682ae32d98211b2ae3f Mon Sep 17 00:00:00 2001 From: Justin Puah Date: Thu, 23 Jul 2026 23:21:13 +1000 Subject: [PATCH 2/2] =?UTF-8?q?fix(tui):=20SingleFlightCache=20=E2=80=94?= =?UTF-8?q?=20dispose=20supersedes,=20CTS=20cancel=20moved=20off=20the=20l?= =?UTF-8?q?ock?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Dispose now acts as the final supersede: both post-fetch guards check _disposed, so a token-ignoring fetch that completes (or faults) after Dispose() can neither publish nor throw out of its returned task. Two new TCS-driven tests pin both paths (each confirmed failing first). CTS Cancel()/Dispose() now run outside _gate in ScheduleAsync and Dispose: Cancel() executes token registrations synchronously, and running arbitrary callbacks under the lock is a lock-inversion hazard for any consumer whose registration takes another lock. Under the lock we only detach the old CTS and bump the stamp/_disposed, so each detached CTS is cancelled exactly once and correctness still rests on the stamp, not the cancel. No new test for this: a deterministic RED for a cross-thread lock inversion would be a hanging test, worse than none — the existing cancellation tests pin the observable behaviour. Also documents the publish-callback constraint: it runs under the cache's internal lock, so it must stay a cheap reference swap. --- src/Cobalt.Tui/Tasks/SingleFlightCache.cs | 30 +++++++++++----- .../Tasks/SingleFlightCacheTests.cs | 34 +++++++++++++++++++ 2 files changed, 56 insertions(+), 8 deletions(-) diff --git a/src/Cobalt.Tui/Tasks/SingleFlightCache.cs b/src/Cobalt.Tui/Tasks/SingleFlightCache.cs index 11dbac2..259c4b3 100644 --- a/src/Cobalt.Tui/Tasks/SingleFlightCache.cs +++ b/src/Cobalt.Tui/Tasks/SingleFlightCache.cs @@ -38,21 +38,31 @@ public long Stamp /// superseded fetch is observed and swallowed, so an abandoned fetch nobody awaits can never /// reach the crash-log hook (ADR 0013; /// precedent: PrDiffViewModel.ObserveFault). + /// runs under the cache's internal lock — that lock is the + /// atomicity guarantee against a concurrent schedule — so it must stay a cheap reference swap + /// (e.g. Published<T>.Publish): never raise events, block, or take other locks. /// public async Task ScheduleAsync( TKey key, Func> fetch, Action publish) { long stamp; CancellationToken token; + CancellationTokenSource? previous; lock (_gate) { ObjectDisposedException.ThrowIf(_disposed, this); stamp = ++_stamp; - _cts?.Cancel(); - _cts?.Dispose(); + previous = _cts; _cts = CancellationTokenSource.CreateLinkedTokenSource(_lifetime); token = _cts.Token; } + // Cancel outside the lock: Cancel() runs token registrations synchronously, and running + // arbitrary callbacks while holding _gate is a lock-inversion hazard. The detached CTS is + // cancelled exactly once, by whoever detached it. Correctness never rested on the cancel + // anyway — the stamp was already bumped under the lock, so the old fetch is superseded + // before its token even trips. + previous?.Cancel(); + previous?.Dispose(); TValue value; try @@ -67,19 +77,19 @@ public async Task ScheduleAsync( { lock (_gate) { - if (stamp == _stamp) + if (!_disposed && stamp == _stamp) { throw; // still the newest — the caller owns surfacing this fault } } - return; // superseded — the fault is observed here (caught) and goes no further + return; // superseded (or disposed, the final supersede) — the fault is observed here and goes no further } lock (_gate) { - if (stamp != _stamp) + if (_disposed || stamp != _stamp) { - return; // superseded while mid-flight — drop the result, even if the fetch ignored its token + return; // superseded or disposed mid-flight — drop the result, even if the fetch ignored its token } publish(key, value); } @@ -89,6 +99,7 @@ public async Task ScheduleAsync( /// dispose throws . public void Dispose() { + CancellationTokenSource? current; lock (_gate) { if (_disposed) @@ -96,9 +107,12 @@ public void Dispose() return; } _disposed = true; - _cts?.Cancel(); - _cts?.Dispose(); + current = _cts; _cts = null; } + // Outside the lock for the same lock-inversion reason as in ScheduleAsync; _disposed is + // already set, so the in-flight fetch is dropped at completion regardless of this cancel. + current?.Cancel(); + current?.Dispose(); } } diff --git a/tests/Cobalt.Tui.Tests/Tasks/SingleFlightCacheTests.cs b/tests/Cobalt.Tui.Tests/Tasks/SingleFlightCacheTests.cs index 2a3dc37..a3a5d98 100644 --- a/tests/Cobalt.Tui.Tests/Tasks/SingleFlightCacheTests.cs +++ b/tests/Cobalt.Tui.Tests/Tasks/SingleFlightCacheTests.cs @@ -136,6 +136,40 @@ await Assert.ThrowsAsync( cache.Dispose(); // double-dispose is a no-op } + [Fact] + public async Task Dispose_Supersedes_A_Token_Ignoring_Fetch_So_Its_Late_Result_Never_Publishes() + { + var cache = new SingleFlightCache(TestContext.Current.CancellationToken); + var published = new List<(string Key, int Value)>(); + var a = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + + // The fake IGNORES its token: only the disposed-guard can stop the late publish. + var task = cache.ScheduleAsync("a", (_, _) => a.Task, Publish(published)); + cache.Dispose(); + a.SetResult(1); + await task; + + Assert.Empty(published); + } + + [Fact] + public async Task Dispose_Supersedes_A_Token_Ignoring_Fetch_So_Its_Late_Fault_Is_Consumed() + { + var cache = new SingleFlightCache(TestContext.Current.CancellationToken); + var published = new List<(string Key, int Value)>(); + var a = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + + var task = cache.ScheduleAsync("a", (_, _) => a.Task, Publish(published)); + cache.Dispose(); + + // The fetch faults AFTER dispose: nobody is listening any more, so the fault must be + // observed and swallowed (dispose is the final supersede), not thrown to the awaiter. + a.SetException(new InvalidOperationException("boom")); + await task; + + Assert.Empty(published); + } + private sealed record Snapshot(string Key, int Value); [Fact]