Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 8 additions & 3 deletions ARCHITECTURE.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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
Expand Down
5 changes: 4 additions & 1 deletion docs/flow-model.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
16 changes: 11 additions & 5 deletions internal/agent/flow.go
Original file line number Diff line number Diff line change
Expand Up @@ -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{
Expand Down Expand Up @@ -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
Expand All @@ -1675,6 +1679,7 @@ func (flow *Flow) writeInputConsumption(
return flow.writeActivity(ctx, AgentEvent{
Kind: EventKindInputConsumed,
Message: strings.Join(messages, " "),
MessageSequence: messageSequence,
InputConsumption: &consumption,
})
}
Expand Down Expand Up @@ -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
Expand All @@ -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
}
}
Expand Down
5 changes: 5 additions & 0 deletions internal/agent/flow_integration_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
Expand Down Expand Up @@ -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)
Expand Down
4 changes: 4 additions & 0 deletions web/packages/superagent-ui/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
33 changes: 33 additions & 0 deletions web/packages/superagent-ui/src/conversationPresentation.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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"),
Expand Down
5 changes: 4 additions & 1 deletion web/packages/superagent-ui/src/conversationPresentation.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}

Expand Down
4 changes: 2 additions & 2 deletions web/src/App.test.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -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);
});

Expand Down
2 changes: 1 addition & 1 deletion web/src/conversation-state.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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: [],
Expand Down
1 change: 1 addition & 0 deletions web/src/conversation-state.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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),
Expand Down
19 changes: 4 additions & 15 deletions web/tests/full-stack.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down Expand Up @@ -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);
Expand Down