From 5f95cd6538184ae05c56fa2045b17fc29893f559 Mon Sep 17 00:00:00 2001 From: Aaron Stannard Date: Fri, 18 Sep 2026 22:37:05 +0000 Subject: [PATCH] Persist accepted session input before acknowledgment --- docs/spec/SPEC-011-daemon-architecture.md | 7 + .../.openspec.yaml | 2 + .../resume-interrupted-sessions/design.md | 104 ++++++ .../resume-interrupted-sessions/proposal.md | 37 +++ .../specs/daemon-container/spec.md | 35 ++ .../specs/session-resume/spec.md | 103 ++++++ .../resume-interrupted-sessions/tasks.md | 26 ++ .../Protocol/SerializationRoundTripTests.cs | 62 ++++ .../Sessions/LlmSessionIntegrationTests.cs | 72 +++++ .../Sessions/SessionStateTests.cs | 32 ++ src/Netclaw.Actors/Channels/TurnContext.cs | 2 +- .../Protocol/SessionSnapshot.cs | 6 + .../Serialization/NetclawProtoMapper.cs | 85 ++++- .../NetclawProtobufSerializer.cs | 8 + .../Protos/netclaw_messages.proto | 22 ++ .../Sessions/LlmSessionActor.cs | 300 +++++++++++------- .../Sessions/SessionProtocol.Commands.cs | 3 + .../Sessions/SessionProtocol.Events.cs | 45 +++ src/Netclaw.Actors/Sessions/SessionState.cs | 38 +++ 19 files changed, 868 insertions(+), 121 deletions(-) create mode 100644 openspec/changes/resume-interrupted-sessions/.openspec.yaml create mode 100644 openspec/changes/resume-interrupted-sessions/design.md create mode 100644 openspec/changes/resume-interrupted-sessions/proposal.md create mode 100644 openspec/changes/resume-interrupted-sessions/specs/daemon-container/spec.md create mode 100644 openspec/changes/resume-interrupted-sessions/specs/session-resume/spec.md create mode 100644 openspec/changes/resume-interrupted-sessions/tasks.md diff --git a/docs/spec/SPEC-011-daemon-architecture.md b/docs/spec/SPEC-011-daemon-architecture.md index d792879ec..d40d0928c 100644 --- a/docs/spec/SPEC-011-daemon-architecture.md +++ b/docs/spec/SPEC-011-daemon-architecture.md @@ -210,6 +210,13 @@ The daemon writes its PID to `~/.netclaw/netclaw.pid` for lifecycle management. shutdown. The daemon handles SIGTERM by draining active sessions and stopping the actor system cleanly. +The session journals each accepted input before it acknowledges the source. +The record retains the text, media, source message ID, and original authority. +A completed reply, a started tool batch, or a terminal failure consumes the +input ID. The actor restores unconsumed records from the journal after a cold +start. A retry with the same stable source message ID does not add a second +record. A source without a stable ID cannot use this deduplication rule. + During drain, a session can stop a tool task that waits only for durable approval prompts. The session waits for the tool task to stop before it acknowledges drain. Its journal retains the prompts and completed sibling diff --git a/openspec/changes/resume-interrupted-sessions/.openspec.yaml b/openspec/changes/resume-interrupted-sessions/.openspec.yaml new file mode 100644 index 000000000..f2cbbe6a6 --- /dev/null +++ b/openspec/changes/resume-interrupted-sessions/.openspec.yaml @@ -0,0 +1,2 @@ +schema: spec-driven +created: 2026-09-18 diff --git a/openspec/changes/resume-interrupted-sessions/design.md b/openspec/changes/resume-interrupted-sessions/design.md new file mode 100644 index 000000000..9cbb7ceaf --- /dev/null +++ b/openspec/changes/resume-interrupted-sessions/design.md @@ -0,0 +1,104 @@ +## Context + +See [proposal.md](proposal.md) for the user outcome. The actor currently acknowledges a model-only request and each buffered request before a journal event stores that input. `TurnRecorded` stores one user message after the reply. A stop clears the actor-local queue. `RestartRecoveryService` warms prior sessions but does not start a model call. The manifest exists only for a configuration restart. + +The [engineering glossary](../../../docs/spec/GLOSSARY.md) defines shared session terms. `TurnContextRecord` already stores the authority fields that a recovered tool approval uses. The channel bindings own output subscribers. A warm session alone cannot deliver a recovered reply to Slack, Discord, or Mattermost. + +## Goals / Non-Goals + +**Goals:** + +- A journal record precedes each accepted input acknowledgment. +- The journal and snapshots retain accepted input until a durable terminal event consumes it. +- A graceful stop confirms a model task stop before it grants a resume candidate. +- A short restart resumes only eligible work under its original authority. +- One model call receives the ordered accepted queue after the original turn ends. + +**Non-Goals:** + +- A crash does not create an automatic resume candidate. +- A tool with an uncertain effect does not run again without new user input. +- A new reminder definition does not represent interrupted work. +- A channel without a confirmed output route does not get an automatic model call. + +## Decisions + +### D1. The journal owns admitted input + +Add `InputAdmitted` with an input ID, a stable source message ID when one exists, text, media, executable text, source IDs, received time, and a `TurnContextRecord`. The actor persists it before the ack. The record excludes actor refs and raw `MessageSource`. Reject a missing or invalid authority record before the model call. The actor uses a source ID only within its channel and session scope for deduplication. Sources without a stable ID cannot claim retry deduplication. + +The actor keeps admitted records in an ordered state list. `TurnRecorded` identifies every input in its model call. `ToolBatchStarted` identifies the inputs before any tool runs. A durable terminal failure event consumes a model-only input. This prevents a failed turn from starting again after a short stop. The live callback and journal replay use the same input IDs. The snapshot stores pending records and a bounded recent source-ID ledger. The actor skips a snapshot while a pending input also appears only in transient history. + +The old `TurnRecorded.UserMessage` remains for journal compatibility. Replay uses admitted records when the new consumed-ID list exists. Replay uses the old field for earlier journal records. This keeps old sessions readable. + +Alternative: Store the input only in a stop manifest. That misses a stop after an ack and before manifest creation. It also duplicates session authority outside the journal. + +### D2. Drain returns a classified result + +`PrepareForDaemonRestart` asks the actor to stop new work. A live model call gets a short grace. If it completes, the actor handles its normal result. If it remains active, the actor cancels its token and awaits the exact `SessionLlmInvoker` task. A call ID rejects stale task results. The actor grants a candidate only when the task stopped, no tool batch started for that turn, and no user-visible text escaped. A confirmed queue can form a candidate after a completed reply. Approval-only work follows the existing durable approval path without a candidate. + +The actor returns a typed drain result with the candidate input IDs or a blocked reason. `SessionDrainHelper` collects results. The daemon writes one manifest after all drain replies. Both the config restart and normal coordinated stop use this path. The manifest write is atomic. It stores an absolute ten-minute deadline per candidate. A timeout or failure produces a warning and no executable candidate. + +Alternative: Infer candidates from actor phase or session IDs. Phase does not prove task cancellation, and a warm session can contain a completed turn. + +### D3. Recovery checks the manifest and the session journal + +The recovery service reconciles the session catalog and warms listed sessions. After channel services start, it prepares the output route for a candidate. It then sends an internal resume request with the candidate IDs and absolute deadline. The actor checks the deadline again. It checks that the IDs still match its pending journal records, no newer turn started, and the recorded authority is valid. The actor resumes the original turn first. It then sends one ordered follow-up model call for accepted queued input. + +The service uses existing channel gateway `StartProactiveThread` messages to rebuild Slack, Discord, and Mattermost bindings. Those messages must confirm the output subscriber before the resume request. A TUI or SignalR session needs a live attachment; otherwise the candidate stays blocked and the service reports it. The service must retain a candidate until it succeeds or expires. It does not reset the deadline after another process start. + +Alternative: Schedule a generic reminder. That creates a new turn, can abandon a parked approval, and may run under authority derived from the reminder rather than the original input. + +### D4. Authority and output safety gate + +Each admitted record carries its original `TurnContextRecord`. The actor reconstructs it with `TurnContext.TryFromRecord`. It never derives a recovery authority from a session ID. A queue with incompatible boundaries remains durable but does not auto-resume. A queue with compatible authority uses the narrowest audience and never widens tool access. The actor rejects a candidate after partial text output because the user might have seen that text. It also rejects a candidate after a tool batch starts because a tool might have an external effect. + +Positive example: A Slack user sends one request. The model call stops before text or tools. The next start restores the Slack binding and resumes that input under its stored personal boundary. + +Negative example: A shell tool starts before a stop. The daemon reports the pending work and does not replay the shell call. + +### D5. Delivery order and failure behavior + +The candidate deadline starts at interruption. The recovery service checks it before each route attempt. The actor checks it before the model call. An expired candidate causes one warning and leaves the journal record available for forensic inspection or a new user turn. A newer input supersedes an old candidate. The manifest remains on disk until each candidate reaches a terminal recovery decision. The service records a blocked route or invalid record with a clear diagnostic. + +## Ordered flow + +This diagram is schematic. It omits the channel ACL and persistence callbacks. + +```text +source -> session: SendUserMessage +session -> journal: InputAdmitted(input ID, authority, content) +journal -> session: persisted +session -> source: CommandAck +session -> model: original turn +stop -> session: PrepareForDaemonRestart +session -> model: cancel if grace ends +model -> session: task stopped +session -> stop: candidate(input IDs) or blocked reason +stop -> manifest: atomic write(deadline) +start -> channel: restore output binding +channel -> start: binding ready +start -> session: ResumeInterruptedTurn(input IDs, deadline) +session -> journal: validate pending input and authority +session -> model: original turn, then one queued call +``` + +## Risks / Trade-offs + +- [Input accepted during a write failure] → The actor sends a nack and starts no model call. +- [Snapshot skips an unconsumed input] → Replay starts before that snapshot, then restores the pending record. +- [Provider ignores cancellation] → Drain reaches its existing deadline and records no candidate. +- [A reply reaches the user before cancellation] → The actor reports a blocked candidate and avoids duplicate text. +- [A tool starts before cancellation] → The actor reports a blocked candidate and avoids duplicate effects. +- [A route is absent after restart] → The service reports the candidate and leaves the agent quiet. +- [Old manifest or a second process start] → The absolute deadline and input IDs prevent a new or duplicate model call. +- [A source has no stable message ID] → The actor records a unique admission ID but cannot deduplicate a source retry. + +## Migration Plan + +1. Add the journal and snapshot fields with new protobuf tags and a new manifest for `InputAdmitted`. +2. Keep old event fields and old journal readers intact. +3. Add actor admission and recovery tests before enabling automatic wakeup. +4. Add the normal stop manifest and recovery request after the actor contract passes. +5. Verify a graceful stop and a short restart in an isolated daemon with a persistent home. +6. Roll back by disabling automatic wakeup in the new binary. Keep the new journal decoder so accepted input remains readable. diff --git a/openspec/changes/resume-interrupted-sessions/proposal.md b/openspec/changes/resume-interrupted-sessions/proposal.md new file mode 100644 index 000000000..009df2dfa --- /dev/null +++ b/openspec/changes/resume-interrupted-sessions/proposal.md @@ -0,0 +1,37 @@ +## Why + +A graceful stop can acknowledge input before the journal stores it. A model call can also delay drain until the global stop limit. A short restart then leaves accepted work quiet or loses queued context. + +Source: [PRD-001 FR-003 and FR-016](../../../docs/prd/PRD-001-netclaw-mvp.md). + +## What Changes + +- Persist each accepted input, its order, media, source ID, delivery context, and authority before the input ack. +- Record the input IDs that a completed turn consumes. Restore accepted input from the journal after cold recovery. +- Stop an eligible model call after a short grace and wait for its task to stop. +- Save one bounded resume candidate after any graceful stop. The deadline expires ten minutes after interruption. +- Resume the original turn under its recorded authority. Deliver the accepted queue in one ordered follow-up model call. +- Leave a completed reply with no accepted queue quiet. Block auto-resume for approvals, uncertain tool effects, and partial replies. +- Report blocked work to the operator. +- Correct the container contract for stop signals and the persistent state volume. + +The MVP change excludes automatic crash recovery, replay of uncertain tool effects, and a new reminder schedule. It does not shorten the global stop limit. + +## Capabilities + +### New Capabilities + +None. + +### Modified Capabilities + +- `session-resume`: Durable input admission and bounded restart recovery for confirmed interruptions. +- `daemon-container`: The entrypoint forwards stop signals, and the state volume retains restart data. + +## Impact + +This change affects the session actor, journal events, protobuf mappings, restart manifest, recovery service, and container contract. It adds no user configuration property. + +### Security and operational impact + +The actor restores the original authority from a durable record. It rejects an incomplete record and does not replay a tool with an uncertain effect. A stale candidate expires without an agent turn. A process stop must complete its journal and manifest writes before exit. diff --git a/openspec/changes/resume-interrupted-sessions/specs/daemon-container/spec.md b/openspec/changes/resume-interrupted-sessions/specs/daemon-container/spec.md new file mode 100644 index 000000000..78d8cd7f3 --- /dev/null +++ b/openspec/changes/resume-interrupted-sessions/specs/daemon-container/spec.md @@ -0,0 +1,35 @@ +## MODIFIED Requirements + +### Requirement: Image entrypoint auto-starts netclawd + +The image SHALL start `tini` as PID 1. Its supervisor SHALL start `netclawd` and forward a container stop signal to it. The supervisor SHALL wait for the daemon to finish graceful drain before it exits. + +#### Scenario: docker run starts the daemon + +- **GIVEN** the image is present locally with valid configuration and identity files +- **WHEN** an operator starts the container +- **THEN** `tini` is PID 1 and the supervisor starts `netclawd` +- **AND** the daemon binds its HTTP port within 60 seconds + +#### Scenario: Pod stop preserves a resume candidate + +- **GIVEN** the container has a persistent operator state volume and an eligible interrupted session +- **WHEN** the pod sends a graceful stop signal with enough termination time +- **THEN** the supervisor forwards the signal and waits for the daemon to exit +- **AND** the state volume retains the candidate for the next container start + +### Requirement: Operator state mounts at /home/netclaw/.netclaw + +The image SHALL declare `VOLUME /home/netclaw/.netclaw`. The volume SHALL hold identity, configuration, session data, and restart candidates. The image SHALL not include operator credentials or identity files. + +#### Scenario: Operator bind-mounts an initialized home + +- **GIVEN** an operator has an initialized Netclaw home on the host +- **WHEN** they mount it at `/home/netclaw/.netclaw` and start the container +- **THEN** the daemon reads identity and configuration from that directory +- **AND** it writes session state and restart candidates to the same directory + +## RENAMED Requirements + +- FROM: `Operator state mounts at /root/.netclaw` +- TO: `Operator state mounts at /home/netclaw/.netclaw` diff --git a/openspec/changes/resume-interrupted-sessions/specs/session-resume/spec.md b/openspec/changes/resume-interrupted-sessions/specs/session-resume/spec.md new file mode 100644 index 000000000..e9e7cf261 --- /dev/null +++ b/openspec/changes/resume-interrupted-sessions/specs/session-resume/spec.md @@ -0,0 +1,103 @@ +## ADDED Requirements + +### Requirement: Accepted input survives a graceful stop + +The session SHALL record each accepted input before it acknowledges that input. The record SHALL retain a stable input ID, source ID, order, text, media, authority, and delivery context. A failed journal write SHALL reject the input. + +Use the [engineering glossary](../../../../../docs/spec/GLOSSARY.md) for shared terms. + +#### Scenario: Input ack follows its journal record + +- **GIVEN** a user sends input to a ready session or a busy session +- **WHEN** the journal confirms the accepted input record +- **THEN** the session acknowledges the input +- **AND** cold recovery restores its text, media, order, and original authority + +#### Scenario: Failed journal write rejects input + +- **GIVEN** the journal cannot store an input record +- **WHEN** the user sends that input +- **THEN** the session rejects the input with a visible error +- **AND** the session does not start a model call for it + +#### Scenario: Source retry does not duplicate accepted input + +- **GIVEN** the journal stores an input with a stable source message ID +- **WHEN** the source retries that same message after it loses the ack +- **THEN** the session acknowledges the existing input +- **AND** the session does not add a second copy to the queue + +### Requirement: Graceful stop records only safe resume candidates + +The daemon SHALL record a resume candidate after any graceful stop only for confirmed interrupted model work or accepted queued input. The candidate SHALL have one absolute deadline ten minutes after interruption. A stopped model task SHALL be confirmed before the session acknowledges drain. + +#### Scenario: Interrupted model call creates a candidate + +- **GIVEN** an admitted turn has a live model call and no tool batch has started in that turn +- **WHEN** a graceful stop interrupts the call after a short completion grace +- **THEN** the session waits for the model task to stop +- **AND** the daemon records a candidate for that original turn + +#### Scenario: Completed reply with queued input creates a candidate + +- **GIVEN** a turn finishes during drain and accepted queued input remains +- **WHEN** the session acknowledges drain +- **THEN** the daemon records a candidate for the accepted queue + +#### Scenario: Completed reply without queued input stays quiet + +- **GIVEN** a turn finishes during drain and no accepted queued input remains +- **WHEN** the daemon starts again +- **THEN** the daemon starts no model call for that session + +#### Scenario: Tool effect or partial reply blocks automatic resume + +- **GIVEN** the interrupted turn has a tool batch or user-visible partial text +- **WHEN** the daemon stops +- **THEN** the daemon does not record an executable model resume candidate +- **AND** it reports the blocked work to the operator + +### Requirement: Eligible restart resumes original work once + +The daemon SHALL resume an eligible candidate only before its deadline and after its output route is ready. The session SHALL use the recorded authority and input IDs. It SHALL not create a new reminder turn or ask the user for stored context. + +#### Scenario: Cold recovery resumes the original turn + +- **GIVEN** a graceful stop confirmed model cancellation and stored the admitted input +- **WHEN** the daemon starts before the candidate deadline +- **THEN** the session resumes one model call for the original turn +- **AND** it uses the recorded requester and trust boundary + +#### Scenario: Accepted queue follows the original turn + +- **GIVEN** a canceled turn and multiple accepted queued messages +- **WHEN** the resumed original turn finishes +- **THEN** the session sends the queued messages in their original order +- **AND** it uses one follow-up model call for that queue + +#### Scenario: Expired candidate stays quiet + +- **GIVEN** the absolute candidate deadline has passed +- **WHEN** the daemon starts or tries to resume the session +- **THEN** it starts no model call from that candidate +- **AND** it reports and removes the expired candidate once + +#### Scenario: Output route is absent + +- **GIVEN** a candidate has no live route that can deliver the resumed output +- **WHEN** the daemon starts before its deadline +- **THEN** it does not run a model call yet +- **AND** it reports the blocked candidate to the operator + +#### Scenario: Newer work supersedes a candidate + +- **GIVEN** a newer user turn starts before the recovery request reaches the session +- **WHEN** the older recovery request arrives +- **THEN** the session rejects the stale request and starts no duplicate call + +#### Scenario: Incompatible queued authority blocks automatic resume + +- **GIVEN** accepted queued messages have incompatible trust boundaries +- **WHEN** the daemon tries to resume that queue +- **THEN** the session keeps the messages in the journal +- **AND** it reports the blocked queue instead of widening tool authority diff --git a/openspec/changes/resume-interrupted-sessions/tasks.md b/openspec/changes/resume-interrupted-sessions/tasks.md new file mode 100644 index 000000000..e12276b62 --- /dev/null +++ b/openspec/changes/resume-interrupted-sessions/tasks.md @@ -0,0 +1,26 @@ +## 1. Durable input admission + +- [x] 1.1 Add `InputAdmitted` and terminal input IDs to protobuf and journal mapping. Verify old and new event round trips. +- [ ] 1.2 Store the ordered pending input ledger in session snapshots. Verify replay across a snapshot boundary. +- [ ] 1.3 Persist ready, busy, and compacting input before each ack. Verify a journal failure rejects input. +- [x] 1.4 Restore source retry deduplication and exact input order. Verify a lost ack does not duplicate a stable source ID. +- [ ] 1.5 Consume admitted IDs in completed, tool-started, and failed turns. Verify replay cannot restart terminal work. + +## 2. Graceful stop classification + +- [ ] 2.1 Retain the active model task and add a short completion grace. Verify task cancellation completes before drain ack. +- [ ] 2.2 Return typed drain results for safe candidates and blocked work. Verify partial output, tool effects, and approvals stay quiet. +- [ ] 2.3 Save candidates from config and normal stop paths with an atomic manifest write. Verify a completed idle session has no candidate. + +## 3. Short restart recovery + +- [ ] 3.1 Restore channel output bindings after channel startup. Verify the binding ack precedes a resumed model call. +- [ ] 3.2 Validate the absolute deadline, input IDs, authority, and newer work inside the actor. Verify stale requests start no call. +- [ ] 3.3 Resume the original turn and one ordered queued call. Verify model input and trust context across cold recovery. +- [ ] 3.4 Retain blocked candidates and emit operator diagnostics. Verify a missing output route and an expired deadline stay quiet. + +## 4. Documentation and gates + +- [ ] 4.1 Update SPEC-011 and the operations skill. Verify the skill version and behavior match the actor contract. +- [ ] 4.2 Correct the daemon container spec and run an isolated persistent-home stop/start test. +- [ ] 4.3 Run actor and daemon suites, Spark2 evals, Slopwatch, header verification, and strict OpenSpec validation. diff --git a/src/Netclaw.Actors.Tests/Protocol/SerializationRoundTripTests.cs b/src/Netclaw.Actors.Tests/Protocol/SerializationRoundTripTests.cs index 8cb27ca01..21499521b 100644 --- a/src/Netclaw.Actors.Tests/Protocol/SerializationRoundTripTests.cs +++ b/src/Netclaw.Actors.Tests/Protocol/SerializationRoundTripTests.cs @@ -14,6 +14,7 @@ using Netclaw.Actors.Reminders; using Netclaw.Actors.Serialization; using Netclaw.Actors.Sessions; +using Netclaw.Configuration; using Netclaw.Tools; using Xunit; using static Netclaw.Actors.Sessions.SessionProtocol; @@ -132,6 +133,67 @@ public void TurnRecorded_round_trips() Assert.Equal(original.RecordedAtMs, result.RecordedAtMs); } + [Fact] + public void Admitted_input_and_terminal_ids_survive_journal_and_snapshot_round_trips() + { + var sessionId = new SessionId("C99999/1708531200.000100"); + var admitted = new InputAdmitted + { + SessionId = sessionId, + InputId = "input-1", + SourceMessageId = "event-1", + UserMessage = new SerializableChatMessage { Role = ChatRole.User, Content = "Continue the task" }, + ExecutableText = "Continue the task", + TurnContext = new TurnContextRecord + { + SessionId = sessionId, + TurnId = "turn-1", + Audience = TrustAudience.Personal, + Boundary = new TrustBoundary("slack:C99999"), + ChannelType = "slack", + RequesterSenderId = new SenderId("U123"), + RequesterPrincipal = PrincipalClassification.Operator + }, + AdmittedAtMs = 1_700_000_000_000 + }; + + var restored = RoundTrip(admitted); + Assert.Equal(admitted.InputId, restored.InputId); + Assert.Equal(admitted.UserMessage.Content, restored.UserMessage.Content); + Assert.Equal(admitted.TurnContext?.Boundary, restored.TurnContext?.Boundary); + Assert.Equal(admitted.TurnContext?.RequesterSenderId, restored.TurnContext?.RequesterSenderId); + + var snapshot = RoundTrip(new SessionSnapshot + { + PendingInputs = [admitted], + RecentSourceMessageKeys = ["slack:event-1"] + }); + Assert.Single(snapshot.PendingInputs); + Assert.Equal("input-1", snapshot.PendingInputs[0].InputId); + Assert.Equal("slack:event-1", Assert.Single(snapshot.RecentSourceMessageKeys)); + + var completed = RoundTrip(new TurnRecorded + { + SessionId = sessionId, + UserMessage = admitted.UserMessage, + AssistantReply = new SerializableChatMessage { Role = ChatRole.Assistant, Content = "Done" }, + ConsumedInputIds = ["input-1"] + }); + Assert.Equal("input-1", Assert.Single(completed.ConsumedInputIds)); + + var closed = RoundTrip(new InputClosed { SessionId = sessionId, InputIds = ["input-1"] }); + Assert.Equal("input-1", Assert.Single(closed.InputIds)); + + var toolStarted = RoundTrip(new ToolBatchStarted + { + SessionId = sessionId, + UserMessage = admitted.UserMessage, + AssistantMessage = new SerializableChatMessage { Role = ChatRole.Assistant }, + ConsumedInputIds = ["input-1"] + }); + Assert.Equal("input-1", Assert.Single(toolStarted.ConsumedInputIds)); + } + [Fact] public void TurnRecorded_round_trips_preserving_value_object_source_ids() { diff --git a/src/Netclaw.Actors.Tests/Sessions/LlmSessionIntegrationTests.cs b/src/Netclaw.Actors.Tests/Sessions/LlmSessionIntegrationTests.cs index 2596f15ca..8a2da5325 100644 --- a/src/Netclaw.Actors.Tests/Sessions/LlmSessionIntegrationTests.cs +++ b/src/Netclaw.Actors.Tests/Sessions/LlmSessionIntegrationTests.cs @@ -6,6 +6,7 @@ using Akka.Actor; using Akka.Event; using Akka.Hosting; +using Akka.Persistence; using Akka.Streams; using System.Text.Json; using Microsoft.Extensions.AI; @@ -48,6 +49,77 @@ public LlmSessionIntegrationTests(ITestOutputHelper output) : base(output) { } + [Fact] + public async Task Input_ack_follows_a_journal_record_with_original_authority() + { + var sessionId = new SessionId("admission/journal-before-ack"); + var manager = ActorRegistry.Get(); + var subscriber = CreateTestProbe("admission-subscriber"); + await JoinSessionAsync(manager, subscriber, sessionId); + + var source = new MessageSource + { + ChannelType = ChannelType.SignalR, + SenderId = new SenderId("operator-1"), + ChannelId = "operator-channel", + MessageId = "source-event-1", + TurnId = new Netclaw.Actors.Protocol.TurnId("source-turn-1"), + Audience = TrustAudience.Personal, + Boundary = TrustBoundary.Personal, + Principal = PrincipalClassification.Operator, + Provenance = new SourceProvenance(TransportAuthenticity.Verified, PayloadTaint.Trusted), + ReceivedAt = _timeProvider.GetUtcNow() + }; + + await manager.Ask(new SendUserMessage + { + SessionId = sessionId, + Content = "Finish the operator task", + Source = source + }, TimeSpan.FromSeconds(3), TestContext.Current.CancellationToken); + + var journalProbe = CreateTestProbe("admission-journal"); + Sys.ActorOf(Props.Create(() => new InputJournalObserver( + $"session-{sessionId.Value}", journalProbe.Ref))); + var admitted = await journalProbe.ExpectMsgAsync( + TimeSpan.FromSeconds(3), cancellationToken: TestContext.Current.CancellationToken); + + Assert.Equal("Finish the operator task", admitted.UserMessage.Content); + Assert.Equal("source-event-1", admitted.SourceMessageId); + Assert.Equal(source.Boundary, admitted.TurnContext?.Boundary); + Assert.Equal(source.SenderId, admitted.TurnContext?.RequesterSenderId); + + await manager.Ask(new SendUserMessage + { + SessionId = sessionId, + Content = "Finish the operator task", + Source = source + }, TimeSpan.FromSeconds(3), TestContext.Current.CancellationToken); + + var retryProbe = CreateTestProbe("admission-retry-journal"); + Sys.ActorOf(Props.Create(() => new InputJournalObserver( + $"session-{sessionId.Value}", retryProbe.Ref))); + await retryProbe.ExpectMsgAsync( + TimeSpan.FromSeconds(3), cancellationToken: TestContext.Current.CancellationToken); + await retryProbe.ExpectMsgAsync( + TimeSpan.FromSeconds(3), cancellationToken: TestContext.Current.CancellationToken); + } + + private sealed record JournalReplayComplete : INoSerializationVerificationNeeded; + + private sealed class InputJournalObserver : ReceivePersistentActor + { + public override string PersistenceId { get; } + + public InputJournalObserver(string persistenceId, IActorRef replyTo) + { + PersistenceId = persistenceId; + Recover(evt => replyTo.Tell(evt)); + Recover(_ => replyTo.Tell(new JournalReplayComplete())); + RecoverAny(_ => { }); + } + } + protected override void ConfigureSessionServices(IServiceCollection services) { services.AddSingleton(new SingleClientProvider(_fakeChatClient)); diff --git a/src/Netclaw.Actors.Tests/Sessions/SessionStateTests.cs b/src/Netclaw.Actors.Tests/Sessions/SessionStateTests.cs index 37a0d1065..1de7315f2 100644 --- a/src/Netclaw.Actors.Tests/Sessions/SessionStateTests.cs +++ b/src/Netclaw.Actors.Tests/Sessions/SessionStateTests.cs @@ -21,6 +21,38 @@ public class SessionStateTests { private static readonly SessionId TestSessionId = new("test/session"); + [Fact] + public void Admitted_input_survives_snapshot_and_closes_only_by_its_id() + { + var first = new InputAdmitted + { + SessionId = TestSessionId, + InputId = "first", + SourceMessageId = "event-1", + TurnContext = new TurnContextRecord + { + SessionId = TestSessionId, + ChannelType = "slack", + RequesterSenderId = new SenderId("user-1") + }, + UserMessage = new SerializableChatMessage { Role = ChatRole.User, Content = "one" } + }; + var second = first with + { + InputId = "second", + SourceMessageId = "event-2", + UserMessage = new SerializableChatMessage { Role = ChatRole.User, Content = "two" } + }; + + var restored = SessionState.FromSnapshot(SessionState.Empty.Apply(first).Apply(second).ToSnapshot()); + Assert.Equal(["first", "second"], restored.PendingInputs.Select(input => input.InputId)); + Assert.Contains("slack:user-1:event-1", restored.RecentSourceMessageKeys); + + var closed = restored.CloseInputs(["first"]); + Assert.Equal("second", Assert.Single(closed.PendingInputs).InputId); + Assert.Contains("slack:user-1:event-1", closed.RecentSourceMessageKeys); + } + [Fact] public void Empty_state_has_no_history() { diff --git a/src/Netclaw.Actors/Channels/TurnContext.cs b/src/Netclaw.Actors/Channels/TurnContext.cs index f8c5f0cf2..4fe90aa74 100644 --- a/src/Netclaw.Actors/Channels/TurnContext.cs +++ b/src/Netclaw.Actors/Channels/TurnContext.cs @@ -100,7 +100,7 @@ public TurnContextRecord ToRecord() RequestedDeliveryTarget = RequestedDeliveryTarget, HasAdoptedContext = HasAdoptedContext, HasThirdPartyAdoptedContext = HasThirdPartyAdoptedContext, - AdoptedSpeakerIds = [.. AdoptedSpeakerIds], + AdoptedSpeakerIds = AdoptedSpeakerIds.ToArray(), SupportsInteractiveApproval = SupportsInteractiveApproval }; } diff --git a/src/Netclaw.Actors/Protocol/SessionSnapshot.cs b/src/Netclaw.Actors/Protocol/SessionSnapshot.cs index d8bdc7f51..bef30a9e2 100644 --- a/src/Netclaw.Actors/Protocol/SessionSnapshot.cs +++ b/src/Netclaw.Actors/Protocol/SessionSnapshot.cs @@ -80,4 +80,10 @@ public sealed record AdoptedContextSnapshotMessage public IReadOnlyList AdoptedContextRecords { get; init; } = Array.Empty(); + + public IReadOnlyList PendingInputs { get; init; } = + Array.Empty(); + + public IReadOnlyList RecentSourceMessageKeys { get; init; } = + Array.Empty(); } diff --git a/src/Netclaw.Actors/Serialization/NetclawProtoMapper.cs b/src/Netclaw.Actors/Serialization/NetclawProtoMapper.cs index 440db781d..530926b4d 100644 --- a/src/Netclaw.Actors/Serialization/NetclawProtoMapper.cs +++ b/src/Netclaw.Actors/Serialization/NetclawProtoMapper.cs @@ -28,6 +28,8 @@ internal static class NetclawProtoMapper SerializableMediaReference v => ToProto(v), SerializableToolCall v => ToProto(v), TurnRecorded v => ToProto(v), + InputAdmitted v => ToProto(v), + InputClosed v => ToProto(v), SessionTitleSet v => ToProto(v), SessionCompacted v => ToProto(v), ToolBatchStarted v => ToProto(v), @@ -166,6 +168,7 @@ internal static Proto.TurnRecordedProto ToProto(TurnRecorded evt) proto.SourceReminderId = reminderId.Value; if (evt.SourceBackgroundJobId is { } backgroundJobId) proto.SourceBackgroundJobId = backgroundJobId.Value; + proto.ConsumedInputIds.AddRange(evt.ConsumedInputIds); return proto; } @@ -176,7 +179,61 @@ internal static Proto.TurnRecordedProto ToProto(TurnRecorded evt) AssistantReply = FromProto(proto.AssistantReply), RecordedAtMs = proto.RecordedAtMs, SourceReminderId = proto.HasSourceReminderId ? new ReminderId(proto.SourceReminderId) : (ReminderId?)null, - SourceBackgroundJobId = proto.HasSourceBackgroundJobId ? new BackgroundJobId(proto.SourceBackgroundJobId) : (BackgroundJobId?)null + SourceBackgroundJobId = proto.HasSourceBackgroundJobId ? new BackgroundJobId(proto.SourceBackgroundJobId) : (BackgroundJobId?)null, + ConsumedInputIds = proto.ConsumedInputIds.ToArray() + }; + + internal static Proto.InputAdmittedProto ToProto(InputAdmitted evt) + { + var proto = new Proto.InputAdmittedProto + { + SessionId = ToProto(evt.SessionId), + InputId = evt.InputId, + UserMessage = ToProto(evt.UserMessage), + AdmittedAtMs = evt.AdmittedAtMs + }; + if (evt.SourceMessageId is not null) + proto.SourceMessageId = evt.SourceMessageId; + if (evt.ExecutableText is not null) + proto.ExecutableText = evt.ExecutableText; + if (evt.TurnContext is not null) + proto.TurnContext = ToProto(evt.TurnContext); + if (evt.SourceReminderId is { } reminderId) + proto.SourceReminderId = reminderId.Value; + if (evt.SourceBackgroundJobId is { } backgroundJobId) + proto.SourceBackgroundJobId = backgroundJobId.Value; + return proto; + } + + internal static InputAdmitted FromProto(Proto.InputAdmittedProto proto) => new() + { + SessionId = FromProto(proto.SessionId), + InputId = proto.InputId, + SourceMessageId = proto.HasSourceMessageId ? proto.SourceMessageId : null, + UserMessage = FromProto(proto.UserMessage), + ExecutableText = proto.HasExecutableText ? proto.ExecutableText : null, + TurnContext = proto.TurnContext is null ? null : FromProto(proto.TurnContext), + SourceReminderId = proto.HasSourceReminderId ? new ReminderId(proto.SourceReminderId) : null, + SourceBackgroundJobId = proto.HasSourceBackgroundJobId ? new BackgroundJobId(proto.SourceBackgroundJobId) : null, + AdmittedAtMs = proto.AdmittedAtMs + }; + + internal static Proto.InputClosedProto ToProto(InputClosed evt) + { + var proto = new Proto.InputClosedProto + { + SessionId = ToProto(evt.SessionId), + ClosedAtMs = evt.ClosedAtMs + }; + proto.InputIds.AddRange(evt.InputIds); + return proto; + } + + internal static InputClosed FromProto(Proto.InputClosedProto proto) => new() + { + SessionId = FromProto(proto.SessionId), + InputIds = proto.InputIds.ToArray(), + ClosedAtMs = proto.ClosedAtMs }; // ── SessionTitleSet ── @@ -224,20 +281,26 @@ internal static Proto.SessionCompactedProto ToProto(SessionCompacted evt) // ── Tool batch / approval events ── - internal static Proto.ToolBatchStartedProto ToProto(ToolBatchStarted evt) => new() + internal static Proto.ToolBatchStartedProto ToProto(ToolBatchStarted evt) { - SessionId = ToProto(evt.SessionId), - UserMessage = ToProto(evt.UserMessage), - AssistantMessage = ToProto(evt.AssistantMessage), - StartedAtMs = evt.StartedAtMs - }; + var proto = new Proto.ToolBatchStartedProto + { + SessionId = ToProto(evt.SessionId), + UserMessage = ToProto(evt.UserMessage), + AssistantMessage = ToProto(evt.AssistantMessage), + StartedAtMs = evt.StartedAtMs + }; + proto.ConsumedInputIds.AddRange(evt.ConsumedInputIds); + return proto; + } internal static ToolBatchStarted FromProto(Proto.ToolBatchStartedProto proto) => new() { SessionId = FromProto(proto.SessionId), UserMessage = FromProto(proto.UserMessage), AssistantMessage = FromProto(proto.AssistantMessage), - StartedAtMs = proto.StartedAtMs + StartedAtMs = proto.StartedAtMs, + ConsumedInputIds = proto.ConsumedInputIds.ToArray() }; internal static Proto.ToolCallRecordedProto ToProto(ToolCallRecorded evt) => new() @@ -513,6 +576,8 @@ internal static Proto.SessionSnapshotProto ToProto(SessionSnapshot snap) proto.History.AddRange(snap.History.Select(ToProto)); proto.ActiveBackgroundJobs.AddRange(snap.ActiveBackgroundJobs.Select(ToProto)); proto.AdoptedContextRecords.AddRange(snap.AdoptedContextRecords.Select(ToAdoptedContextSnapshotRecord)); + proto.PendingInputs.AddRange(snap.PendingInputs.Select(ToProto)); + proto.RecentSourceMessageKeys.AddRange(snap.RecentSourceMessageKeys); return proto; } @@ -526,7 +591,9 @@ internal static Proto.SessionSnapshotProto ToProto(SessionSnapshot snap) WorkingContext = proto.WorkingContext is not null ? FromProto(proto.WorkingContext) : null, History = proto.History.Select(FromProto).ToArray(), ActiveBackgroundJobs = proto.ActiveBackgroundJobs.Select(FromProto).ToArray(), - AdoptedContextRecords = proto.AdoptedContextRecords.Select(FromAdoptedContextSnapshotRecord).ToArray() + AdoptedContextRecords = proto.AdoptedContextRecords.Select(FromAdoptedContextSnapshotRecord).ToArray(), + PendingInputs = proto.PendingInputs.Select(FromProto).ToArray(), + RecentSourceMessageKeys = proto.RecentSourceMessageKeys.ToArray() }; private static Proto.SessionSnapshotProto.Types.AdoptedContextSnapshotRecord ToAdoptedContextSnapshotRecord( diff --git a/src/Netclaw.Actors/Serialization/NetclawProtobufSerializer.cs b/src/Netclaw.Actors/Serialization/NetclawProtobufSerializer.cs index 5cf1b2e19..1893f5a24 100644 --- a/src/Netclaw.Actors/Serialization/NetclawProtobufSerializer.cs +++ b/src/Netclaw.Actors/Serialization/NetclawProtobufSerializer.cs @@ -27,6 +27,8 @@ public sealed class NetclawProtobufSerializer : SerializerWithStringManifest private const string SerializableMediaReferenceManifest = "smr-v1"; private const string SerializableToolCallManifest = "stc-v1"; private const string TurnRecordedManifest = "tr-v1"; + private const string InputAdmittedManifest = "ia-v1"; + private const string InputClosedManifest = "ic-v1"; private const string SessionTitleSetManifest = "sts-v1"; private const string SessionCompactedManifest = "sc-v1"; private const string SessionSnapshotManifest = "ss-v1"; @@ -55,6 +57,8 @@ public sealed class NetclawProtobufSerializer : SerializerWithStringManifest [typeof(SerializableMediaReference)] = SerializableMediaReferenceManifest, [typeof(SerializableToolCall)] = SerializableToolCallManifest, [typeof(TurnRecorded)] = TurnRecordedManifest, + [typeof(InputAdmitted)] = InputAdmittedManifest, + [typeof(InputClosed)] = InputClosedManifest, [typeof(SessionTitleSet)] = SessionTitleSetManifest, [typeof(SessionCompacted)] = SessionCompactedManifest, [typeof(SessionSnapshot)] = SessionSnapshotManifest, @@ -112,6 +116,10 @@ public override object FromBinary(byte[] bytes, string manifest) Proto.SerializableToolCallProto.Parser.ParseFrom(bytes)), TurnRecordedManifest => NetclawProtoMapper.FromProto( Proto.TurnRecordedProto.Parser.ParseFrom(bytes)), + InputAdmittedManifest => NetclawProtoMapper.FromProto( + Proto.InputAdmittedProto.Parser.ParseFrom(bytes)), + InputClosedManifest => NetclawProtoMapper.FromProto( + Proto.InputClosedProto.Parser.ParseFrom(bytes)), SessionTitleSetManifest => NetclawProtoMapper.FromProto( Proto.SessionTitleSetProto.Parser.ParseFrom(bytes)), SessionCompactedManifest => NetclawProtoMapper.FromProto( diff --git a/src/Netclaw.Actors/Serialization/Protos/netclaw_messages.proto b/src/Netclaw.Actors/Serialization/Protos/netclaw_messages.proto index 320406aa8..7bd134df7 100644 --- a/src/Netclaw.Actors/Serialization/Protos/netclaw_messages.proto +++ b/src/Netclaw.Actors/Serialization/Protos/netclaw_messages.proto @@ -94,6 +94,25 @@ message TurnRecordedProto { int64 recorded_at_ms = 4; optional string source_reminder_id = 5; optional string source_background_job_id = 6; + repeated string consumed_input_ids = 7; +} + +message InputAdmittedProto { + SessionIdProto session_id = 1; + string input_id = 2; + optional string source_message_id = 3; + SerializableChatMessageProto user_message = 4; + optional string executable_text = 5; + ToolApprovalRequestedProto.TurnContextRecordProto turn_context = 6; + optional string source_reminder_id = 7; + optional string source_background_job_id = 8; + int64 admitted_at_ms = 9; +} + +message InputClosedProto { + SessionIdProto session_id = 1; + repeated string input_ids = 2; + int64 closed_at_ms = 3; } message SessionTitleSetProto { @@ -117,6 +136,7 @@ message ToolBatchStartedProto { SerializableChatMessageProto user_message = 2; SerializableChatMessageProto assistant_message = 3; int64 started_at_ms = 4; + repeated string consumed_input_ids = 5; } message ToolCallRecordedProto { @@ -272,6 +292,8 @@ message SessionSnapshotProto { repeated ActiveJobInfoProto active_background_jobs = 7; repeated AdoptedContextSnapshotRecord adopted_context_records = 8; reserved 9; + repeated InputAdmittedProto pending_inputs = 10; + repeated string recent_source_message_keys = 11; } // ── Session state ── diff --git a/src/Netclaw.Actors/Sessions/LlmSessionActor.cs b/src/Netclaw.Actors/Sessions/LlmSessionActor.cs index 2a97ea01b..2560bd99c 100644 --- a/src/Netclaw.Actors/Sessions/LlmSessionActor.cs +++ b/src/Netclaw.Actors/Sessions/LlmSessionActor.cs @@ -77,6 +77,7 @@ public sealed class LlmSessionActor : ReceivePersistentActor, IWithTimers // Transient state (not persisted) private readonly List<(SendUserMessage Message, bool IsReplay)> _buffer = []; + private readonly List _activeInputIds = []; // In-flight reminder/background-job dedup (transient; rebuilt from journal on recovery). private readonly InFlightTurnDedup _inFlightDedup = new(); private readonly SessionSubscriberManager _subscribers = new(); @@ -280,6 +281,7 @@ public LlmSessionActor( // ── Recovery handlers ── Recover(evt => { + RestoreConsumedInputs(evt.ConsumedInputIds, evt.UserMessage); ApplyTurnRecorded(evt); _toolApprovals.ClearCalls(); ClearApprovalTurnState(); @@ -294,7 +296,13 @@ public LlmSessionActor( ClearApprovalTurnState(); ClearActiveToolBatchTracking(); }); - Recover(ApplyToolBatchStarted); + Recover(evt => + { + RestoreConsumedInputs(evt.ConsumedInputIds, evt.UserMessage); + ApplyToolBatchStarted(evt); + }); + Recover(evt => _state = _state.Apply(evt)); + Recover(evt => _state = _state.CloseInputs(evt.InputIds)); Recover(ApplyToolCallRecorded); Recover(ApplyToolApprovalRequested); Recover(ApplyToolApprovalResolved); @@ -506,10 +514,13 @@ private void Processing() return; } - _deliveryRetry.Clear(); - _log.Info("Buffering user message (LLM call in progress)"); - _buffer.Add((cmd, false)); - TryReplyAck(); + AdmitInput(cmd, admitted => + { + _deliveryRetry.Clear(); + _log.Info("Buffering user message (LLM call in progress)"); + _buffer.Add((admitted, false)); + TryReplyAck(); + }); }); Command(msg => @@ -1161,9 +1172,12 @@ private void Compacting() return; } - _log.Info("Buffering user message (compaction in progress)"); - _buffer.Add((cmd, false)); - TryReplyAck(); + AdmitInput(cmd, admitted => + { + _log.Info("Buffering user message (compaction in progress)"); + _buffer.Add((admitted, false)); + TryReplyAck(); + }); }); // Buffer an approval response during compaction — compaction rewrites @@ -1879,9 +1893,12 @@ private void HandleToolCallResponse( SessionId = _sessionId, UserMessage = userMsg, AssistantMessage = assistantMsg, - StartedAtMs = NowMs() + StartedAtMs = NowMs(), + ConsumedInputIds = _activeInputIds.ToArray() }, evt => { + _state = _state.CloseInputs(evt.ConsumedInputIds); + _activeInputIds.Clear(); ApplyToolBatchStarted(evt); EmitAndDispatchToolBatch( lastMessage, @@ -2132,13 +2149,16 @@ private void HandleTextResponse( AssistantReply = reply, RecordedAtMs = NowMs(), SourceReminderId = _currentTurnSource?.ReminderId, - SourceBackgroundJobId = _currentTurnSource?.BackgroundJobId + SourceBackgroundJobId = _currentTurnSource?.BackgroundJobId, + ConsumedInputIds = _activeInputIds.ToArray() }; Persist(turnEvent, evt => { _inFlightDedup.CompleteReminder(evt.SourceReminderId); _inFlightDedup.CompleteBackgroundJob(evt.SourceBackgroundJobId); + _state = _state.CloseInputs(evt.ConsumedInputIds); + _activeInputIds.Clear(); var processed = _state.ProcessedReminderIds; if (evt.SourceReminderId is { } reminderId && !string.IsNullOrEmpty(reminderId.Value)) @@ -2212,6 +2232,8 @@ private bool DrainBufferedUserMessages() { var refs = message.MediaReferences.Count > 0 ? message.MediaReferences : null; _state = _state.AddUserMessage(message.Content, refs); + if (message.AdmittedInputId is { } id && !_activeInputIds.Contains(id)) + _activeInputIds.Add(id); } _buffer.Clear(); @@ -2292,9 +2314,6 @@ private void HandleIncomingUserMessage(SendUserMessage cmd) if (TryRejectIncompatibleInput(cmd.MediaReferences, cmd.Source)) return; - _inFlightDedup.ReserveReminder(reminderId); - _inFlightDedup.ReserveBackgroundJob(bgJobId); - // A new inbound message while a tool batch is still parked on an // approval gate means the user abandoned that approval. This is only // reachable on a cold-recovered session — a live session parked on @@ -2310,7 +2329,7 @@ private void HandleIncomingUserMessage(SendUserMessage cmd) Persist(abandoned, evt => { ApplyToolBatchAbandoned(evt); - ContinueIncomingUserMessage(cmd); + AdmitInput(cmd, ContinueIncomingUserMessage); }); return; } @@ -2321,7 +2340,7 @@ private void HandleIncomingUserMessage(SendUserMessage cmd) Persist(abandoned, evt => { ApplyToolBatchAbandoned(evt); - ContinueIncomingUserMessage(cmd); + AdmitInput(cmd, ContinueIncomingUserMessage); }); return; } @@ -2332,16 +2351,93 @@ private void HandleIncomingUserMessage(SendUserMessage cmd) Persist(abandoned, evt => { ApplyToolBatchAbandoned(evt); - ContinueIncomingUserMessage(cmd); + AdmitInput(cmd, ContinueIncomingUserMessage); }); return; } - ContinueIncomingUserMessage(cmd); + AdmitInput(cmd, ContinueIncomingUserMessage); + } + + private void AdmitInput(SendUserMessage cmd, Action onAccepted) + { + if (cmd.AdmittedInputId is not null) + { + onAccepted(cmd); + return; + } + + var context = TurnContext.FromMessageSource( + _sessionId, + cmd.Source?.TurnId ?? new Protocol.TurnId(IdGen.ShortId()), + cmd.Source); + var evt = new InputAdmitted + { + SessionId = _sessionId, + InputId = Guid.NewGuid().ToString("N"), + SourceMessageId = cmd.Source?.MessageId + ?? cmd.Source?.ReminderId?.Value + ?? cmd.Source?.BackgroundJobId?.Value, + UserMessage = new SerializableChatMessage + { + Role = Protocol.ChatRole.User, + Content = cmd.Content ?? string.Empty, + MediaReferences = cmd.MediaReferences.ToArray() + }, + ExecutableText = cmd.Source?.ExecutableText, + TurnContext = context.ToRecord(), + SourceReminderId = cmd.Source?.ReminderId, + SourceBackgroundJobId = cmd.Source?.BackgroundJobId, + AdmittedAtMs = NowMs() + }; + + if (!TurnContext.TryFromRecord(evt.TurnContext, out _, out var reason)) + { + _log.Error("Rejecting input with invalid durable authority: {Reason}", reason); + TryReplyNack($"Input authority cannot be stored: {reason}"); + return; + } + + var sourceKey = SessionState.SourceMessageKey(evt); + if (sourceKey is not null && _state.RecentSourceMessageKeys.Contains(sourceKey)) + { + TryReplyAck(); + return; + } + + Persist(evt, admitted => + { + _state = _state.Apply(admitted); + _inFlightDedup.ReserveReminder(admitted.SourceReminderId); + _inFlightDedup.ReserveBackgroundJob(admitted.SourceBackgroundJobId); + onAccepted(cmd with { AdmittedInputId = admitted.InputId }); + }); + } + + private void RestoreConsumedInputs(IReadOnlyList inputIds, SerializableChatMessage lastUserMessage) + { + if (inputIds.Count == 0) + return; + + var pendingById = _state.PendingInputs.ToDictionary(evt => evt.InputId, StringComparer.Ordinal); + for (var index = 0; index < inputIds.Count; index++) + { + var id = inputIds[index]; + if (!pendingById.TryGetValue(id, out var input)) + throw new InvalidOperationException($"Journal input {id} is missing before its terminal event."); + + var message = index == inputIds.Count - 1 ? lastUserMessage : input.UserMessage; + _state = _state with { History = _state.History.Add(message) }; + } + + _state = _state.CloseInputs(inputIds); } private void ContinueIncomingUserMessage(SendUserMessage cmd) { + _activeInputIds.Clear(); + if (cmd.AdmittedInputId is { } admittedId) + _activeInputIds.Add(admittedId); _sessionManagedTemporaryCorrections.Clear(); _deliveryRetry.Clear(); _currentTurnSource = cmd.Source; @@ -2605,6 +2701,22 @@ protected override void PostStop() base.PostStop(); } + protected override void OnPersistFailure(Exception cause, object @event, long sequenceNr) + { + if (@event is InputAdmitted) + TryReplyNack("The input journal failed. The session did not start a model call."); + + base.OnPersistFailure(cause, @event, sequenceNr); + } + + protected override void OnPersistRejected(Exception cause, object @event, long sequenceNr) + { + if (@event is InputAdmitted) + TryReplyNack("The input journal rejected the message. The session did not start a model call."); + + base.OnPersistRejected(cause, @event, sequenceNr); + } + private void CancelAndDisposeLlmCts() { _activeLlmCts?.Cancel(); @@ -3111,21 +3223,7 @@ private bool TryHandleSlashCommand(string userContent, IReadOnlyList + { + _state = _state.CloseInputs(evt.InputIds); + _activeInputIds.Clear(); + Reply(); }); - TryReplyAck(); return true; } @@ -3173,19 +3292,7 @@ private bool HandleInlineSlashCommand(SkillEntry skill, string remainder, IReadO { _log.Warning("Failed to read skill file for slash command /{SkillName}: {Error}", skill.Name, ex.Message); - EmitOutput(new TextOutput($"Failed to load skill /{skill.Name}: {ex.Message}\n\nThe skill file may be missing or corrupted.") - { - SessionId = _sessionId - }, OutputFilter.Text); - EmitOutput(new TurnCompleted - { - SessionId = _sessionId, - TurnNumber = new TurnNumber(_state.TurnCount), - Outcome = TurnOutcome.Skipped, - SourceReminderId = _currentTurnSource?.ReminderId - }); - TryReplyAck(); - return true; + return RejectSlashCommand($"Failed to load skill /{skill.Name}: {ex.Message}\n\nThe skill file may be missing or corrupted."); } _sessionMetrics?.RecordSkillLoaded(skill.Name, SkillLoadMethod.SlashCommand); @@ -3212,58 +3319,16 @@ private bool TryHandleRoutedSlashCommand(SkillEntry skill, string remainder, IRe return RejectSlashCommand($"Skill /{skill.Name} cannot use file-based routed dispatch."); if (_subAgentRegistry is null || _subAgentSpawner is null) - { - EmitOutput(new TextOutput($"Skill '/{skill.Name}' routes to subagent '{routedSubagent}', but subagent routing is not available in this runtime.") - { - SessionId = _sessionId - }, OutputFilter.Text); - EmitOutput(new TurnCompleted - { - SessionId = _sessionId, - TurnNumber = new TurnNumber(_state.TurnCount), - Outcome = TurnOutcome.Skipped, - SourceReminderId = _currentTurnSource?.ReminderId - }); - TryReplyAck(); - return true; - } + return RejectSlashCommand($"Skill '/{skill.Name}' routes to subagent '{routedSubagent}', but subagent routing is not available in this runtime."); _subAgentLoader?.SyncInto(_subAgentRegistry); var profile = _subAgentRegistry.TryGetByName(routedSubagent); if (profile is null) - { - EmitOutput(new TextOutput(SkillActivationRouter.UnknownTargetError(skill.Name, routedSubagent)) - { - SessionId = _sessionId - }, OutputFilter.Text); - EmitOutput(new TurnCompleted - { - SessionId = _sessionId, - TurnNumber = new TurnNumber(_state.TurnCount), - Outcome = TurnOutcome.Skipped, - SourceReminderId = _currentTurnSource?.ReminderId - }); - TryReplyAck(); - return true; - } + return RejectSlashCommand(SkillActivationRouter.UnknownTargetError(skill.Name, routedSubagent)); if (profile.Visibility != SubAgentVisibility.UserFacing) - { - EmitOutput(new TextOutput(SkillActivationRouter.InternalTargetError(skill.Name, routedSubagent)) - { - SessionId = _sessionId - }, OutputFilter.Text); - EmitOutput(new TurnCompleted - { - SessionId = _sessionId, - TurnNumber = new TurnNumber(_state.TurnCount), - Outcome = TurnOutcome.Skipped, - SourceReminderId = _currentTurnSource?.ReminderId - }); - TryReplyAck(); - return true; - } + return RejectSlashCommand(SkillActivationRouter.InternalTargetError(skill.Name, routedSubagent)); string skillBody; try @@ -3275,19 +3340,7 @@ private bool TryHandleRoutedSlashCommand(SkillEntry skill, string remainder, IRe { _log.Warning("Failed to read skill file for routed slash command /{SkillName}: {Error}", skill.Name, ex.Message); - EmitOutput(new TextOutput($"Failed to load skill /{skill.Name}: {ex.Message}\n\nThe skill file may be missing or corrupted.") - { - SessionId = _sessionId - }, OutputFilter.Text); - EmitOutput(new TurnCompleted - { - SessionId = _sessionId, - TurnNumber = new TurnNumber(_state.TurnCount), - Outcome = TurnOutcome.Skipped, - SourceReminderId = _currentTurnSource?.ReminderId - }); - TryReplyAck(); - return true; + return RejectSlashCommand($"Failed to load skill /{skill.Name}: {ex.Message}\n\nThe skill file may be missing or corrupted."); } _sessionMetrics?.RecordSkillLoaded(skill.Name, SkillLoadMethod.SlashCommand); @@ -3390,13 +3443,16 @@ private void HandleRoutedSkillExecutionCompleted(RoutedSkillExecutionCompleted m }, RecordedAtMs = NowMs(), SourceReminderId = _currentTurnSource?.ReminderId, - SourceBackgroundJobId = _currentTurnSource?.BackgroundJobId + SourceBackgroundJobId = _currentTurnSource?.BackgroundJobId, + ConsumedInputIds = _activeInputIds.ToArray() }; Persist(turnEvent, evt => { _inFlightDedup.CompleteReminder(evt.SourceReminderId); _inFlightDedup.CompleteBackgroundJob(evt.SourceBackgroundJobId); + _state = _state.CloseInputs(evt.ConsumedInputIds); + _activeInputIds.Clear(); var processed = _state.ProcessedReminderIds; if (evt.SourceReminderId is { } reminderId && !string.IsNullOrEmpty(reminderId.Value)) @@ -3730,9 +3786,10 @@ private void SaveSnapshotIfSafe() // skip the journal event that rehydrates pending approval context. if (_toolApprovals.PendingCount > 0 || _toolApprovals.ResolvedCount > 0 + || _state.PendingInputs.Count > 0 || ParkedToolBatchHistory.FindRedrivableAssistantMessage(_state.History, null) is not null) { - _log.Info("Skipping snapshot while approval-paused tool batch is still unresolved"); + _log.Info("Skipping snapshot while a durable input or approval batch is unresolved"); return; } @@ -4495,6 +4552,27 @@ private ToolBatchAbandoned BuildToolBatchAbandonedEvent(string resultContent) } private void FailCurrentTurn(string errorMessage, Exception cause, ErrorCategory category = ErrorCategory.Unknown) + { + if (_activeInputIds.Count > 0) + { + Persist(new InputClosed + { + SessionId = _sessionId, + InputIds = _activeInputIds.ToArray(), + ClosedAtMs = NowMs() + }, evt => + { + _state = _state.CloseInputs(evt.InputIds); + _activeInputIds.Clear(); + FinishFailedTurn(errorMessage, cause, category); + }); + return; + } + + FinishFailedTurn(errorMessage, cause, category); + } + + private void FinishFailedTurn(string errorMessage, Exception cause, ErrorCategory category) { _inFlightDedup.CompleteReminder(_currentTurnSource?.ReminderId); _inFlightDedup.CompleteBackgroundJob(_currentTurnSource?.BackgroundJobId); diff --git a/src/Netclaw.Actors/Sessions/SessionProtocol.Commands.cs b/src/Netclaw.Actors/Sessions/SessionProtocol.Commands.cs index 440cad541..21fac79b7 100644 --- a/src/Netclaw.Actors/Sessions/SessionProtocol.Commands.cs +++ b/src/Netclaw.Actors/Sessions/SessionProtocol.Commands.cs @@ -39,6 +39,9 @@ public sealed record SendUserMessage : ISessionCommand, INetclawSerializableMess /// Ephemeral channel metadata for ACL/audit. Not persisted. /// public MessageSource? Source { get; init; } + + /// Actor-local correlation for an input that the journal accepted. + public string? AdmittedInputId { get; init; } } /// diff --git a/src/Netclaw.Actors/Sessions/SessionProtocol.Events.cs b/src/Netclaw.Actors/Sessions/SessionProtocol.Events.cs index 5daf844ce..f0531d553 100644 --- a/src/Netclaw.Actors/Sessions/SessionProtocol.Events.cs +++ b/src/Netclaw.Actors/Sessions/SessionProtocol.Events.cs @@ -49,11 +49,54 @@ public sealed record TurnRecorded : ISessionEvent /// public BackgroundJobId? SourceBackgroundJobId { get; init; } + public IReadOnlyList ConsumedInputIds { get; init; } = []; + public DateTimeOffset RecordedAt => DateTimeOffset.FromUnixTimeMilliseconds(RecordedAtMs); public DateTimeOffset Timestamp => RecordedAt; } + /// + /// Records accepted input and its original authority before the input ack. + /// Runtime actor references remain outside the journal. + /// + public sealed record InputAdmitted : ISessionEvent + { + public SessionId SessionId { get; init; } + + public string InputId { get; init; } = string.Empty; + + public string? SourceMessageId { get; init; } + + public SerializableChatMessage UserMessage { get; init; } = new(); + + public string? ExecutableText { get; init; } + + public TurnContextRecord? TurnContext { get; init; } + + public ReminderId? SourceReminderId { get; init; } + + public BackgroundJobId? SourceBackgroundJobId { get; init; } + + public long AdmittedAtMs { get; init; } + + public DateTimeOffset Timestamp => DateTimeOffset.FromUnixTimeMilliseconds(AdmittedAtMs); + } + + /// + /// Closes admitted input when a turn ends without a recorded model reply. + /// + public sealed record InputClosed : ISessionEvent + { + public SessionId SessionId { get; init; } + + public IReadOnlyList InputIds { get; init; } = []; + + public long ClosedAtMs { get; init; } + + public DateTimeOffset Timestamp => DateTimeOffset.FromUnixTimeMilliseconds(ClosedAtMs); + } + /// /// Persisted event recording an assistant tool-call batch before any tool is /// executed or approval prompt is emitted. This makes in-flight tool history @@ -70,6 +113,8 @@ public sealed record ToolBatchStarted : ISessionEvent public long StartedAtMs { get; init; } + public IReadOnlyList ConsumedInputIds { get; init; } = []; + public DateTimeOffset Timestamp => DateTimeOffset.FromUnixTimeMilliseconds(StartedAtMs); } diff --git a/src/Netclaw.Actors/Sessions/SessionState.cs b/src/Netclaw.Actors/Sessions/SessionState.cs index d45a27eaa..3b20cf877 100644 --- a/src/Netclaw.Actors/Sessions/SessionState.cs +++ b/src/Netclaw.Actors/Sessions/SessionState.cs @@ -46,6 +46,10 @@ public sealed record AdoptedContextAuditMessage( public ImmutableList History { get; init; } = []; + public ImmutableList PendingInputs { get; init; } = []; + + public ImmutableList RecentSourceMessageKeys { get; init; } = []; + public int TurnCount { get; init; } public string? Title { get; init; } @@ -94,6 +98,36 @@ public sealed record AdoptedContextAuditMessage( // ── Event application (pure functions) ── + public SessionState Apply(InputAdmitted evt) + { + var sourceKey = SourceMessageKey(evt); + var keys = sourceKey is null || RecentSourceMessageKeys.Contains(sourceKey) + ? RecentSourceMessageKeys + : RecentSourceMessageKeys.Add(sourceKey); + if (keys.Count > 256) + keys = keys.RemoveRange(0, keys.Count - 256); + + return this with + { + PendingInputs = PendingInputs.Add(evt), + RecentSourceMessageKeys = keys + }; + } + + public SessionState CloseInputs(IReadOnlyList inputIds) + { + if (inputIds.Count == 0) + return this; + + var ids = inputIds.ToHashSet(StringComparer.Ordinal); + return this with { PendingInputs = PendingInputs.RemoveAll(evt => ids.Contains(evt.InputId)) }; + } + + public static string? SourceMessageKey(InputAdmitted evt) + => string.IsNullOrWhiteSpace(evt.SourceMessageId) + ? null + : $"{evt.TurnContext?.ChannelType ?? string.Empty}:{evt.TurnContext?.RequesterSenderId?.Value ?? string.Empty}:{evt.SourceMessageId}"; + public SessionState Apply(TurnRecorded evt) { var processedReminders = ProcessedReminderIds; @@ -432,6 +466,8 @@ public SessionSnapshot ToSnapshot() return new SessionSnapshot { History = new List(History), + PendingInputs = PendingInputs.ToArray(), + RecentSourceMessageKeys = RecentSourceMessageKeys.ToArray(), TurnCount = TurnCount, Title = Title, WorkingContext = WorkingContext.IsEmpty ? null : WorkingContext, @@ -492,6 +528,8 @@ [.. record.Messages return new SessionState { History = ImmutableList.CreateRange(snapshot.History), + PendingInputs = ImmutableList.CreateRange(snapshot.PendingInputs), + RecentSourceMessageKeys = ImmutableList.CreateRange(snapshot.RecentSourceMessageKeys), TurnCount = snapshot.TurnCount, Title = snapshot.Title, WorkingContext = snapshot.WorkingContext ?? WorkingContext.Empty,