diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index 0176268..54edcc9 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -209,7 +209,10 @@ becomes the next watermark before requesting Snapshot. Every consumed queued message, steered message, or Plan execution request emits one `input_consumed` Activity event with exact application IDs or revision and -no user content. The reducer removes only +no user content. Message consumption carries the latest durable user message +sequence so conversation presentation associates it with that turn. Plan-only +consumption has no message sequence and remains a hidden reconciliation hint. +The reducer removes only matching visible inputs and closes the transient `isWaitingForInput` gate. It projects payloads already known from Snapshot into temporary user bubbles without adding content to the Stream. @@ -219,8 +222,10 @@ A later authoritative Snapshot at a real waiting boundary reopens the gate. Stream loss is corrected by Snapshot. Consumed user bubbles are anchored after the latest explicit message sequence -seen in the Activity stream. This preserves causal ordering while Snapshot is -temporarily behind the Stream. +seen before the consumption event. The event's own message sequence associates +the hint with its durable turn and does not advance that temporary projection +boundary. This preserves causal ordering while Snapshot is temporarily behind +the Stream. The approval, manual recovery, and Timer target Steps emit a hidden `snapshot_required` Activity control from `WaitFor`, after the preceding Step diff --git a/docs/flow-model.md b/docs/flow-model.md index 7b93a83..ff95f0a 100644 --- a/docs/flow-model.md +++ b/docs/flow-model.md @@ -296,7 +296,10 @@ Every Execute that consumes a queued or steered message emits one for a consumed Plan request, including a request that became stale in a race. The payload contains only `queuedMessageIds`, `steeredMessageIds`, and nullable `planExecutionRevision`. The display message contains counts or the revision, -never user content. +never user content. Queued consumption carries its new user message sequence; +batched steering carries the latest new user message sequence. Plan-only +consumption has no associated message sequence and is not presented as +conversation history. The browser removes only matching IDs and clears only the matching Plan revision. It temporarily projects consumed payloads already present in Snapshot diff --git a/internal/agent/flow.go b/internal/agent/flow.go index 0d15899..f22a24f 100644 --- a/internal/agent/flow.go +++ b/internal/agent/flow.go @@ -1120,12 +1120,15 @@ func (flow *Flow) beginSteeredTurn(ctx dex.Context, messages []PendingUserMessag return err } } + var latestMessageSequence *Sequence for _, message := range messages { - if _, err := flow.beginUserTurn(ctx, message.Value); err != nil { + sequence, err := flow.beginUserTurn(ctx, message.Value) + if err != nil { return err } + latestMessageSequence = &sequence } - if err := flow.writeInputConsumption(ctx, nil, messages, nil); err != nil { + if err := flow.writeInputConsumption(ctx, nil, messages, nil, latestMessageSequence); err != nil { return err } return flow.writeActivity(ctx, AgentEvent{ @@ -1650,6 +1653,7 @@ func (flow *Flow) writeInputConsumption( queued []PendingUserMessage, steered []PendingUserMessage, planExecutionRevision *PlanRevision, + messageSequence *Sequence, ) error { if len(queued) == 0 && len(steered) == 0 && planExecutionRevision == nil { return nil @@ -1675,6 +1679,7 @@ func (flow *Flow) writeInputConsumption( return flow.writeActivity(ctx, AgentEvent{ Kind: EventKindInputConsumed, Message: strings.Join(messages, " "), + MessageSequence: messageSequence, InputConsumption: &consumption, }) } @@ -2285,10 +2290,11 @@ func (step awaitUserStep) Execute(ctx dex.Context, _ dex.None) (*dex.StepDecisio return nil, err } if len(queued) > 0 { - if _, beginErr := step.flow.beginUserTurn(ctx, queued[0].Value); beginErr != nil { + sequence, beginErr := step.flow.beginUserTurn(ctx, queued[0].Value) + if beginErr != nil { return nil, beginErr } - if consumptionErr := step.flow.writeInputConsumption(ctx, []PendingUserMessage{queued[0]}, nil, nil); consumptionErr != nil { + if consumptionErr := step.flow.writeInputConsumption(ctx, []PendingUserMessage{queued[0]}, nil, nil, &sequence); consumptionErr != nil { return nil, consumptionErr } return dex.GoTo(checkSteeredStep{flow: step.flow}, continueCompactContext), nil @@ -2306,7 +2312,7 @@ func (step awaitUserStep) Execute(ctx dex.Context, _ dex.None) (*dex.StepDecisio } if len(executions) > 0 { revision := executions[0].Revision - if consumptionErr := step.flow.writeInputConsumption(ctx, nil, nil, &revision); consumptionErr != nil { + if consumptionErr := step.flow.writeInputConsumption(ctx, nil, nil, &revision, nil); consumptionErr != nil { return nil, consumptionErr } } diff --git a/internal/agent/flow_integration_test.go b/internal/agent/flow_integration_test.go index ee25b83..e839eb6 100644 --- a/internal/agent/flow_integration_test.go +++ b/internal/agent/flow_integration_test.go @@ -1041,6 +1041,8 @@ func TestAgentFlowDurabilityIntegration(t *testing.T) { len(event.InputConsumption.QueuedMessageIDs) == 1 }) if consumed.Activity.Message != "Consumed 1 queued user message." || + consumed.Activity.MessageSequence == nil || + *consumed.Activity.MessageSequence != state.LastSequence-1 || strings.Contains(consumed.Activity.Message, "hello") { t.Fatalf("queued input consumption Activity = %#v", consumed.Activity) } @@ -2746,6 +2748,9 @@ func assertSteeredInputConsumption( if event.Activity.Message != consumedMessagesDescription(len(consumption.SteeredMessageIDs), "steered") { t.Fatalf("steered input consumption Activity = %#v", event.Activity) } + if event.Activity.MessageSequence == nil { + t.Fatalf("steered input consumption Activity has no user message sequence: %#v", event.Activity) + } for _, messageID := range consumption.SteeredMessageIDs { if _, expected := remaining[messageID]; !expected { t.Fatalf("unexpected steered message ID %q in %#v", messageID, event.Activity) diff --git a/web/packages/superagent-ui/README.md b/web/packages/superagent-ui/README.md index a6c1b24..16689a4 100644 --- a/web/packages/superagent-ui/README.md +++ b/web/packages/superagent-ui/README.md @@ -64,6 +64,10 @@ the presentation does not require a reasoning Stream. Every turn has one collapsed Work log, so losing retained Stream events never destroys the durable conversation structure. +Activity with a durable `messageSequence`, including input-consumption hints, +belongs to that message's turn instead of the global historical section. +Unanchored input-consumption hints remain reconciliation-only and are not shown. + `ConversationTurn.operations` preserves causal order: each model operation is followed by the tools it requested. Adjacent identical calls in one model batch may collapse, but repeated calls separated by another model stay in place. diff --git a/web/packages/superagent-ui/src/conversationPresentation.test.ts b/web/packages/superagent-ui/src/conversationPresentation.test.ts index c0a4358..6b48ebc 100644 --- a/web/packages/superagent-ui/src/conversationPresentation.test.ts +++ b/web/packages/superagent-ui/src/conversationPresentation.test.ts @@ -221,6 +221,39 @@ describe("conversation presentation", () => { ]); }); + it("associates input consumption with its durable user turn", () => { + const messages = [ + message(1, "user", "first"), + message(2, "assistant", "first answer"), + message(3, "user", "second"), + ]; + const activity: TimelineActivityEntry = { + resumeToken: "input-consumed", + source: "await-user", + createdAt: iso(3_000), + value: { + kind: "input_consumed", + message: "Consumed 1 queued user message.", + messageSequence: 3, + }, + }; + + const view = buildConversationPresentation({ + messages, + activities: [ + { + ...activity, + resumeToken: "legacy-input-consumed", + value: { ...activity.value, messageSequence: null }, + }, + activity, + ], + }); + + expect(view.earlierActivity).toHaveLength(0); + expect(required(view.turns[1]).activities).toContain(activity); + }); + it("keeps current and resolved waits separate from tool execution", () => { const messages = [ message(1, "user", "build"), diff --git a/web/packages/superagent-ui/src/conversationPresentation.ts b/web/packages/superagent-ui/src/conversationPresentation.ts index fe2c8a9..4d1acb4 100644 --- a/web/packages/superagent-ui/src/conversationPresentation.ts +++ b/web/packages/superagent-ui/src/conversationPresentation.ts @@ -271,7 +271,10 @@ export function buildConversationPresentation( turn.activities.push(activity); if (HISTORICAL_ACTIVITY.has(activity.value.kind)) turn.historicalActivities.push(activity); - } else if (!MODEL_LIFECYCLE.has(activity.value.kind)) + } else if ( + !MODEL_LIFECYCLE.has(activity.value.kind) && + activity.value.kind !== "input_consumed" + ) earlierActivity.push(activity); } diff --git a/web/src/App.test.tsx b/web/src/App.test.tsx index 41fc68e..9416dbb 100644 --- a/web/src/App.test.tsx +++ b/web/src/App.test.tsx @@ -886,8 +886,8 @@ describe("App", () => { ).toBeDisabled(); }); expect( - screen.getByText("Consumed plan execution request for revision 7."), - ).toBeInTheDocument(); + screen.queryByText("Consumed plan execution request for revision 7."), + ).not.toBeInTheDocument(); expect(getAgentSnapshot).toHaveBeenCalledTimes(1); }); diff --git a/web/src/conversation-state.test.ts b/web/src/conversation-state.test.ts index ba831fc..166fc4c 100644 --- a/web/src/conversation-state.test.ts +++ b/web/src/conversation-state.test.ts @@ -360,7 +360,7 @@ describe("conversationReducer", () => { message: "Consumed 1 queued user message.", callId: null, toolName: null, - messageSequence: null, + messageSequence: 4, inputConsumption: { queuedMessageIds: ["queued-1"], steeredMessageIds: [], diff --git a/web/src/conversation-state.ts b/web/src/conversation-state.ts index 29e4496..7343e0f 100644 --- a/web/src/conversation-state.ts +++ b/web/src/conversation-state.ts @@ -790,6 +790,7 @@ function applyInputConsumption( ); const consumedAfterSequence = state.activities.reduce( (latestSequence, activity) => + activity.resumeToken === update.resumeToken || activity.value.messageSequence === null ? latestSequence : Math.max(latestSequence, activity.value.messageSequence), diff --git a/web/tests/full-stack.spec.ts b/web/tests/full-stack.spec.ts index a18902b..6af37e1 100644 --- a/web/tests/full-stack.spec.ts +++ b/web/tests/full-stack.spec.ts @@ -108,13 +108,7 @@ test("renders chronological transient activity and durable queue interactions", "/reason Checked the constraints | Durable answer", ), ).toHaveCount(0); - const historicalActivity = history.locator("details.sa-earlier-activity"); - await expect(historicalActivity.locator("summary")).toContainText( - "Historical activity ยท 1 event", - ); - await expect(historicalActivity).toContainText( - "Consumed 1 queued user message.", - ); + await expect(history.locator("details.sa-earlier-activity")).toHaveCount(0); expect(snapshots.length).toBeGreaterThanOrEqual(2); const workLog = history.locator("details.sa-work-log").last(); @@ -349,14 +343,9 @@ test("removes a consumed queued message before the next Snapshot completes", asy expect(steer.status()).toBe(200); await expect(queuedMessage).toHaveCount(0); - const recoveredActivity = page - .getByRole("region", { name: "Conversation history" }) - .locator("details.sa-earlier-activity"); - await expect(recoveredActivity).toContainText( - "Consumed 1 steered user message.", - ); - const consumedMessage = page - .getByRole("region", { name: "Conversation history" }) + const history = page.getByRole("region", { name: "Conversation history" }); + await expect(history).not.toContainText("Consumed 1 steered user message."); + const consumedMessage = history .locator(".message-bubble.user") .filter({ hasText: "consume this queued message" }); await expect(consumedMessage).toHaveCount(1);