From 98da332260aea9ac881422b4358172335a46affd Mon Sep 17 00:00:00 2001 From: creatixchu Date: Wed, 26 Aug 2026 11:48:33 +0800 Subject: [PATCH] =?UTF-8?q?feat(session-controller):=20=E5=AE=A2=E6=88=B7?= =?UTF-8?q?=E7=AB=AF=E6=9C=AC=E5=9C=B0=E6=8F=90=E4=BA=A4=E5=9B=9E=E6=98=BE?= =?UTF-8?q?=E4=B8=8E=20rpcId=20=E5=85=B3=E8=81=94?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit beginSubmission 在 prompt 之前同步把本地提交回显写入 SessionSnapshot.pendingSubmissions; durable user/message(source.rpcId)或队列投影(SessionQueuedItem.rpcId)到达后延迟一帧退休, prompt 失败与放弃立即退休并回调 onRetire。fixture 的 prompt 同步回显 requestId。 --- .../src/client/contract/session.ts | 45 +++- .../src/client/contract/snapshot.ts | 34 +++ .../session-controller/src/client/index.ts | 11 +- .../src/client/sessions/queue-mirror.ts | 1 + .../src/client/sessions/session.ts | 134 +++++++++- .../api/session-controller/src/control.ts | 8 + packages/api/session-controller/src/types.ts | 2 + ...session-pending-submissions.client.spec.ts | 238 ++++++++++++++++++ .../client/connection/src/client/fixture.ts | 9 +- .../tests/ui-session.client.spec.ts | 1 + .../ui-trajectory/tests/views.client.spec.tsx | 1 + .../user-questions-composer.client.spec.tsx | 1 + .../client-runtime/src/fixtures.ts | 1 + .../client-runtime/src/sessions.ts | 19 +- 14 files changed, 496 insertions(+), 9 deletions(-) create mode 100644 packages/api/session-controller/tests/session-pending-submissions.client.spec.ts diff --git a/packages/api/session-controller/src/client/contract/session.ts b/packages/api/session-controller/src/client/contract/session.ts index 6ac4182bb9..9b8ed3f7ec 100644 --- a/packages/api/session-controller/src/client/contract/session.ts +++ b/packages/api/session-controller/src/client/contract/session.ts @@ -12,9 +12,37 @@ import type { MessageId } from '@deepseek-ai/dsh-llm/brand' import type { SessionId } from '@deepseek-ai/dsh-session/types' import type { RemoteResult } from '@deepseek-ai/dsh-typert-protocol' import type { ObservableSnapshot } from '@deepseek-ai/dsh-client-store' -import type { PromptContentPart, QueueAction } from '../../types.ts' +import type { PromptContentPart, QueueAction, SessionRequestId } from '../../types.ts' import type { ClientResult } from './result.ts' -import type { SessionSnapshot } from './snapshot.ts' +import type { PendingSubmissionImage, SessionSnapshot } from './snapshot.ts' + +/** + * Why a local submission echo left the snapshot: `observed` when its durable + * `user/message` event or host queue occurrence arrived (with the admitted + * image references in prompt order), `failed` when the prompt was rejected, + * threw, or was aborted before acceptance. + */ +export type PendingSubmissionRetirement = + | { readonly reason: 'observed'; readonly attachments: readonly ImageAttachmentRef[] } + | { readonly reason: 'failed' } + +/** Input registering one local submission echo ahead of its prompt call. */ +export interface BeginSubmissionInput { + /** Prompt text exactly as the upcoming prompt will send it. */ + readonly text: string + /** Ordered image previews matching the upcoming prompt's image parts. */ + readonly images: readonly PendingSubmissionImage[] + /** Settlement callback fired exactly once when the echo retires. */ + readonly onRetire?: (retirement: PendingSubmissionRetirement) => void +} + +/** One registered submission echo: the identity its prompt must carry, and the pre-prompt escape hatch. */ +export interface SubmissionHandle { + /** The prompt RPC identity; pass it to {@link ISession.prompt}. */ + readonly requestId: SessionRequestId + /** Retire the echo as failed when the caller cannot reach prompt() (serialization failure); no-op after any other settlement. */ + abandon(): void +} /** Key-addressed projection read face (the useProjection resolution path; see ProjectionValueStore). */ export interface ProjectionsFace { @@ -33,16 +61,29 @@ export interface ISession { readonly sessionId: SessionId /** Host-computed projection values by key (the useProjection seat). */ readonly projections: ProjectionsFace + /** + * Register one local submission echo in `snapshot.pendingSubmissions`, + * synchronously, before the caller serializes and sends the prompt. The + * echo retires when a durable `user/message` event or queue occurrence + * carrying the returned identity arrives, or when the identified prompt + * call fails. + * @param input - echo content and the optional settlement callback. + * @returns the minted identity for {@link prompt} plus the pre-prompt abandon path. + */ + beginSubmission(input: BeginSubmissionInput): SubmissionHandle /** * Send a prompt into the session. * @param content - text plus browser-owned temporary image uploads. * @param mode - 'queue' appends a turn; 'steer' interrupts the running one. + * @param signal - optional caller cancellation for the complete admission round-trip. + * @param requestId - identity from {@link beginSubmission}; a failed identified prompt retires its echo. * @returns acceptance, or the business error (also mirrored into snapshot.promptError). */ prompt( content: PromptContentPart[], mode: 'queue' | 'steer', signal?: AbortSignal, + requestId?: SessionRequestId, ): Promise> /** * Resolve one durable image referenced by this session. diff --git a/packages/api/session-controller/src/client/contract/snapshot.ts b/packages/api/session-controller/src/client/contract/snapshot.ts index dff2317359..85b5a3013e 100644 --- a/packages/api/session-controller/src/client/contract/snapshot.ts +++ b/packages/api/session-controller/src/client/contract/snapshot.ts @@ -3,6 +3,7 @@ import type { ContentBlock } from '@deepseek-ai/dsh-llm/types' import type { MessageId } from '@deepseek-ai/dsh-llm/brand' import type { SessionId } from '@deepseek-ai/dsh-session/types' import type { SubagentAddress } from '@deepseek-ai/dsh-client-connection/client' +import type { SessionRequestId } from '../../types.ts' import type { ClientFailure } from './result.ts' /** One transient inbox occurrence from the authoritative queue snapshot. */ @@ -10,11 +11,42 @@ export interface QueuedMessage { readonly id: MessageId readonly messageId: MessageId readonly placement: 'queued' | 'steering' | 'context' + /** Prompt-RPC identity of a browser-submitted occurrence; correlates the local submission echo. */ + readonly rpcId?: SessionRequestId readonly content: readonly ContentBlock[] readonly preview: string readonly text: string | null } +/** One image displayed by a local submission echo before durable admission. */ +export interface PendingSubmissionImage { + /** Browser-owned preview URL; its lifecycle belongs to the submitter, never this snapshot. */ + readonly previewUrl: string + /** Browser file name, when the file had one. */ + readonly name?: string + /** Intrinsic pixel width, when the submitter has probed it. */ + readonly width?: number + /** Intrinsic pixel height, when the submitter has probed it. */ + readonly height?: number +} + +/** + * One local prompt-submission echo: inserted synchronously when a submission + * begins, so the conversation can show the message before serialization, + * transport, and durable admission complete. Client-memory only — reload and + * reconnect rebuild the conversation from durable events alone. + */ +export interface PendingSubmission { + /** The prompt RPC identity; the durable `user/message` source echoes it as `rpcId`. */ + readonly requestId: SessionRequestId + /** Client wall-clock ms when the submission began. */ + readonly time: number + /** Prompt text exactly as it will be sent (one text block). */ + readonly text: string + /** Ordered image previews matching the prompt's image parts. */ + readonly images: readonly PendingSubmissionImage[] +} + /** History-open lifecycle of a Session event window. */ export type OpenState = 'cold' | 'loading' | 'open' | 'error' @@ -28,6 +60,8 @@ export interface PromptError { export interface SessionSnapshot { readonly sessionId: SessionId readonly queue: readonly QueuedMessage[] + /** Local prompt-submission echoes not yet observed as durable events or queue occurrences. */ + readonly pendingSubmissions: readonly PendingSubmission[] readonly running: boolean readonly subagent: { readonly address: SubagentAddress diff --git a/packages/api/session-controller/src/client/index.ts b/packages/api/session-controller/src/client/index.ts index 7fa6ede0d0..965b5e2bb8 100644 --- a/packages/api/session-controller/src/client/index.ts +++ b/packages/api/session-controller/src/client/index.ts @@ -40,7 +40,14 @@ export type { SessionProjectionMap, UseProjection, } from './sessions/projection-store.ts' -export type { ISession, ProjectionsFace, SessionFace } from './contract/session.ts' +export type { + BeginSubmissionInput, + ISession, + PendingSubmissionRetirement, + ProjectionsFace, + SessionFace, + SubmissionHandle, +} from './contract/session.ts' export type { ISessions } from './contract/sessions.ts' export { MutableSessionEventSource } from './contract/events.ts' export type { @@ -53,6 +60,8 @@ export type { } from './contract/events.ts' export type { OpenState, + PendingSubmission, + PendingSubmissionImage, PromptError, QueuedMessage, SessionSnapshot, diff --git a/packages/api/session-controller/src/client/sessions/queue-mirror.ts b/packages/api/session-controller/src/client/sessions/queue-mirror.ts index 209349af2c..2a9f274b28 100644 --- a/packages/api/session-controller/src/client/sessions/queue-mirror.ts +++ b/packages/api/session-controller/src/client/sessions/queue-mirror.ts @@ -43,6 +43,7 @@ export class SessionQueueMirror { id: item.id, messageId: item.message.id, placement: item.placement, + ...(item.rpcId === undefined ? {} : { rpcId: item.rpcId }), content, preview: previewOf(content), text: textOf(content), diff --git a/packages/api/session-controller/src/client/sessions/session.ts b/packages/api/session-controller/src/client/sessions/session.ts index 182d8514d9..acf671b9e9 100644 --- a/packages/api/session-controller/src/client/sessions/session.ts +++ b/packages/api/session-controller/src/client/sessions/session.ts @@ -24,9 +24,11 @@ import type { } from '../../types.ts' import type { ClientFailure, ClientResult } from '../contract/result.ts' import { transportResult } from '../contract/result.ts' -import type { SessionFace } from '../contract/session.ts' import type { - OpenState, PromptError, SessionSnapshot, + BeginSubmissionInput, PendingSubmissionRetirement, SessionFace, SubmissionHandle, +} from '../contract/session.ts' +import type { + OpenState, PendingSubmission, PromptError, SessionSnapshot, } from '../contract/snapshot.ts' import { MutableSessionEventSource } from '../contract/events.ts' import type { @@ -101,6 +103,14 @@ export class Session implements SessionFace { private removed = false private promptError: PromptError | null = null private lastAgentError: string | null = null + /** Local submission echoes, insertion-ordered (see SessionSnapshot.pendingSubmissions). */ + private pendingSubmissions: readonly PendingSubmission[] = [] + /** Per-echo settlement state; `retiring` latches the first observation so a + * queue frame and its durable event cannot both retire one echo. */ + private readonly submissionSettlements = new Map void) | undefined + retiring: boolean + }>() /** Owns the addressed page/follow lifecycle while this Session is open. */ private events: SessionEventStream | undefined @@ -172,16 +182,42 @@ export class Session implements SessionFace { // ---- Operations ---- + /** + * Register one local submission echo (see the ISession declaration). + * Synchronous through markDirty: the echo is in the very next snapshot, so + * the conversation can paint it before the caller starts serializing. + * @param input - echo content and the optional settlement callback. + * @returns the minted identity for {@link prompt} plus the pre-prompt abandon path. + */ + beginSubmission(input: BeginSubmissionInput): SubmissionHandle { + const requestId = randomUUID() as SessionRequestId + this.pendingSubmissions = [...this.pendingSubmissions, { + requestId, + time: Date.now(), + text: input.text, + images: input.images, + }] + this.submissionSettlements.set(requestId, { onRetire: input.onRetire, retiring: false }) + // The blank → engaging edge flips here, ahead of prompt(): the composer + // docks and the echo renders on the click's own frame. + this.promptAttempted = true + this.notifier.markDirty() + return { requestId, abandon: () => { this.retireFailedSubmission(requestId) } } + } + /** * Send (queue/steer passed through 1:1); failures land in the snapshot's promptError. * @param content - text plus browser-owned temporary image uploads. * @param mode - queue appends after the current turn; steer interrupts it. + * @param signal - optional caller cancellation for the complete admission round-trip. + * @param requestId - identity from {@link beginSubmission}; a failed identified prompt retires its echo. * @returns the prompt result (also mirrored into promptError on failure). */ async prompt( content: PromptContentPart[], mode: 'queue' | 'steer', signal?: AbortSignal, + requestId?: SessionRequestId, ): Promise> { this.promptError = null this.lastAgentError = null @@ -196,7 +232,7 @@ export class Session implements SessionFace { if (this.address === undefined) { const clientTimeZone = resolvedClientTimeZone() result = toSessionResult(await this.remote.session.prompt({ - requestId: randomUUID() as SessionRequestId, + requestId: requestId ?? randomUUID() as SessionRequestId, sessionId: this.sessionId, mode, content, @@ -236,6 +272,7 @@ export class Session implements SessionFace { result = transportResult(error) } if (!result.ok) { + if (requestId !== undefined) this.retireFailedSubmission(requestId) this.promptError = { op: 'send', error: result.error } this.notifier.markDirty() return result @@ -434,6 +471,7 @@ export class Session implements SessionFace { */ replaceControl(queue: readonly SessionQueuedItem[]): void { this.queueMirror.replace(queue) + this.observeSubmissionQueue(queue) this.notifier.markDirty() } @@ -443,6 +481,7 @@ export class Session implements SessionFace { */ handleControlFrame(frame: Extract): void { this.queueMirror.replace(frame.items) + this.observeSubmissionQueue(frame.items) this.notifier.markDirty() } @@ -523,6 +562,12 @@ export class Session implements SessionFace { * @returns when the Remote iterator has completed teardown. */ async dispose(): Promise { + // Unsettled echoes retire as failed so their owners can restore or + // release browser resources; echoes already scheduled as observed keep + // that settlement. + for (const requestId of [...this.submissionSettlements.keys()]) { + this.retireFailedSubmission(requestId) + } this.openGeneration++ const events = this.events this.events = undefined @@ -581,6 +626,7 @@ export class Session implements SessionFace { if (entries.some(entry => entry.event.type === 'turn/start')) this.firstPromptPendingTurn = false if (projections !== undefined) this.projections.seed(projections) this.eventSource.replace(entries, hasMore) + for (const entry of entries) this.observeSubmissionEvent(entry.event) this.notifier.markDirty() } @@ -598,9 +644,70 @@ export class Session implements SessionFace { if (event.type === 'turn/start') this.firstPromptPendingTurn = false const queueChanged = this.queueMirror.acceptDurable(event) this.eventSource.append(entry) + // After the feed append: the conversation assembly's animation frame is + // registered by the feed subscribers above, so the echo-retirement frame + // scheduled here always runs after the durable node became renderable. + this.observeSubmissionEvent(event) return queueChanged || awaitingFirstTurn !== this.firstPromptPendingTurn } + /** Retire the matching echo when a durable browser-prompt `user/message` becomes visible. */ + private observeSubmissionEvent(event: { readonly type: string; readonly data?: unknown }): void { + if (this.submissionSettlements.size === 0 || event.type !== 'user/message') return + // Structural read: window entries may be compact history records, so the + // fields are narrowed rather than trusted (same posture as Conversation + // assembly matchers). + const data = event.data as { readonly source?: unknown; readonly content?: unknown } | undefined + const source = data?.source as { readonly kind?: unknown; readonly rpcId?: unknown } | undefined + if (source?.kind !== 'user' || typeof source.rpcId !== 'string') return + this.scheduleObservedRetirement(source.rpcId as SessionRequestId, imageRefsIn(data?.content)) + } + + /** Retire echoes whose prompts landed in the host inbox instead of the log (running-turn submissions). */ + private observeSubmissionQueue(items: readonly SessionQueuedItem[]): void { + if (this.submissionSettlements.size === 0) return + for (const item of items) { + if (item.rpcId !== undefined) { + this.scheduleObservedRetirement(item.rpcId, imageRefsIn(item.message.content)) + } + } + } + + /** + * Latch one observed settlement and remove the echo an animation frame + * later. The delay keeps the echo in the snapshot until the frame in which + * the durable node (whose assembly frame was registered first) is + * renderable; the render-time rpcId dedupe hides the one-frame overlap. + */ + private scheduleObservedRetirement( + requestId: SessionRequestId, + attachments: readonly ImageAttachmentRef[], + ): void { + const settlement = this.submissionSettlements.get(requestId) + if (settlement === undefined || settlement.retiring) return + settlement.retiring = true + scheduleFrame(() => { this.finishSubmission(requestId, { reason: 'observed', attachments }) }) + } + + /** Remove one unsettled echo immediately (prompt rejection, abort, or disposal). */ + private retireFailedSubmission(requestId: SessionRequestId): void { + const settlement = this.submissionSettlements.get(requestId) + if (settlement === undefined || settlement.retiring) return + settlement.retiring = true + this.finishSubmission(requestId, { reason: 'failed' }) + } + + /** Single removal point: drop the echo, publish, then notify the owner. */ + private finishSubmission(requestId: SessionRequestId, retirement: PendingSubmissionRetirement): void { + const settlement = this.submissionSettlements.get(requestId) + /* v8 ignore next -- retiring latches before every schedule, so one settlement never finishes twice. */ + if (settlement === undefined) return + this.submissionSettlements.delete(requestId) + this.pendingSubmissions = this.pendingSubmissions.filter(echo => echo.requestId !== requestId) + this.notifier.markDirty() + settlement.onRetire?.(retirement) + } + /** Publish a terminal background failure only while this stream still owns the Session. */ private failEventStream(events: SessionEventStream, generation: number, error: unknown): void { if (generation !== this.openGeneration || this.events !== events) return @@ -617,6 +724,7 @@ export class Session implements SessionFace { return { sessionId: this.sessionId, queue: this.queueMirror.snapshot(), + pendingSubmissions: this.pendingSubmissions, running: this.running, subagent: this.address === undefined ? null @@ -644,6 +752,26 @@ export class Session implements SessionFace { } } +/** Run one callback on the next animation frame, or a macrotask where no frame clock exists. */ +function scheduleFrame(fn: () => void): void { + if (typeof requestAnimationFrame === 'function') requestAnimationFrame(() => { fn() }) + else setTimeout(fn, 0) +} + +/** Image attachment references in one structurally-read content block list, in block order. */ +function imageRefsIn(content: unknown): readonly ImageAttachmentRef[] { + if (!Array.isArray(content)) return [] + const refs: ImageAttachmentRef[] = [] + for (const block of content) { + if (typeof block !== 'object' || block === null) continue + const candidate = block as { readonly type?: unknown; readonly attachment?: unknown } + if (candidate.type === 'image' && typeof candidate.attachment === 'object' && candidate.attachment !== null) { + refs.push(candidate.attachment as ImageAttachmentRef) + } + } + return refs +} + /** Convert a terminal Session stream failure to the Client error vocabulary. */ function openFailure(error: unknown): ClientFailure { const failure = sessionStreamFailure(error) diff --git a/packages/api/session-controller/src/control.ts b/packages/api/session-controller/src/control.ts index 4068b5536d..624c9a9493 100644 --- a/packages/api/session-controller/src/control.ts +++ b/packages/api/session-controller/src/control.ts @@ -188,16 +188,24 @@ function queueItems( ...project('next-turn').map(message => ({ id: message.id, placement: 'queued' as const, + ...promptRpcId(message), message: { id: message.id, content: message.content as unknown as JsonValue[] }, })), ...project('next-step').map(message => ({ id: message.id, placement: message.source.kind === 'user' ? 'steering' as const : 'context' as const, + ...promptRpcId(message), message: { id: message.id, content: message.content as unknown as JsonValue[] }, })), ] } +/** Prompt-RPC identity carried by a browser-submitted message's user source. */ +function promptRpcId(message: UserMessage): Pick { + const source = message.source + return source.kind === 'user' && 'rpcId' in source ? { rpcId: source.rpcId } : {} +} + function jobView(job: JobSnapshot): SessionJob { return { id: job.id, diff --git a/packages/api/session-controller/src/types.ts b/packages/api/session-controller/src/types.ts index c045e00707..e9e416777c 100644 --- a/packages/api/session-controller/src/types.ts +++ b/packages/api/session-controller/src/types.ts @@ -432,6 +432,8 @@ export type SessionFollowFrame = export interface SessionQueuedItem { readonly id: MessageId readonly placement: 'queued' | 'steering' | 'context' + /** Prompt-RPC identity from the queued message's user source; clients retire the matching local submission echo on it. */ + readonly rpcId?: SessionRequestId /** JSON-safe message fields consumed by pending-queue presentation. */ readonly message: { readonly id: MessageId diff --git a/packages/api/session-controller/tests/session-pending-submissions.client.spec.ts b/packages/api/session-controller/tests/session-pending-submissions.client.spec.ts new file mode 100644 index 0000000000..fa0e69d89a --- /dev/null +++ b/packages/api/session-controller/tests/session-pending-submissions.client.spec.ts @@ -0,0 +1,238 @@ +/** Local submission echoes: synchronous insertion, observed/failed retirement, and settlement callbacks. */ + +import { afterEach, describe, expect, it, vi } from 'vitest' +import { createUserMessage } from '@deepseek-ai/dsh-llm' +import type { MessageSource } from '@deepseek-ai/dsh-llm' +import type { ImageAttachmentRef } from '@deepseek-ai/dsh-attachment' +import type { SessionEvent, SessionId } from '@deepseek-ai/dsh-session/types' +import { Session } from '../src/client/sessions/session.ts' +import type { PendingSubmissionRetirement } from '../src/client/contract/session.ts' +import type { SessionQueuedItem, SessionRequestId } from '../src/types.ts' +import { FakeApiClient, err, fakeRemote, ok } from './fake-api.client.ts' +import { historyValue } from './event-script.client.ts' + +const SID = 'fk-s1' as SessionId + +afterEach(() => { + vi.unstubAllGlobals() +}) + +function makeSession(api = new FakeApiClient()): { api: FakeApiClient; session: Session } { + return { api, session: new Session(SID, api, fakeRemote(api)) } +} + +function imageRef(id: string): ImageAttachmentRef { + return { + attachmentId: id, + mediaType: 'image/png', + bytes: 1, + width: 2, + height: 2, + } as unknown as ImageAttachmentRef +} + +/** A durable browser-prompt user/message whose source echoes `rpcId`. */ +function promptEvent(seq: number, rpcId: SessionRequestId, refs: readonly ImageAttachmentRef[] = []): SessionEvent { + return { + seq, + time: 1_700_000_000_000 + seq, + type: 'user/message', + surfaceOp: 'append', + data: createUserMessage({ + content: [ + ...refs.map(attachment => ({ type: 'image' as const, attachment })), + { type: 'text' as const, text: '发送' }, + ], + source: { kind: 'user', rpcId } as MessageSource, + }), + } as unknown as SessionEvent +} + +function queuedItem(rpcId: SessionRequestId, refs: readonly ImageAttachmentRef[] = []): SessionQueuedItem { + return { + id: 'm-queued' as SessionQueuedItem['id'], + placement: 'queued', + rpcId, + message: { + id: 'm-queued' as SessionQueuedItem['id'], + content: refs.map(attachment => ({ type: 'image', attachment })) as unknown as SessionQueuedItem['message']['content'], + }, + } +} + +/** Let the frame-delayed retirement (setTimeout fallback in this node environment) run. */ +function settleFrames(): Promise { + return new Promise(resolve => setTimeout(resolve, 0)) +} + +describe('beginSubmission', () => { + it('inserts the echo synchronously and flips the engaging edge before any prompt call', () => { + const { session } = makeSession() + expect(session.getSnapshot()).toMatchObject({ pendingSubmissions: [], promptAttempted: false }) + const handle = session.beginSubmission({ + text: '你好', + images: [{ previewUrl: 'blob:p1', name: 'a.png', width: 4, height: 3 }], + }) + expect(session.getSnapshot().promptAttempted).toBe(true) + expect(session.getSnapshot().pendingSubmissions).toMatchObject([{ + requestId: handle.requestId, + text: '你好', + images: [{ previewUrl: 'blob:p1', name: 'a.png', width: 4, height: 3 }], + }]) + }) + + it('abandon retires the echo as failed exactly once', () => { + const { session } = makeSession() + const retirements: PendingSubmissionRetirement[] = [] + const handle = session.beginSubmission({ + text: '放弃', + images: [], + onRetire: retirement => retirements.push(retirement), + }) + handle.abandon() + handle.abandon() + expect(session.getSnapshot().pendingSubmissions).toEqual([]) + expect(retirements).toEqual([{ reason: 'failed' }]) + }) +}) + +describe('prompt-coupled retirement', () => { + it('a rejected identified prompt retires its echo immediately alongside promptError', async () => { + const { api, session } = makeSession() + api.onPrompt = () => Promise.resolve(err({ code: 'agent-busy', message: '忙', details: {} })) + const retirements: PendingSubmissionRetirement[] = [] + const handle = session.beginSubmission({ + text: '失败的', + images: [], + onRetire: retirement => retirements.push(retirement), + }) + const result = await session.prompt([{ type: 'text', text: '失败的' }], 'queue', undefined, handle.requestId) + expect(result.ok).toBe(false) + expect(session.getSnapshot().pendingSubmissions).toEqual([]) + expect(session.getSnapshot().promptError).toMatchObject({ op: 'send' }) + expect(retirements).toEqual([{ reason: 'failed' }]) + }) + + it('sends the echo identity as the prompt requestId', async () => { + const { api, session } = makeSession() + const handle = session.beginSubmission({ text: '带 id', images: [] }) + await session.prompt([{ type: 'text', text: '带 id' }], 'queue', undefined, handle.requestId) + expect(api.callsOf('session.prompt')).toMatchObject([{ requestId: handle.requestId }]) + }) + + it('an unidentified prompt failure leaves registered echoes alone', async () => { + const { api, session } = makeSession() + api.onPrompt = () => Promise.resolve(err({ code: 'agent-busy', message: '忙', details: {} })) + session.beginSubmission({ text: '还在', images: [] }) + await session.prompt([{ type: 'text', text: '另一个' }], 'queue') + expect(session.getSnapshot().pendingSubmissions).toHaveLength(1) + }) +}) + +describe('observed retirement', () => { + it('a live durable event carrying the rpcId retires the echo one frame later with the admitted refs', async () => { + const { api, session } = makeSession() + api.onHistory = () => Promise.resolve(ok(historyValue([]))) + await session.open() + const retirements: PendingSubmissionRetirement[] = [] + const handle = session.beginSubmission({ + text: '发送', + images: [{ previewUrl: 'blob:p1' }], + onRetire: retirement => retirements.push(retirement), + }) + const refs = [imageRef('att-1')] + await api.pushFollow(SID, { type: 'event', event: promptEvent(0, handle.requestId, refs) as never }) + // Synchronously after the append the echo is still in the snapshot; the + // render-time dedupe owns the overlap frame. + expect(session.getSnapshot().pendingSubmissions).toHaveLength(1) + await settleFrames() + expect(session.getSnapshot().pendingSubmissions).toEqual([]) + expect(retirements).toEqual([{ reason: 'observed', attachments: refs }]) + }) + + it('a queue occurrence carrying the rpcId retires the echo (running-turn submissions)', async () => { + const { session } = makeSession() + const retirements: PendingSubmissionRetirement[] = [] + const handle = session.beginSubmission({ + text: '排队', + images: [{ previewUrl: 'blob:p1' }], + onRetire: retirement => retirements.push(retirement), + }) + const refs = [imageRef('att-q')] + session.handleControlFrame({ type: 'queue', sessionId: SID, items: [queuedItem(handle.requestId, refs)] }) + await settleFrames() + expect(session.getSnapshot().pendingSubmissions).toEqual([]) + expect(retirements).toEqual([{ reason: 'observed', attachments: refs }]) + // The queue projection keeps the correlation id for render-time dedupe. + expect(session.getSnapshot().queue).toMatchObject([{ rpcId: handle.requestId }]) + }) + + it('a full-window install (reconnect resync) retires echoes observed in the window', async () => { + const { api, session } = makeSession() + const handle = session.beginSubmission({ text: '重连', images: [] }) + api.onHistory = () => Promise.resolve(ok(historyValue([promptEvent(12, handle.requestId)]))) + await session.open() + await settleFrames() + expect(session.getSnapshot().pendingSubmissions).toEqual([]) + }) + + it('the first observation wins: a later prompt failure cannot re-retire an observed echo', async () => { + const { api, session } = makeSession() + api.onHistory = () => Promise.resolve(ok(historyValue([]))) + await session.open() + const retirements: PendingSubmissionRetirement[] = [] + const handle = session.beginSubmission({ + text: '先观察', + images: [], + onRetire: retirement => retirements.push(retirement), + }) + await api.pushFollow(SID, { type: 'event', event: promptEvent(0, handle.requestId) as never }) + handle.abandon() + await settleFrames() + expect(retirements).toEqual([{ reason: 'observed', attachments: [] }]) + }) + + it('uses requestAnimationFrame for the retirement delay when the runtime provides one', async () => { + const frames: FrameRequestCallback[] = [] + vi.stubGlobal('requestAnimationFrame', (fn: FrameRequestCallback) => { + frames.push(fn) + return frames.length + }) + const { api, session } = makeSession() + api.onHistory = () => Promise.resolve(ok(historyValue([]))) + await session.open() + const handle = session.beginSubmission({ text: '帧', images: [] }) + await api.pushFollow(SID, { type: 'event', event: promptEvent(0, handle.requestId) as never }) + expect(session.getSnapshot().pendingSubmissions).toHaveLength(1) + expect(frames).toHaveLength(1) + frames[0]?.(0) + expect(session.getSnapshot().pendingSubmissions).toEqual([]) + }) +}) + +describe('disposal', () => { + it('retires unsettled echoes as failed and preserves an already-observed settlement', async () => { + const { api, session } = makeSession() + api.onHistory = () => Promise.resolve(ok(historyValue([]))) + await session.open() + const retirements: { text: string; retirement: PendingSubmissionRetirement }[] = [] + const observed = session.beginSubmission({ + text: '已观察', + images: [], + onRetire: retirement => retirements.push({ text: '已观察', retirement }), + }) + session.beginSubmission({ + text: '未settle', + images: [], + onRetire: retirement => retirements.push({ text: '未settle', retirement }), + }) + await api.pushFollow(SID, { type: 'event', event: promptEvent(0, observed.requestId) as never }) + await session.dispose() + await settleFrames() + expect(retirements).toEqual([ + { text: '未settle', retirement: { reason: 'failed' } }, + { text: '已观察', retirement: { reason: 'observed', attachments: [] } }, + ]) + expect(session.getSnapshot().pendingSubmissions).toEqual([]) + }) +}) diff --git a/packages/client/connection/src/client/fixture.ts b/packages/client/connection/src/client/fixture.ts index ce2fcdae6a..d3be8087a6 100644 --- a/packages/client/connection/src/client/fixture.ts +++ b/packages/client/connection/src/client/fixture.ts @@ -2757,9 +2757,14 @@ function createFixtureWorld(options: FixtureOptions): FixtureWorld { attachments.set(String(attachment.attachmentId), { attachment, data: block.data }) return { type: 'image', attachment } }) + // The host echoes the prompt's requestId as the user source's rpcId; + // the Session object retires its local submission echo on it. The + // user-rpc source member is declared by dsh-api-session-controller, + // which this standalone fixture does not import — hence the assertion. + const promptSource = { kind: 'user', rpcId: request.requestId } as MessageSource if (mode === 'steer' && replays.has(id)) { // Steering: the durable user/message lands inside the current turn; the replay continues. - append(id, { type: 'user/message', surfaceOp: 'append', data: userMessage(durable) }) + append(id, { type: 'user/message', surfaceOp: 'append', data: userMessage(durable, promptSource) }) return sessionOk({ accepted: true as const }) } const turn = nextTurn.get(id) ?? 0 @@ -2772,7 +2777,7 @@ function createFixtureWorld(options: FixtureOptions): FixtureWorld { if (plan.wanted !== null && plan.wanted !== plan.active) { append(id, { type: 'plan/mode', data: { active: plan.wanted } }) } - append(id, { type: 'user/message', surfaceOp: 'append', data: userMessage(durable) }) + append(id, { type: 'user/message', surfaceOp: 'append', data: userMessage(durable, promptSource) }) // Capacity parallel of the host token-meter's request/context record: // log-only, appended inside the open turn, and deduplicated against the // route already recorded (the fixture never varies contextWindow). diff --git a/packages/client/ui-session/tests/ui-session.client.spec.ts b/packages/client/ui-session/tests/ui-session.client.spec.ts index 97946f4c1d..ce6d154c0c 100644 --- a/packages/client/ui-session/tests/ui-session.client.spec.ts +++ b/packages/client/ui-session/tests/ui-session.client.spec.ts @@ -76,6 +76,7 @@ function createSessionsBench(_ctx: Context): SessionsBench { const snapshot = createSnapshotStore({ sessionId: id, queue: [], + pendingSubmissions: [], running: false, subagent: null, removed: false, diff --git a/packages/client/ui-trajectory/tests/views.client.spec.tsx b/packages/client/ui-trajectory/tests/views.client.spec.tsx index 4aa01e9ca6..4bf4793de0 100644 --- a/packages/client/ui-trajectory/tests/views.client.spec.tsx +++ b/packages/client/ui-trajectory/tests/views.client.spec.tsx @@ -111,6 +111,7 @@ function sessionSnapshot(nodes: LegacyConversationSlice['nodes']): SessionSnapsh return { sessionId: SID, queue: [], + pendingSubmissions: [], running: false, subagent: null, removed: false, diff --git a/packages/client/ui-user-questions/tests/user-questions-composer.client.spec.tsx b/packages/client/ui-user-questions/tests/user-questions-composer.client.spec.tsx index a650dd6a5b..c2ee66b3b0 100644 --- a/packages/client/ui-user-questions/tests/user-questions-composer.client.spec.tsx +++ b/packages/client/ui-user-questions/tests/user-questions-composer.client.spec.tsx @@ -25,6 +25,7 @@ type AttentionState = Parameters {}, + } + } + + private submissionSeq = 0 + /** * Fail-loud stub; supply `readAttachment` on the fixture's session face to exercise it. * @param _attachmentId - opaque durable attachment id.