From 61b631597f99056a5b51d968f9799f87e3fa4f03 Mon Sep 17 00:00:00 2001 From: Aaron Stannard Date: Wed, 23 Sep 2026 23:20:53 +0000 Subject: [PATCH 1/3] Resume interrupted sessions through reminders --- docs/spec/SPEC-011-daemon-architecture.md | 13 +- .../.system/files/netclaw-operations/SKILL.md | 13 +- .../resume-interrupted-sessions/design.md | 13 +- .../specs/session-resume/spec.md | 6 +- .../resume-interrupted-sessions/tasks.md | 18 +- .../Reminders/ReminderManagerActorTests.cs | 7 +- .../Sessions/LlmSessionIntegrationTests.cs | 144 ++++++- .../Reminders/ReminderManagerActor.cs | 9 +- .../Reminders/ReminderProtocol.cs | 2 +- .../Sessions/LlmSessionActor.cs | 356 +++++++++++++++++- .../Sessions/SessionProtocol.Responses.cs | 10 +- .../Services/DaemonRestartCoordinatorTests.cs | 36 +- .../Services/RestartRecoveryServiceTests.cs | 80 ++-- src/Netclaw.Daemon/Program.cs | 13 + .../Services/DaemonRestartCoordinator.cs | 7 +- .../Services/RestartManifestStore.cs | 23 +- .../Services/RestartRecoveryService.cs | 106 ++++-- .../Services/SessionDrainHelper.cs | 28 +- 18 files changed, 743 insertions(+), 141 deletions(-) diff --git a/docs/spec/SPEC-011-daemon-architecture.md b/docs/spec/SPEC-011-daemon-architecture.md index 66d113a46..42f2f3be8 100644 --- a/docs/spec/SPEC-011-daemon-architecture.md +++ b/docs/spec/SPEC-011-daemon-architecture.md @@ -234,6 +234,15 @@ The pipeline reports cancellation through its canceled task state. The daemon keeps the existing approval prompt. Button and text responses can resume the recovered turn. A channel UI can temporarily lag the session state after restart. +During any graceful stop, the actor gives an active model call a two-second +completion grace. It then cancels an eligible call and waits for its task. +The actor returns a standard `current_session` reminder for durable pending +input. The reminder expires ten minutes after the interruption. Startup +registers each fresh reminder through the reminder manager. The session +restores the pending input under its recorded authority when the reminder +arrives. A completed turn, partial text, or possible tool effect produces no +restart reminder. + `netclaw daemon status` checks the PID file and verifies the process is alive. Reports: running/stopped, PID, uptime, port, number of active sessions. @@ -360,8 +369,8 @@ not execute tools. 4. **Valid config**: close daemon-managed ingress, enumerate live session actors, ask them to drain, persist a restart manifest, and request coordinated daemon restart -5. **After restart**: warm the sessions that were active when restart began and - inject a continuity notice for the next turn +5. **After restart**: register fresh restart reminders. The normal reminder + route activates each target session. 6. **Invalid config**: log warning with validation errors, preserve previous config ### What Changes Take Effect After Restart diff --git a/feeds/skills/.system/files/netclaw-operations/SKILL.md b/feeds/skills/.system/files/netclaw-operations/SKILL.md index be3d434ef..d1eb69083 100644 --- a/feeds/skills/.system/files/netclaw-operations/SKILL.md +++ b/feeds/skills/.system/files/netclaw-operations/SKILL.md @@ -3,7 +3,7 @@ name: netclaw-operations description: "REQUIRED when the user asks about scheduling, reminders, cron jobs, timers, background jobs, diagnostics, troubleshooting, MCP tools, daemon health, identity updates, or Netclaw capabilities and self-maintenance." metadata: author: netclaw - version: "2.75.1" + version: "2.75.2" --- # Netclaw Operations @@ -461,11 +461,16 @@ stale button, ask them to re-issue the request. During a graceful stop, the session can stop a tool task that waits only for journaled approval prompts. The session waits for that task to stop before it -acknowledges drain. The daemon keeps the existing approval prompt. The original requester can -still approve through its button or a text response after restart. The UI +acknowledges drain. The daemon keeps the existing approval prompt. The original +requester can approve it after restart. They can use its button or a text response. The UI can temporarily lag the session state. Shutdown cancellation does not mean the approval expired. -An active tool or accepted buffered input keeps the current bounded drain path. +An active tool with a possible external effect keeps the bounded drain path. + +An interrupted model call can create a short-lived restart reminder. The +session restores accepted input and its original authority from the journal. +The reminder expires ten minutes after the interruption. A completed turn, +partial reply, or possible tool effect does not create this reminder. **Why you may not see a prompt at all.** If the user invokes a read-only verb (say `grep`) with a path argument under a tree the operator has previously diff --git a/openspec/changes/resume-interrupted-sessions/design.md b/openspec/changes/resume-interrupted-sessions/design.md index eca0fbc61..6826896c0 100644 --- a/openspec/changes/resume-interrupted-sessions/design.md +++ b/openspec/changes/resume-interrupted-sessions/design.md @@ -38,15 +38,15 @@ The actor creates a standard one shot `ReminderDefinition` only when these condi - no tool batch started; - no partial text reached a subscriber; - the stored turn context is valid; -- the existing reminder path supports the session channel. +- the stored turn has a channel type for current-session delivery. -The reminder expires ten minutes after the interruption. The restart manifest stores the definition with the active session list. +The reminder expires ten minutes after the interruption. The restart manifest stores only restart reminder definitions. ### D3. The reminder manager owns wakeup and delivery Startup registers each fresh definition through `SaveReminderCommand`. The reminder uses `DeliveryKind.CurrentSession` and the existing gateway path. -The daemon adds no route binder, channel state, retry loop, or resume candidate protocol. The reminder manager owns persistence, delivery retries, and deduplication. +The daemon adds no route binder, channel state, retry loop, or resume candidate protocol. The reminder manager owns route resolution, persistence, delivery retries, and deduplication. ### D4. The reminder is a trigger @@ -67,8 +67,8 @@ journal -> session: stored session -> source: CommandAck stop -> session: PrepareForDaemonRestart session -> model: cancel and await stop -session -> stop: ReminderDefinition or no reminder -stop -> manifest: active sessions and reminders +session -> stop: DaemonRestartPrepared with a reminder or no reminder +stop -> manifest: restart reminders start -> reminder manager: SaveReminderCommand reminder manager -> existing gateway: current_session reminder gateway -> session: SendUserMessage @@ -82,4 +82,5 @@ session -> model: resume prior work - A tool can have an uncertain effect. `ToolBatchStarted` closes its input before execution and blocks this path. - A partial reply can repeat text. The actor records transient text emission and does not create a reminder. - A reminder can register twice after a process failure. Its stored ID makes `CreateOnly` registration idempotent. -- A channel can lack `current_session` support. The actor creates no reminder for that session. +- A stored turn can lack a channel type. The actor then creates no reminder. +- The reminder system owns gateway resolution for a stored channel type. diff --git a/openspec/changes/resume-interrupted-sessions/specs/session-resume/spec.md b/openspec/changes/resume-interrupted-sessions/specs/session-resume/spec.md index 3ab4c0424..9a3838fb1 100644 --- a/openspec/changes/resume-interrupted-sessions/specs/session-resume/spec.md +++ b/openspec/changes/resume-interrupted-sessions/specs/session-resume/spec.md @@ -78,9 +78,9 @@ The reminder manager SHALL deliver a fresh restart reminder through its existing - **THEN** it does not register or deliver that reminder - **AND** it logs one warning -#### Scenario: A channel lacks current session delivery +#### Scenario: Stored authority has no channel type -- **GIVEN** an interrupted session uses a channel that the reminder manager cannot address +- **GIVEN** an interrupted session has no stored channel type - **WHEN** graceful drain classifies the session - **THEN** the actor creates no restart reminder -- **AND** no channel adapter is added by this change +- **AND** the actor does not contain a channel-specific route list diff --git a/openspec/changes/resume-interrupted-sessions/tasks.md b/openspec/changes/resume-interrupted-sessions/tasks.md index 60c9e872f..a6668dd96 100644 --- a/openspec/changes/resume-interrupted-sessions/tasks.md +++ b/openspec/changes/resume-interrupted-sessions/tasks.md @@ -7,18 +7,18 @@ ## 2. Graceful drain -- [ ] 2.1 Retain and cancel the active model task after a short grace. -- [ ] 2.2 Return one standard reminder definition for eligible pending input. -- [ ] 2.3 Exclude approvals, tool work, partial replies, and unsupported channels. +- [x] 2.1 Retain and cancel the active model task after a short grace. +- [x] 2.2 Return one standard reminder definition for eligible pending input. +- [x] 2.3 Exclude approvals, tool work, partial replies, and input without a channel type. ## 3. Existing reminder path -- [ ] 3.1 Store reminder definitions in the restart manifest. -- [ ] 3.2 Register fresh reminders through the reminder manager after startup. -- [ ] 3.3 Restore pending input under its original context when the reminder arrives. -- [ ] 3.4 Verify expiration, duplicate registration, and a cold session wakeup. +- [x] 3.1 Store reminder definitions in the restart manifest. +- [x] 3.2 Register fresh reminders through the reminder manager after startup. +- [x] 3.3 Restore pending input under its original context when the reminder arrives. +- [x] 3.4 Verify expiration, duplicate registration, and a cold session wakeup. ## 4. Verification -- [ ] 4.1 Update SPEC-011 and the operations skill. -- [ ] 4.2 Run actor and daemon tests, evals, Slopwatch, headers, and OpenSpec validation. +- [x] 4.1 Update SPEC-011 and the operations skill. +- [x] 4.2 Run actor and daemon tests, evals, Slopwatch, headers, and OpenSpec validation. diff --git a/src/Netclaw.Actors.Tests/Reminders/ReminderManagerActorTests.cs b/src/Netclaw.Actors.Tests/Reminders/ReminderManagerActorTests.cs index d5765f084..35e1a7ddd 100644 --- a/src/Netclaw.Actors.Tests/Reminders/ReminderManagerActorTests.cs +++ b/src/Netclaw.Actors.Tests/Reminders/ReminderManagerActorTests.cs @@ -448,7 +448,7 @@ public async Task Save_rejects_missing_authorization_context() } [Fact] - public async Task Save_rejects_expiration_for_oneshot_reminders() + public async Task Save_accepts_expiration_for_oneshot_reminders() { var manager = await GetManagerAsync(); var now = TimeProvider.System.GetUtcNow(); @@ -469,9 +469,8 @@ public async Task Save_rejects_expiration_for_oneshot_reminders() TimeSpan.FromSeconds(5), TestContext.Current.CancellationToken); - Assert.False(response.Success); - Assert.Equal(ReminderSaveError.Validation, response.Error); - Assert.Contains("one-shot", response.ErrorMessage, StringComparison.OrdinalIgnoreCase); + Assert.True(response.Success, response.ErrorMessage); + Assert.Equal(definition.ExpiresAt, _definitionStore.Get(definition.Id)?.ExpiresAt); } [Fact] diff --git a/src/Netclaw.Actors.Tests/Sessions/LlmSessionIntegrationTests.cs b/src/Netclaw.Actors.Tests/Sessions/LlmSessionIntegrationTests.cs index cd4722461..1a216c9ec 100644 --- a/src/Netclaw.Actors.Tests/Sessions/LlmSessionIntegrationTests.cs +++ b/src/Netclaw.Actors.Tests/Sessions/LlmSessionIntegrationTests.cs @@ -1747,7 +1747,7 @@ await sessionManager.Ask(new SendUserMessage var child = await Sys.ActorSelection($"/user/session-manager/{escapedId}").ResolveOne(TimeSpan.FromSeconds(3), TestContext.Current.CancellationToken); Watch(child); - var drainTask = sessionManager.Ask(new PrepareForDaemonRestart(sessionId, "config-reload"), TimeSpan.FromSeconds(5), cancellationToken: TestContext.Current.CancellationToken); + var drainTask = sessionManager.Ask(new PrepareForDaemonRestart(sessionId, "config-reload"), TimeSpan.FromSeconds(5), cancellationToken: TestContext.Current.CancellationToken); var nack = await sessionManager.Ask(new SendUserMessage { @@ -1770,7 +1770,9 @@ await sessionManager.Ask(new SendUserMessage await subscriber.ExpectMsgAsync(TimeSpan.FromSeconds(3), cancellationToken: TestContext.Current.CancellationToken); await subscriber.ExpectMsgAsync(TimeSpan.FromSeconds(3), cancellationToken: TestContext.Current.CancellationToken); - Assert.Equal(sessionId, (await drainTask).SessionId); + var drainAck = await drainTask; + Assert.Equal(sessionId, drainAck.SessionId); + Assert.Null(drainAck.RestartReminder); await ExpectTerminatedAsync(child, TimeSpan.FromSeconds(5), cancellationToken: TestContext.Current.CancellationToken); Assert.DoesNotContain(_fakeChatClient.ReceivedMessages, conversation => conversation.Any(msg => @@ -1778,6 +1780,133 @@ msg.Text is not null && msg.Text.Contains("msg_too_long", StringComparison.OrdinalIgnoreCase))); } + [Fact] + public async Task Restart_reminder_restores_an_interrupted_input_without_a_second_user_prompt() + { + var sessionId = new SessionId("signalr/restart-resume"); + var sessionManager = ActorRegistry.Get(); + var responseGate = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + _fakeChatClient.NextResponseGate = responseGate; + var source = new MessageSource + { + ChannelType = ChannelType.SignalR, + SenderId = new SenderId("operator-1"), + MessageId = "restart-input-1", + TurnId = new Netclaw.Actors.Protocol.TurnId("restart-turn-1"), + Audience = TrustAudience.Personal, + Boundary = TrustBoundary.Personal, + Principal = PrincipalClassification.Operator, + Provenance = new SourceProvenance(TransportAuthenticity.Verified, PayloadTaint.Trusted), + ReceivedAt = _timeProvider.GetUtcNow() + }; + + await sessionManager.Ask(new SendUserMessage + { + SessionId = sessionId, + Content = "Finish the interrupted task", + Source = source + }, TimeSpan.FromSeconds(3), TestContext.Current.CancellationToken); + await _fakeChatClient.FirstCallEntered.Task.WaitAsync(TestContext.Current.CancellationToken); + + var child = await Sys.ActorSelection($"/user/session-manager/{Uri.EscapeDataString(sessionId.Value)}") + .ResolveOne(TimeSpan.FromSeconds(3), TestContext.Current.CancellationToken); + Watch(child); + + var drainAck = await sessionManager.Ask( + new PrepareForDaemonRestart(sessionId, "daemon-stop"), + TimeSpan.FromSeconds(8), + TestContext.Current.CancellationToken); + var reminder = Assert.IsType(drainAck.RestartReminder); + Assert.Equal(DeliveryKind.CurrentSession, reminder.Delivery.Kind); + Assert.Equal(TrustAudience.Personal, reminder.Audience); + Assert.Equal(TrustBoundary.Personal, reminder.Boundary); + Assert.Equal(TimeSpan.FromMinutes(10), reminder.ExpiresAt - reminder.CreatedAt); + await ExpectTerminatedAsync( + child, + TimeSpan.FromSeconds(3), + cancellationToken: TestContext.Current.CancellationToken); + + var subscriber = CreateTestProbe("restart-resume-subscriber"); + await JoinSessionAsync(sessionManager, subscriber, sessionId); + + await sessionManager.Ask(new SendUserMessage + { + SessionId = sessionId, + Content = reminder.Instructions, + Source = new MessageSource + { + ChannelType = ChannelType.SignalR, + SenderId = new SenderId("reminder-system"), + MessageId = $"{reminder.Id}:1", + TurnId = new Netclaw.Actors.Protocol.TurnId($"{reminder.Id}:1"), + Audience = reminder.Audience, + Boundary = reminder.Boundary, + Principal = PrincipalClassification.VerifiedAutomation, + Provenance = new SourceProvenance( + TransportAuthenticity.LocalProcess, + PayloadTaint.Trusted), + ReceivedAt = _timeProvider.GetUtcNow(), + ReminderId = new ReminderId($"{reminder.Id}:1") + } + }, TimeSpan.FromSeconds(3), TestContext.Current.CancellationToken); + + await subscriber.ExpectMsgAsync( + TimeSpan.FromSeconds(3), + cancellationToken: TestContext.Current.CancellationToken); + await subscriber.ExpectMsgAsync( + TimeSpan.FromSeconds(3), + cancellationToken: TestContext.Current.CancellationToken); + + var resumedCall = _fakeChatClient.ReceivedMessages.Last(conversation => + conversation.Any(message => string.Equals( + message.Text, + "Finish the interrupted task", + StringComparison.Ordinal))); + Assert.DoesNotContain(resumedCall, message => string.Equals( + message.Text, + reminder.Instructions, + StringComparison.Ordinal)); + } + + [Fact] + public async Task Restart_drain_does_not_remind_after_partial_text() + { + var sessionId = new SessionId("signalr/restart-partial-text"); + var sessionManager = ActorRegistry.Get(); + var subscriber = CreateTestProbe("restart-partial-text-subscriber"); + await JoinSessionAsync(sessionManager, subscriber, sessionId, OutputFilter.TextStreaming); + _fakeChatClient.StreamCompletionGate = new TaskCompletionSource( + TaskCreationOptions.RunContinuationsAsynchronously); + + await sessionManager.Ask(new SendUserMessage + { + SessionId = sessionId, + Content = "Do not repeat a partial reply", + Source = new MessageSource + { + ChannelType = ChannelType.SignalR, + SenderId = new SenderId("operator-1"), + Audience = TrustAudience.Personal, + Boundary = TrustBoundary.Personal, + Principal = PrincipalClassification.Operator, + Provenance = new SourceProvenance( + TransportAuthenticity.Verified, + PayloadTaint.Trusted), + ReceivedAt = _timeProvider.GetUtcNow() + } + }, TimeSpan.FromSeconds(3), TestContext.Current.CancellationToken); + await subscriber.ExpectMsgAsync( + TimeSpan.FromSeconds(3), + cancellationToken: TestContext.Current.CancellationToken); + + var ack = await sessionManager.Ask( + new PrepareForDaemonRestart(sessionId, "daemon-stop"), + TimeSpan.FromSeconds(8), + TestContext.Current.CancellationToken); + + Assert.Null(ack.RestartReminder); + } + [Fact] public async Task Compacting_session_rejects_new_work_and_passivates_after_compaction_finishes() { @@ -2506,6 +2635,8 @@ public IReadOnlyList ReceivedOptions public TaskCompletionSource? NextResponseGate { get; set; } + public TaskCompletionSource? StreamCompletionGate { get; set; } + public async Task GetResponseAsync( IEnumerable messages, ChatOptions? options = null, @@ -2694,6 +2825,15 @@ private async IAsyncEnumerable CreateStreamingUpdatesAsync( [System.Runtime.CompilerServices.EnumeratorCancellation] CancellationToken cancellationToken) { var response = await GetResponseAsync(messages, options, cancellationToken); + if (StreamCompletionGate is { } streamGate) + { + yield return new ChatResponseUpdate(Microsoft.Extensions.AI.ChatRole.Assistant, "partial "); + yield return new ChatResponseUpdate(Microsoft.Extensions.AI.ChatRole.Assistant, "reply"); + using var registration = cancellationToken.Register(() => streamGate.TrySetCanceled(cancellationToken)); + await streamGate.Task; + yield break; + } + foreach (var update in response.ToChatResponseUpdates()) { cancellationToken.ThrowIfCancellationRequested(); diff --git a/src/Netclaw.Actors/Reminders/ReminderManagerActor.cs b/src/Netclaw.Actors/Reminders/ReminderManagerActor.cs index 66100a169..bbf1384c8 100644 --- a/src/Netclaw.Actors/Reminders/ReminderManagerActor.cs +++ b/src/Netclaw.Actors/Reminders/ReminderManagerActor.cs @@ -267,12 +267,6 @@ static ReminderSavedResponse ValidationFailure(ReminderId id, string title, stri : cmd.Definition.CreatedBy }; - if (normalized.Schedule.Type == ReminderScheduleType.OneShot && normalized.ExpiresAt is not null) - { - replyTo.Tell(ValidationFailure(id, title, "expires_at is not applicable to one-shot reminders.")); - return; - } - if (exists) { var existing = _definitionStore.Get(id); @@ -570,8 +564,7 @@ private async Task HandleReminderFiredAsync(ReminderEnvelope en return; } - if (definition.Schedule.Type is not ReminderScheduleType.OneShot - && definition.ExpiresAt is { } expiresAt + if (definition.ExpiresAt is { } expiresAt && expiresAt <= _timeProvider.GetUtcNow()) { _log.Info("Reminder '{0}' has expired (expiresAt={1}), disabling", reminderId.Value, expiresAt); diff --git a/src/Netclaw.Actors/Reminders/ReminderProtocol.cs b/src/Netclaw.Actors/Reminders/ReminderProtocol.cs index b43716ba9..ec95c2745 100644 --- a/src/Netclaw.Actors/Reminders/ReminderProtocol.cs +++ b/src/Netclaw.Actors/Reminders/ReminderProtocol.cs @@ -239,7 +239,7 @@ public sealed record ReminderDefinition public long UpdatedAtMs { get; set; } /// - /// Optional expiration for recurring reminders. When set, the reminder + /// Optional expiration for reminders. When set, the reminder /// auto-disables on next fire after this time without executing. /// Null means no expiration (default for backwards compatibility). /// diff --git a/src/Netclaw.Actors/Sessions/LlmSessionActor.cs b/src/Netclaw.Actors/Sessions/LlmSessionActor.cs index dac9d26ec..a8e6aa0b6 100644 --- a/src/Netclaw.Actors/Sessions/LlmSessionActor.cs +++ b/src/Netclaw.Actors/Sessions/LlmSessionActor.cs @@ -144,6 +144,7 @@ public sealed class LlmSessionActor : ReceivePersistentActor, IWithTimers // is the authoritative timeout — this CTS just propagates cancellation to the // HTTP layer so timed-out connections are released. private CancellationTokenSource? _activeLlmCts; + private Task? _activeLlmWorkTask; // Actor-owned CTS for active tool execution. Cancels direct approval waits // and tool calls when the session stops, restarts, or fails the turn. @@ -165,6 +166,7 @@ public sealed class LlmSessionActor : ReceivePersistentActor, IWithTimers // (the generous prefill budget before the first token vs the tighter inter-delta // budget after). private bool _anyContentStreamed; + private bool _turnEmittedText; // Per-turn diagnostic correlation (ephemeral) private Protocol.TurnId? _activeTurnId; @@ -192,12 +194,19 @@ public sealed class LlmSessionActor : ReceivePersistentActor, IWithTimers private readonly Telemetry.ISessionMetrics? _sessionMetrics; private bool _restartDrainRequested; + private bool _llmRestartStopRequested; + private DateTimeOffset? _restartInterruptionAt; private bool _passivationCompleted; private bool _passivationFinalStopScheduled; private sealed record ToolPipelineStoppedForRestart(Task WorkTask, Exception? Failure) : INoSerializationVerificationNeeded; + private sealed record RestartModelGraceExpired(long CallId) : INoSerializationVerificationNeeded; + + private sealed record ModelStoppedForRestart(long CallId, Exception? Failure) + : INoSerializationVerificationNeeded; + // Reap-on-passivation handshake: while a KillJobsForSession ask is in // flight, the final snapshot is deferred so it captures the reaped marks. // _jobReapEpoch is bumped per reap request so a late reply from a @@ -643,6 +652,8 @@ private void Processing() Command(msg => { if (msg.CallId != _activeCallId) return; // stale failure from cancelled call + if (_llmRestartStopRequested) return; + _activeLlmWorkTask = null; _watchdog.Stop(Timers); CancelAndDisposeLlmCts(); @@ -765,7 +776,10 @@ private void HandleLlmResponseReceived(LlmResponseReceived msg) { if (msg.CallId != _activeCallId) return; + if (_llmRestartStopRequested) + return; + _activeLlmWorkTask = null; _watchdog.Stop(Timers); CancelAndDisposeLlmCts(); @@ -823,7 +837,7 @@ private void HandleLlmResponseReceived(LlmResponseReceived msg) private void HandleLlmResponseDeltaReceived(LlmResponseDeltaReceived msg) { - if (msg.CallId != _activeCallId) + if (msg.CallId != _activeCallId || _llmRestartStopRequested) return; // Two-phase watchdog (shared with the sub-agent path): keep the generous @@ -841,6 +855,7 @@ private void HandleLlmResponseDeltaReceived(LlmResponseDeltaReceived msg) switch (msg.Content) { case TextContent text when !string.IsNullOrEmpty(text.Text): + _turnEmittedText = true; EmitOutput(new TextDeltaOutput(text.Text) { SessionId = _sessionId @@ -1398,6 +1413,7 @@ private void DrainBufferOrReady() if (_restartDrainRequested) { _deferredApprovalResponse = null; + MarkRestartReminderEligible(); ClearBufferedMessagesForRestartDrain(); _resumeToolLoopAfterCompaction = false; TransitionTo(SessionPhase.Ready); @@ -1442,6 +1458,12 @@ private void DrainBufferOrReady() } } + private const string RestartReminderInstruction = + "Resume the work that was interrupted by the daemon restart."; + private const string RestartReminderIdPrefix = "restart-resume-"; + private static readonly TimeSpan RestartModelGracePeriod = TimeSpan.FromSeconds(2); + private static readonly TimeSpan RestartReminderLifetime = TimeSpan.FromMinutes(10); + private static readonly object RestartModelGraceTimerKey = new(); private static readonly TimeSpan PassivationGracePeriod = TimeSpan.FromSeconds(5); // Bounds the reap-on-passivation handshake; the Ask always resolves within @@ -1598,8 +1620,13 @@ private void Passivating() HandleDeliveryFailedWhenReady(msg); }); - // Request final distillation from observer, or stop immediately - if (_observerActor is not null) + // Restart drain already preserves accepted input in the journal. + // Skip the optional memory pass because ingress is closed. + if (_restartDrainRequested) + { + CompletePassivation(); + } + else if (_observerActor is not null) { _observerActor.Tell(new DistillMemories(), Self); Timers.StartSingleTimer(PassivationTimerKey, new PassivationTimeout(), PassivationGracePeriod); @@ -1632,6 +1659,13 @@ private void CompletePassivation() _passivationFinalStopScheduled = true; SaveSnapshotIfSafe(); + + if (_restartDrainRequested) + { + FinalizePassivation(); + return; + } + Timers.StartSingleTimer( PassivationFinalStopTimerKey, new PassivationFinalStop(), @@ -1669,7 +1703,9 @@ private void FinalizePassivation() _passivationCompleted = true; _lifecycleObserver?.OnSessionDeactivated(_sessionId); - _restartDrainReplyTo?.Tell(CommandAck.For(_sessionId)); + _restartDrainReplyTo?.Tell(new DaemonRestartPrepared( + _sessionId, + BuildRestartReminder())); _restartDrainReplyTo = null; Context.Stop(Self); } @@ -1892,6 +1928,7 @@ private void HandleToolCallResponse( { _state = _state.CloseInputs(evt.ConsumedInputIds); _activeInputIds.Clear(); + _turnEmittedText = false; ApplyToolBatchStarted(evt); EmitAndDispatchToolBatch( lastMessage, @@ -2152,6 +2189,7 @@ private void HandleTextResponse( _inFlightDedup.CompleteBackgroundJob(evt.SourceBackgroundJobId); _state = _state.CloseInputs(evt.ConsumedInputIds); _activeInputIds.Clear(); + _turnEmittedText = false; var processed = _state.ProcessedReminderIds; if (evt.SourceReminderId is { } reminderId && !string.IsNullOrEmpty(reminderId.Value)) @@ -2191,6 +2229,7 @@ private void DrainBufferedMessagesOrBecomeReady() { if (_restartDrainRequested) { + MarkRestartReminderEligible(); ClearBufferedMessagesForRestartDrain(); TransitionTo(SessionPhase.Ready); TransitionTo(SessionPhase.Passivating); @@ -2219,6 +2258,7 @@ private bool DrainBufferedUserMessages() // window only after the prior batch or response completes. _turnState.ResetForNewTurn(); _recallManager.ResetForNewTurn(); + _turnEmittedText = false; } foreach (var (message, _) in _buffer) @@ -2304,6 +2344,9 @@ private void HandleIncomingUserMessage(SendUserMessage cmd) return; } + if (TryResumeInterruptedInput(cmd)) + return; + if (TryRejectIncompatibleInput(cmd.MediaReferences, cmd.Source)) return; @@ -2429,6 +2472,7 @@ private void RestoreConsumedInputs(IReadOnlyList inputIds, Serializable private void ContinueIncomingUserMessage(SendUserMessage cmd) { + _turnEmittedText = false; _activeInputIds.Clear(); if (cmd.AdmittedInputId is { } admittedId) _activeInputIds.Add(admittedId); @@ -2484,6 +2528,103 @@ private void ContinueIncomingUserMessage(SendUserMessage cmd) TransitionTo(SessionPhase.Processing); } + private bool TryResumeInterruptedInput(SendUserMessage cmd) + { + var reminderId = cmd.Source?.ReminderId?.Value; + if (reminderId is null + || !reminderId.StartsWith(RestartReminderIdPrefix, StringComparison.Ordinal) + || !string.Equals(cmd.Content, RestartReminderInstruction, StringComparison.Ordinal)) + return false; + + if (_state.PendingInputs.Count == 0) + { + _log.Warning( + "Restart reminder {ReminderId} found no pending input for session {SessionId}.", + reminderId, + _sessionId.Value); + TryReplyNack("The interrupted input is no longer pending."); + return true; + } + + if (!TryGetRestartContext(out var firstContext, out var reason)) + { + _log.Warning( + "Restart reminder {ReminderId} cannot restore session {SessionId}: {Reason}.", + reminderId, + _sessionId.Value, + reason ?? "invalid stored authority"); + TryReplyNack("The interrupted input authority is invalid."); + return true; + } + + var pending = _state.PendingInputs.ToArray(); + _restartInterruptionAt = null; + _turnEmittedText = false; + _activeInputIds.Clear(); + _currentTurnSource = cmd.Source; + _currentTurnContext = firstContext; + BindTurnTelemetry(firstContext); + _toolApprovals.StartTurn(firstContext); + _currentTrustContext = _trustContextDeriver?.DeriveFromTurnContext(firstContext); + _inFlightDedup.ReserveReminder(cmd.Source!.ReminderId); + SetSystemPrompt(); + _turnState.ResetForNewTurn(); + _discoveredToolCache.PrepareForNewTurn( + _config.Tuning.DiscoveredToolRetentionTurns, + _config.Tuning.DiscoveredToolMaxCount, + _fullRegistry); + + if (TryResumeSlashCommand(pending)) + return true; + + _activeInputIds.AddRange(pending.Select(static item => item.InputId)); + foreach (var item in pending) + _state = _state with { History = _state.History.Add(item.UserMessage) }; + + _recallManager.ResetForNewTurn(); + _compactionOverflowRetryCount = 0; + + TryReplyAck(); + var recallQuery = string.Join( + '\n', + pending.Select(static item => item.ExecutableText ?? item.UserMessage.Content ?? string.Empty)); + FireInitialTurnLlmCall(recallQuery); + TransitionTo(SessionPhase.Processing); + return true; + } + + private bool TryResumeSlashCommand(IReadOnlyList pending) + { + var first = pending[0]; + if (first.ExecutableText is not { } executableText + || string.IsNullOrWhiteSpace(executableText) + || executableText.TrimStart()[0] != '/') + return false; + + _activeInputIds.Add(first.InputId); + if (!TryHandleSlashCommand(executableText, first.UserMessage.MediaReferences)) + { + _activeInputIds.Clear(); + return false; + } + + if (_phase.Current != SessionPhase.Processing) + return true; + + foreach (var item in pending.Skip(1)) + { + _buffer.Add((new SendUserMessage + { + SessionId = _sessionId, + Content = item.UserMessage.Content ?? string.Empty, + MediaReferences = item.UserMessage.MediaReferences, + AdmittedInputId = item.InputId + }, false)); + } + + return true; + } + private bool IsReminderDedupHit(ReminderId? reminderId, bool includeBuffered) { if (reminderId is not { } id || string.IsNullOrEmpty(id.Value)) @@ -2529,6 +2670,8 @@ private void CommandCommonMessages() { Command(_ => RequestRestartDrain()); Command(HandleToolPipelineStoppedForRestart); + Command(HandleRestartModelGraceExpired); + Command(HandleModelStoppedForRestart); Command(HandleWorkingContextSnapshotReady); Command(msg => @@ -3022,7 +3165,14 @@ private void ContinueFireLlmCall(bool forceNoTools) forceNoTools, _activeCallId); - _ = SessionLlmInvoker.InvokeAsync(client, messages, options, self, _activeCallId, _sessionId, _activeLlmCts!.Token); + _activeLlmWorkTask = SessionLlmInvoker.InvokeAsync( + client, + messages, + options, + self, + _activeCallId, + _sessionId, + _activeLlmCts!.Token); } private async Task CreateWorkingContextContinuationAsync( @@ -3085,6 +3235,9 @@ private void HandleWorkingContextSnapshotFailed(WorkingContextSnapshotFailed mes private void HandleWorkingContextSnapshotReady(WorkingContextSnapshotReady message) { + if (_llmRestartStopRequested) + return; + if (!ShouldApplyWorkingContextSnapshot( message.Generation, _workingContextGeneration, @@ -5055,19 +5208,31 @@ private void RequestRestartDrain() if (_phase.Current == SessionPhase.Ready) TransitionTo(SessionPhase.Passivating); else if (_phase.Current == SessionPhase.Processing) - TryStopDurableApprovalWaits(); + { + if (TryStopDurableApprovalWaits()) + return; + + if (_activeLlmCts is not null && _activeInputIds.Count > 0) + { + Timers.StartSingleTimer( + RestartModelGraceTimerKey, + new RestartModelGraceExpired(_activeCallId), + RestartModelGracePeriod); + } + } } - private void TryStopDurableApprovalWaits() + private bool TryStopDurableApprovalWaits() { if (!CanStopToolsForRestart()) - return; + return false; var task = _activeToolWorkTask!; _toolRestartStopRequested = true; _log.Info("Stopping a tool batch that waits only for durable approvals before restart drain"); _activeToolExecutionCts!.Cancel(); _ = ReportToolPipelineStopAsync(task, Self); + return true; } private bool CanStopToolsForRestart() @@ -5087,6 +5252,71 @@ private bool CanStopToolsForRestart() _toolApprovals.HasRecoverablePending); } + private void HandleRestartModelGraceExpired(RestartModelGraceExpired expired) + { + if (!_restartDrainRequested || _phase.Current != SessionPhase.Processing + || expired.CallId != _activeCallId || _activeLlmCts is null + || _activeInputIds.Count == 0) + return; + + _llmRestartStopRequested = true; + _watchdog.Stop(Timers); + _workingContextGeneration++; + _activeLlmCts.Cancel(); + + if (_activeLlmWorkTask is { IsCompleted: false } task) + { + _ = ReportModelStopAsync(task, Self, expired.CallId); + return; + } + + Self.Tell(new ModelStoppedForRestart(expired.CallId, null)); + } + + private static async Task ReportModelStopAsync(Task task, IActorRef actor, long callId) + { + try + { + await task.ConfigureAwait(false); + actor.Tell(new ModelStoppedForRestart(callId, null)); + } + catch (Exception ex) + { + actor.Tell(new ModelStoppedForRestart(callId, ex)); + } + } + + private void HandleModelStoppedForRestart(ModelStoppedForRestart stopped) + { + if (!_llmRestartStopRequested || stopped.CallId != _activeCallId) + return; + + _activeCallId++; + _activeLlmWorkTask = null; + CancelAndDisposeLlmCts(); + _llmRestartStopRequested = false; + + if (stopped.Failure is { } failure) + { + _log.Error(failure, "The model task failed during restart drain"); + } + + if (stopped.Failure is null && !_turnEmittedText) + { + CompleteModelStopForRestart(); + return; + } + + CloseActiveInputs(CompleteModelStopForRestart); + } + + private void CompleteModelStopForRestart() + { + MarkRestartReminderEligible(); + ClearBufferedMessagesForRestartDrain(); + TransitionTo(SessionPhase.Passivating); + } + private static async Task ReportToolPipelineStopAsync(Task task, IActorRef actor) { try @@ -5121,13 +5351,121 @@ private void HandleToolPipelineStoppedForRestart(ToolPipelineStoppedForRestart s TransitionTo(SessionPhase.Passivating); } + private void MarkRestartReminderEligible() + { + if (_state.PendingInputs.Count == 0) + return; + + if (_turnEmittedText && _activeInputIds.Count > 0) + { + _log.Warning( + "Restart drain cannot resume session {SessionId} because the interrupted turn emitted text.", + _sessionId.Value); + return; + } + + _restartInterruptionAt = _timeProvider.GetUtcNow(); + } + + private ReminderDefinition? BuildRestartReminder() + { + if (_restartInterruptionAt is not { } interruptedAt || _state.PendingInputs.Count == 0) + return null; + + if (!TryGetRestartContext(out var context, out var reason)) + { + _log.Warning( + "Restart drain cannot resume session {SessionId}: {Reason}.", + _sessionId.Value, + reason ?? "invalid stored authority"); + return null; + } + + if (context.ChannelType is not { } channelType) + { + _log.Warning( + "Restart drain cannot resume session {SessionId}: the turn has no channel type.", + _sessionId.Value); + return null; + } + + return new ReminderDefinition + { + Id = ReminderIdGenerator.Generate("restart resume"), + Title = "Resume after daemon restart", + Instructions = RestartReminderInstruction, + Schedule = new ReminderSchedule + { + Type = ReminderScheduleType.OneShot, + FireAt = interruptedAt + }, + Delivery = new ReminderDelivery + { + Kind = DeliveryKind.CurrentSession, + SessionId = _sessionId.Value, + OriginChannelType = channelType + }, + DeliveryRequired = true, + Audience = context.Audience, + Boundary = context.Boundary, + CreatedBy = "system", + CreatedAt = interruptedAt, + UpdatedAt = interruptedAt, + ExpiresAt = interruptedAt + RestartReminderLifetime + }; + } + + private bool TryGetRestartContext(out TurnContext context, out string? reason) + { + context = null!; + TurnContext? first = null; + reason = "no pending input"; + foreach (var pending in _state.PendingInputs) + { + if (!TurnContext.TryFromRecord(pending.TurnContext, out var candidate, out reason) + || candidate is null) + return false; + + if (first is null) + { + first = candidate; + continue; + } + + if (!HasSameRestartAuthority(first, candidate)) + { + reason = "pending inputs have different authority"; + context = first; + return false; + } + } + + context = first!; + return first is not null; + } + + private static bool HasSameRestartAuthority(TurnContext left, TurnContext right) + => left.SessionId == right.SessionId + && left.Audience == right.Audience + && left.Boundary == right.Boundary + && left.ChannelType == right.ChannelType + && left.RequesterSenderId == right.RequesterSenderId + && left.RequesterPrincipal == right.RequesterPrincipal + && left.Provenance == right.Provenance + && left.DefaultDeliveryTarget == right.DefaultDeliveryTarget + && left.RequestedDeliveryTarget == right.RequestedDeliveryTarget + && left.HasAdoptedContext == right.HasAdoptedContext + && left.HasThirdPartyAdoptedContext == right.HasThirdPartyAdoptedContext + && left.AdoptedSpeakerIds.SequenceEqual(right.AdoptedSpeakerIds, StringComparer.Ordinal) + && left.SupportsInteractiveApproval == right.SupportsInteractiveApproval; + private void ClearBufferedMessagesForRestartDrain() { if (_buffer.Count == 0) return; _log.Warning( - "Dropping {BufferCount} buffered message(s) because coordinated restart drain is completing.", + "Clearing {BufferCount} transient buffered message copy or copies; the journal keeps each accepted input.", _buffer.Count); _buffer.Clear(); } diff --git a/src/Netclaw.Actors/Sessions/SessionProtocol.Responses.cs b/src/Netclaw.Actors/Sessions/SessionProtocol.Responses.cs index 8d1d7a18b..135043854 100644 --- a/src/Netclaw.Actors/Sessions/SessionProtocol.Responses.cs +++ b/src/Netclaw.Actors/Sessions/SessionProtocol.Responses.cs @@ -4,6 +4,7 @@ // // ----------------------------------------------------------------------- using Netclaw.Actors.Protocol; +using Netclaw.Actors.Reminders; namespace Netclaw.Actors.Sessions; @@ -15,11 +16,18 @@ public static partial class SessionProtocol /// Acknowledged receipt of a command by the session actor. /// The command has been accepted and will be processed. /// - public sealed record CommandAck(SessionId SessionId) : ISessionResponse + public record CommandAck(SessionId SessionId) : ISessionResponse { public static CommandAck For(SessionId sessionId) => new(sessionId); } + /// + /// Reports that a session stopped and supplies any reminder that startup must register. + /// + public sealed record DaemonRestartPrepared( + SessionId SessionId, + ReminderDefinition? RestartReminder) : CommandAck(SessionId); + /// /// Negative acknowledgement — the command was rejected. /// diff --git a/src/Netclaw.Daemon.Tests/Services/DaemonRestartCoordinatorTests.cs b/src/Netclaw.Daemon.Tests/Services/DaemonRestartCoordinatorTests.cs index 0e548af4f..dcd5abe2d 100644 --- a/src/Netclaw.Daemon.Tests/Services/DaemonRestartCoordinatorTests.cs +++ b/src/Netclaw.Daemon.Tests/Services/DaemonRestartCoordinatorTests.cs @@ -13,6 +13,7 @@ using Netclaw.Actors.Channels; using Netclaw.Actors.Hosting; using Netclaw.Actors.Protocol; +using Netclaw.Actors.Reminders; using Netclaw.Configuration; using Netclaw.Daemon.Services; using Netclaw.Tests.Utilities; @@ -47,6 +48,26 @@ public async Task RequestConfigRestartAsync_drains_active_sessions_and_requests_ var (coordinator, drain) = CreateCoordinator( ["slack/C123.1", "slack/C123.2"], timeProvider: time); + drain.RestartReminder = new ReminderDefinition + { + Id = new ReminderId("restart-resume-test"), + Title = "Resume after daemon restart", + Instructions = "Resume the work that was interrupted by the daemon restart.", + Schedule = new ReminderSchedule + { + Type = ReminderScheduleType.OneShot, + FireAt = time.GetUtcNow() + }, + Delivery = new ReminderDelivery + { + Kind = DeliveryKind.CurrentSession, + SessionId = "slack/C123.1", + OriginChannelType = ChannelType.Slack + }, + Audience = TrustAudience.Team, + Boundary = TrustBoundary.Team, + ExpiresAt = time.GetUtcNow().AddMinutes(10) + }; var restart = coordinator.RequestConfigRestartAsync(CancellationToken.None); await drain.AllRequestsObserved; @@ -60,8 +81,7 @@ public async Task RequestConfigRestartAsync_drains_active_sessions_and_requests_ var manifest = await new RestartManifestStore(_paths).ReadAsync(CancellationToken.None); Assert.NotNull(manifest); - Assert.Equal(["slack/C123.1", "slack/C123.2"], manifest!.SessionIds); - Assert.Empty(manifest.TimedOutSessionIds); + Assert.Equal("restart-resume-test", Assert.Single(manifest!.RestartReminders).Id.Value); var alert = Assert.Single(_sink.Alerts); Assert.Equal("drained", alert.Context!["drainOutcome"]); @@ -89,9 +109,7 @@ public async Task RequestConfigRestartAsync_records_timed_out_sessions() Assert.True(_restartSignal.RestartRequested); Assert.True(_appLifetime.StopRequested); - var manifest = await new RestartManifestStore(_paths).ReadAsync(CancellationToken.None); - Assert.NotNull(manifest); - Assert.Equal(["slack/C123.2"], manifest!.TimedOutSessionIds); + Assert.Null(await new RestartManifestStore(_paths).ReadAsync(CancellationToken.None)); var alert = Assert.Single(_sink.Alerts); Assert.Equal("timeout", alert.Context!["drainOutcome"]); @@ -312,6 +330,8 @@ private sealed class DrainControl private int _requestCount; private int _acknowledgementCount; + public ReminderDefinition? RestartReminder { get; set; } + public DrainControl( IReadOnlyList activeSessionIds, IReadOnlyList timedOutSessionIds) @@ -346,7 +366,11 @@ public void AcknowledgeAll() if (_timedOutSessionIds.Contains(sessionId)) continue; - replyTo.Tell(CommandAck.For(new SessionId(sessionId))); + replyTo.Tell(new DaemonRestartPrepared( + new SessionId(sessionId), + RestartReminder?.Delivery.SessionId == sessionId + ? RestartReminder + : null)); if (Interlocked.Increment(ref _acknowledgementCount) == _expectedAcknowledgementCount) _allAcknowledgementsSent.TrySetResult(); } diff --git a/src/Netclaw.Daemon.Tests/Services/RestartRecoveryServiceTests.cs b/src/Netclaw.Daemon.Tests/Services/RestartRecoveryServiceTests.cs index 4a4f0b8f8..76ee4b111 100644 --- a/src/Netclaw.Daemon.Tests/Services/RestartRecoveryServiceTests.cs +++ b/src/Netclaw.Daemon.Tests/Services/RestartRecoveryServiceTests.cs @@ -7,14 +7,17 @@ using Akka.Actor; using Akka.Hosting; using Microsoft.Extensions.Logging.Abstractions; +using Microsoft.Extensions.Time.Testing; using Netclaw.Actors.Hosting; using Netclaw.Actors.Protocol; +using Netclaw.Actors.Reminders; +using Netclaw.Actors.Channels; using Netclaw.Configuration; using Netclaw.Daemon.Gateway; using Netclaw.Daemon.Services; using Netclaw.Tests.Utilities; using Xunit; -using static Netclaw.Actors.Sessions.SessionProtocol; +using static Netclaw.Actors.Reminders.ReminderProtocol; namespace Netclaw.Daemon.Tests.Services; @@ -32,44 +35,41 @@ public RestartRecoveryServiceTests() } [Fact] - public async Task StartAsync_warms_manifest_sessions_and_marks_catalog_active() + public async Task StartAsync_registers_only_fresh_reminders_and_accepts_an_existing_definition() { - var warmedSessions = new ConcurrentQueue(); - var actor = _system.ActorOf(Props.Create(() => new WarmSessionActor(warmedSessions))); + var time = new FakeTimeProvider(new DateTimeOffset(2026, 9, 23, 12, 0, 0, TimeSpan.Zero)); + var saved = new ConcurrentQueue(); + var reminderActor = _system.ActorOf(Props.Create(() => new ReminderActor(saved))); var manifestStore = new RestartManifestStore(_paths); var catalog = new SessionCatalogService( _paths, - TimeProvider.System, + time, new TestSessionStorageResolver(_paths), NullLogger.Instance); - var sessionId = new SessionId("slack/C123/1710000000.000001"); - - catalog.OnSessionActivated(sessionId, Netclaw.Actors.Channels.ChannelType.Slack); - catalog.OnSessionDeactivated(sessionId); await manifestStore.WriteAsync(new RestartManifest { - Reason = "config-reload", - RequestedAt = TimeProvider.System.GetUtcNow(), - SessionIds = [sessionId.Value], - TimedOutSessionIds = [sessionId.Value] + RestartReminders = + [ + CreateReminder("fresh", time.GetUtcNow().AddMinutes(10)), + CreateReminder("stale", time.GetUtcNow().AddSeconds(-1)) + ] }, CancellationToken.None); var sut = new RestartRecoveryService( manifestStore, - new StubRequiredActor(actor), + new StubRequiredActor(reminderActor), catalog, + time, NullLogger.Instance); await sut.StartAsync(CancellationToken.None); - Assert.True(warmedSessions.TryDequeue(out var warmed)); - Assert.Equal(sessionId, warmed!.SessionId); - Assert.Contains("last durable checkpoint", warmed.RestartNotice, StringComparison.Ordinal); + var command = Assert.Single(saved); + Assert.Equal("fresh", command.Definition.Id.Value); + Assert.True(command.Definition.Schedule.FireAt > time.GetUtcNow()); + Assert.Equal(TrustAudience.Personal, command.Authorization?.SourceAudience); Assert.Null(await manifestStore.ReadAsync(CancellationToken.None)); - - var entry = Assert.Single(catalog.ListRecent()); - Assert.Equal("active", entry.Status); } public void Dispose() @@ -79,7 +79,29 @@ public void Dispose() _dir.Dispose(); } - private sealed class StubRequiredActor : IRequiredActor + private static ReminderDefinition CreateReminder(string id, DateTimeOffset expiresAt) + => new() + { + Id = new ReminderId(id), + Title = "Resume after daemon restart", + Instructions = "Resume the work that was interrupted by the daemon restart.", + Schedule = new ReminderSchedule + { + Type = ReminderScheduleType.OneShot, + FireAt = expiresAt.AddMinutes(-10) + }, + Delivery = new ReminderDelivery + { + Kind = DeliveryKind.CurrentSession, + SessionId = "signalr/restart", + OriginChannelType = ChannelType.SignalR + }, + Audience = TrustAudience.Personal, + Boundary = TrustBoundary.Personal, + ExpiresAt = expiresAt + }; + + private sealed class StubRequiredActor : IRequiredActor { public StubRequiredActor(IActorRef actorRef) { @@ -92,14 +114,20 @@ public Task GetAsync(CancellationToken cancellationToken = default) => Task.FromResult(ActorRef); } - private sealed class WarmSessionActor : ReceiveActor + private sealed class ReminderActor : ReceiveActor { - public WarmSessionActor(ConcurrentQueue warmedSessions) + public ReminderActor(ConcurrentQueue saved) { - Receive(msg => + Receive(command => { - warmedSessions.Enqueue(msg); - Sender.Tell(CommandAck.For(msg.SessionId)); + saved.Enqueue(command); + Sender.Tell(new ReminderSavedResponse( + command.Definition.Id, + command.Definition.Title, + Success: false, + NextFire: null, + Error: ReminderSaveError.Conflict, + ErrorMessage: "The reminder already exists.")); }); } } diff --git a/src/Netclaw.Daemon/Program.cs b/src/Netclaw.Daemon/Program.cs index 19c9ca26b..4a5b05010 100644 --- a/src/Netclaw.Daemon/Program.cs +++ b/src/Netclaw.Daemon/Program.cs @@ -1126,6 +1126,7 @@ static void ConfigureDaemonServices( var sessionManager = registry.Get(); var ingressGate = sp.GetRequiredService(); var lifecycleNotifier = sp.GetRequiredService(); + var restartManifestStore = sp.GetRequiredService(); var drainLogger = sp.GetRequiredService().CreateLogger("Netclaw.Daemon.SessionDrain"); cs.AddTask(CoordinatedShutdown.PhaseBeforeServiceUnbind, "drain-llm-sessions", async () => @@ -1154,6 +1155,18 @@ static void ConfigureDaemonServices( drainDeadlineCts.Token, CancellationToken.None); + if (drainResult.RestartReminders.Count == 0) + { + await restartManifestStore.DeleteAsync(); + } + else + { + await restartManifestStore.WriteAsync(new RestartManifest + { + RestartReminders = [.. drainResult.RestartReminders] + }, CancellationToken.None); + } + lifecycleNotifier.NotifyShutdown("daemon-stop", drainResult.ToNotificationContext()); } catch (Exception ex) diff --git a/src/Netclaw.Daemon/Services/DaemonRestartCoordinator.cs b/src/Netclaw.Daemon/Services/DaemonRestartCoordinator.cs index ecf60bf36..fbddd06be 100644 --- a/src/Netclaw.Daemon/Services/DaemonRestartCoordinator.cs +++ b/src/Netclaw.Daemon/Services/DaemonRestartCoordinator.cs @@ -84,13 +84,10 @@ public async Task RequestConfigRestartAsync(CancellationToken cancellationToken) var manifest = new RestartManifest { - Reason = "config-reload", - RequestedAt = _timeProvider.GetUtcNow(), - SessionIds = [.. drainResult.AllSessionIds.Select(static id => id.Value)], - TimedOutSessionIds = [.. drainResult.TimedOutSessionIds.Select(static id => id.Value)] + RestartReminders = [.. drainResult.RestartReminders] }; - if (manifest.SessionIds.Count == 0) + if (manifest.RestartReminders.Count == 0) await _manifestStore.DeleteAsync(); else await _manifestStore.WriteAsync(manifest, cancellationToken); diff --git a/src/Netclaw.Daemon/Services/RestartManifestStore.cs b/src/Netclaw.Daemon/Services/RestartManifestStore.cs index e85cf014f..33a3d1b75 100644 --- a/src/Netclaw.Daemon/Services/RestartManifestStore.cs +++ b/src/Netclaw.Daemon/Services/RestartManifestStore.cs @@ -4,23 +4,18 @@ // // ----------------------------------------------------------------------- using System.Text.Json; +using Netclaw.Actors.Reminders; using Netclaw.Configuration; namespace Netclaw.Daemon.Services; public sealed record RestartManifest { - public required string Reason { get; init; } - - public required DateTimeOffset RequestedAt { get; init; } - - public required List SessionIds { get; init; } - - public List TimedOutSessionIds { get; init; } = []; + public List RestartReminders { get; init; } = []; } /// -/// Persists short-lived restart recovery state across a coordinated in-process host restart. +/// Persists short-lived reminders across a coordinated daemon restart. /// public sealed class RestartManifestStore { @@ -41,8 +36,16 @@ public async Task WriteAsync(RestartManifest manifest, CancellationToken cancell ArgumentNullException.ThrowIfNull(manifest); _paths.EnsureDirectoriesExist(); - await using var stream = File.Create(_paths.RestartManifestPath); - await JsonSerializer.SerializeAsync(stream, manifest, JsonOptions, cancellationToken); + var json = JsonSerializer.Serialize(manifest, JsonOptions); + await AtomicFile.WriteAllTextAsync( + _paths.RestartManifestPath, + json, + static path => + { + if (!OperatingSystem.IsWindows()) + File.SetUnixFileMode(path, UnixFileMode.UserRead | UnixFileMode.UserWrite); + }, + cancellationToken); } public async Task ReadAsync(CancellationToken cancellationToken) diff --git a/src/Netclaw.Daemon/Services/RestartRecoveryService.cs b/src/Netclaw.Daemon/Services/RestartRecoveryService.cs index 8470b1ccb..378161f3b 100644 --- a/src/Netclaw.Daemon/Services/RestartRecoveryService.cs +++ b/src/Netclaw.Daemon/Services/RestartRecoveryService.cs @@ -9,80 +9,114 @@ using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; using Netclaw.Actors.Hosting; -using Netclaw.Actors.Protocol; +using Netclaw.Actors.Reminders; using Netclaw.Daemon.Gateway; -using static Netclaw.Actors.Sessions.SessionProtocol; +using static Netclaw.Actors.Reminders.ReminderProtocol; namespace Netclaw.Daemon.Services; /// -/// Rehydrates the sessions that were active before coordinated restart began. +/// Registers short-lived reminders for work that a graceful stop interrupted. /// public sealed class RestartRecoveryService : IHostedService { - private const string RestartNotice = "The daemon restarted due to a configuration change. Recovery resumed from the last durable checkpoint."; + private static readonly TimeSpan AskTimeout = TimeSpan.FromSeconds(10); + private static readonly TimeSpan ReminderStartDelay = TimeSpan.FromMilliseconds(250); private readonly RestartManifestStore _manifestStore; - private readonly IRequiredActor _sessionManagerProvider; + private readonly IRequiredActor _reminderManagerProvider; private readonly SessionCatalogService _sessionCatalog; + private readonly TimeProvider _timeProvider; private readonly ILogger _logger; public RestartRecoveryService( RestartManifestStore manifestStore, - IRequiredActor sessionManagerProvider, + IRequiredActor reminderManagerProvider, SessionCatalogService sessionCatalog, + TimeProvider timeProvider, ILogger logger) { _manifestStore = manifestStore; - _sessionManagerProvider = sessionManagerProvider; + _reminderManagerProvider = reminderManagerProvider; _sessionCatalog = sessionCatalog; + _timeProvider = timeProvider; _logger = logger; } public async Task StartAsync(CancellationToken cancellationToken) { - // Reconcile stale 'active' sessions from the previous daemon lifetime. - // Must run before reading the manifest so that re-warmed sessions get - // a clean 'inactive' → 'active' transition via MarkSessionActive(). + // The reminder path activates only the sessions that have work to resume. _sessionCatalog.ReconcileStaleActiveSessions(); var manifest = await _manifestStore.ReadAsync(cancellationToken); if (manifest is null) return; - try + if (manifest.RestartReminders.Count == 0) { - if (manifest.SessionIds.Count == 0) - return; + await _manifestStore.DeleteAsync(); + return; + } + + var reminderManager = await _reminderManagerProvider.GetAsync(cancellationToken); + var now = _timeProvider.GetUtcNow(); + var results = await Task.WhenAll(manifest.RestartReminders.Select( + reminder => RegisterAsync(reminderManager, reminder, now, cancellationToken))); + + if (results.All(static result => result)) + await _manifestStore.DeleteAsync(); + + _logger.LogInformation( + "Restart recovery read {ReminderCount} restart reminder(s).", + manifest.RestartReminders.Count); + } - var sessionManager = await _sessionManagerProvider.GetAsync(cancellationToken); - foreach (var sessionIdValue in manifest.SessionIds) - { - var sessionId = new SessionId(sessionIdValue); + private async Task RegisterAsync( + IActorRef reminderManager, + ReminderDefinition stored, + DateTimeOffset now, + CancellationToken cancellationToken) + { + if (stored.ExpiresAt is not { } expiresAt || expiresAt <= now + ReminderStartDelay) + { + _logger.LogWarning( + "Restart reminder {ReminderId} expired before startup recovery; the session stays quiet.", + stored.Id.Value); + return true; + } - try - { - await sessionManager.Ask( - new WarmSession(sessionId, RestartNotice), - timeout: TimeSpan.FromSeconds(10), - cancellationToken: cancellationToken); + var reminder = stored with + { + Schedule = stored.Schedule with { FireAt = now + ReminderStartDelay } + }; + try + { + var result = await reminderManager.Ask( + new SaveReminderCommand( + reminder, + ReminderWriteMode.CreateOnly, + new ReminderAudienceAuthorizationContext( + reminder.Audience, + "restart manifest")), + timeout: AskTimeout, + cancellationToken: cancellationToken); - _sessionCatalog.MarkSessionActive(sessionId); - } - catch (Exception ex) - { - _logger.LogWarning(ex, "Failed to warm session {SessionId} during restart recovery.", sessionIdValue); - } - } + if (result.Success || result.Error == ReminderSaveError.Conflict) + return true; - _logger.LogInformation( - "Restart recovery warmed {SessionCount} session(s); {TimedOutCount} were previously timed out during drain.", - manifest.SessionIds.Count, - manifest.TimedOutSessionIds.Count); + _logger.LogWarning( + "Restart reminder {ReminderId} could not register: {Reason}", + reminder.Id.Value, + result.ErrorMessage ?? result.Error.ToString()); + return false; } - finally + catch (AskTimeoutException ex) { - await _manifestStore.DeleteAsync(); + _logger.LogWarning( + ex, + "Restart reminder {ReminderId} timed out during registration.", + reminder.Id.Value); + return false; } } diff --git a/src/Netclaw.Daemon/Services/SessionDrainHelper.cs b/src/Netclaw.Daemon/Services/SessionDrainHelper.cs index b9dbdc6b0..fc9bf52fa 100644 --- a/src/Netclaw.Daemon/Services/SessionDrainHelper.cs +++ b/src/Netclaw.Daemon/Services/SessionDrainHelper.cs @@ -8,6 +8,7 @@ using Microsoft.Extensions.Logging; using Netclaw.Actors.Hosting; using Netclaw.Actors.Protocol; +using Netclaw.Actors.Reminders; using static Netclaw.Actors.Sessions.SessionProtocol; namespace Netclaw.Daemon.Services; @@ -57,7 +58,7 @@ public static async Task DrainAsync( { try { - var ack = await sessionManager.Ask( + var ack = await sessionManager.Ask( new PrepareForDaemonRestart(sessionId, reason), timeout: Timeout.InfiniteTimeSpan, cancellationToken: operationCancellationToken); @@ -72,11 +73,11 @@ public static async Task DrainAsync( reason); } - return new DrainOutcome(sessionId, drained); + return new DrainOutcome(sessionId, drained, drained ? ack.RestartReminder : null); } catch (OperationCanceledException) when (!callerCancellationToken.IsCancellationRequested) { - return new DrainOutcome(sessionId, false); + return new DrainOutcome(sessionId, false, null); } catch (OperationCanceledException) { @@ -85,7 +86,7 @@ public static async Task DrainAsync( catch (Exception ex) { logger.LogWarning(ex, "Failed to drain session {SessionId} before shutdown.", sessionId.Value); - return new DrainOutcome(sessionId, false); + return new DrainOutcome(sessionId, false, null); } }).ToArray(); @@ -106,24 +107,33 @@ public static async Task DrainAsync( string.Join(", ", timedOut.Select(static id => id.Value))); } - return new DrainResult(sessionIds, drained, timedOut); + var reminders = outcomes + .Where(static outcome => outcome.RestartReminder is not null) + .Select(static outcome => outcome.RestartReminder!) + .ToArray(); + return new DrainResult(sessionIds, drained, timedOut, reminders); } - internal sealed record DrainOutcome(SessionId SessionId, bool Drained); + internal sealed record DrainOutcome( + SessionId SessionId, + bool Drained, + ReminderDefinition? RestartReminder); internal sealed record DrainResult( IReadOnlyList AllSessionIds, IReadOnlyList DrainedSessionIds, - IReadOnlyList TimedOutSessionIds) + IReadOnlyList TimedOutSessionIds, + IReadOnlyList RestartReminders) { - public static readonly DrainResult Empty = new([], [], []); + public static readonly DrainResult Empty = new([], [], [], []); public Dictionary ToNotificationContext() => new() { ["drainOutcome"] = TimedOutSessionIds.Count == 0 ? "drained" : "timeout", ["activeSessions"] = AllSessionIds.Count.ToString(CultureInfo.InvariantCulture), ["drainedSessions"] = DrainedSessionIds.Count.ToString(CultureInfo.InvariantCulture), - ["timedOutSessions"] = TimedOutSessionIds.Count.ToString(CultureInfo.InvariantCulture) + ["timedOutSessions"] = TimedOutSessionIds.Count.ToString(CultureInfo.InvariantCulture), + ["restartReminders"] = RestartReminders.Count.ToString(CultureInfo.InvariantCulture) }; } } From 2b0430251dedbefc49994324108dd55ccbb3596f Mon Sep 17 00:00:00 2001 From: Aaron Stannard Date: Thu, 24 Sep 2026 04:02:35 +0000 Subject: [PATCH 2/3] Make the compaction drain test deterministic --- .../Sessions/LlmSessionIntegrationTests.cs | 52 ++++++++++++++----- 1 file changed, 40 insertions(+), 12 deletions(-) diff --git a/src/Netclaw.Actors.Tests/Sessions/LlmSessionIntegrationTests.cs b/src/Netclaw.Actors.Tests/Sessions/LlmSessionIntegrationTests.cs index 1a216c9ec..648ef697b 100644 --- a/src/Netclaw.Actors.Tests/Sessions/LlmSessionIntegrationTests.cs +++ b/src/Netclaw.Actors.Tests/Sessions/LlmSessionIntegrationTests.cs @@ -1913,13 +1913,6 @@ public async Task Compacting_session_rejects_new_work_and_passivates_after_compa var sessionId = new SessionId("test-channel/restart-drain-compacting"); var sessionManager = ActorRegistry.Get(); var subscriber = CreateTestProbe("restart-drain-compacting-sub"); - _fakeChatClient.UsageOverride = new UsageDetails - { - InputTokenCount = 200_000, - OutputTokenCount = 1, - TotalTokenCount = 200_001 - }; - _fakeChatClient.HangingObservationCallsRemaining = 1; await sessionManager.Ask(new JoinSession(subscriber) { @@ -1928,6 +1921,26 @@ await sessionManager.Ask(new JoinSession(subscriber) }, TimeSpan.FromSeconds(3), cancellationToken: TestContext.Current.CancellationToken); await subscriber.ExpectMsgAsync(cancellationToken: TestContext.Current.CancellationToken); + for (var turn = 1; turn <= 3; turn++) + { + await sessionManager.Ask(new SendUserMessage + { + SessionId = sessionId, + Content = $"Build compaction history {turn}" + }, TimeSpan.FromSeconds(3), cancellationToken: TestContext.Current.CancellationToken); + + await subscriber.ExpectMsgAsync(TimeSpan.FromSeconds(3), cancellationToken: TestContext.Current.CancellationToken); + await subscriber.ExpectMsgAsync(TimeSpan.FromSeconds(3), cancellationToken: TestContext.Current.CancellationToken); + } + + _fakeChatClient.UsageOverride = new UsageDetails + { + InputTokenCount = 200_000, + OutputTokenCount = 1, + TotalTokenCount = 200_001 + }; + _fakeChatClient.HangingObservationCallsRemaining = 1; + await sessionManager.Ask(new SendUserMessage { SessionId = sessionId, @@ -1937,17 +1950,29 @@ await sessionManager.Ask(new SendUserMessage await subscriber.ExpectMsgAsync(TimeSpan.FromSeconds(3), cancellationToken: TestContext.Current.CancellationToken); await subscriber.ExpectMsgAsync(TimeSpan.FromSeconds(3), cancellationToken: TestContext.Current.CancellationToken); + await AwaitAssertAsync(() => + { + Assert.Contains(_fakeChatClient.ReceivedMessages, conversation => + conversation.Any(message => + message.Role == Microsoft.Extensions.AI.ChatRole.System + && message.Text?.Contains("You are a session summarizer", StringComparison.Ordinal) == true)); + return Task.CompletedTask; + }, TimeSpan.FromSeconds(3), TimeSpan.FromMilliseconds(50), cancellationToken: TestContext.Current.CancellationToken); + var escapedId = Uri.EscapeDataString(sessionId.Value); var child = await Sys.ActorSelection($"/user/session-manager/{escapedId}").ResolveOne(TimeSpan.FromSeconds(3), TestContext.Current.CancellationToken); Watch(child); - var drainTask = sessionManager.Ask(new PrepareForDaemonRestart(sessionId, "config-reload"), TimeSpan.FromSeconds(5), cancellationToken: TestContext.Current.CancellationToken); - - var nack = await sessionManager.Ask(new SendUserMessage + sessionManager.Tell(new PrepareForDaemonRestart(sessionId, "config-reload"), TestActor); + sessionManager.Tell(new SendUserMessage { SessionId = sessionId, Content = "Should be rejected during compaction" - }, TimeSpan.FromSeconds(3), cancellationToken: TestContext.Current.CancellationToken); + }, TestActor); + + var nack = await ExpectMsgAsync( + TimeSpan.FromSeconds(3), + cancellationToken: TestContext.Current.CancellationToken); Assert.Equal(SessionIngressGate.RestartInProgressMessage, nack.Reason); @@ -1956,7 +1981,10 @@ await sessionManager.Ask(new SendUserMessage Cause = new InvalidOperationException("test compaction completion") }); - Assert.Equal(sessionId, (await drainTask).SessionId); + var drainAck = await ExpectMsgAsync( + TimeSpan.FromSeconds(5), + cancellationToken: TestContext.Current.CancellationToken); + Assert.Equal(sessionId, drainAck.SessionId); await ExpectTerminatedAsync(child, TimeSpan.FromSeconds(5), cancellationToken: TestContext.Current.CancellationToken); Assert.Contains(sessionId.Value, _lifecycleObserver.DeactivatedSessionIds); } From 1264be5589318ca27b4a336ebe129f68a0d82e29 Mon Sep 17 00:00:00 2001 From: Aaron Stannard Date: Thu, 24 Sep 2026 12:37:17 +0000 Subject: [PATCH 3/3] Batch shell mutation targets by project --- TOOLING.md | 3 +- .../run-shell-command-analysis-mutations.sh | 166 +++++------------- 2 files changed, 44 insertions(+), 125 deletions(-) diff --git a/TOOLING.md b/TOOLING.md index 4cfc102e8..102f9b063 100644 --- a/TOOLING.md +++ b/TOOLING.md @@ -222,7 +222,8 @@ Two status-parameter mutants test the rule that only bare `$?` can preserve reusable candidates. The focused test rejects other unknown output data. The new target took 49 seconds after package restore. -The local run on 2026-09-18 took about 22 minutes. CI allows 30 minutes for +The script groups targets by source project. Stryker analyzes each source project once. +The local run on 2026-09-24 took under four minutes. CI allows 30 minutes for hosted-runner variance and report upload. The report directory is `artifacts/stryker/shell-command-analysis`. diff --git a/scripts/run-shell-command-analysis-mutations.sh b/scripts/run-shell-command-analysis-mutations.sh index e399817e8..0219b1831 100755 --- a/scripts/run-shell-command-analysis-mutations.sh +++ b/scripts/run-shell-command-analysis-mutations.sh @@ -28,19 +28,23 @@ find_span() { ' "$context_marker" "$start_marker" "$end_marker" "$source_file" } -run_target() { +run_group() { local config_file="$1" - local source_name="$2" - local span_start="$3" - local span_end="$4" - local target_output="$5" - local expected_count="$6" + local target_output="$2" + local expected_count="$3" + shift 3 + + local mutate_args=() + local mutation + for mutation in "$@"; do + mutate_args+=(--mutate "$mutation") + done ( cd "$test_project" dotnet stryker \ --config-file "$config_file" \ - --mutate "$source_name{$span_start..$span_end}" \ + "${mutate_args[@]}" \ --output "$target_output" \ --skip-version-check ) @@ -59,6 +63,8 @@ run_target() { fi } +security_mutations=() + analysis_file="$repo_root/src/Netclaw.Security/ShellCommandAnalysis.cs" read -r region_start region_end < <( find_span \ @@ -67,13 +73,7 @@ read -r region_start region_end < <( "=> argument.Argument.Kind == ArgKind.DynamicSkip" \ "&& accountedRegionArguments.Contains(argument.Element);" ) -run_target \ - "stryker-shell-command-analysis.json" \ - "ShellCommandAnalysis.cs" \ - "$region_start" \ - "$region_end" \ - "$output_path/execution-region" \ - 2 +security_mutations+=("ShellCommandAnalysis.cs{$region_start..$region_end}") policy_file="$repo_root/src/Netclaw.Security/ShellCommandPolicy.cs" read -r gate_start gate_end < <( @@ -83,13 +83,7 @@ read -r gate_start gate_end < <( "var denyOnlyDecision = EvaluateDenyOnlyClauses" \ "return denyOnlyDecision;" ) -run_target \ - "stryker-shell-command-analysis.json" \ - "ShellCommandPolicy.cs" \ - "$gate_start" \ - "$gate_end" \ - "$output_path/deny-only-gate" \ - 1 +security_mutations+=("ShellCommandPolicy.cs{$gate_start..$gate_end}") read -r trust_start trust_end < <( find_span \ @@ -98,13 +92,7 @@ read -r trust_start trust_end < <( "if (!tokens[i].IsKnown)" \ "return false;" ) -run_target \ - "stryker-shell-command-analysis.json" \ - "ShellCommandPolicy.cs" \ - "$trust_start" \ - "$trust_end" \ - "$output_path/deny-only-token-trust" \ - 2 +security_mutations+=("ShellCommandPolicy.cs{$trust_start..$trust_end}") tree_policy_file="$repo_root/src/Netclaw.Security/ShellFileSystemTreeAccessPolicy.cs" read -r tree_decision_start tree_decision_end < <( @@ -114,13 +102,7 @@ read -r tree_decision_start tree_decision_end < <( "var accesses = command.FileSystemTreeAccesses;" \ "return !CanUseReusableApproval(command, accesses[0]);" ) -run_target \ - "stryker-shell-command-analysis.json" \ - "ShellFileSystemTreeAccessPolicy.cs" \ - "$tree_decision_start" \ - "$tree_decision_end" \ - "$output_path/tree-decision" \ - 7 +security_mutations+=("ShellFileSystemTreeAccessPolicy.cs{$tree_decision_start..$tree_decision_end}") read -r tree_root_start tree_root_end < <( find_span \ @@ -129,13 +111,7 @@ read -r tree_root_start tree_root_end < <( "if (!IsReusableTraversal(access.Traversal)" \ "&& string.Equals(root.Value, cwd.Value, StringComparison.Ordinal);" ) -run_target \ - "stryker-shell-command-analysis.json" \ - "ShellFileSystemTreeAccessPolicy.cs" \ - "$tree_root_start" \ - "$tree_root_end" \ - "$output_path/tree-root" \ - 10 +security_mutations+=("ShellFileSystemTreeAccessPolicy.cs{$tree_root_start..$tree_root_end}") read -r root_match_start root_match_end < <( find_span \ @@ -144,13 +120,7 @@ read -r root_match_start root_match_end < <( "=> root switch" \ "_ => false" ) -run_target \ - "stryker-shell-command-analysis.json" \ - "ShellFileSystemTreeAccessPolicy.cs" \ - "$root_match_start" \ - "$root_match_end" \ - "$output_path/tree-root-correspondence" \ - 3 +security_mutations+=("ShellFileSystemTreeAccessPolicy.cs{$root_match_start..$root_match_end}") read -r leaf_call_start leaf_call_end < <( find_span \ @@ -159,13 +129,7 @@ read -r leaf_call_start leaf_call_end < <( "!IsReusableLeafPattern(pattern)" \ "!IsReusableLeafPattern(pattern)" ) -run_target \ - "stryker-shell-command-analysis.json" \ - "ShellFileSystemTreeAccessPolicy.cs" \ - "$leaf_call_start" \ - "$leaf_call_end" \ - "$output_path/tree-root-leaf-call" \ - 1 +security_mutations+=("ShellFileSystemTreeAccessPolicy.cs{$leaf_call_start..$leaf_call_end}") read -r leaf_shape_start leaf_shape_end < <( find_span \ @@ -174,13 +138,7 @@ read -r leaf_shape_start leaf_shape_end < <( "internal static bool IsReusableLeafPattern(" \ "return separator != 0 && !IsIncompleteUncLeafPattern(pattern, separator);" ) -run_target \ - "stryker-shell-command-analysis.json" \ - "ShellFileSystemTreeAccessPolicy.cs" \ - "$leaf_shape_start" \ - "$leaf_shape_end" \ - "$output_path/tree-root-leaf-shape" \ - 15 +security_mutations+=("ShellFileSystemTreeAccessPolicy.cs{$leaf_shape_start..$leaf_shape_end}") read -r traversal_start traversal_end < <( find_span \ @@ -189,13 +147,7 @@ read -r traversal_start traversal_end < <( "=> Enum.IsDefined(traversal)" \ "or ShellTreeTraversalMode.RecursiveWithoutFollowingLinks;" ) -run_target \ - "stryker-shell-command-analysis.json" \ - "ShellFileSystemTreeAccessPolicy.cs" \ - "$traversal_start" \ - "$traversal_end" \ - "$output_path/tree-traversal" \ - 2 +security_mutations+=("ShellFileSystemTreeAccessPolicy.cs{$traversal_start..$traversal_end}") read -r nonfile_start nonfile_end < <( find_span \ @@ -204,13 +156,7 @@ read -r nonfile_start nonfile_end < <( "if (argument.Argument.IsPath" \ "ShellValueDomain.Unknown => HasMatchingIntegerRange(argument)," ) -run_target \ - "stryker-shell-command-analysis.json" \ - "ShellCommandAnalysis.cs" \ - "$nonfile_start" \ - "$nonfile_end" \ - "$output_path/audited-nonfilesystem" \ - 5 +security_mutations+=("ShellCommandAnalysis.cs{$nonfile_start..$nonfile_end}") read -r range_start range_end < <( find_span \ @@ -219,13 +165,7 @@ read -r range_start range_end < <( "=> argument.Value is" \ "< MaximumReviewedIntegerRangeCardinality;" ) -run_target \ - "stryker-shell-command-analysis.json" \ - "ShellCommandAnalysis.cs" \ - "$range_start" \ - "$range_end" \ - "$output_path/integer-range" \ - 11 +security_mutations+=("ShellCommandAnalysis.cs{$range_start..$range_end}") read -r status_start status_end < <( find_span \ @@ -234,13 +174,7 @@ read -r status_start status_end < <( "argument.Argument.Raw == \"\$?\"" \ "argument.Argument.Raw == \"\$?\"" ) -run_target \ - "stryker-shell-command-analysis.json" \ - "ShellCommandAnalysis.cs" \ - "$status_start" \ - "$status_end" \ - "$output_path/status-parameter" \ - 2 +security_mutations+=("ShellCommandAnalysis.cs{$status_start..$status_end}") matcher_file="$repo_root/src/Netclaw.Security/IToolApprovalMatcher.cs" read -r candidate_start candidate_end < <( @@ -250,13 +184,7 @@ read -r candidate_start candidate_end < <( "if (!result.IsResolved" \ "return [];" ) -run_target \ - "stryker-shell-command-analysis.json" \ - "IToolApprovalMatcher.cs" \ - "$candidate_start" \ - "$candidate_end" \ - "$output_path/reusable-candidates" \ - 5 +security_mutations+=("IToolApprovalMatcher.cs{$candidate_start..$candidate_end}") read -r messy_start messy_end < <( find_span \ @@ -265,13 +193,15 @@ read -r messy_start messy_end < <( "if (!analysis.IsResolved" \ "return true;" ) -run_target \ +security_mutations+=("IToolApprovalMatcher.cs{$messy_start..$messy_end}") + +run_group \ "stryker-shell-command-analysis.json" \ - "IToolApprovalMatcher.cs" \ - "$messy_start" \ - "$messy_end" \ - "$output_path/messy-analysis" \ - 6 + "$output_path/security" \ + 72 \ + "${security_mutations[@]}" + +actor_mutations=() tool_policy_file="$repo_root/src/Netclaw.Actors/Tools/ToolAccessPolicy.cs" read -r mode_start mode_end < <( @@ -281,13 +211,7 @@ read -r mode_start mode_end < <( "=> configuredMode == ToolApprovalMode.Auto" \ ": configuredMode;" ) -run_target \ - "stryker-config.json" \ - "Tools/ToolAccessPolicy.cs" \ - "$mode_start" \ - "$mode_end" \ - "$output_path/exact-tree-mode" \ - 4 +actor_mutations+=("Tools/ToolAccessPolicy.cs{$mode_start..$mode_end}") path_facts_file="$repo_root/src/Netclaw.Actors/Tools/ShellPolicyPathFacts.cs" read -r path_fact_start path_fact_end < <( @@ -297,13 +221,7 @@ read -r path_fact_start path_fact_end < <( "facts.Add(CreateFact(" \ "ShellPathShape.Unknown));" ) -run_target \ - "stryker-config.json" \ - "Tools/ShellPolicyPathFacts.cs" \ - "$path_fact_start" \ - "$path_fact_end" \ - "$output_path/tree-path-fact" \ - 1 +actor_mutations+=("Tools/ShellPolicyPathFacts.cs{$path_fact_start..$path_fact_end}") reviewed_file="$repo_root/src/Netclaw.Actors/Tools/ReviewedSafeShellPolicy.cs" read -r reviewed_start reviewed_end < <( @@ -313,10 +231,10 @@ read -r reviewed_start reviewed_end < <( "if (ShellFileSystemTreeAccessPolicy.RequiresExactApproval(" \ "return false;" ) -run_target \ +actor_mutations+=("Tools/ReviewedSafeShellPolicy.cs{$reviewed_start..$reviewed_end}") + +run_group \ "stryker-config.json" \ - "Tools/ReviewedSafeShellPolicy.cs" \ - "$reviewed_start" \ - "$reviewed_end" \ - "$output_path/reviewed-safe-tree" \ - 4 + "$output_path/actors" \ + 9 \ + "${actor_mutations[@]}"