From 9bbf6ba5793bcafca6ae121d756cb7e52405107f Mon Sep 17 00:00:00 2001 From: Aaron Stannard Date: Wed, 23 Sep 2026 22:35:27 +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 | 85 +++++ .../resume-interrupted-sessions/proposal.md | 28 ++ .../specs/daemon-container/spec.md | 35 +++ .../specs/session-resume/spec.md | 86 +++++ .../resume-interrupted-sessions/tasks.md | 24 ++ .../Protocol/SerializationRoundTripTests.cs | 63 ++++ .../Sessions/LlmSessionIntegrationTests.cs | 49 +++ .../Sessions/SessionStateTests.cs | 34 ++ .../Protocol/SessionSnapshot.cs | 6 + .../Serialization/NetclawProtoMapper.cs | 80 ++++- .../NetclawProtobufSerializer.cs | 8 + .../Protos/netclaw_messages.proto | 21 ++ .../Sessions/LlmSessionActor.cs | 294 +++++++++++------- .../Sessions/SessionProtocol.Commands.cs | 3 + .../Sessions/SessionProtocol.Events.cs | 42 +++ .../Sessions/SessionProtocol.cs | 10 + src/Netclaw.Actors/Sessions/SessionState.cs | 38 +++ 19 files changed, 793 insertions(+), 122 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 3c7611a21..66d113a46 100644 --- a/docs/spec/SPEC-011-daemon-architecture.md +++ b/docs/spec/SPEC-011-daemon-architecture.md @@ -216,6 +216,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..eca0fbc61 --- /dev/null +++ b/openspec/changes/resume-interrupted-sessions/design.md @@ -0,0 +1,85 @@ +## Context + +The actor stores completed turns and tool work. It does not store a model only request before its acknowledgment. Its buffer is also actor local. + +The reminder manager already persists schedules, routes a `current_session` turn, deduplicates delivery, and records the result. + +## Goals + +- Store accepted input before acknowledgment. +- Stop an eligible model call during graceful drain. +- Wake the session within ten minutes through the reminder manager. +- Restore pending input under its recorded authority. +- Add no channel adapter or second delivery path. + +## Non Goals + +- Resume after an ungraceful crash. +- Replay a turn after a tool starts or partial text reaches a user. +- Create a daemon restart command or a configuration property. + +## Decisions + +### D1. The session journal owns accepted input + +`InputAdmitted` stores an `InputId`, content, media, source ID, executable text, and `TurnContextRecord`. The actor persists it before acknowledgment. + +`TurnRecorded` and `ToolBatchStarted` close the input IDs that they consume. `InputClosed` closes input after a terminal path without either event. + +`SessionState` keeps the ordered pending input ledger. A bounded source ID ledger rejects a retry after a lost acknowledgment. + +### D2. Drain produces a standard reminder + +The actor gives the model call a short completion grace. It cancels the call and waits for its task to stop when the grace ends. + +The actor creates a standard one shot `ReminderDefinition` only when these conditions hold: + +- pending input exists; +- 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 reminder expires ten minutes after the interruption. The restart manifest stores the definition with the active session list. + +### 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. + +### D4. The reminder is a trigger + +The reminder text is generic: `Resume the work that was interrupted by the daemon restart.` + +When this internal reminder arrives, the actor restores its pending input and original `TurnContextRecord`. The model sees the stored input and the restart notice. + +The reminder does not replace the original authority. A missing or incompatible context causes a visible operator warning and no model call. + +## Ordered Flow + +This flow is schematic. It omits persistence callbacks and reminder delivery acknowledgments. + +```text +input -> session: SendUserMessage +session -> journal: InputAdmitted +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 +start -> reminder manager: SaveReminderCommand +reminder manager -> existing gateway: current_session reminder +gateway -> session: SendUserMessage +session -> journal: restore pending input and authority +session -> model: resume prior work +``` + +## Risks + +- A stale reminder can start old work. The one shot definition has an absolute expiration. +- 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. diff --git a/openspec/changes/resume-interrupted-sessions/proposal.md b/openspec/changes/resume-interrupted-sessions/proposal.md new file mode 100644 index 000000000..095aa6aa6 --- /dev/null +++ b/openspec/changes/resume-interrupted-sessions/proposal.md @@ -0,0 +1,28 @@ +## Why + +A session acknowledges user input before the journal stores it. A graceful stop can therefore lose an active request or its queued input. + +Source: [PRD-001 FR-003 and FR-016](../../../docs/prd/PRD-001-netclaw-mvp.md). + +## What Changes + +- Persist each accepted input before its acknowledgment. +- Retain its content, order, source ID, media, and original authority. +- Cancel an eligible model call during graceful drain. +- Put a short lived `current_session` reminder in the restart manifest. +- Register that reminder through the existing reminder manager after startup. +- Restore the pending input under its original authority when the reminder arrives. +- Keep approvals, partial replies, and turns with possible tool effects quiet. + +This change adds no channel code and no configuration property. It excludes crash recovery and replay of uncertain tool effects. + +## Capabilities + +### Modified Capabilities + +- `session-resume`: Durable input admission and a bounded restart reminder. +- `daemon-container`: The state volume retains the restart manifest across a graceful pod replacement. + +## Impact + +This change affects session persistence, graceful drain, the restart manifest, and reminder registration. Existing channel gateways deliver the reminder without new adapters. 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..2ab01c828 --- /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 restart reminder + +- **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 reminder 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 reminders. 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 reminders 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..3ab4c0424 --- /dev/null +++ b/openspec/changes/resume-interrupted-sessions/specs/session-resume/spec.md @@ -0,0 +1,86 @@ +## ADDED Requirements + +### Requirement: Accepted input survives a graceful stop + +The session SHALL store each accepted input before acknowledgment. The record SHALL retain its identity, order, content, media, source identity, and original authority. + +Use the [engineering glossary](../../../../../docs/spec/GLOSSARY.md) for shared terms. + +#### Scenario: Input acknowledgment follows its journal record + +- **GIVEN** a session receives user input +- **WHEN** the journal stores its admission record +- **THEN** the session acknowledges the input +- **AND** cold recovery restores the pending input and its authority + +#### Scenario: Journal failure rejects input + +- **GIVEN** the journal cannot store an admission record +- **WHEN** the session receives input +- **THEN** the session rejects that input +- **AND** it starts no model call for that input + +#### Scenario: A lost acknowledgment does not duplicate input + +- **GIVEN** the journal stores input with a stable source ID +- **WHEN** the source retries that input +- **THEN** the session acknowledges the stored input +- **AND** it does not add a second pending record + +### Requirement: Graceful drain creates only a safe restart reminder + +The session SHALL create a restart reminder only after an eligible model task stops. It SHALL use the existing reminder definition and `current_session` delivery contract. + +#### Scenario: An interrupted model call creates a reminder + +- **GIVEN** a model call has pending admitted input +- **AND** no tool batch or partial reply exists +- **WHEN** graceful drain cancels the call and confirms its task stopped +- **THEN** the restart manifest stores one reminder for that session +- **AND** the reminder expires ten minutes after interruption + +#### Scenario: A completed turn stays quiet + +- **GIVEN** a model call completes during drain +- **AND** no admitted input remains pending +- **WHEN** the daemon starts again +- **THEN** it registers no restart reminder for that session + +#### Scenario: A possible effect blocks the reminder + +- **GIVEN** a tool batch started or partial text reached a subscriber +- **WHEN** graceful drain stops the session +- **THEN** the manifest contains no restart reminder for that turn +- **AND** the daemon reports the blocked session + +### Requirement: A fresh restart reminder resumes stored work + +The reminder manager SHALL deliver a fresh restart reminder through its existing `current_session` path. The session SHALL restore pending input under its recorded authority. + +#### Scenario: A fresh reminder resumes the pending input + +- **GIVEN** the restart manifest contains a reminder that has not expired +- **WHEN** the daemon starts +- **THEN** startup registers the reminder through `SaveReminderCommand` +- **AND** the session resumes the stored input without a user prompt + +#### Scenario: The original authority remains in force + +- **GIVEN** a restart reminder wakes a session with pending input +- **WHEN** the session starts the model call +- **THEN** it restores the recorded requester, audience, and trust boundary +- **AND** reminder automation authority does not replace that context + +#### Scenario: An expired reminder stays quiet + +- **GIVEN** the reminder expiration is in the past +- **WHEN** startup reads the restart manifest +- **THEN** it does not register or deliver that reminder +- **AND** it logs one warning + +#### Scenario: A channel lacks current session delivery + +- **GIVEN** an interrupted session uses a channel that the reminder manager cannot address +- **WHEN** graceful drain classifies the session +- **THEN** the actor creates no restart reminder +- **AND** no channel adapter is added by this change diff --git a/openspec/changes/resume-interrupted-sessions/tasks.md b/openspec/changes/resume-interrupted-sessions/tasks.md new file mode 100644 index 000000000..60c9e872f --- /dev/null +++ b/openspec/changes/resume-interrupted-sessions/tasks.md @@ -0,0 +1,24 @@ +## 1. Durable input + +- [x] 1.1 Persist accepted input before acknowledgment. +- [x] 1.2 Store pending input and recent source IDs in snapshots. +- [x] 1.3 Close input from completed, failed, and tool started turns. +- [x] 1.4 Verify journal, snapshot, order, and source retry behavior. + +## 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. + +## 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. + +## 4. Verification + +- [ ] 4.1 Update SPEC-011 and the operations skill. +- [ ] 4.2 Run actor and daemon tests, evals, Slopwatch, headers, and OpenSpec validation. diff --git a/src/Netclaw.Actors.Tests/Protocol/SerializationRoundTripTests.cs b/src/Netclaw.Actors.Tests/Protocol/SerializationRoundTripTests.cs index 8cb27ca01..51e3c1074 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,68 @@ 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 inputId = new InputId("input-1"); + var admitted = new InputAdmitted + { + SessionId = sessionId, + InputId = inputId, + 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(inputId, 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 = [inputId] + }); + Assert.Equal(inputId, Assert.Single(completed.ConsumedInputIds)); + + var closed = RoundTrip(new InputClosed { SessionId = sessionId, InputIds = [inputId] }); + Assert.Equal(inputId, Assert.Single(closed.InputIds)); + + var toolStarted = RoundTrip(new ToolBatchStarted + { + SessionId = sessionId, + UserMessage = admitted.UserMessage, + AssistantMessage = new SerializableChatMessage { Role = ChatRole.Assistant }, + ConsumedInputIds = [inputId] + }); + Assert.Equal(inputId, 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..cd4722461 100644 --- a/src/Netclaw.Actors.Tests/Sessions/LlmSessionIntegrationTests.cs +++ b/src/Netclaw.Actors.Tests/Sessions/LlmSessionIntegrationTests.cs @@ -48,6 +48,55 @@ public LlmSessionIntegrationTests(ITestOutputHelper output) : base(output) { } + [Fact] + public async Task Duplicate_source_message_is_acknowledged_without_a_second_turn() + { + 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); + + await subscriber.FishForMessageAsync( + _ => true, + TimeSpan.FromSeconds(3), + cancellationToken: TestContext.Current.CancellationToken); + var callCount = _fakeChatClient.CallCount; + + var duplicateAck = await manager.Ask(new SendUserMessage + { + SessionId = sessionId, + Content = "Finish the operator task", + Source = source + }, TimeSpan.FromSeconds(3), TestContext.Current.CancellationToken); + + Assert.Equal(sessionId, duplicateAck.SessionId); + await subscriber.ExpectNoMsgAsync( + TimeSpan.FromMilliseconds(250), + TestContext.Current.CancellationToken); + Assert.Equal(callCount, _fakeChatClient.CallCount); + } + 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..db6a786ab 100644 --- a/src/Netclaw.Actors.Tests/Sessions/SessionStateTests.cs +++ b/src/Netclaw.Actors.Tests/Sessions/SessionStateTests.cs @@ -21,6 +21,40 @@ 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 firstId = new InputId("first"); + var secondId = new InputId("second"); + var first = new InputAdmitted + { + SessionId = TestSessionId, + InputId = firstId, + 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 = secondId, + SourceMessageId = "event-2", + UserMessage = new SerializableChatMessage { Role = ChatRole.User, Content = "two" } + }; + + var restored = SessionState.FromSnapshot(SessionState.Empty.Apply(first).Apply(second).ToSnapshot()); + Assert.Equal([firstId, secondId], restored.PendingInputs.Select(input => input.InputId)); + Assert.Contains("slack:user-1:event-1", restored.RecentSourceMessageKeys); + + var closed = restored.CloseInputs([firstId]); + Assert.Equal(secondId, 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/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..c3967bf2c 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.Select(static id => id.Value)); return proto; } @@ -176,7 +179,56 @@ 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.Select(static id => new InputId(id)).ToArray() + }; + + internal static Proto.InputAdmittedProto ToProto(InputAdmitted evt) + { + var proto = new Proto.InputAdmittedProto + { + SessionId = ToProto(evt.SessionId), + InputId = evt.InputId.Value, + 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; + proto.TurnContext = ToProto(evt.TurnContext); + return proto; + } + + internal static InputAdmitted FromProto(Proto.InputAdmittedProto proto) => new() + { + SessionId = FromProto(proto.SessionId), + InputId = new InputId(proto.InputId), + SourceMessageId = proto.HasSourceMessageId ? proto.SourceMessageId : null, + UserMessage = FromProto(proto.UserMessage), + ExecutableText = proto.HasExecutableText ? proto.ExecutableText : null, + TurnContext = proto.TurnContext is null + ? throw new InvalidDataException("An admitted input has no turn context.") + : FromProto(proto.TurnContext), + 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.Select(static id => id.Value)); + return proto; + } + + internal static InputClosed FromProto(Proto.InputClosedProto proto) => new() + { + SessionId = FromProto(proto.SessionId), + InputIds = proto.InputIds.Select(static id => new InputId(id)).ToArray(), + ClosedAtMs = proto.ClosedAtMs }; // ── SessionTitleSet ── @@ -224,20 +276,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.Select(static id => id.Value)); + 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.Select(static id => new InputId(id)).ToArray() }; internal static Proto.ToolCallRecordedProto ToProto(ToolCallRecorded evt) => new() @@ -513,6 +571,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 +586,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..2609e00d4 100644 --- a/src/Netclaw.Actors/Serialization/Protos/netclaw_messages.proto +++ b/src/Netclaw.Actors/Serialization/Protos/netclaw_messages.proto @@ -94,6 +94,24 @@ 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; + int64 admitted_at_ms = 9; + reserved 7, 8; +} + +message InputClosedProto { + SessionIdProto session_id = 1; + repeated string input_ids = 2; + int64 closed_at_ms = 3; } message SessionTitleSetProto { @@ -117,6 +135,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 +291,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 33b2c91cf..dac9d26ec 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(); @@ -282,6 +283,7 @@ public LlmSessionActor( // ── Recovery handlers ── Recover(evt => { + RestoreConsumedInputs(evt.ConsumedInputIds, evt.UserMessage); ApplyTurnRecorded(evt); _toolApprovals.ClearCalls(); ClearApprovalTurnState(); @@ -296,7 +298,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 => @@ -1156,9 +1167,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 @@ -1872,9 +1886,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, @@ -2125,13 +2142,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)) @@ -2205,6 +2225,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(); @@ -2285,9 +2307,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 @@ -2303,7 +2322,7 @@ private void HandleIncomingUserMessage(SendUserMessage cmd) Persist(abandoned, evt => { ApplyToolBatchAbandoned(evt); - ContinueIncomingUserMessage(cmd); + AdmitInput(cmd, ContinueIncomingUserMessage); }); return; } @@ -2314,7 +2333,7 @@ private void HandleIncomingUserMessage(SendUserMessage cmd) Persist(abandoned, evt => { ApplyToolBatchAbandoned(evt); - ContinueIncomingUserMessage(cmd); + AdmitInput(cmd, ContinueIncomingUserMessage); }); return; } @@ -2325,16 +2344,94 @@ 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 = InputId.New(), + 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() with + { + AdoptedSpeakerIds = context.AdoptedSpeakerIds.ToArray() + }, + 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(cmd.Source?.ReminderId); + _inFlightDedup.ReserveBackgroundJob(cmd.Source?.BackgroundJobId); + 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); + 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; @@ -2601,6 +2698,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(); @@ -3107,21 +3220,7 @@ private bool TryHandleSlashCommand(string userContent, IReadOnlyList { _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)) @@ -3728,9 +3769,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; } @@ -4493,6 +4535,32 @@ private ToolBatchAbandoned BuildToolBatchAbandonedEvent(string resultContent) } private void FailCurrentTurn(string errorMessage, Exception cause, ErrorCategory category = ErrorCategory.Unknown) + { + CloseActiveInputs(() => FinishFailedTurn(errorMessage, cause, category)); + } + + private void CloseActiveInputs(Action onClosed) + { + if (_activeInputIds.Count == 0) + { + onClosed(); + return; + } + + Persist(new InputClosed + { + SessionId = _sessionId, + InputIds = _activeInputIds.ToArray(), + ClosedAtMs = NowMs() + }, evt => + { + _state = _state.CloseInputs(evt.InputIds); + _activeInputIds.Clear(); + onClosed(); + }); + } + + 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..6b89b6902 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. + internal InputId? AdmittedInputId { get; init; } } /// diff --git a/src/Netclaw.Actors/Sessions/SessionProtocol.Events.cs b/src/Netclaw.Actors/Sessions/SessionProtocol.Events.cs index 5daf844ce..d61c7771c 100644 --- a/src/Netclaw.Actors/Sessions/SessionProtocol.Events.cs +++ b/src/Netclaw.Actors/Sessions/SessionProtocol.Events.cs @@ -49,11 +49,51 @@ 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 InputId InputId { get; init; } + + public string? SourceMessageId { get; init; } + + public required SerializableChatMessage UserMessage { get; init; } + + /// The executable form when it differs from the user-visible content. + public string? ExecutableText { get; init; } + + public required TurnContextRecord TurnContext { 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 +110,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/SessionProtocol.cs b/src/Netclaw.Actors/Sessions/SessionProtocol.cs index 760ca6d91..4fe80ef5e 100644 --- a/src/Netclaw.Actors/Sessions/SessionProtocol.cs +++ b/src/Netclaw.Actors/Sessions/SessionProtocol.cs @@ -25,6 +25,16 @@ namespace Netclaw.Actors.Sessions; /// public static partial class SessionProtocol { + /// + /// Identifies one input that the session journal accepted. + /// + public readonly record struct InputId(string Value) + { + public static InputId New() => new(Guid.NewGuid().ToString("N")); + + public override string ToString() => Value; + } + /// Marker for an imperative request received by the session actor. public interface ISessionCommand : IWithSessionId { diff --git a/src/Netclaw.Actors/Sessions/SessionState.cs b/src/Netclaw.Actors/Sessions/SessionState.cs index d45a27eaa..7ae3bde7c 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(); + 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,