diff --git a/docs/ACTIVITY_DOTS.md b/docs/ACTIVITY_DOTS.md index 8771b6e..29bcad5 100644 --- a/docs/ACTIVITY_DOTS.md +++ b/docs/ACTIVITY_DOTS.md @@ -14,8 +14,9 @@ state in memory. - A completion animation after the final active turn finishes The widget accepts only the lifecycle event type and the session and turn identifiers -provided by Codex. It does not collect prompts, responses, transcript contents, transcript -paths, or model output. +provided by Codex hooks. While activity is present, it also checks turn completion metadata +through the official local app-server protocol, with turn items omitted. It does not collect +prompts, responses, transcript contents, transcript paths, or model output. ## Set up in the widget @@ -73,10 +74,21 @@ Each session owns at most one active turn. A later `UserPromptSubmit` replaces a turn in that session, and a late `Stop` for the old turn cannot clear the new one. Duplicate events are harmless. `SessionEnd` removes only the matching session. -If Codex terminates without sending a final lifecycle event, a later turn in the same session -replaces the stale turn. Restarting the widget also clears all in-memory activity state. The -widget does not use an arbitrary timeout because legitimate Codex tasks can run for a long -time. +While activity is present, the widget checks every 15 seconds whether each tracked turn has +finished. It uses `thread/turns/list` with `itemsView: "notLoaded"`, follows pagination, and +clears a turn only when its exact identifier has a terminal status and an explicit completion +timestamp. A late result cannot clear a newer turn in the same session. No checks are sent +while idle, and each turn's check has a five-second request timeout. + +This requires a Codex CLI that supports `thread/turns/list` and completion timestamps, verified +with CLI 0.154.0. Unsupported requests, unavailable history, missing turns, and failed checks +leave hook activity unchanged. In particular, a separate app-server can reconstruct a running +turn as `interrupted` without a completion timestamp; that is not evidence of completion. + +If Codex terminates without recording completion or sending a final lifecycle event, a later +turn in the same session replaces the stale turn. Restarting the widget also clears all +in-memory activity state. The widget never expires activity solely because a task has run +for a long time. If the widget is closed, the hook bridge exits successfully after a short connection attempt and Codex continues normally. diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index c4eb360..82b6494 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -87,8 +87,12 @@ tests/CodexUsageWidget.Tests/ Unit tests for parsing, formatting and persistence Codex remains the owner of hook trust; the widget only reads trust state and opens the interactive CLI for the user's explicit `/hooks` approval. - Activity state is not persisted or reconstructed with private transcript/database polling. - A later turn in the same session recovers missing cleanup; a hard Codex termination with no - later lifecycle event is cleared by restarting the widget. + While active, the monitor checks tracked turns every 15 seconds through + `CodexTurnCompletionReader`, using official `thread/turns/list` metadata with items omitted. + Only an exact turn match with a terminal status and completion timestamp clears activity; + missing or unavailable evidence preserves it. Each check is bounded to five seconds and + shutdown cancels pending checks before disposing the shared app-server session. A hard + termination with no recorded completion still requires a later lifecycle event or restart. - Unhandled exceptions and CLI diagnostics are recorded locally for support. - Publish trimming is disabled because WPF is not a safe trimming boundary. diff --git a/src/CodexUsageWidget/App.xaml.cs b/src/CodexUsageWidget/App.xaml.cs index dd66e38..6497db2 100644 --- a/src/CodexUsageWidget/App.xaml.cs +++ b/src/CodexUsageWidget/App.xaml.cs @@ -72,7 +72,9 @@ protected override void OnStartup(StartupEventArgs e) usageMonitor.DiagnosticMessage += (_, message) => _logger.Info(message); var resetUseCase = new RateLimitResetUseCase(resetConsumer, usageMonitor); - activityMonitor = new CodexActivityMonitor(new CodexActivityPipeSignalSource()); + activityMonitor = new CodexActivityMonitor( + new CodexActivityPipeSignalSource(), + new CodexTurnCompletionReader(appServerSession)); activityMonitor.DiagnosticMessage += (_, message) => _logger.Info(message); var processPath = Environment.ProcessPath ?? throw new InvalidOperationException("Cannot determine the widget executable path."); diff --git a/src/CodexUsageWidget/Application/CodexActivityMonitor.cs b/src/CodexUsageWidget/Application/CodexActivityMonitor.cs index c5edbb7..21da734 100644 --- a/src/CodexUsageWidget/Application/CodexActivityMonitor.cs +++ b/src/CodexUsageWidget/Application/CodexActivityMonitor.cs @@ -5,13 +5,92 @@ public sealed class CodexActivityMonitor : IAsyncDisposable private readonly object _stateLock = new(); private readonly object _transitionLock = new(); private readonly ICodexActivitySignalSource _source; + private readonly ICodexTurnCompletionReader? _completionReader; + private readonly TimeSpan _reconciliationInterval; + private readonly TimeSpan _requestTimeout; + private readonly CancellationTokenSource _lifetime = new(); + private readonly SemaphoreSlim _reconciliationGate = new(1, 1); + private Task? _reconciliationTask; + private int _disposed; + private bool _completionCheckFailed; private readonly Dictionary _activeTurnsBySession = new(StringComparer.Ordinal); private bool _started; - public CodexActivityMonitor(ICodexActivitySignalSource source) + public CodexActivityMonitor( + ICodexActivitySignalSource source, + ICodexTurnCompletionReader? completionReader = null, + TimeSpan? reconciliationInterval = null, + TimeSpan? requestTimeout = null) { _source = source; + _completionReader = completionReader; + _reconciliationInterval = reconciliationInterval ?? TimeSpan.FromSeconds(15); + _requestTimeout = requestTimeout ?? TimeSpan.FromSeconds(5); + ArgumentOutOfRangeException.ThrowIfLessThanOrEqual(_reconciliationInterval, TimeSpan.Zero); + ArgumentOutOfRangeException.ThrowIfLessThanOrEqual(_requestTimeout, TimeSpan.Zero); + } + + public async Task ReconcileAsync(CancellationToken cancellationToken = default) + { + ObjectDisposedException.ThrowIf(_disposed != 0, this); + if (_completionReader is null || + !await _reconciliationGate.WaitAsync(0, cancellationToken).ConfigureAwait(false)) + { + return; + } + + try + { + await ReconcileWithGateHeldAsync(cancellationToken).ConfigureAwait(false); + } + finally + { + _reconciliationGate.Release(); + } + } + + private async Task ReconcileWithGateHeldAsync(CancellationToken cancellationToken) + { + KeyValuePair[] turns; + lock (_stateLock) + { + turns = _activeTurnsBySession.ToArray(); + } + + var checkFailed = false; + foreach (var turn in turns) + { + try + { + using var timeout = CancellationTokenSource.CreateLinkedTokenSource( + cancellationToken, _lifetime.Token); + timeout.CancelAfter(_requestTimeout); + if (await _completionReader!.IsCompletedAsync(turn.Key, turn.Value, timeout.Token) + .ConfigureAwait(false)) + { + timeout.Token.ThrowIfCancellationRequested(); + SourceOnSignalReceived(new(CodexActivitySignalKind.TurnStopped, turn.Key, turn.Value)); + } + } + catch (OperationCanceledException) when ( + cancellationToken.IsCancellationRequested || _lifetime.IsCancellationRequested) + { + throw; + } + catch (Exception) + { + checkFailed = true; + } + } + + if (checkFailed && !_completionCheckFailed) + { + DiagnosticMessage?.Invoke(this, + "Codex activity completion check unavailable; keeping hook activity until completion is confirmed."); + } + + _completionCheckFailed = checkFailed; } public event Action? ActivityChanged; @@ -31,6 +110,7 @@ public bool IsActive public async Task StartAsync(CancellationToken cancellationToken = default) { + ObjectDisposedException.ThrowIf(_disposed != 0, this); if (_started) { return; @@ -39,6 +119,30 @@ public async Task StartAsync(CancellationToken cancellationToken = default) _started = true; _source.SignalReceived += SourceOnSignalReceived; await _source.StartAsync(cancellationToken).ConfigureAwait(false); + if (_completionReader is not null) + { + _reconciliationTask = RunReconciliationAsync(_lifetime.Token); + } + } + + private async Task RunReconciliationAsync(CancellationToken cancellationToken) + { + using var timer = new PeriodicTimer(_reconciliationInterval); + try + { + while (await timer.WaitForNextTickAsync(cancellationToken).ConfigureAwait(false)) + { + await ReconcileAsync(cancellationToken).ConfigureAwait(false); + } + } + catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested) + { + } + // Disposal marks the monitor before cancelling the timer. A tick in that + // interval is normal shutdown, even if cancellation is not visible yet. + catch (ObjectDisposedException) when (Volatile.Read(ref _disposed) != 0) + { + } } private void SourceOnSignalReceived(CodexActivitySignal signal) @@ -95,7 +199,22 @@ private void SourceOnSignalReceived(CodexActivitySignal signal) public async ValueTask DisposeAsync() { + if (Interlocked.Exchange(ref _disposed, 1) != 0) + { + return; + } + _source.SignalReceived -= SourceOnSignalReceived; + await _lifetime.CancelAsync().ConfigureAwait(false); + if (_reconciliationTask is not null) + { + await _reconciliationTask.ConfigureAwait(false); + } + + await _reconciliationGate.WaitAsync().ConfigureAwait(false); + _reconciliationGate.Release(); await _source.DisposeAsync().ConfigureAwait(false); + _lifetime.Dispose(); + _reconciliationGate.Dispose(); } } diff --git a/src/CodexUsageWidget/Application/ICodexTurnCompletionReader.cs b/src/CodexUsageWidget/Application/ICodexTurnCompletionReader.cs new file mode 100644 index 0000000..b3a2774 --- /dev/null +++ b/src/CodexUsageWidget/Application/ICodexTurnCompletionReader.cs @@ -0,0 +1,6 @@ +namespace CodexUsageWidget.Application; + +public interface ICodexTurnCompletionReader +{ + Task IsCompletedAsync(string sessionId, string turnId, CancellationToken cancellationToken); +} diff --git a/src/CodexUsageWidget/Infrastructure/Codex/CodexAppServerSession.cs b/src/CodexUsageWidget/Infrastructure/Codex/CodexAppServerSession.cs index 0d7df46..af8beeb 100644 --- a/src/CodexUsageWidget/Infrastructure/Codex/CodexAppServerSession.cs +++ b/src/CodexUsageWidget/Infrastructure/Codex/CodexAppServerSession.cs @@ -88,6 +88,7 @@ await connection.RequestAsync( }, capabilities = new { + experimentalApi = true, optOutNotificationMethods = Array.Empty() } }, diff --git a/src/CodexUsageWidget/Infrastructure/Codex/CodexTurnCompletionReader.cs b/src/CodexUsageWidget/Infrastructure/Codex/CodexTurnCompletionReader.cs new file mode 100644 index 0000000..a1a018b --- /dev/null +++ b/src/CodexUsageWidget/Infrastructure/Codex/CodexTurnCompletionReader.cs @@ -0,0 +1,51 @@ +using System.Text.Json; +using CodexUsageWidget.Application; + +namespace CodexUsageWidget.Infrastructure.Codex; + +public sealed class CodexTurnCompletionReader(ICodexAppServerSession session) : ICodexTurnCompletionReader +{ + public async Task IsCompletedAsync( + string sessionId, string turnId, CancellationToken cancellationToken) + { + string? cursor = null; + var seenCursors = new HashSet(StringComparer.Ordinal); + do + { + cancellationToken.ThrowIfCancellationRequested(); + var result = await session.RequestAsync( + "thread/turns/list", + new { threadId = sessionId, cursor, limit = 100, sortDirection = "desc", itemsView = "notLoaded" }, + cancellationToken).ConfigureAwait(false); + if (result.ValueKind != JsonValueKind.Object || + !result.TryGetProperty("data", out var turns) || turns.ValueKind != JsonValueKind.Array) + { + return false; + } + + foreach (var turn in turns.EnumerateArray()) + { + if (turn.ValueKind == JsonValueKind.Object && + turn.TryGetProperty("id", out var id) && id.ValueKind == JsonValueKind.String && + string.Equals(id.GetString(), turnId, StringComparison.Ordinal)) + { + // A foreign running turn can be reconstructed as interrupted without an end time. + // Never infer completion from that status alone, or from an absent turn. + return turn.TryGetProperty("status", out var status) && + status.ValueKind == JsonValueKind.String && + status.GetString() is "completed" or "interrupted" or "failed" && + turn.TryGetProperty("completedAt", out var completedAt) && + completedAt.ValueKind == JsonValueKind.Number && + completedAt.TryGetInt64(out var timestamp) && timestamp > 0; + } + } + + cursor = result.TryGetProperty("nextCursor", out var next) && next.ValueKind == JsonValueKind.String + ? next.GetString() + : null; + } + while (!string.IsNullOrEmpty(cursor) && seenCursors.Add(cursor)); + + return false; + } +} diff --git a/tests/CodexUsageWidget.Tests/CodexActivityMonitorTests.cs b/tests/CodexUsageWidget.Tests/CodexActivityMonitorTests.cs index 9a42421..9602820 100644 --- a/tests/CodexUsageWidget.Tests/CodexActivityMonitorTests.cs +++ b/tests/CodexUsageWidget.Tests/CodexActivityMonitorTests.cs @@ -4,6 +4,227 @@ namespace CodexUsageWidget.Tests; public sealed class CodexActivityMonitorTests { + [Fact] + public async Task DisposalDuringTimerTickStillDisposesSignalSource() + { + var source = new FakeSignalSource(); + var checkedTurn = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var reader = new FakeCompletionReader((_, _, _) => + { + checkedTurn.TrySetResult(); + return Task.FromResult(false); + }); + var monitor = new CodexActivityMonitor(source, reader, TimeSpan.FromMilliseconds(1)); + await monitor.StartAsync(); + source.Publish(new(CodexActivitySignalKind.TurnStarted, "session", "turn")); + await checkedTurn.Task.WaitAsync(TimeSpan.FromSeconds(2)); + + // Hold unsubscription open so timer ticks overlap shutdown before source cleanup. + source.OnUnsubscribe = () => Task.Delay(200).GetAwaiter().GetResult(); + await monitor.DisposeAsync().AsTask().WaitAsync(TimeSpan.FromSeconds(2)); + + Assert.True(source.IsDisposed); + } + + [Fact] + public async Task CompletionOfOneSessionDoesNotHideAnotherUnconfirmedSession() + { + await using var source = new FakeSignalSource(); + var reader = new FakeCompletionReader((session, _, _) => Task.FromResult(session == "finished")); + await using var monitor = new CodexActivityMonitor(source, reader); + await monitor.StartAsync(); + source.Publish(new(CodexActivitySignalKind.TurnStarted, "finished", "turn-a")); + source.Publish(new(CodexActivitySignalKind.TurnStarted, "running", "turn-b")); + + await monitor.ReconcileAsync(); + await monitor.ReconcileAsync(); + Assert.True(monitor.IsActive); + source.Publish(new(CodexActivitySignalKind.TurnStopped, "running", "turn-b")); + Assert.False(monitor.IsActive); + } + + [Fact] + public async Task LateCompletionCheckCannotClearANewerTurn() + { + await using var source = new FakeSignalSource(); + var result = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var reader = new FakeCompletionReader((_, _, _) => result.Task); + await using var monitor = new CodexActivityMonitor(source, reader); + await monitor.StartAsync(); + source.Publish(new(CodexActivitySignalKind.TurnStarted, "session", "old")); + var reconciliation = monitor.ReconcileAsync(); + source.Publish(new(CodexActivitySignalKind.TurnStarted, "session", "new")); + result.SetResult(true); + + await reconciliation; + Assert.True(monitor.IsActive); + source.Publish(new(CodexActivitySignalKind.TurnStopped, "session", "new")); + Assert.False(monitor.IsActive); + } + + [Fact] + public async Task TimedOutCheckKeepsActivityAndRetriesLater() + { + await using var source = new FakeSignalSource(); + var attempts = 0; + var reader = new FakeCompletionReader(async (_, _, token) => + { + if (++attempts == 1) + { + await Task.Delay(Timeout.InfiniteTimeSpan, token); + } + + return true; + }); + await using var monitor = new CodexActivityMonitor(source, reader, + requestTimeout: TimeSpan.FromMilliseconds(30)); + await monitor.StartAsync(); + source.Publish(new(CodexActivitySignalKind.TurnStarted, "session", "turn")); + + await monitor.ReconcileAsync().WaitAsync(TimeSpan.FromSeconds(2)); + Assert.True(monitor.IsActive); + await monitor.ReconcileAsync(); + Assert.False(monitor.IsActive); + } + + [Fact] + public async Task DisposalCancelsPendingCheckWithoutClearingActivity() + { + await using var source = new FakeSignalSource(); + var reader = new FakeCompletionReader(async (_, _, token) => + { + await Task.Delay(Timeout.InfiniteTimeSpan, token); + return true; + }); + await using var monitor = new CodexActivityMonitor(source, reader); + await monitor.StartAsync(); + source.Publish(new(CodexActivitySignalKind.TurnStarted, "session", "turn")); + var reconciliation = monitor.ReconcileAsync(); + + await monitor.DisposeAsync().AsTask().WaitAsync(TimeSpan.FromSeconds(2)); + await Assert.ThrowsAnyAsync(() => reconciliation); + Assert.True(monitor.IsActive); + } + + [Fact] + public async Task IdleMonitorDoesNotQueryCompletion() + { + await using var source = new FakeSignalSource(); + var reads = 0; + var reader = new FakeCompletionReader((_, _, _) => + { + reads++; + return Task.FromResult(false); + }); + await using var monitor = new CodexActivityMonitor(source, reader); + await monitor.StartAsync(); + await monitor.ReconcileAsync(); + Assert.Equal(0, reads); + } + + [Fact] + public async Task PeriodicCheckRecoversMissingStopWithoutAnotherHook() + { + await using var source = new FakeSignalSource(); + var reader = new FakeCompletionReader((_, _, _) => Task.FromResult(true)); + await using var monitor = new CodexActivityMonitor(source, reader, TimeSpan.FromMilliseconds(20)); + var stopped = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + monitor.ActivityChanged += active => + { + if (!active) + { + stopped.TrySetResult(); + } + }; + await monitor.StartAsync(); + source.Publish(new(CodexActivitySignalKind.TurnStarted, "session", "turn")); + + await stopped.Task.WaitAsync(TimeSpan.FromSeconds(2)); + Assert.False(monitor.IsActive); + } + + [Fact] + public async Task FailedCompletionCheckPreservesActivityAndCanRecoverOnRetry() + { + await using var source = new FakeSignalSource(); + var attempts = 0; + var reader = new FakeCompletionReader((_, _, _) => ++attempts == 1 + ? Task.FromException(new InvalidOperationException("Unavailable")) + : Task.FromResult(true)); + await using var monitor = new CodexActivityMonitor(source, reader); + await monitor.StartAsync(); + source.Publish(new(CodexActivitySignalKind.TurnStarted, "session", "turn")); + + await monitor.ReconcileAsync(); + Assert.True(monitor.IsActive); + await monitor.ReconcileAsync(); + Assert.False(monitor.IsActive); + } + + [Fact] + public async Task ConfirmedCompletionRecoversMissingStopAfterAnotherSessionFinishes() + { + await using var source = new FakeSignalSource(); + var reader = new FakeCompletionReader((_, _, _) => Task.FromResult(true)); + await using var monitor = new CodexActivityMonitor(source, reader); + var changes = new List(); + monitor.ActivityChanged += changes.Add; + await monitor.StartAsync(); + + source.Publish(new(CodexActivitySignalKind.TurnStarted, "session-a", "turn-a")); + source.Publish(new(CodexActivitySignalKind.TurnStarted, "session-b", "turn-b")); + source.Publish(new(CodexActivitySignalKind.TurnStopped, "session-b", "turn-b")); + Assert.True(monitor.IsActive); + + await monitor.ReconcileAsync(); + + Assert.False(monitor.IsActive); + Assert.Equal([true, false], changes); + } + + [Theory] + [InlineData(true, false)] + [InlineData(false, true)] + public async Task MissingStopLeavesActivityActiveAfterAnotherSessionFinishes( + bool deliverFirstStop, + bool expectedActivity) + { + await using var source = new FakeSignalSource(); + await using var monitor = new CodexActivityMonitor(source); + var changes = new List(); + monitor.ActivityChanged += changes.Add; + await monitor.StartAsync(); + + source.Publish(new(CodexActivitySignalKind.TurnStarted, "session-a", "turn-a")); + + // Session A has finished externally. Simulate either delivery or loss of its Stop. + // No SessionEnd is delivered, as the Codex session can remain open between turns. + if (deliverFirstStop) + { + source.Publish(new(CodexActivitySignalKind.TurnStopped, "session-a", "turn-a")); + } + + source.Publish(new(CodexActivitySignalKind.TurnStarted, "session-b", "turn-b")); + source.Publish(new(CodexActivitySignalKind.TurnStopped, "session-b", "turn-b")); + source.Publish(new(CodexActivitySignalKind.SessionEnded, "session-b")); + + // Characterizes the current limitation, rather than asserting recovery exists. + Assert.Equal(expectedActivity, monitor.IsActive); + if (deliverFirstStop) + { + Assert.Equal([true, false, true, false], changes); + } + else + { + Assert.Equal([true], changes); + } + + // A final event for the stranded session clears the activity immediately. + source.Publish(new(CodexActivitySignalKind.SessionEnded, "session-a")); + Assert.False(monitor.IsActive); + Assert.False(changes[^1]); + } + [Fact] public async Task FirstStartEnablesActivityAndDuplicateStartIsIdempotent() { @@ -169,14 +390,38 @@ await Assert.ThrowsAsync(() => Assert.False(monitor.IsActive); } + private sealed class FakeCompletionReader( + Func> read) : ICodexTurnCompletionReader + { + public Task IsCompletedAsync( + string sessionId, string turnId, CancellationToken cancellationToken) => + read(sessionId, turnId, cancellationToken); + } + private sealed class FakeSignalSource : ICodexActivitySignalSource { - public event Action? SignalReceived; + private Action? _signalReceived; + public Action? OnUnsubscribe { get; set; } + public bool IsDisposed { get; private set; } + + public event Action? SignalReceived + { + add => _signalReceived += value; + remove + { + OnUnsubscribe?.Invoke(); + _signalReceived -= value; + } + } public Task StartAsync(CancellationToken cancellationToken = default) => Task.CompletedTask; - public void Publish(CodexActivitySignal signal) => SignalReceived?.Invoke(signal); + public void Publish(CodexActivitySignal signal) => _signalReceived?.Invoke(signal); - public ValueTask DisposeAsync() => ValueTask.CompletedTask; + public ValueTask DisposeAsync() + { + IsDisposed = true; + return ValueTask.CompletedTask; + } } } diff --git a/tests/CodexUsageWidget.Tests/CodexTurnCompletionReaderTests.cs b/tests/CodexUsageWidget.Tests/CodexTurnCompletionReaderTests.cs new file mode 100644 index 0000000..95797f6 --- /dev/null +++ b/tests/CodexUsageWidget.Tests/CodexTurnCompletionReaderTests.cs @@ -0,0 +1,82 @@ +using System.Text.Json; +using CodexUsageWidget.Infrastructure.Codex; + +namespace CodexUsageWidget.Tests; + +public sealed class CodexTurnCompletionReaderTests +{ + [Theory] + [InlineData("null")] + [InlineData("{}")] + [InlineData("{\"data\":[]}")] + [InlineData("{\"data\":[{\"id\":\"other\",\"status\":\"completed\",\"completedAt\":1789463298}]}")] + [InlineData("{\"data\":[{\"id\":\"target-turn\",\"status\":\"completed\",\"completedAt\":\"invalid\"}]}")] + public async Task MissingOrMalformedCompletionIsNotEvidenceOfAnIdleTurn(string json) + { + await using var session = new FakeSession(json); + var reader = new CodexTurnCompletionReader(session); + Assert.False(await reader.IsCompletedAsync("session", "target-turn", CancellationToken.None)); + } + + [Fact] + public async Task FindsCompletedTurnOnAnOlderPage() + { + await using var session = new FakeSession( + """{ "data": [{"id":"newer", "status":"inProgress"}], "nextCursor":"older" }""", + """{ "data": [{"id":"target-turn", "status":"completed", "completedAt":1789463298}], "nextCursor":null }"""); + var reader = new CodexTurnCompletionReader(session); + Assert.True(await reader.IsCompletedAsync("session", "target-turn", CancellationToken.None)); + Assert.Equal("older", session.LastParameters.GetProperty("cursor").GetString()); + } + + [Theory] + [InlineData("completed", "1789463298", true)] + [InlineData("interrupted", "1789463298", true)] + [InlineData("failed", "1789463298", true)] + [InlineData("interrupted", "null", false)] + [InlineData("completed", "null", false)] + [InlineData("inProgress", "null", false)] + [InlineData("inProgress", "1789463298", false)] + [InlineData("unknown", "1789463298", false)] + public async Task RequiresExplicitCompletionForTheExactTurn( + string status, string completedAt, bool expected) + { + await using var session = new FakeSession($$""" + { "data": [ + { "id": "another-turn", "status": "completed", "completedAt": 1789463298 }, + { "id": "target-turn", "status": "{{status}}", "completedAt": {{completedAt}}, + "itemsView": "notLoaded", "items": [] } + ], "nextCursor": null } + """); + var reader = new CodexTurnCompletionReader(session); + + Assert.Equal(expected, await reader.IsCompletedAsync("session", "target-turn", CancellationToken.None)); + } + + private sealed class FakeSession(params string[] results) : ICodexAppServerSession + { + private readonly Queue _results = new(results); + public JsonElement LastParameters { get; private set; } + public event EventHandler? NotificationReceived; + public event EventHandler? DiagnosticMessage; + + public Task RequestAsync(string method, object? parameters, CancellationToken cancellationToken) + { + Assert.Equal("thread/turns/list", method); + var request = JsonSerializer.SerializeToElement(parameters); + LastParameters = request; + Assert.Equal("session", request.GetProperty("threadId").GetString()); + Assert.Equal("notLoaded", request.GetProperty("itemsView").GetString()); + Assert.Equal("desc", request.GetProperty("sortDirection").GetString()); + using var document = JsonDocument.Parse(_results.Dequeue()); + return Task.FromResult(document.RootElement.Clone()); + } + + public ValueTask DisposeAsync() + { + GC.KeepAlive(NotificationReceived); + GC.KeepAlive(DiagnosticMessage); + return ValueTask.CompletedTask; + } + } +}