From 5d8605ab0f049b6487164ab46b6b75f77af63584 Mon Sep 17 00:00:00 2001 From: me2seeks Date: Thu, 20 Aug 2026 18:39:10 +0800 Subject: [PATCH 1/3] fix(runtime-host): forward the durable steering echo to session subscribers The TUI never rendered a consumed steering message: the host event pump dropped steering_message (not a transient tool/text event), leaving the transient in-flight queue observation as the only live render path, which the coalesced canonical refresh can skip entirely (queued -> consumed). Forward the durable steering_message on the session-event frame so every subscriber renders the interjection in place, keep the queue in-flight synthesis for the attach/rejoin window, and dedup both paths per message. Bumps the session continuity wire schema to 5. Fixes apache/maka#3304 --- .../__tests__/execution-host-message.test.ts | 58 ++++++++ .../src/__tests__/protocol.test.ts | 33 ++++- .../session-continuity-coordinator.test.ts | 33 +++++ .../src/__tests__/session-projector.test.ts | 125 ++++++++++++++++++ .../src/adapter/session-projector.ts | 30 ++++- .../src/protocol/session-continuity.ts | 53 +++++++- .../src/server/root-turn-coordinator.ts | 1 + .../server/session-continuity-coordinator.ts | 23 +++- 8 files changed, 346 insertions(+), 10 deletions(-) diff --git a/packages/runtime-host/src/__tests__/execution-host-message.test.ts b/packages/runtime-host/src/__tests__/execution-host-message.test.ts index 9f69b0c2a5..a7800c6c18 100644 --- a/packages/runtime-host/src/__tests__/execution-host-message.test.ts +++ b/packages/runtime-host/src/__tests__/execution-host-message.test.ts @@ -95,6 +95,64 @@ import { withTimeout, } from './fixtures/execution-host-suite.js'; +test('subscribed Clients receive the durable steering echo as a session event', async () => { + await withExecutionRoot(async (fixture) => { + const host = await fixture.startHost(); + const client = await connectClient(fixture.root, 'tui'); + const subscription = await client.openSessionSubscription({ + sessionId: fixture.sessionId, + transcript: { kind: 'none' }, + }); + const probe = new SubscriptionProbe(subscription); + + const turnId = randomUUID(); + await client.startTurn({ + sessionId: fixture.sessionId, + turnId, + content: { text: FAKE_WAIT_FOR_STEERING_PROMPT }, + }); + const steeringId = randomUUID(); + const steeringContent = { + text: 'steer mid-turn', + displayText: 'steer mid-turn', + }; + const submitted = await client.request('turn.message.submit', { + originHostEpoch: host.hostEpoch, + sessionId: fixture.sessionId, + messageId: steeringId, + content: steeringContent, + placement: 'current_turn', + }); + assert.equal(submitted.disposition, 'steering'); + + // apache/maka#3304: the steering render must not depend on observing the + // transient in-flight queue state; the durable echo is forwarded verbatim. + const echoed = await probe.waitFor( + (frame) => + frame.kind === 'subscription.session_event' && frame.event.type === 'steering_message', + 'continuity did not forward the durable steering echo', + ); + assert.equal(echoed.kind, 'subscription.session_event'); + if (echoed.kind === 'subscription.session_event') { + assert.equal(echoed.event.type, 'steering_message'); + if (echoed.event.type === 'steering_message') { + assert.equal(echoed.event.turnId, turnId); + assert.equal(echoed.event.messageId, steeringId); + assert.deepEqual(echoed.event.content, steeringContent); + } + } + + assert.equal( + (await waitForTerminalTurn(client, fixture.sessionId, turnId)).status, + 'completed', + ); + await subscription.close(); + await probe.done; + await client.close(); + await fixture.stopHost(host); + }); +}); + test('steering becomes durable and ordered followups automatically start the next root', async () => { await withExecutionRoot(async (fixture) => { const host = await fixture.startHost(); diff --git a/packages/runtime-host/src/__tests__/protocol.test.ts b/packages/runtime-host/src/__tests__/protocol.test.ts index a1a7b3810a..6fcbbcaa64 100644 --- a/packages/runtime-host/src/__tests__/protocol.test.ts +++ b/packages/runtime-host/src/__tests__/protocol.test.ts @@ -106,7 +106,7 @@ describe('Runtime Host bootstrap protocol', () => { }); test('keeps the subscription queue Epoch correlated', () => { - assert.equal(SESSION_CONTINUITY_SCHEMA_VERSION, 4); + assert.equal(SESSION_CONTINUITY_SCHEMA_VERSION, 5); const opened = { requestId: 'open-1', operation: 'subscription.open', @@ -305,6 +305,37 @@ describe('Runtime Host bootstrap protocol', () => { ]) { assert.throws(() => decodeHostFrame({ ...envelope, event }), isInvalidFrame); } + + // The durable steering echo shares the session-event frame without a + // toolUseId; unknown keys stay rejected. + const steering = { + type: 'steering_message' as const, + id: 'steering-event-1', + turnId: 'turn-1', + ts: 7, + messageId: 'steering-message-1', + content: { text: 'steer the turn' }, + }; + const decodedSteering = decodeHostFrame({ ...envelope, event: steering }); + assert.ok('kind' in decodedSteering); + if ('kind' in decodedSteering) { + assert.equal(decodedSteering.kind, 'subscription.session_event'); + if (decodedSteering.kind === 'subscription.session_event') { + assert.deepEqual(decodedSteering.event, steering); + } + } + assert.throws( + () => decodeHostFrame({ ...envelope, event: { ...steering, toolUseId: 'tool-1' } }), + isInvalidFrame, + ); + assert.throws( + () => + decodeHostFrame({ + ...envelope, + event: { ...steering, content: { text: 'x'.repeat(49 * 1024) } }, + }), + isInvalidFrame, + ); assert.throws( () => decodeHostFrame({ diff --git a/packages/runtime-host/src/__tests__/session-continuity-coordinator.test.ts b/packages/runtime-host/src/__tests__/session-continuity-coordinator.test.ts index 9dc5e30a56..120453b28a 100644 --- a/packages/runtime-host/src/__tests__/session-continuity-coordinator.test.ts +++ b/packages/runtime-host/src/__tests__/session-continuity-coordinator.test.ts @@ -69,6 +69,39 @@ test('open is an inactive publication barrier and live sequence starts at nextSe coordinator.close(); }); +test('forwards the durable steering echo to subscribers as a session event', async () => { + const sink = new RecordingSink(); + const coordinator = new SessionContinuityCoordinator( + HOST_EPOCH, + async () => canonical(), + new SessionAdmissionGate(), + ); + const connection = coordinator.attachConnection('connection-1', sink); + const opened = await open(coordinator, 'connection-1'); + connection.activate(opened.subscriptionId); + await delayImmediate(); + sink.frames.length = 0; + + const steering = { + type: 'steering_message' as const, + id: 'steering-event-1', + turnId: 'turn-1', + ts: 7, + messageId: 'steering-message-1', + content: { text: 'steer the turn' }, + }; + await coordinator.acceptRuntimeEvent(SESSION_ID, 'run-1', steering); + + assert.equal(sink.frames.length, 1); + const frame = sink.frames[0]; + assert.equal(frame?.kind, 'subscription.session_event'); + if (frame?.kind !== 'subscription.session_event') return; + assert.deepEqual(frame.event, steering); + + connection.abort(opened.subscriptionId); + coordinator.close(); +}); + test('open snapshot includes pending Interactions from the canonical projection', async () => { const pending = pendingInteraction(); const coordinator = new SessionContinuityCoordinator( diff --git a/packages/runtime-host/src/__tests__/session-projector.test.ts b/packages/runtime-host/src/__tests__/session-projector.test.ts index a3b0e462e6..8079c5f6e2 100644 --- a/packages/runtime-host/src/__tests__/session-projector.test.ts +++ b/packages/runtime-host/src/__tests__/session-projector.test.ts @@ -8,6 +8,7 @@ import { import { SESSION_CONTINUITY_SCHEMA_VERSION, type SessionContinuitySnapshot, + type SteeringMessageSnapshot, type SubscriptionFrame, } from '../protocol/index.js'; @@ -243,6 +244,130 @@ test('does not replay settled transcript steps when the active step reaches term ); }); +test('projects the durable steering echo even when the in-flight queue state was never observed', () => { + // Regression for apache/maka#3304: the coalesced canonical refresh can jump + // the queue straight from queued to consumed, so the in-flight synthesis + // never fires. The forwarded steering_message event must render the message. + const projector = new RuntimeHostSessionProjector( + snapshot({ queue: queue(2, [steeringEntry('queued')]) }), + createRuntimeHostSessionProjectionSeed([], snapshot()), + () => 10, + ); + + const skipped = projector.accept({ + kind: 'subscription.session_projection', + hostEpoch: 'host-1', + subscriptionId: 'subscription-1', + sequence: 1, + snapshot: snapshot({ queue: queue(4, []) }), + }); + assert.deepEqual( + skipped.events.map((event) => event.type), + ['queue_update'], + ); + + const echoed = projector.accept(steeringFrame(2)).events; + assert.equal(echoed.length, 1); + assert.deepEqual(echoed[0], { + type: 'steering_message', + id: 'steering-event-1', + turnId: 'turn-1', + ts: 10, + messageId: 'steering-message-1', + content: { text: 'steer the turn' }, + }); +}); + +test('projects a steering message exactly once across both authoritative paths', () => { + // The queue in-flight synthesis and the durable session-event echo race; + // whichever projects the message first suppresses the other. + const inFlightFirst = new RuntimeHostSessionProjector( + snapshot({ queue: queue(2, [steeringEntry('queued')]) }), + createRuntimeHostSessionProjectionSeed([], snapshot()), + () => 10, + ); + const synthesized = inFlightFirst.accept({ + kind: 'subscription.session_projection', + hostEpoch: 'host-1', + subscriptionId: 'subscription-1', + sequence: 1, + snapshot: snapshot({ queue: queue(3, [steeringEntry('in_flight')]) }), + }); + assert.deepEqual( + synthesized.events.map((event) => event.type), + ['steering_message', 'queue_update'], + ); + assert.deepEqual(inFlightFirst.accept(steeringFrame(2)).events, []); + + const echoFirst = new RuntimeHostSessionProjector( + snapshot({ queue: queue(2, [steeringEntry('queued')]) }), + createRuntimeHostSessionProjectionSeed([], snapshot()), + () => 10, + ); + assert.equal(echoFirst.accept(steeringFrame(1)).events.length, 1); + const suppressed = echoFirst.accept({ + kind: 'subscription.session_projection', + hostEpoch: 'host-1', + subscriptionId: 'subscription-1', + sequence: 2, + snapshot: snapshot({ queue: queue(3, [steeringEntry('in_flight')]) }), + }); + assert.deepEqual( + suppressed.events.map((event) => event.type), + ['queue_update'], + ); +}); + +test('seeds an unrendered in-flight steering message once on rejoin', () => { + const projector = new RuntimeHostSessionProjector( + snapshot({ queue: queue(3, [steeringEntry('in_flight')]) }), + createRuntimeHostSessionProjectionSeed([], snapshot()), + () => 10, + ); + assert.deepEqual( + projector.seedActive(false).map((event) => event.type), + ['steering_message', 'queue_update'], + ); + // A live echo of the same message arriving after the seed is the duplicate. + assert.deepEqual(projector.accept(steeringFrame(1)).events, []); +}); + +function steeringEntry(state: 'queued' | 'in_flight'): SteeringMessageSnapshot { + return { + entryId: 'entry-1', + messageId: 'steering-message-1', + content: { text: 'steer the turn' }, + placement: 'current_turn', + state, + }; +} + +function queue( + queueRevision: number, + steering: readonly SteeringMessageSnapshot[], +): SessionContinuitySnapshot['queue'] { + return { hostEpoch: 'host-1', queueRevision, steering, followup: [] }; +} + +function steeringFrame(sequence: number): SubscriptionFrame { + return { + kind: 'subscription.session_event', + hostEpoch: 'host-1', + subscriptionId: 'subscription-1', + sequence, + sessionId: 'session-1', + runId: 'run-1', + event: { + type: 'steering_message', + id: 'steering-event-1', + turnId: 'turn-1', + ts: 10, + messageId: 'steering-message-1', + content: { text: 'steer the turn' }, + }, + }; +} + function deltaFrame( sequence: number, startOffset: number, diff --git a/packages/runtime-host/src/adapter/session-projector.ts b/packages/runtime-host/src/adapter/session-projector.ts index 6a37dc2457..cd9885e21f 100644 --- a/packages/runtime-host/src/adapter/session-projector.ts +++ b/packages/runtime-host/src/adapter/session-projector.ts @@ -65,6 +65,12 @@ export class RuntimeHostSessionProjector { #snapshot: SessionContinuitySnapshot; readonly #now: () => number; readonly #transcriptIds: Set; + /** + * One render per steering message, whichever authoritative source projects + * it first: the durable session-event echo (in stream order) or the queue + * in-flight synthesis (which can win the race or cover a rejoin window). + */ + readonly #renderedSteeringMessageIds = new Set(); readonly #accumulators = new Map(); constructor( @@ -148,6 +154,8 @@ export class RuntimeHostSessionProjector { } for (const entry of rootQueueInFlight(this.#snapshot.queue)) { if (this.#transcriptIds.has(entry.messageId)) continue; + if (this.#renderedSteeringMessageIds.has(entry.messageId)) continue; + this.#renderedSteeringMessageIds.add(entry.messageId); events.push({ type: 'steering_message', id: `host-queue:${this.#snapshot.queue.hostEpoch}:${this.#snapshot.queue.queueRevision}:${entry.entryId}`, @@ -321,7 +329,15 @@ export class RuntimeHostSessionProjector { } if (frame.kind === 'subscription.session_event') { const event = projectToolEvent(frame); - if (event) events.push(event); + if (event) { + if (event.type === 'steering_message') { + if (this.#renderedSteeringMessageIds.has(event.messageId)) { + return emptyUpdate(events); + } + this.#renderedSteeringMessageIds.add(event.messageId); + } + events.push(event); + } return emptyUpdate(events); } if (frame.kind !== 'subscription.session_projection') return emptyUpdate(events); @@ -340,6 +356,8 @@ export class RuntimeHostSessionProjector { const root = next.rootTurn; if (root && queueChanged(previousSnapshot.queue, next.queue)) { for (const entry of newlyInFlight(previousSnapshot.queue, next.queue)) { + if (this.#renderedSteeringMessageIds.has(entry.messageId)) continue; + this.#renderedSteeringMessageIds.add(entry.messageId); events.push({ type: 'steering_message', id: `host-queue:${next.queue.hostEpoch}:${next.queue.queueRevision}:${entry.entryId}`, @@ -462,6 +480,16 @@ function projectToolEvent( frame: Extract, ): SessionEvent | undefined { const event = frame.event; + if (event.type === 'steering_message') { + return { + type: 'steering_message', + id: event.id, + turnId: event.turnId, + ts: event.ts, + messageId: event.messageId, + content: structuredClone(event.content), + }; + } const base = { id: event.id, turnId: event.turnId, diff --git a/packages/runtime-host/src/protocol/session-continuity.ts b/packages/runtime-host/src/protocol/session-continuity.ts index cb713f3c88..d310cfb610 100644 --- a/packages/runtime-host/src/protocol/session-continuity.ts +++ b/packages/runtime-host/src/protocol/session-continuity.ts @@ -22,7 +22,12 @@ import { type SessionMessageQueueProjection, } from './message.js'; import { defineOperation } from './operation-spec.js'; -import { decodeTurnSnapshot, type TurnSnapshot } from './turn.js'; +import { + decodeMessageContent, + decodeTurnSnapshot, + type MessageContent, + type TurnSnapshot, +} from './turn.js'; import { decodeGoalProjection, type GoalProjection } from './goal.js'; import { decodeRuntimeResourceRef } from './runtime-resource.js'; import { @@ -31,7 +36,7 @@ import { type SessionTranscriptBootstrap, } from './session-transcript.js'; -export const SESSION_CONTINUITY_SCHEMA_VERSION = 4 as const; +export const SESSION_CONTINUITY_SCHEMA_VERSION = 5 as const; export const SESSION_CONTINUITY_SNAPSHOT_MAX_BYTES = 56 * 1024; // Leave transport headroom for the response envelope and request correlation. export const SUBSCRIPTION_OPEN_RESULT_MAX_BYTES = 92 * 1024; @@ -170,11 +175,26 @@ export type SessionToolEvent = content: ToolResultPreviewContent; }); +/** + * The durable mid-turn user interjection (steering), forwarded verbatim from the + * run's event stream. Unlike tool events it has no toolUseId; it shares the + * frame so subscribers render the interjection in place without depending on + * observing the transient in-flight queue state. + */ +export interface SessionSteeringEvent { + type: 'steering_message'; + id: string; + turnId: string; + ts: number; + messageId: string; + content: MessageContent; +} + export interface SessionEventFrame extends SubscriptionEnvelope { kind: 'subscription.session_event'; sessionId: string; runId: string; - event: SessionToolEvent; + event: SessionToolEvent | SessionSteeringEvent; } export interface SessionTranscriptAdvancedFrame extends SubscriptionEnvelope { @@ -339,7 +359,7 @@ export function decodeSubscriptionFrame(value: unknown): SubscriptionFrame { ...envelope, sessionId: requireEntityId(record.sessionId, 'sessionId'), runId: requireEntityId(record.runId, 'runId'), - event: decodeSessionToolEvent(record.event), + event: decodeSessionFrameEvent(record.event), }; } else if (record.kind === 'subscription.transcript_advanced') { assertExactKeys(record, 'Session transcript advanced frame', [ @@ -684,6 +704,31 @@ function decodeAssistantDelta(value: unknown): SessionAssistantDelta { }; } +function decodeSessionFrameEvent(value: unknown): SessionToolEvent | SessionSteeringEvent { + const record = requireRecord(value, 'Session event'); + if (record.type === 'steering_message') return decodeSessionSteeringEvent(record); + return decodeSessionToolEvent(record); +} + +function decodeSessionSteeringEvent(record: Record): SessionSteeringEvent { + assertExactKeys(record, 'Session steering event', [ + 'type', + 'id', + 'turnId', + 'ts', + 'messageId', + 'content', + ]); + return { + type: 'steering_message', + id: requireId(record.id, 'Session steering event id'), + turnId: requireEntityId(record.turnId, 'turnId'), + ts: requireCount(record.ts, 'Session steering event timestamp'), + messageId: requireEntityId(record.messageId, 'messageId'), + content: decodeMessageContent(record.content), + }; +} + function decodeSessionToolEvent(value: unknown): SessionToolEvent { const record = requireRecord(value, 'Session tool event'); const identity = { diff --git a/packages/runtime-host/src/server/root-turn-coordinator.ts b/packages/runtime-host/src/server/root-turn-coordinator.ts index b9afbd4fd0..d0dd484a41 100644 --- a/packages/runtime-host/src/server/root-turn-coordinator.ts +++ b/packages/runtime-host/src/server/root-turn-coordinator.ts @@ -2787,6 +2787,7 @@ function isRuntimeSessionTransientEvent( event.type === 'tool_progress' || event.type === 'tool_result_preview' || event.type === 'tool_result' || + event.type === 'steering_message' || event.type === 'provider_retry' ); } diff --git a/packages/runtime-host/src/server/session-continuity-coordinator.ts b/packages/runtime-host/src/server/session-continuity-coordinator.ts index 02d50c741b..f9e36eda31 100644 --- a/packages/runtime-host/src/server/session-continuity-coordinator.ts +++ b/packages/runtime-host/src/server/session-continuity-coordinator.ts @@ -19,6 +19,7 @@ import { type SessionDomainChangedFrame, type SessionEventFrame, type SessionRuntimeResourcePtyDataFrame, + type SessionSteeringEvent, type SessionToolEvent, type SessionTranscriptAdvancedFrame, type SessionTranscriptPageInput, @@ -73,6 +74,7 @@ export type RuntimeSessionTransientEvent = Extract< | 'tool_progress' | 'tool_result_preview' | 'tool_result' + | 'steering_message' | 'provider_retry'; } >; @@ -714,7 +716,7 @@ export class SessionContinuityCoordinator implements SessionContinuityService { } else if (event.type === 'tool_result') { state.toolResultPreviews.delete(event.toolUseId); } - const projected = projectToolEvent(event); + const projected = projectSessionEvent(event); for (const subscriber of state.subscribers.values()) { const frame: SessionEventFrame = { kind: 'subscription.session_event', @@ -948,7 +950,7 @@ export class SessionContinuityCoordinator implements SessionContinuityService { sequence: subscriber.nextSequence, sessionId, runId: rootTurn.runId, - event: projectToolEvent(preview), + event: projectSessionEvent(preview), }; this.#enqueue(subscriber, frame); } @@ -1899,7 +1901,7 @@ function jsonStringContentBytes(value: string): number { return Buffer.byteLength(encoded.slice(1, -1), 'utf8'); } -function projectToolEvent( +function projectSessionEvent( event: Exclude< RuntimeSessionTransientEvent, { @@ -1911,7 +1913,20 @@ function projectToolEvent( | 'provider_retry'; } >, -): SessionToolEvent { +): SessionToolEvent | SessionSteeringEvent { + if (event.type === 'steering_message') { + // The durable steering echo: forwarded verbatim so subscribers render the + // interjection in place instead of depending on observing the transient + // in-flight queue state. + return { + type: 'steering_message', + id: event.id, + turnId: event.turnId, + ts: event.ts, + messageId: event.messageId, + content: structuredClone(event.content), + }; + } const identity = { id: event.id, turnId: event.turnId, From d6010baa02f4889f6cc407406fe524e2b1bc918d Mon Sep 17 00:00:00 2001 From: me2seeks Date: Thu, 20 Aug 2026 20:11:29 +0800 Subject: [PATCH 2/3] refactor(runtime-host): name the forwarding rule, not the durability class MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit isRuntimeSessionForwardedEvent / RuntimeSessionForwardedEvent describe what the predicate decides — forward live to subscribers — now that the durable steering_message belongs to it. Rename the client projectSessionEvent to match its Host-side counterpart. Addresses review on #3316. --- .../runtime-host/src/adapter/session-projector.ts | 4 ++-- .../runtime-host/src/server/root-turn-coordinator.ts | 12 ++++++++---- .../src/server/session-continuity-coordinator.ts | 10 +++++----- 3 files changed, 15 insertions(+), 11 deletions(-) diff --git a/packages/runtime-host/src/adapter/session-projector.ts b/packages/runtime-host/src/adapter/session-projector.ts index cd9885e21f..13755ffc21 100644 --- a/packages/runtime-host/src/adapter/session-projector.ts +++ b/packages/runtime-host/src/adapter/session-projector.ts @@ -328,7 +328,7 @@ export class RuntimeHostSessionProjector { return emptyUpdate(events); } if (frame.kind === 'subscription.session_event') { - const event = projectToolEvent(frame); + const event = projectSessionEvent(frame); if (event) { if (event.type === 'steering_message') { if (this.#renderedSteeringMessageIds.has(event.messageId)) { @@ -476,7 +476,7 @@ export function projectRuntimeHostInteractionRequest( return []; } -function projectToolEvent( +function projectSessionEvent( frame: Extract, ): SessionEvent | undefined { const event = frame.event; diff --git a/packages/runtime-host/src/server/root-turn-coordinator.ts b/packages/runtime-host/src/server/root-turn-coordinator.ts index d0dd484a41..11842217b2 100644 --- a/packages/runtime-host/src/server/root-turn-coordinator.ts +++ b/packages/runtime-host/src/server/root-turn-coordinator.ts @@ -73,7 +73,7 @@ import type { ConnectionContext, TurnOperationHandlerMap } from './operation-dis import { RootAdmissionOwner } from './root-admission-owner.js'; import { type SessionAdmissionLease, SessionAdmissionGate } from './session-admission-gate.js'; import { - type RuntimeSessionTransientEvent, + type RuntimeSessionForwardedEvent, SessionContinuityCoordinator, } from './session-continuity-coordinator.js'; import type { @@ -2096,7 +2096,7 @@ export class RootTurnCoordinator implements HostedExecutionAuthority { // Presentation observers do not participate in execution authority. } } - if (isRuntimeSessionTransientEvent(event)) { + if (isRuntimeSessionForwardedEvent(event)) { await this.continuity.acceptRuntimeEvent(input.sessionId, active.runId, event); } else if (isInteractionAnswerAck(event)) { await this.continuity.refreshCanonical(input.sessionId); @@ -2774,9 +2774,13 @@ function isStoppedInteractionAdmission( ); } -function isRuntimeSessionTransientEvent( +// Membership answers one question: forward this event live to subscribers via +// the continuity coordinator instead of letting the canonical refresh carry +// it. Persistence is orthogonal — it happens upstream in the run's own event +// stream, which is why the durable steering_message belongs here. +function isRuntimeSessionForwardedEvent( event: SessionEvent, -): event is RuntimeSessionTransientEvent { +): event is RuntimeSessionForwardedEvent { return ( event.type === 'text_delta' || event.type === 'text_complete' || diff --git a/packages/runtime-host/src/server/session-continuity-coordinator.ts b/packages/runtime-host/src/server/session-continuity-coordinator.ts index f9e36eda31..d6664204d9 100644 --- a/packages/runtime-host/src/server/session-continuity-coordinator.ts +++ b/packages/runtime-host/src/server/session-continuity-coordinator.ts @@ -61,7 +61,7 @@ const MAX_SUBSCRIBER_QUEUED_BYTES = 256 * 1024; export type { CanonicalSessionProjection } from './canonical-session-projection.js'; -export type RuntimeSessionTransientEvent = Extract< +export type RuntimeSessionForwardedEvent = Extract< SessionEvent, { type: @@ -96,7 +96,7 @@ interface SessionProjectionState { */ toolResultPreviews: Map< string, - Extract + Extract >; terminalPublicationFence?: TerminalPublicationFence; } @@ -599,7 +599,7 @@ export class SessionContinuityCoordinator implements SessionContinuityService { async acceptRuntimeEvent( sessionId: string, runId: string, - event: RuntimeSessionTransientEvent, + event: RuntimeSessionForwardedEvent, ): Promise { if ( (event.type === 'text_delta' || event.type === 'thinking_delta') && @@ -1432,7 +1432,7 @@ export class SessionContinuityCoordinator implements SessionContinuityService { subscriber: Subscriber, sessionId: string, runId: string, - event: Extract, + event: Extract, kind: SessionAssistantDelta['kind'], startOffset: number, ): void { @@ -1903,7 +1903,7 @@ function jsonStringContentBytes(value: string): number { function projectSessionEvent( event: Exclude< - RuntimeSessionTransientEvent, + RuntimeSessionForwardedEvent, { type: | 'text_delta' From 6397a567776876817957f6cff62a373b2de1903a Mon Sep 17 00:00:00 2001 From: me2seeks Date: Thu, 20 Aug 2026 22:09:24 +0800 Subject: [PATCH 3/3] fix(runtime-host): suppress the steering echo when the bootstrap already rendered it subscription.open can bootstrap the durable steering message and install the subscriber before the Host's forwarded echo arrives; seed the render-dedup set from steering messages already durable in the transcript so the bootstrapped render stays the only one. Also rebased onto main and updated the new e2e test to the post-#3277 single-argument connectClient. Addresses review on apache/maka#3316 --- .../__tests__/execution-host-message.test.ts | 2 +- .../src/__tests__/session-projector.test.ts | 43 ++++++++++++++++++- .../src/adapter/session-projector.ts | 18 ++++++-- 3 files changed, 57 insertions(+), 6 deletions(-) diff --git a/packages/runtime-host/src/__tests__/execution-host-message.test.ts b/packages/runtime-host/src/__tests__/execution-host-message.test.ts index a7800c6c18..9e1e84ae02 100644 --- a/packages/runtime-host/src/__tests__/execution-host-message.test.ts +++ b/packages/runtime-host/src/__tests__/execution-host-message.test.ts @@ -98,7 +98,7 @@ import { test('subscribed Clients receive the durable steering echo as a session event', async () => { await withExecutionRoot(async (fixture) => { const host = await fixture.startHost(); - const client = await connectClient(fixture.root, 'tui'); + const client = await connectClient(fixture.root); const subscription = await client.openSessionSubscription({ sessionId: fixture.sessionId, transcript: { kind: 'none' }, diff --git a/packages/runtime-host/src/__tests__/session-projector.test.ts b/packages/runtime-host/src/__tests__/session-projector.test.ts index 8079c5f6e2..d10729d041 100644 --- a/packages/runtime-host/src/__tests__/session-projector.test.ts +++ b/packages/runtime-host/src/__tests__/session-projector.test.ts @@ -332,6 +332,31 @@ test('seeds an unrendered in-flight steering message once on rejoin', () => { assert.deepEqual(projector.accept(steeringFrame(1)).events, []); }); +test('suppresses the live echo for a steering message already durable in the bootstrap', () => { + // subscription.open can bootstrap the durable steering message and install + // the subscriber while the Host's forwarded echo for it is still pending: + // the bootstrapped render must stay the only one (apache/maka#3316 review). + const inFlight = snapshot({ queue: queue(3, [steeringEntry('in_flight')]) }); + const projector = new RuntimeHostSessionProjector( + inFlight, + createRuntimeHostSessionProjectionSeed( + [userSteering('steering-message-1', 'steering-event-1')], + inFlight, + ), + () => 10, + ); + + // Durable and in-flight: no synthesis seed… + assert.deepEqual( + projector.seedActive(false).map((event) => event.type), + ['queue_update'], + ); + // …and the late echo of the same message is the duplicate. + assert.deepEqual(projector.accept(steeringFrame(1)).events, []); + // A different steering message still renders normally. + assert.equal(projector.accept(steeringFrame(2, 'steering-message-2')).events.length, 1); +}); + function steeringEntry(state: 'queued' | 'in_flight'): SteeringMessageSnapshot { return { entryId: 'entry-1', @@ -349,7 +374,7 @@ function queue( return { hostEpoch: 'host-1', queueRevision, steering, followup: [] }; } -function steeringFrame(sequence: number): SubscriptionFrame { +function steeringFrame(sequence: number, messageId = 'steering-message-1'): SubscriptionFrame { return { kind: 'subscription.session_event', hostEpoch: 'host-1', @@ -362,12 +387,26 @@ function steeringFrame(sequence: number): SubscriptionFrame { id: 'steering-event-1', turnId: 'turn-1', ts: 10, - messageId: 'steering-message-1', + messageId, content: { text: 'steer the turn' }, }, }; } +function userSteering( + id: string, + steeringEventId: string, +): Extract { + return { + type: 'user', + id, + turnId: 'turn-1', + ts: 1, + text: 'steer the turn', + steeringEventId, + }; +} + function deltaFrame( sequence: number, startOffset: number, diff --git a/packages/runtime-host/src/adapter/session-projector.ts b/packages/runtime-host/src/adapter/session-projector.ts index 13755ffc21..d979f79a83 100644 --- a/packages/runtime-host/src/adapter/session-projector.ts +++ b/packages/runtime-host/src/adapter/session-projector.ts @@ -24,6 +24,13 @@ interface AssistantAccumulator { export interface RuntimeHostSessionProjectionSeed { readonly durableInFlightMessageIds: readonly string[]; + /** + * messageIds of steering interjections already durable in the transcript. + * `subscription.open` can bootstrap message X and install the subscriber + * before the Host's live echo for X arrives; these ids suppress that late + * echo so the bootstrapped render stays the only one. + */ + readonly durableSteeringMessageIds: readonly string[]; readonly activeAssistantMessages: readonly Extract[]; } @@ -38,6 +45,9 @@ export function createRuntimeHostSessionProjectionSeed( durableInFlightMessageIds: transcript .filter((message) => inFlightMessageIds.has(message.id)) .map((message) => message.id), + durableSteeringMessageIds: transcript + .filter((message) => message.type === 'user' && message.steeringEventId !== undefined) + .map((message) => message.id), activeAssistantMessages: snapshot.rootTurn === null ? [] @@ -67,10 +77,11 @@ export class RuntimeHostSessionProjector { readonly #transcriptIds: Set; /** * One render per steering message, whichever authoritative source projects - * it first: the durable session-event echo (in stream order) or the queue - * in-flight synthesis (which can win the race or cover a rejoin window). + * it first: the transcript bootstrap (seeded here), the durable session-event + * echo (in stream order), or the queue in-flight synthesis (which can win the + * race or cover a rejoin window). */ - readonly #renderedSteeringMessageIds = new Set(); + readonly #renderedSteeringMessageIds: Set; readonly #accumulators = new Map(); constructor( @@ -82,6 +93,7 @@ export class RuntimeHostSessionProjector { this.#snapshot = structuredClone(snapshot); this.#now = now; this.#transcriptIds = new Set(seed.durableInFlightMessageIds); + this.#renderedSteeringMessageIds = new Set(seed.durableSteeringMessageIds); const root = snapshot.rootTurn; if (!root) return; for (const message of seed.activeAssistantMessages) {