From 63027e2584ad517541ac6597ac71c36557654694 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=9F=B3=E6=95=AC=E8=8D=A3?= <15920017+itjingrong@user.noreply.gitee.com> Date: Fri, 11 Sep 2026 12:12:56 +0800 Subject: [PATCH 01/14] feat(core): add Stream Event Protocol v1 types (ADR 0001) Introduce the pure-type foundation for the structured streaming event protocol (issue #27): versioned envelope, 12 typed events, and the idempotency/cancel/reconnect contracts as types. Zero runtime change. Adds ADR 0001 documenting context, decision, migration path and open questions for maintainer review. --- docs/adr/0001-stream-event-protocol.md | 79 +++++++++++++++++++++++ packages/core/src/index.ts | 3 + packages/core/src/stream-events.test.ts | 68 ++++++++++++++++++++ packages/core/src/stream-events.ts | 83 +++++++++++++++++++++++++ 4 files changed, 233 insertions(+) create mode 100644 docs/adr/0001-stream-event-protocol.md create mode 100644 packages/core/src/stream-events.test.ts create mode 100644 packages/core/src/stream-events.ts diff --git a/docs/adr/0001-stream-event-protocol.md b/docs/adr/0001-stream-event-protocol.md new file mode 100644 index 0000000..9811c29 --- /dev/null +++ b/docs/adr/0001-stream-event-protocol.md @@ -0,0 +1,79 @@ +# ADR 0001: Stream Event Protocol v1 + +- **Status:** Proposed (awaiting review on #27) +- **Date:** 2026-09-11 +- **Deciders:** helsome/folio maintainers, contributor for #27 + +## Context + +The Copilot answer path currently streams through an implicit `AgentEvent` protocol +(`packages/core/src/index.ts`): 8 event types +(`run_started / message_started / message_delta / message_completed / tool_started / tool_completed / run_completed / run_failed`) +carried over one IPC channel (`agent:event` in `apps/electron`). The envelope already has +`id / sessionId / runId / timestamp / sequence`, but the protocol lacks: + +- a protocol version (no forward-compat contract); +- a monotonic `sequence` contract that all producers honor (idempotency / resume); +- an independent `cancelled` event (UI currently infers cancellation from `run_failed`); +- typed payloads for `tool_progress`, `citation_added`, and `status`. + +Issue #27 asks for a stable, versioned streaming event protocol with cancel and +reconnect support. + +## Decision + +Introduce **Stream Event Protocol v1** as pure type/enum definitions in +`packages/core/src/stream-events.ts` (re-exported from `@finagent/core`). It is a +protocol-layer upgrade that coexists with `AgentEvent` during migration and +gradually replaces it. + +### Envelope + +```ts +interface StreamEventEnvelope { + protocolVersion: 1; // forward-compat gate + runId: string; + messageId: string; // equals runId in v1; split later if one message spans runs + sequence: number; // monotonic per run; idempotency key = runId + messageId + sequence + type: T; + timestamp: string; // ISO 8601 UTC; display only, never identity + payload: StreamEventTypeToPayload[T]; +} +``` + +### Event types (12) + +`run_started · message_started · text_delta · tool_started · tool_progress · tool_result · +citation_added · status · error · cancelled · message_completed · run_completed` + +Key behavior change: `message_delta` (full-answer snapshot) becomes `text_delta` +(incremental, append-only text). + +### Cross-cutting contracts + +| Capability | Contract | +|---|---| +| Idempotency | Consumers dedupe on `runId + messageId + sequence`; replayed events never re-insert text, tool cards, or citations | +| Cancel | renderer `cancelRun` → runtime propagates → runtime **explicitly emits `cancelled`**; partial answer preserved | +| Reconnect | client reconnects with `lastSequence`; unconsumed gap is replayed (resume), consumed events are not re-applied | +| Final-state parity | UI `stopReason` and persisted `run.status` come from the same single final state in run-manager (aligns with #18) | +| Security (#19) | `status` never exposes chain-of-thought; tool payloads pass redaction before reaching the UI; renderer never executes model-returned code | + +## Migration path + +1. **Land types first (this ADR + `stream-events.ts`)**: zero runtime change, reviewable alone. +2. Runtime event loop (`run-manager`) emits the new typed events. +3. Transport upgrade (`kernelHost` + `preload`); renderer consumes the new protocol. +4. **Compat window**: keep `message_delta` as a "resume snapshot backfill" event — + transport sends a full snapshot once after reconnect, then only `text_delta` increments. + +## Consequences + +- **Positive:** versioned contract; deterministic idempotency; explicit cancellation + state; groundwork for reconnect/resume and for #21 immutable run manifests. +- **Negative:** dual event families during migration; consumers must handle both. +- **Open questions for reviewers:** + 1. Resume data source: v1 keeps events in memory for the run's lifetime and defers + persisting an event log (ties into #21) — acceptable? + 2. Is `status.phase` of `thinking / searching / working` sufficient? + 3. Keep `messageId` merged with `runId` in v1 and split only when needed? \ No newline at end of file diff --git a/packages/core/src/index.ts b/packages/core/src/index.ts index 8b749ac..24001b9 100644 --- a/packages/core/src/index.ts +++ b/packages/core/src/index.ts @@ -5,6 +5,9 @@ import type { FinancialEvidenceEnvelope } from './financial-evidence.ts'; export type { SupportedLocale, LocalePreference } from './locale.ts'; +// Stream Event Protocol v1 (issue #27, docs/adr/0001-stream-event-protocol.md) +export * from './stream-events.ts'; + export interface Quote { symbol: string; /** Folio canonical instrument id when the quote was resolved through the catalog. */ diff --git a/packages/core/src/stream-events.test.ts b/packages/core/src/stream-events.test.ts new file mode 100644 index 0000000..c9f6110 --- /dev/null +++ b/packages/core/src/stream-events.test.ts @@ -0,0 +1,68 @@ +// Stream Event Protocol v1 — 类型层完整性测试 +// 纯类型交付的自检:事件类型数量/无重复、payload 判别映射、envelope 构造约束。 + +import { describe, expect, it } from 'bun:test'; +import { + STREAM_EVENT_PROTOCOL_VERSION, + STREAM_EVENT_TYPES, + type StreamEvent, + type StreamEventEnvelope, + type StreamEventType, + type StreamEventTypeToPayload, +} from './stream-events.ts'; + +describe('Stream Event Protocol v1', () => { + it('枚举 12 种协议事件类型且无重复', () => { + expect(STREAM_EVENT_TYPES).toHaveLength(12); + expect(new Set(STREAM_EVENT_TYPES).size).toBe(12); + }); + + it('每个事件类型都有对应的 typed payload(编译期覆盖)', () => { + // 编译期强制约束:所有 StreamEventType 必须在映射表中存在。 + type Coverage = { [K in StreamEventType]: StreamEventTypeToPayload[K] }; + const coverage: Coverage = {} as Coverage; + expect(coverage).toBeDefined(); + }); + + it('protocol version 恒为 1', () => { + expect(STREAM_EVENT_PROTOCOL_VERSION).toBe(1); + }); + + it('envelope 按 type 判别 payload(编译期约束)', () => { + // 通过类型构造示例事件流;completion 事件可用新事件类型 + const runStarted: StreamEvent = { + protocolVersion: 1, + runId: 'run-1', + messageId: 'run-1', + sequence: 1, + type: 'run_started', + timestamp: '2026-09-11T00:00:00.000Z', + payload: { input: 'show me AAPL.US', startedAt: '2026-09-11T00:00:00.000Z' }, + }; + const delta: StreamEvent = { + ...runStarted, + protocolVersion: STREAM_EVENT_PROTOCOL_VERSION, + sequence: 2, + type: 'text_delta', + payload: { text: 'Apple ' }, + }; + const cancelled: StreamEvent = { + ...runStarted, + protocolVersion: STREAM_EVENT_PROTOCOL_VERSION, + sequence: 9, + type: 'cancelled', + payload: { reason: 'user', partial: { text: 'Apple ' } }, + }; + const done: StreamEvent = { + ...runStarted, + protocolVersion: STREAM_EVENT_PROTOCOL_VERSION, + sequence: 10, + type: 'run_completed', + payload: { stopReason: 'cancelled' }, + }; + + const events: StreamEvent[] = [runStarted, delta, cancelled, done]; + expect(events).toHaveLength(4); + expect(events[2].type).toBe('cancelled'); + }); +}); \ No newline at end of file diff --git a/packages/core/src/stream-events.ts b/packages/core/src/stream-events.ts new file mode 100644 index 0000000..63b1e8e --- /dev/null +++ b/packages/core/src/stream-events.ts @@ -0,0 +1,83 @@ +// Stream Event Protocol v1 +// +// 结构化流式事件协议的类型定义(issue #27)。 +// 这是协议层的"纯类型 + 枚举"交付,无任何运行时行为变更。 +// 设计文档:docs/adr/0001-stream-event-protocol.md +// +// 与现有 AgentEvent(同文件 index.ts)的关系: +// - AgentEvent 是内部 8 事件隐式协议;本模块是其协议化升级版(12 事件 + 版本 + 单调 seq)。 +// - 迁移期间两者并存,AgentEvent 逐步被取代(见 ADR "Migration" 一节)。 + +export const STREAM_EVENT_PROTOCOL_VERSION = 1 as const; + +export type StreamStatusPhase = 'thinking' | 'searching' | 'working'; +export type StreamCancelReason = 'user' | 'budget' | 'runtime'; +export type StreamStopReason = 'completed' | 'cancelled' | 'error' | 'budget'; + +/** 协议事件类型全集(12 种)。 */ +export type StreamEventType = + | 'run_started' + | 'message_started' + | 'text_delta' + | 'tool_started' + | 'tool_progress' + | 'tool_result' + | 'citation_added' + | 'status' + | 'error' + | 'cancelled' + | 'message_completed' + | 'run_completed'; + +/** 类型 -> payload 映射。envelope.type 作为唯一判别字段,payload 不再重复 type。 */ +export interface StreamEventTypeToPayload { + run_started: { input: string; startedAt: string }; + message_started: Record; + text_delta: { text: string }; + tool_started: { callId: string; name: string; input?: unknown }; + tool_progress: { callId: string; progress?: unknown }; + tool_result: { callId: string; name: string; result: unknown }; + citation_added: { citationId: string; sourceId: string }; + status: { phase: StreamStatusPhase; detail?: string }; + error: { code: string; message: string; retryable: boolean }; + cancelled: { reason: StreamCancelReason; partial: { text: string } }; + message_completed: Record; + run_completed: { stopReason: StreamStopReason }; +} + +export type StreamEventPayload = StreamEventTypeToPayload[StreamEventType]; + +/** + * 统一事件信封。 + * - 幂等键:runId + messageId + sequence。 + * - sequence 为 run 内单调递增;reconnect 以它为游标(lastSequence 补发)。 + * - timestamp 仅用于展示/排序,不作为身份。 + */ +export interface StreamEventEnvelope { + protocolVersion: typeof STREAM_EVENT_PROTOCOL_VERSION; + runId: string; + /** v1 与 runId 相同;预留“一条 message 跨多次 run”时拆分。 */ + messageId: string; + sequence: number; + type: T; + timestamp: string; + payload: StreamEventTypeToPayload[T]; +} + +export type StreamEvent = StreamEventEnvelope; + +/** 供完整性检查/测试用的枚举列表,必须与 StreamEventType 一一对应。 */ +export const STREAM_EVENT_TYPES = [ + 'run_started', + 'message_started', + 'text_delta', + 'tool_started', + 'tool_progress', + 'tool_result', + 'citation_added', + 'status', + 'error', + 'cancelled', + 'message_completed', + 'run_completed', +] as const satisfies readonly StreamEventType[]; \ No newline at end of file From 60022849dc844f5c594b9a7e963b163f44dd29a6 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=9F=B3=E6=95=AC=E8=8D=A3?= <15920017+itjingrong@user.noreply.gitee.com> Date: Fri, 11 Sep 2026 12:20:55 +0800 Subject: [PATCH 02/14] feat(shared): add parallel Stream Event v1 channel to RunManager Emit Stream Event Protocol v1 events alongside the existing AgentEvent stream (issue #27, ADR 0001 migration step 2). Adds toStreamEvents mapping (8 AgentEvent types -> 12 protocol events, cancel normalized to 'cancelled') and RunManager.subscribeStream. Existing AgentEvent consumers are untouched; the parallel channel only activates when a stream subscriber is registered. --- packages/shared/src/kernel/run-manager.ts | 20 ++++ .../src/kernel/stream-event-adapter.test.ts | 111 ++++++++++++++++++ .../shared/src/kernel/stream-event-adapter.ts | 77 ++++++++++++ 3 files changed, 208 insertions(+) create mode 100644 packages/shared/src/kernel/stream-event-adapter.test.ts create mode 100644 packages/shared/src/kernel/stream-event-adapter.ts diff --git a/packages/shared/src/kernel/run-manager.ts b/packages/shared/src/kernel/run-manager.ts index 5d0c8c4..c3dd1a7 100644 --- a/packages/shared/src/kernel/run-manager.ts +++ b/packages/shared/src/kernel/run-manager.ts @@ -7,6 +7,7 @@ import type { Message, Run, SessionMeta, + StreamEvent, TokenUsage, ToolCall, ToolCallRecord, @@ -37,6 +38,7 @@ import { type RunawayState, } from './runaway-detector.ts'; import { buildFinancialEvidence } from '../evidence/financial-evidence.ts'; +import { toStreamEvents } from './stream-event-adapter.ts'; export interface RunManagerOptions { sessions: SessionManager; @@ -87,6 +89,7 @@ export class RunManager { private readonly searchToolPatterns: readonly string[]; private readonly runawayPolicy: Partial; private readonly listeners = new Set<(event: AgentEvent) => void>(); + private readonly streamListeners = new Set<(sessionId: string, event: StreamEvent) => void>(); private activeRun: ActiveRun | null = null; constructor(options: RunManagerOptions) { @@ -104,6 +107,16 @@ export class RunManager { return () => this.listeners.delete(listener); } + /** + * Stream Event Protocol v1 channel (issue #27). Parallel to `subscribe`; + * every AgentEvent is additionally mapped to StreamEvents. The sessionId is + * passed along since the envelope intentionally does not carry it. + */ + subscribeStream(listener: (sessionId: string, event: StreamEvent) => void): () => void { + this.streamListeners.add(listener); + return () => this.streamListeners.delete(listener); + } + /** Whether a run is currently executing (Pi runtime executes one at a time). */ isRunning(): boolean { return this.activeRun !== null; @@ -408,6 +421,13 @@ export class RunManager { for (const listener of this.listeners) { listener(event); } + if (this.streamListeners.size > 0) { + for (const mapped of toStreamEvents(event)) { + for (const listener of this.streamListeners) { + listener(event.sessionId, mapped); + } + } + } } } diff --git a/packages/shared/src/kernel/stream-event-adapter.test.ts b/packages/shared/src/kernel/stream-event-adapter.test.ts new file mode 100644 index 0000000..cbcb579 --- /dev/null +++ b/packages/shared/src/kernel/stream-event-adapter.test.ts @@ -0,0 +1,111 @@ +// Stream Event Protocol v1 — AgentEvent→StreamEvent 映射完整性测试。 + +import { describe, expect, it } from 'bun:test'; +import type { AgentEvent } from '@finagent/core'; +import { toStreamEvents } from './stream-event-adapter.ts'; + +function makeEvent(partial: Partial & { type: AgentEvent['type']; payload: AgentEvent['payload'] }): AgentEvent { + return { + id: 'id-1', + sessionId: 'sess-1', + runId: 'run-1', + timestamp: 1726000000000, + sequence: 5, + ...partial, + } as AgentEvent; +} + +describe('toStreamEvents', () => { + it('映射 run_started(含 input 与 ISO 时间戳)', () => { + const [ev] = toStreamEvents( + makeEvent({ + type: 'run_started', + payload: { + run: { id: 'run-1', sessionId: 'sess-1', status: 'running', input: 'AAPL.US', startedAt: 1726000000000 }, + userMessage: { id: 'm1', role: 'user', content: 'AAPL.US', timestamp: 1726000000000 }, + }, + }) + ); + expect(ev.type).toBe('run_started'); + if (ev.type === 'run_started') { + expect(ev.protocolVersion).toBe(1); + expect(ev.messageId).toBe('run-1'); + expect(ev.sequence).toBe(5); + expect(ev.payload.input).toBe('AAPL.US'); + expect(ev.payload.startedAt).toBe('2024-09-10T20:26:40.000Z'); + expect(ev.timestamp).toBe('2024-09-10T20:26:40.000Z'); + } + }); + + it('映射 message_delta → text_delta(增量字段无损)', () => { + const [ev] = toStreamEvents( + makeEvent({ type: 'message_delta', payload: { delta: 'Apple ', answer: 'Apple Inc.' } }) + ); + expect(ev.type).toBe('text_delta'); + if (ev.type === 'text_delta') { + expect(ev.payload.text).toBe('Apple '); + } + }); + + it('映射 tool 事件为 tool_started / tool_result', () => { + const toolCall = { + id: 'tc-1', + toolName: 'get_quote', + args: { symbol: 'AAPL.US' }, + startedAt: 1, + status: 'success' as const, + result: { lastPrice: 220 }, + }; + const [started] = toStreamEvents(makeEvent({ type: 'tool_started', payload: { toolCall } })); + const [result] = toStreamEvents(makeEvent({ type: 'tool_completed', payload: { toolCall } })); + if (started.type === 'tool_started') { + expect(started.payload.callId).toBe('tc-1'); + expect(started.payload.name).toBe('get_quote'); + } + if (result.type === 'tool_result') { + expect(result.payload.result).toEqual({ lastPrice: 220 }); + } + }); + + it('映射 run_completed → stopReason=completed / message_completed', () => { + const [completed] = toStreamEvents( + makeEvent({ type: 'message_completed', payload: { answer: 'done' } }) + ); + const [done] = toStreamEvents( + makeEvent({ type: 'run_completed', payload: { answer: 'done', toolCalls: [] } }) + ); + expect(completed.type).toBe('message_completed'); + if (done.type === 'run_completed') { + expect(done.payload.stopReason).toBe('completed'); + } + }); + + it('任务失败映射为 error(保留 code/message)', () => { + const [ev] = toStreamEvents( + makeEvent({ + type: 'run_failed', + payload: { error: { code: 'TOOL_ERROR', message: 'provider timeout' } }, + }) + ); + expect(ev.type).toBe('error'); + if (ev.type === 'error') { + expect(ev.payload.code).toBe('TOOL_ERROR'); + expect(ev.payload.message).toBe('provider timeout'); + expect(ev.payload.retryable).toBe(false); + } + }); + + it('用户取消映射为 cancelled(reason=user)', () => { + const [ev] = toStreamEvents( + makeEvent({ + type: 'run_failed', + payload: { error: { code: 'RUN_CANCELLED', message: 'Run cancelled by user.' } }, + }) + ); + expect(ev.type).toBe('cancelled'); + if (ev.type === 'cancelled') { + expect(ev.payload.reason).toBe('user'); + expect(ev.payload.partial).toEqual({ text: '' }); + } + }); +}); \ No newline at end of file diff --git a/packages/shared/src/kernel/stream-event-adapter.ts b/packages/shared/src/kernel/stream-event-adapter.ts new file mode 100644 index 0000000..2e0d254 --- /dev/null +++ b/packages/shared/src/kernel/stream-event-adapter.ts @@ -0,0 +1,77 @@ +// AgentEvent → Stream Event Protocol v1 转换器(仅运行期映射,无副作用)。 +// 供 RunManager 在现有 AgentEvent 广播旁并行产出协议事件(issue #27, +// docs/adr/0001-stream-event-protocol.md,Migration step 2)。 + +import type { AgentEvent, StreamEvent } from '@finagent/core'; +import { STREAM_EVENT_PROTOCOL_VERSION } from '@finagent/core'; + +/** + * 把单个 AgentEvent 映射为一条或多条 StreamEvent。 + * - 时间戳:AgentEvent 用 epoch 毫秒 number,协议层用 ISO 8601 UTC string。 + * - messageId 与 runId 合并(v1,见 ADR)。 + * - run_failed(code=RUN_CANCELLED) 归一为 cancelled 事件;其余失败归一为 error。 + */ +export function toStreamEvents(event: AgentEvent): StreamEvent[] { + const base = { + protocolVersion: STREAM_EVENT_PROTOCOL_VERSION, + runId: event.runId, + messageId: event.runId, + sequence: event.sequence, + timestamp: new Date(event.timestamp).toISOString(), + }; + + switch (event.type) { + case 'run_started': + return [ + { + ...base, + type: 'run_started', + payload: { + input: event.payload.run.input, + startedAt: new Date(event.payload.run.startedAt).toISOString(), + }, + }, + ]; + case 'message_started': + return [{ ...base, type: 'message_started', payload: {} }]; + case 'message_delta': + return [{ ...base, type: 'text_delta', payload: { text: event.payload.delta } }]; + case 'tool_started': + return [ + { + ...base, + type: 'tool_started', + payload: { + callId: event.payload.toolCall.id, + name: event.payload.toolCall.toolName, + input: event.payload.toolCall.args, + }, + }, + ]; + case 'tool_completed': + return [ + { + ...base, + type: 'tool_result', + payload: { + callId: event.payload.toolCall.id, + name: event.payload.toolCall.toolName, + result: event.payload.toolCall.result, + }, + }, + ]; + case 'message_completed': + return [{ ...base, type: 'message_completed', payload: {} }]; + case 'run_completed': + return [{ ...base, type: 'run_completed', payload: { stopReason: 'completed' } }]; + case 'run_failed': { + const { error } = event.payload; + if (error.code === 'RUN_CANCELLED') { + return [{ ...base, type: 'cancelled', payload: { reason: 'user', partial: { text: '' } } }]; + } + return [ + { ...base, type: 'error', payload: { code: error.code, message: error.message, retryable: false } }, + ]; + } + } +} \ No newline at end of file From 9275dece8398893546329c11d272b3fe2af7e394 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=9F=B3=E6=95=AC=E8=8D=A3?= <15920017+itjingrong@user.noreply.gitee.com> Date: Fri, 11 Sep 2026 12:28:45 +0800 Subject: [PATCH 03/14] fix(core): make StreamEvent a true discriminated union StreamEvent was a single indexed-union instantiation (Tagged StreamEventEnvelope), so payload could not be narrowed by type in switch/if. Rewrite as a distributive mapped union; on-disk type shape is unchanged. Adjusts adapter unit-test helper accordingly. --- packages/core/src/stream-events.ts | 10 ++++- .../src/kernel/stream-event-adapter.test.ts | 44 +++++++------------ 2 files changed, 25 insertions(+), 29 deletions(-) diff --git a/packages/core/src/stream-events.ts b/packages/core/src/stream-events.ts index 63b1e8e..bbe9e27 100644 --- a/packages/core/src/stream-events.ts +++ b/packages/core/src/stream-events.ts @@ -48,7 +48,7 @@ export interface StreamEventTypeToPayload { export type StreamEventPayload = StreamEventTypeToPayload[StreamEventType]; /** - * 统一事件信封。 + * 统一事件信封(参数化版本,供内部组合用)。 * - 幂等键:runId + messageId + sequence。 * - sequence 为 run 内单调递增;reconnect 以它为游标(lastSequence 补发)。 * - timestamp 仅用于展示/排序,不作为身份。 @@ -64,7 +64,13 @@ export interface StreamEventEnvelope = StreamEventEnvelope; +/** + * 可判别事件联合:type 与 payload 强关联,switch/if 收窄后可直接访问 + * 对应 payload 字段(非简单的 envelope 索引联合)。 + */ +export type StreamEvent = { + [P in T]: StreamEventEnvelope

; +}[T]; /** 供完整性检查/测试用的枚举列表,必须与 StreamEventType 一一对应。 */ export const STREAM_EVENT_TYPES = [ diff --git a/packages/shared/src/kernel/stream-event-adapter.test.ts b/packages/shared/src/kernel/stream-event-adapter.test.ts index cbcb579..0a0f748 100644 --- a/packages/shared/src/kernel/stream-event-adapter.test.ts +++ b/packages/shared/src/kernel/stream-event-adapter.test.ts @@ -4,26 +4,28 @@ import { describe, expect, it } from 'bun:test'; import type { AgentEvent } from '@finagent/core'; import { toStreamEvents } from './stream-event-adapter.ts'; -function makeEvent(partial: Partial & { type: AgentEvent['type']; payload: AgentEvent['payload'] }): AgentEvent { +/** + * 测试辅助:拼装一个满足 AgentEventBase 的 AgentEvent。 + * 输入从宽(type + payload),输出强类型由 toStreamEvents 的返回类型把关。 + */ +function makeEvent(type: AgentEvent['type'], payload: unknown): AgentEvent { return { id: 'id-1', sessionId: 'sess-1', runId: 'run-1', timestamp: 1726000000000, sequence: 5, - ...partial, + type, + payload, } as AgentEvent; } describe('toStreamEvents', () => { it('映射 run_started(含 input 与 ISO 时间戳)', () => { const [ev] = toStreamEvents( - makeEvent({ - type: 'run_started', - payload: { - run: { id: 'run-1', sessionId: 'sess-1', status: 'running', input: 'AAPL.US', startedAt: 1726000000000 }, - userMessage: { id: 'm1', role: 'user', content: 'AAPL.US', timestamp: 1726000000000 }, - }, + makeEvent('run_started', { + run: { id: 'run-1', sessionId: 'sess-1', status: 'running', input: 'AAPL.US', startedAt: 1726000000000 }, + userMessage: { id: 'm1', role: 'user', content: 'AAPL.US', timestamp: 1726000000000 }, }) ); expect(ev.type).toBe('run_started'); @@ -38,9 +40,7 @@ describe('toStreamEvents', () => { }); it('映射 message_delta → text_delta(增量字段无损)', () => { - const [ev] = toStreamEvents( - makeEvent({ type: 'message_delta', payload: { delta: 'Apple ', answer: 'Apple Inc.' } }) - ); + const [ev] = toStreamEvents(makeEvent('message_delta', { delta: 'Apple ', answer: 'Apple Inc.' })); expect(ev.type).toBe('text_delta'); if (ev.type === 'text_delta') { expect(ev.payload.text).toBe('Apple '); @@ -56,8 +56,8 @@ describe('toStreamEvents', () => { status: 'success' as const, result: { lastPrice: 220 }, }; - const [started] = toStreamEvents(makeEvent({ type: 'tool_started', payload: { toolCall } })); - const [result] = toStreamEvents(makeEvent({ type: 'tool_completed', payload: { toolCall } })); + const [started] = toStreamEvents(makeEvent('tool_started', { toolCall })); + const [result] = toStreamEvents(makeEvent('tool_completed', { toolCall })); if (started.type === 'tool_started') { expect(started.payload.callId).toBe('tc-1'); expect(started.payload.name).toBe('get_quote'); @@ -68,12 +68,8 @@ describe('toStreamEvents', () => { }); it('映射 run_completed → stopReason=completed / message_completed', () => { - const [completed] = toStreamEvents( - makeEvent({ type: 'message_completed', payload: { answer: 'done' } }) - ); - const [done] = toStreamEvents( - makeEvent({ type: 'run_completed', payload: { answer: 'done', toolCalls: [] } }) - ); + const [completed] = toStreamEvents(makeEvent('message_completed', { answer: 'done' })); + const [done] = toStreamEvents(makeEvent('run_completed', { answer: 'done', toolCalls: [] })); expect(completed.type).toBe('message_completed'); if (done.type === 'run_completed') { expect(done.payload.stopReason).toBe('completed'); @@ -82,10 +78,7 @@ describe('toStreamEvents', () => { it('任务失败映射为 error(保留 code/message)', () => { const [ev] = toStreamEvents( - makeEvent({ - type: 'run_failed', - payload: { error: { code: 'TOOL_ERROR', message: 'provider timeout' } }, - }) + makeEvent('run_failed', { error: { code: 'TOOL_ERROR', message: 'provider timeout' } }) ); expect(ev.type).toBe('error'); if (ev.type === 'error') { @@ -97,10 +90,7 @@ describe('toStreamEvents', () => { it('用户取消映射为 cancelled(reason=user)', () => { const [ev] = toStreamEvents( - makeEvent({ - type: 'run_failed', - payload: { error: { code: 'RUN_CANCELLED', message: 'Run cancelled by user.' } }, - }) + makeEvent('run_failed', { error: { code: 'RUN_CANCELLED', message: 'Run cancelled by user.' } }) ); expect(ev.type).toBe('cancelled'); if (ev.type === 'cancelled') { From 903c544a759274ed6c7cf44cc7552312b571b340 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=9F=B3=E6=95=AC=E8=8D=A3?= <15920017+itjingrong@user.noreply.gitee.com> Date: Fri, 11 Sep 2026 12:28:45 +0800 Subject: [PATCH 04/14] feat(electron): forward Stream Event v1 channel to renderer KernelHost subscribes RunManager.subscribeStream (issue #27) and forwards { sessionId, event } over IPC channel 'agent:stream'; preload exposes electronAPI.kernel.onStreamEvent. Legacy 'agent:event' delivery untouched. Transport only; renderer consumption follows. --- apps/electron/src/main/kernelHost.ts | 9 +++++++++ apps/electron/src/preload/index.ts | 9 +++++++++ 2 files changed, 18 insertions(+) diff --git a/apps/electron/src/main/kernelHost.ts b/apps/electron/src/main/kernelHost.ts index 2b4bed1..15b66bb 100644 --- a/apps/electron/src/main/kernelHost.ts +++ b/apps/electron/src/main/kernelHost.ts @@ -283,6 +283,7 @@ export class AgentKernelHost { private instrumentResolver: InstrumentResolver; private activeLogin: { cancel: () => void } | null = null; private unsubscribe: (() => void) | null = null; + private streamUnsubscribe: (() => void) | null = null; private connectionsUnsubscribe: (() => void) | null = null; private window: BrowserWindow | null = null; @@ -479,6 +480,14 @@ export class AgentKernelHost { window.webContents.send('agent:event', event); } }); + // Stream Event Protocol v1 (issue #27): parallel transport for the + // protocol channel, alongside the legacy agent:event delivery. + this.streamUnsubscribe?.(); + this.streamUnsubscribe = this.kernel.runs.subscribeStream((sessionId, event) => { + if (!window.isDestroyed()) { + window.webContents.send('agent:stream', { sessionId, event }); + } + }); this.connectionsUnsubscribe?.(); this.connectionsUnsubscribe = this.connectionStore.subscribe(() => { void this.pushConnections(); diff --git a/apps/electron/src/preload/index.ts b/apps/electron/src/preload/index.ts index 4a0d283..2d6c0b1 100644 --- a/apps/electron/src/preload/index.ts +++ b/apps/electron/src/preload/index.ts @@ -28,6 +28,7 @@ export interface ElectronAPI { }) => Promise; cancelRun: (input: { sessionId: string; runId: string }) => Promise; onAgentEvent: (callback: (event: unknown) => void) => () => void; + onStreamEvent: (callback: (payload: { sessionId: string; event: unknown }) => void) => () => void; }; agent: { getTools: () => Promise; @@ -214,6 +215,14 @@ const electronAPI: ElectronAPI = { ipcRenderer.removeListener('agent:event', listener); }; }, + onStreamEvent: (callback: (payload: { sessionId: string; event: unknown }) => void) => { + const listener = (_event: Electron.IpcRendererEvent, payload: { sessionId: string; event: unknown }) => + callback(payload); + ipcRenderer.on('agent:stream', listener); + return () => { + ipcRenderer.removeListener('agent:stream', listener); + }; + }, }, agent: { getTools: () => ipcRenderer.invoke('agent:getTools'), From 9ca42c837887fa448e691e11292892a705a50107 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=9F=B3=E6=95=AC=E8=8D=A3?= <15920017+itjingrong@user.noreply.gitee.com> Date: Fri, 11 Sep 2026 12:34:15 +0800 Subject: [PATCH 05/14] feat(ui): consume Stream Event v1 channel (idempotent log) Renderer-side data layer for issue #27: KernelBridge subscribes client.kernel.onStreamEvent into a parallel StreamEvent log. reduceStreamLog keeps per-run events ordered by sequence, dedupes replays (drops), and flags gaps/out-of-order (anomalies) as a protocol health signal. FinagentClient and preload.cjs wire onStreamEvent; FinagentClient adds onStreamEvent contract with fallback noop. No visual change: existing AgentEvent rendering untouched. --- apps/electron/src/preload/index.cjs | 7 +++ apps/electron/src/renderer/finagentClient.ts | 5 ++ packages/ui/src/atoms/streamAtoms.test.ts | 58 +++++++++++++++++++ packages/ui/src/atoms/streamAtoms.ts | 57 ++++++++++++++++++ packages/ui/src/client.tsx | 3 + .../ui/src/components/kernel/KernelBridge.tsx | 12 +++- 6 files changed, 141 insertions(+), 1 deletion(-) create mode 100644 packages/ui/src/atoms/streamAtoms.test.ts create mode 100644 packages/ui/src/atoms/streamAtoms.ts diff --git a/apps/electron/src/preload/index.cjs b/apps/electron/src/preload/index.cjs index fd46cb4..0a09a75 100644 --- a/apps/electron/src/preload/index.cjs +++ b/apps/electron/src/preload/index.cjs @@ -48,6 +48,13 @@ var electronAPI = { return () => { import_electron.ipcRenderer.removeListener("agent:event", listener); }; + }, + onStreamEvent: (callback) => { + const listener = (_event, payload) => callback(payload); + import_electron.ipcRenderer.on("agent:stream", listener); + return () => { + import_electron.ipcRenderer.removeListener("agent:stream", listener); + }; } }, agent: { diff --git a/apps/electron/src/renderer/finagentClient.ts b/apps/electron/src/renderer/finagentClient.ts index 4ef1101..9263c9e 100644 --- a/apps/electron/src/renderer/finagentClient.ts +++ b/apps/electron/src/renderer/finagentClient.ts @@ -3,6 +3,7 @@ import type { AlertTriggerEvent, ApiResult, CustomProviderConfig, + StreamEvent, ThesisImpact, WorkspaceContext, } from '@finagent/core'; @@ -36,6 +37,10 @@ function createElectronClient(): FinagentClient { ipcResult(window.electronAPI.kernel.cancelRun({ sessionId, runId })), onAgentEvent: (callback: (event: AgentEvent) => void) => window.electronAPI.kernel.onAgentEvent((event) => callback(event as AgentEvent)), + onStreamEvent: (callback: (payload: { sessionId: string; event: StreamEvent }) => void) => + window.electronAPI.kernel.onStreamEvent((payload) => + callback({ sessionId: payload.sessionId, event: payload.event as StreamEvent }) + ), }, agent: { getTools: () => ipcResult(window.electronAPI.agent.getTools()), diff --git a/packages/ui/src/atoms/streamAtoms.test.ts b/packages/ui/src/atoms/streamAtoms.test.ts new file mode 100644 index 0000000..2917d14 --- /dev/null +++ b/packages/ui/src/atoms/streamAtoms.test.ts @@ -0,0 +1,58 @@ +// Stream Event Protocol v1 — renderer 事件缓冲 reducer 单测。 + +import { describe, expect, it } from 'bun:test'; +import type { StreamEvent } from '@finagent/core'; +import { reduceStreamLog, type StreamLogState } from './streamAtoms.ts'; + +function make(input: { runId?: string; sequence: number; type: StreamEvent['type'] }): StreamEvent { + return { + protocolVersion: 1, + runId: input.runId ?? 'run-1', + messageId: input.runId ?? 'run-1', + sequence: input.sequence, + type: input.type, + timestamp: '2026-09-11T00:00:00.000Z', + payload: {} as never, + }; +} + +const s0: StreamLogState = { byRun: new Map(), drops: 0, anomalies: 0 }; + +describe('reduceStreamLog', () => { + it('按 sequence 顺序累积同 run 事件', () => { + const s1 = reduceStreamLog(s0, { sessionId: 's1', event: make({ sequence: 1, type: 'run_started' }) }); + const s2 = reduceStreamLog(s1, { sessionId: 's1', event: make({ sequence: 2, type: 'text_delta' }) }); + expect(s2.byRun.get('run-1')?.map((e) => e.type)).toEqual(['run_started', 'text_delta']); + expect(s2.anomalies).toBe(0); + expect(s2.drops).toBe(0); + }); + + it('重复 sequence 幂等丢弃(重放不重复投递)', () => { + const s1 = reduceStreamLog(s0, { sessionId: 's1', event: make({ sequence: 1, type: 'run_started' }) }); + const s2 = reduceStreamLog(s1, { sessionId: 's1', event: make({ sequence: 1, type: 'run_started' }) }); + expect(s2.byRun.get('run-1')).toHaveLength(1); + expect(s2.drops).toBe(1); + expect(s2.anomalies).toBe(0); + }); + + it('乱序事件会重排并标记 anomaly', () => { + const s1 = reduceStreamLog(s0, { sessionId: 's1', event: make({ sequence: 2, type: 'text_delta' }) }); + const s2 = reduceStreamLog(s1, { sessionId: 's1', event: make({ sequence: 1, type: 'run_started' }) }); + expect(s2.byRun.get('run-1')?.map((e) => e.sequence)).toEqual([1, 2]); + expect(s2.anomalies).toBe(2); + }); + + it('不同 run 互不干扰', () => { + const s1 = reduceStreamLog(s0, { sessionId: 's1', event: make({ runId: 'run-a', sequence: 1, type: 'run_started' }) }); + const s2 = reduceStreamLog(s1, { sessionId: 's1', event: make({ runId: 'run-b', sequence: 1, type: 'run_started' }) }); + expect(s2.byRun.size).toBe(2); + expect(s2.byRun.get('run-a')).toHaveLength(1); + expect(s2.byRun.get('run-b')).toHaveLength(1); + }); + + it('保留最近一条事件', () => { + const s1 = reduceStreamLog(s0, { sessionId: 's1', event: make({ sequence: 1, type: 'run_started' }) }); + expect(s1.last?.sessionId).toBe('s1'); + expect(s1.last?.event.type).toBe('run_started'); + }); +}); \ No newline at end of file diff --git a/packages/ui/src/atoms/streamAtoms.ts b/packages/ui/src/atoms/streamAtoms.ts new file mode 100644 index 0000000..084f32a --- /dev/null +++ b/packages/ui/src/atoms/streamAtoms.ts @@ -0,0 +1,57 @@ +// Stream Event Protocol v1 — renderer 侧事件缓冲与幂等(issue #27,PR-D)。 +// +// 目标:在不改变现有渲染路径的前提下,把 onStreamEvent 的协议事件按 +// run 有序缓存 + 幂等去重,供后续渲染切换使用(text_delta 增量渲染、 +// cancelled 显式状态等)。纯函数 reducer,可单测。 + +import { atom } from 'jotai'; +import type { StreamEvent } from '@finagent/core'; + +export interface StreamLogInput { + sessionId: string; + event: StreamEvent; +} + +export interface StreamLogState { + /** runId → 该 run 内已收到的事件(按 sequence 有序)。 */ + byRun: ReadonlyMap>; + /** 因重复 sequence 被幂等丢弃的事件数。 */ + drops: number; + /** sequence 出现 gap 或乱序的次数(协议健康指标)。 */ + anomalies: number; + /** 最近一条事件。 */ + last?: StreamLogInput; +} + +const EMPTY_STATE: StreamLogState = { byRun: new Map(), drops: 0, anomalies: 0 }; + +export function reduceStreamLog(state: StreamLogState, input: StreamLogInput): StreamLogState { + const { event } = input; + const prev = state.byRun.get(event.runId) ?? []; + + // 幂等:同一 run 内相同 sequence 的事件直接丢弃(重放/重复投递)。 + if (prev.some((e) => e.sequence === event.sequence)) { + return { ...state, drops: state.drops + 1, last: input }; + } + + const lastSeq = prev.length > 0 ? prev[prev.length - 1].sequence : 0; + const isAnomaly = event.sequence <= lastSeq || event.sequence !== lastSeq + 1; + const next = isAnomaly + ? [...prev, event].sort((a, b) => a.sequence - b.sequence) + : [...prev, event]; + + return { + byRun: new Map(state.byRun).set(event.runId, next), + drops: state.drops, + anomalies: state.anomalies + (isAnomaly ? 1 : 0), + last: input, + }; +} + +/** 运行中的协议事件日志(调试与渐进渲染的数据源)。 */ +export const streamLogAtom = atom(EMPTY_STATE); + +/** 由 KernelBridge 订阅 onStreamEvent 后喂入。 */ +export const applyStreamEventAtom = atom(null, (_get, set, input: StreamLogInput) => { + set(streamLogAtom, reduceStreamLog(_get(streamLogAtom), input)); +}); \ No newline at end of file diff --git a/packages/ui/src/client.tsx b/packages/ui/src/client.tsx index 1d2afff..c339679 100644 --- a/packages/ui/src/client.tsx +++ b/packages/ui/src/client.tsx @@ -47,6 +47,7 @@ import type { Skill, SkillReadiness, StaticInfo, + StreamEvent, ThesisImpact, ToolDefinition, WorkspaceContext, @@ -188,6 +189,7 @@ export interface FinagentClient { ) => Promise>; cancelRun: (sessionId: string, runId: string) => Promise>; onAgentEvent: (callback: (event: AgentEvent) => void) => () => void; + onStreamEvent: (callback: (payload: { sessionId: string; event: StreamEvent }) => void) => () => void; }; agent: { getTools: () => Promise>; @@ -317,6 +319,7 @@ export const fallbackClient: FinagentClient = { startRun: missingClient('kernel.startRun'), cancelRun: missingClient('kernel.cancelRun'), onAgentEvent: () => () => undefined, + onStreamEvent: () => () => undefined, }, agent: { getTools: missingClient('agent.getTools'), diff --git a/packages/ui/src/components/kernel/KernelBridge.tsx b/packages/ui/src/components/kernel/KernelBridge.tsx index b3e0667..0133656 100644 --- a/packages/ui/src/components/kernel/KernelBridge.tsx +++ b/packages/ui/src/components/kernel/KernelBridge.tsx @@ -8,18 +8,22 @@ import { loadedSessionIdsAtom, } from '../../atoms/sessionAtoms'; import { applyAgentEventAtom } from '../../atoms/runAtoms'; +import { applyStreamEventAtom } from '../../atoms/streamAtoms'; /** * Bridges the kernel's `agent:event` stream and persistence into Jotai state. * * Rendered once under the client provider: hydrates the session list from the * kernel, loads messages for the active session (and on session switches), and - * feeds every kernel agent event through the run reducer. + * feeds every kernel agent event through the run reducer. Also subscribes to + * the Stream Event Protocol v1 channel (issue #27) into a parallel idempotent + * log for progressive renderer adoption. */ export const KernelBridge: React.FC<{ client: FinagentClient }> = ({ client }) => { const hydrate = useSetAtom(hydrateSessionsAtom); const loadMessages = useSetAtom(loadMessagesAtom); const applyEvent = useSetAtom(applyAgentEventAtom); + const applyStreamEvent = useSetAtom(applyStreamEventAtom); const [activeSessionId] = useAtom(activeSessionIdAtom); const [loadedSessionIds] = useAtom(loadedSessionIdsAtom); @@ -31,6 +35,12 @@ export const KernelBridge: React.FC<{ client: FinagentClient }> = ({ client }) = }); }, [client, hydrate, applyEvent]); + useEffect(() => { + return client.kernel.onStreamEvent((payload) => { + applyStreamEvent({ sessionId: payload.sessionId, event: payload.event }); + }); + }, [client, applyStreamEvent]); + useEffect(() => { if (activeSessionId && !loadedSessionIds.has(activeSessionId)) { void loadMessages(client, activeSessionId); From cc7c888db3e55a9614ca5e0c0a24486a6a251192 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=9F=B3=E6=95=AC=E8=8D=A3?= <15920017+itjingrong@user.noreply.gitee.com> Date: Fri, 11 Sep 2026 15:36:30 +0800 Subject: [PATCH 06/14] feat: decouple message/run identity and add replay history Addresses #43 review: messageId is no longer fused with runId (issue #34 - message/generation/run are distinct identities). message-level events carry the real assistant messageId, pre-assigned per run in RunManager; run-level events omit it; idempotency key is now runId + sequence. Adds StreamEventHistory: in-memory per-run tail used by RunManager.replayStream(runId, lastSequence) for reconnect resume, with an explicit recoverable:false path when the run is unknown or the tail is non-contiguous (eviction). --- packages/core/src/stream-events.ts | 13 +++- packages/shared/src/kernel/run-manager.ts | 30 +++++++- .../src/kernel/stream-event-adapter.test.ts | 19 ++++- .../shared/src/kernel/stream-event-adapter.ts | 26 +++++-- .../shared/src/kernel/stream-history.test.ts | 63 ++++++++++++++++ packages/shared/src/kernel/stream-history.ts | 72 +++++++++++++++++++ 6 files changed, 209 insertions(+), 14 deletions(-) create mode 100644 packages/shared/src/kernel/stream-history.test.ts create mode 100644 packages/shared/src/kernel/stream-history.ts diff --git a/packages/core/src/stream-events.ts b/packages/core/src/stream-events.ts index bbe9e27..7a1258e 100644 --- a/packages/core/src/stream-events.ts +++ b/packages/core/src/stream-events.ts @@ -49,15 +49,22 @@ export type StreamEventPayload = StreamEventTypeToPayload[StreamEventType]; /** * 统一事件信封(参数化版本,供内部组合用)。 - * - 幂等键:runId + messageId + sequence。 + * - 幂等键:runId + sequence(messageId 不作为幂等身份的一部分)。 * - sequence 为 run 内单调递增;reconnect 以它为游标(lastSequence 补发)。 * - timestamp 仅用于展示/排序,不作为身份。 + * + * 身份语义(对齐 issue #34 / 维护者 #43 review): + * - runId:一次 generation 的唯一身份,所有事件必填。 + * - messageId:仅 message 级事件(message_started / text_delta / + * message_completed / cancelled 的部分文本归属)携带到真实、稳定的 + * message id;run 级事件(run_started / run_completed / tool_* / + * citation_added / status / error)可省略。message 与 run 是两个身份, + * 允许一条 message 跨多次 run(edit/regenerate/fork,见 #34)。 */ export interface StreamEventEnvelope { protocolVersion: typeof STREAM_EVENT_PROTOCOL_VERSION; runId: string; - /** v1 与 runId 相同;预留“一条 message 跨多次 run”时拆分。 */ - messageId: string; + messageId?: string; sequence: number; type: T; timestamp: string; diff --git a/packages/shared/src/kernel/run-manager.ts b/packages/shared/src/kernel/run-manager.ts index c3dd1a7..e6a0813 100644 --- a/packages/shared/src/kernel/run-manager.ts +++ b/packages/shared/src/kernel/run-manager.ts @@ -39,6 +39,7 @@ import { } from './runaway-detector.ts'; import { buildFinancialEvidence } from '../evidence/financial-evidence.ts'; import { toStreamEvents } from './stream-event-adapter.ts'; +import { StreamEventHistory, type StreamReplayResult } from './stream-history.ts'; export interface RunManagerOptions { sessions: SessionManager; @@ -91,6 +92,10 @@ export class RunManager { private readonly listeners = new Set<(event: AgentEvent) => void>(); private readonly streamListeners = new Set<(sessionId: string, event: StreamEvent) => void>(); private activeRun: ActiveRun | null = null; + /** 当前 run 对应的 assistant message id(贯穿 message 级事件,见 #34 身份模型)。 */ + private currentMessageId: string | null = null; + /** run 级内存事件历史:为 reconnect replay 提供数据源(ADR 0001)。 */ + private readonly streamHistory = new StreamEventHistory(); constructor(options: RunManagerOptions) { this.sessions = options.sessions; @@ -117,6 +122,14 @@ export class RunManager { return () => this.streamListeners.delete(listener); } + /** + * Reconnect/重连补发(ADR 0001 §Reconnect):基于运行期内存历史返回 + * lastSequence 之后的连续段;段缺失或运行未知时明确返回不可恢复。 + */ + replayStream(runId: string, lastSequence: number): StreamReplayResult { + return this.streamHistory.replay(runId, lastSequence); + } + /** Whether a run is currently executing (Pi runtime executes one at a time). */ isRunning(): boolean { return this.activeRun !== null; @@ -227,6 +240,11 @@ export class RunManager { const toolCalls: ToolCall[] = []; let sawTerminal = false; + // 一次 run = 一次 assistant generation(#34):在 run 启动时确定稳定的 + // assistant message id,贯穿 message 级 stream 事件,并作为最终写入的 + // assistantMessage.id(避免流中与落库两套 id)。 + this.currentMessageId = randomUUID(); + try { await this.runtime.ensureSession({ id: run.sessionId, @@ -296,7 +314,7 @@ export class RunManager { const isInfraFailure = run.status === 'failed' && isRuntimeInfraCode(run.error?.code); if (!isInfraFailure) { const assistantMessage: Message = { - id: randomUUID(), + id: this.currentMessageId ?? randomUUID(), role: 'assistant', content: answer || (run.status === 'failed' ? run.error?.message ?? 'Run failed.' : ''), timestamp: now, @@ -316,6 +334,7 @@ export class RunManager { // The run is fully settled (persisted) only now; only then allow the next run. this.activeRun = null; + this.currentMessageId = null; // Adapters emit the terminal event themselves; synthesize it only when the // stream failed before producing one (e.g. runtime spawn failure), so the @@ -421,10 +440,15 @@ export class RunManager { for (const listener of this.listeners) { listener(event); } + const mapped = toStreamEvents(event, { messageId: this.currentMessageId ?? undefined }); + // 无论是否有实时订阅者,都先记录进内存历史,保证 replay 有数据源。 + for (const streamEvent of mapped) { + this.streamHistory.append(streamEvent); + } if (this.streamListeners.size > 0) { - for (const mapped of toStreamEvents(event)) { + for (const streamEvent of mapped) { for (const listener of this.streamListeners) { - listener(event.sessionId, mapped); + listener(event.sessionId, streamEvent); } } } diff --git a/packages/shared/src/kernel/stream-event-adapter.test.ts b/packages/shared/src/kernel/stream-event-adapter.test.ts index 0a0f748..f497bd2 100644 --- a/packages/shared/src/kernel/stream-event-adapter.test.ts +++ b/packages/shared/src/kernel/stream-event-adapter.test.ts @@ -31,7 +31,8 @@ describe('toStreamEvents', () => { expect(ev.type).toBe('run_started'); if (ev.type === 'run_started') { expect(ev.protocolVersion).toBe(1); - expect(ev.messageId).toBe('run-1'); + expect(ev.runId).toBe('run-1'); + expect(ev.messageId).toBeUndefined(); expect(ev.sequence).toBe(5); expect(ev.payload.input).toBe('AAPL.US'); expect(ev.payload.startedAt).toBe('2024-09-10T20:26:40.000Z'); @@ -88,6 +89,22 @@ describe('toStreamEvents', () => { } }); + it('message 级事件携带真实 messageId,run 级事件不携带', () => { + const [delta] = toStreamEvents(makeEvent('message_delta', { delta: 'x', answer: 'x' }), { + messageId: 'msg-9', + }); + const [done] = toStreamEvents(makeEvent('run_completed', { answer: 'x', toolCalls: [] }), { + messageId: 'msg-9', + }); + if (delta.type === 'text_delta') { + expect(delta.messageId).toBe('msg-9'); + } + if (done.type === 'run_completed') { + expect(done.messageId).toBeUndefined(); + } + expect(done.runId).toBe('run-1'); + }); + it('用户取消映射为 cancelled(reason=user)', () => { const [ev] = toStreamEvents( makeEvent('run_failed', { error: { code: 'RUN_CANCELLED', message: 'Run cancelled by user.' } }) diff --git a/packages/shared/src/kernel/stream-event-adapter.ts b/packages/shared/src/kernel/stream-event-adapter.ts index 2e0d254..30947ad 100644 --- a/packages/shared/src/kernel/stream-event-adapter.ts +++ b/packages/shared/src/kernel/stream-event-adapter.ts @@ -8,17 +8,22 @@ import { STREAM_EVENT_PROTOCOL_VERSION } from '@finagent/core'; /** * 把单个 AgentEvent 映射为一条或多条 StreamEvent。 * - 时间戳:AgentEvent 用 epoch 毫秒 number,协议层用 ISO 8601 UTC string。 - * - messageId 与 runId 合并(v1,见 ADR)。 + * - 身份(对齐 issue #34):runId 恒有;messageId 仅注入到 message 级事件 + * (message_started / text_delta / message_completed / cancelled), + * 取 run 对应的真实 assistant message id;run 级事件不携带 messageId。 * - run_failed(code=RUN_CANCELLED) 归一为 cancelled 事件;其余失败归一为 error。 */ -export function toStreamEvents(event: AgentEvent): StreamEvent[] { +export function toStreamEvents( + event: AgentEvent, + opts?: { messageId?: string } +): StreamEvent[] { const base = { protocolVersion: STREAM_EVENT_PROTOCOL_VERSION, runId: event.runId, - messageId: event.runId, sequence: event.sequence, timestamp: new Date(event.timestamp).toISOString(), }; + const messageish = opts?.messageId ? { messageId: opts.messageId } : {}; switch (event.type) { case 'run_started': @@ -33,9 +38,9 @@ export function toStreamEvents(event: AgentEvent): StreamEvent[] { }, ]; case 'message_started': - return [{ ...base, type: 'message_started', payload: {} }]; + return [{ ...base, ...messageish, type: 'message_started', payload: {} }]; case 'message_delta': - return [{ ...base, type: 'text_delta', payload: { text: event.payload.delta } }]; + return [{ ...base, ...messageish, type: 'text_delta', payload: { text: event.payload.delta } }]; case 'tool_started': return [ { @@ -61,13 +66,20 @@ export function toStreamEvents(event: AgentEvent): StreamEvent[] { }, ]; case 'message_completed': - return [{ ...base, type: 'message_completed', payload: {} }]; + return [{ ...base, ...messageish, type: 'message_completed', payload: {} }]; case 'run_completed': return [{ ...base, type: 'run_completed', payload: { stopReason: 'completed' } }]; case 'run_failed': { const { error } = event.payload; if (error.code === 'RUN_CANCELLED') { - return [{ ...base, type: 'cancelled', payload: { reason: 'user', partial: { text: '' } } }]; + return [ + { + ...base, + ...messageish, + type: 'cancelled', + payload: { reason: 'user', partial: { text: '' } }, + }, + ]; } return [ { ...base, type: 'error', payload: { code: error.code, message: error.message, retryable: false } }, diff --git a/packages/shared/src/kernel/stream-history.test.ts b/packages/shared/src/kernel/stream-history.test.ts new file mode 100644 index 0000000..0cb8490 --- /dev/null +++ b/packages/shared/src/kernel/stream-history.test.ts @@ -0,0 +1,63 @@ +// Stream Event Protocol v1 — 内存 replay 历史单测。 + +import { describe, expect, it } from 'bun:test'; +import type { StreamEvent } from '@finagent/core'; +import { StreamEventHistory } from './stream-history.ts'; + +function make(runId: string, sequence: number, type: StreamEvent['type']): StreamEvent { + return { + protocolVersion: 1, + runId, + sequence, + type, + timestamp: '2026-09-11T00:00:00.000Z', + payload: {} as never, + }; +} + +function fill(history: StreamEventHistory, runId: string, from: number, to: number, type: StreamEvent['type'] = 'text_delta') { + for (let s = from; s <= to; s += 1) history.append(make(runId, s, type)); +} + +describe('StreamEventHistory', () => { + it('从 0 开始补发整段连续事件', () => { + const h = new StreamEventHistory(); + fill(h, 'run-1', 1, 3); + h.append(make('run-1', 4, 'run_completed')); + const result = h.replay('run-1', 0); + expect(result.recoverable).toBe(true); + expect(result.events.map((e) => e.sequence)).toEqual([1, 2, 3, 4]); + expect(result.atEnd).toBe(true); + }); + + it('只补发 lastSequence 之后的连续段', () => { + const h = new StreamEventHistory(); + fill(h, 'run-1', 1, 5); + const result = h.replay('run-1', 2); + expect(result.recoverable).toBe(true); + expect(result.events.map((e) => e.sequence)).toEqual([3, 4, 5]); + }); + + it('lastSequence 已是末尾 → 空补发', () => { + const h = new StreamEventHistory(); + fill(h, 'run-1', 1, 2); + const result = h.replay('run-1', 2); + expect(result.recoverable).toBe(true); + expect(result.events).toEqual([]); + }); + + it('未知 run → 明确不可恢复', () => { + const h = new StreamEventHistory(); + const result = h.replay('never-seen', 0); + expect(result.recoverable).toBe(false); + expect(result.events).toEqual([]); + }); + + it('缺失段(客户端断在 5,历史只剩 ≥99)→ 明确不可恢复', () => { + const h = new StreamEventHistory(); + fill(h, 'run-1', 99, 102); + const result = h.replay('run-1', 5); + expect(result.recoverable).toBe(false); + expect(result.events).toEqual([]); + }); +}); \ No newline at end of file diff --git a/packages/shared/src/kernel/stream-history.ts b/packages/shared/src/kernel/stream-history.ts new file mode 100644 index 0000000..559fb5e --- /dev/null +++ b/packages/shared/src/kernel/stream-history.ts @@ -0,0 +1,72 @@ +// Stream Event Protocol v1 — 运行期内存事件历史(reconnect/重连补发)。 +// +// 依据维护者 #43 review 与 ADR 0001:reconnect 需要一个明确的 replay +// source + “无法恢复”的失败路径。v1 采用 run 级内存缓冲: +// - append 时按 run 累积(每 run 上限 MAX_EVENTS_PER_RUN,超出截断最旧); +// - replay(runId, lastSequence) 只补发“连续且已存在”的段; +// - 截断/未知 run/段不连续 → recoverable=false,调用方应呈现明确失败。 + +import type { StreamEvent } from '@finagent/core'; + +/** 保留的最大 run 数(最早的新增 run 被淘汰)。 */ +const MAX_RUNS = 32; +/** 单 run 内存缓冲的事件上限(截断后无法从头恢复 → 明确不可恢复)。 */ +const MAX_EVENTS_PER_RUN = 2000; + +export interface StreamReplayResult { + /** true:补发成功;false:内存中已无连续段/未知 run,明确不可恢复。 */ + recoverable: boolean; + /** 连续补发的事件(lastSequence 之后)。不可恢复时为 []。 */ + events: StreamEvent[]; + /** 该 run 是否已结束(最后事件为终端事件:run_completed / cancelled / error)。 */ + atEnd: boolean; +} + +const TERMINAL_TYPES = new Set(['run_completed', 'cancelled', 'error']); + +export class StreamEventHistory { + private readonly runs = new Map(); + /** FIFO LRU 顺序,用于淘汰最久未更新的 run。 */ + private readonly order: string[] = []; + + append(event: StreamEvent): void { + const list = this.runs.get(event.runId); + if (list) { + list.push(event); + if (list.length > MAX_EVENTS_PER_RUN) { + list.splice(0, list.length - MAX_EVENTS_PER_RUN); + } + return; + } + this.runs.set(event.runId, [event]); + this.order.push(event.runId); + while (this.order.length > MAX_RUNS) { + const evicted = this.order.shift(); + if (evicted) this.runs.delete(evicted); + } + } + + /** 从 lastSequence 之后的位置补发;段缺失/未知 run 判为不可恢复。 */ + replay(runId: string, lastSequence: number): StreamReplayResult { + const list = this.runs.get(runId); + if (!list || list.length === 0) { + return { recoverable: false, events: [], atEnd: true }; + } + + const tail = list.filter((e) => e.sequence > lastSequence); + const contiguous = + tail.length === 0 || + (tail[0].sequence === lastSequence + 1 && + tail.every((e, i) => i === 0 || e.sequence === tail[i - 1].sequence + 1)); + if (!contiguous) { + return { recoverable: false, events: [], atEnd: false }; + } + + const last = list[list.length - 1]; + return { + recoverable: true, + events: tail, + atEnd: TERMINAL_TYPES.has(last.type), + }; + } +} \ No newline at end of file From 7c923a3a1bc910aae0c63ec859d7065f787e313c Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=9F=B3=E6=95=AC=E8=8D=A3?= <15920017+itjingrong@user.noreply.gitee.com> Date: Fri, 11 Sep 2026 15:36:31 +0800 Subject: [PATCH 07/14] feat(electron,ui): expose stream replay over IPC Adds 'runs:stream-replay' handler (KernelHost.streamReplay -> RunManager.replayStream), preload (ts + cjs) streamReplay, finagentClient wiring and FinagentClient contract with fallback noop. --- apps/electron/src/main/index.ts | 3 +++ apps/electron/src/main/kernelHost.ts | 13 +++++++++++++ apps/electron/src/preload/index.cjs | 1 + apps/electron/src/preload/index.ts | 3 +++ apps/electron/src/renderer/finagentClient.ts | 2 ++ packages/ui/src/client.tsx | 2 ++ 6 files changed, 24 insertions(+) diff --git a/apps/electron/src/main/index.ts b/apps/electron/src/main/index.ts index 9ae724e..423564b 100644 --- a/apps/electron/src/main/index.ts +++ b/apps/electron/src/main/index.ts @@ -134,6 +134,9 @@ ipcMain.handle('runs:start', async (_event, input: unknown) => ipcMain.handle('runs:cancel', async (_event, input: unknown) => toIpcResult(() => agentKernelHost.cancelRun(input)) ); +ipcMain.handle('runs:stream-replay', async (_event, input: unknown) => + toIpcResult(() => agentKernelHost.streamReplay(input)) +); ipcMain.handle('agent:getTools', async () => toIpcResult(() => agentKernelHost.getTools()) diff --git a/apps/electron/src/main/kernelHost.ts b/apps/electron/src/main/kernelHost.ts index 15b66bb..55235bb 100644 --- a/apps/electron/src/main/kernelHost.ts +++ b/apps/electron/src/main/kernelHost.ts @@ -572,6 +572,19 @@ export class AgentKernelHost { ); } + /** Stream Event replay(ADR 0001 §Reconnect):按 lastSequence 补发或明确不可恢复。 */ + streamReplay(input: unknown) { + const request = requireObject(input); + const lastSequence = request.lastSequence; + if (typeof lastSequence !== 'number' || !Number.isFinite(lastSequence)) { + throw createCodeError('INVALID_ARGUMENT', 'lastSequence must be a finite number.'); + } + return this.kernel.runs.replayStream( + requireString(request.runId, 'runId'), + lastSequence + ); + } + getTools(): Promise> { return this.kernel.getTools(); } diff --git a/apps/electron/src/preload/index.cjs b/apps/electron/src/preload/index.cjs index 0a09a75..81ed063 100644 --- a/apps/electron/src/preload/index.cjs +++ b/apps/electron/src/preload/index.cjs @@ -42,6 +42,7 @@ var electronAPI = { listRuns: (sessionId) => import_electron.ipcRenderer.invoke("sessions:listRuns", sessionId), startRun: (input) => import_electron.ipcRenderer.invoke("runs:start", input), cancelRun: (input) => import_electron.ipcRenderer.invoke("runs:cancel", input), + streamReplay: (input) => import_electron.ipcRenderer.invoke("runs:stream-replay", input), onAgentEvent: (callback) => { const listener = (_event, agentEvent) => callback(agentEvent); import_electron.ipcRenderer.on("agent:event", listener); diff --git a/apps/electron/src/preload/index.ts b/apps/electron/src/preload/index.ts index 2d6c0b1..d4e018c 100644 --- a/apps/electron/src/preload/index.ts +++ b/apps/electron/src/preload/index.ts @@ -27,6 +27,7 @@ export interface ElectronAPI { workspaceContext?: unknown; }) => Promise; cancelRun: (input: { sessionId: string; runId: string }) => Promise; + streamReplay: (input: { runId: string; lastSequence: number }) => Promise; onAgentEvent: (callback: (event: unknown) => void) => () => void; onStreamEvent: (callback: (payload: { sessionId: string; event: unknown }) => void) => () => void; }; @@ -208,6 +209,8 @@ const electronAPI: ElectronAPI = { workspaceContext?: unknown; }) => ipcRenderer.invoke('runs:start', input), cancelRun: (input: { sessionId: string; runId: string }) => ipcRenderer.invoke('runs:cancel', input), + streamReplay: (input: { runId: string; lastSequence: number }) => + ipcRenderer.invoke('runs:stream-replay', input), onAgentEvent: (callback: (event: unknown) => void) => { const listener = (_event: Electron.IpcRendererEvent, agentEvent: unknown) => callback(agentEvent); ipcRenderer.on('agent:event', listener); diff --git a/apps/electron/src/renderer/finagentClient.ts b/apps/electron/src/renderer/finagentClient.ts index 9263c9e..456172b 100644 --- a/apps/electron/src/renderer/finagentClient.ts +++ b/apps/electron/src/renderer/finagentClient.ts @@ -35,6 +35,8 @@ function createElectronClient(): FinagentClient { ipcResult(window.electronAPI.kernel.startRun({ sessionId, content, workspaceContext })), cancelRun: (sessionId: string, runId: string) => ipcResult(window.electronAPI.kernel.cancelRun({ sessionId, runId })), + streamReplay: (input: { runId: string; lastSequence: number }) => + ipcResult(window.electronAPI.kernel.streamReplay(input)), onAgentEvent: (callback: (event: AgentEvent) => void) => window.electronAPI.kernel.onAgentEvent((event) => callback(event as AgentEvent)), onStreamEvent: (callback: (payload: { sessionId: string; event: StreamEvent }) => void) => diff --git a/packages/ui/src/client.tsx b/packages/ui/src/client.tsx index c339679..1876307 100644 --- a/packages/ui/src/client.tsx +++ b/packages/ui/src/client.tsx @@ -188,6 +188,7 @@ export interface FinagentClient { workspaceContext?: WorkspaceContext ) => Promise>; cancelRun: (sessionId: string, runId: string) => Promise>; + streamReplay: (input: { runId: string; lastSequence: number }) => Promise>; onAgentEvent: (callback: (event: AgentEvent) => void) => () => void; onStreamEvent: (callback: (payload: { sessionId: string; event: StreamEvent }) => void) => () => void; }; @@ -318,6 +319,7 @@ export const fallbackClient: FinagentClient = { listRuns: missingClient('kernel.listRuns'), startRun: missingClient('kernel.startRun'), cancelRun: missingClient('kernel.cancelRun'), + streamReplay: missingClient('kernel.streamReplay'), onAgentEvent: () => () => undefined, onStreamEvent: () => () => undefined, }, From 5084cea6e0e28213be6cd4a064a67c2620f42593 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=9F=B3=E6=95=AC=E8=8D=A3?= <15920017+itjingrong@user.noreply.gitee.com> Date: Fri, 11 Sep 2026 15:36:31 +0800 Subject: [PATCH 08/14] docs(adr): align 0001 with #34 identity model and reconnect implementation Envelope messageId is now optional (message-level events only); idempotency key runId+sequence; Reconnect row documents implemented StreamEventHistory + IPC with explicit unrecoverable path; open question 3 resolved. --- docs/adr/0001-stream-event-protocol.md | 27 ++++++++++++++++++-------- 1 file changed, 19 insertions(+), 8 deletions(-) diff --git a/docs/adr/0001-stream-event-protocol.md b/docs/adr/0001-stream-event-protocol.md index 9811c29..ac73c18 100644 --- a/docs/adr/0001-stream-event-protocol.md +++ b/docs/adr/0001-stream-event-protocol.md @@ -32,15 +32,24 @@ gradually replaces it. ```ts interface StreamEventEnvelope { protocolVersion: 1; // forward-compat gate - runId: string; - messageId: string; // equals runId in v1; split later if one message spans runs - sequence: number; // monotonic per run; idempotency key = runId + messageId + sequence + runId: string; // one run = one generation (identity, see #34) + messageId?: string; // message-level events only: stable message id (#34); + // run-level events omit it. message and run are distinct identities + sequence: number; // monotonic per run; idempotency key = runId + sequence type: T; timestamp: string; // ISO 8601 UTC; display only, never identity payload: StreamEventTypeToPayload[T]; } ``` +Identity model (aligned with #34): a message has a stable id and can span +multiple runs (edit / regenerate / fork); an assistant generation binds to +exactly one run. Message-level events (`message_started / text_delta / +message_completed / cancelled`) carry the real assistant `messageId`; run-level +events (`run_started / run_completed / tool_* / citation_added / status / +error`) carry only `runId`. Implemented in `run-manager` (pre-assigned +assistant message id per run) and `stream-event-adapter`. + ### Event types (12) `run_started · message_started · text_delta · tool_started · tool_progress · tool_result · @@ -53,9 +62,9 @@ Key behavior change: `message_delta` (full-answer snapshot) becomes `text_delta` | Capability | Contract | |---|---| -| Idempotency | Consumers dedupe on `runId + messageId + sequence`; replayed events never re-insert text, tool cards, or citations | +| Idempotency | Consumers dedupe on `runId + sequence`; replayed events never re-insert text, tool cards, or citations | | Cancel | renderer `cancelRun` → runtime propagates → runtime **explicitly emits `cancelled`**; partial answer preserved | -| Reconnect | client reconnects with `lastSequence`; unconsumed gap is replayed (resume), consumed events are not re-applied | +| Reconnect | **Implemented (v1, in-memory):** `RunManager.replayStream(runId, lastSequence)` returns the contiguous tail from `StreamEventHistory` (`packages/shared/src/kernel/stream-history.ts`); when the run is unknown or the tail is non-contiguous (buffer eviction), it returns `recoverable: false` — an explicit unrecoverable path the renderer must surface. IPC: `runs:stream-replay` | | Final-state parity | UI `stopReason` and persisted `run.status` come from the same single final state in run-manager (aligns with #18) | | Security (#19) | `status` never exposes chain-of-thought; tool payloads pass redaction before reaching the UI; renderer never executes model-returned code | @@ -73,7 +82,9 @@ Key behavior change: `message_delta` (full-answer snapshot) becomes `text_delta` state; groundwork for reconnect/resume and for #21 immutable run manifests. - **Negative:** dual event families during migration; consumers must handle both. - **Open questions for reviewers:** - 1. Resume data source: v1 keeps events in memory for the run's lifetime and defers - persisting an event log (ties into #21) — acceptable? + 1. Resume data source: v1 keeps events in memory (`StreamEventHistory`); persisting + an event log is deferred (ties into #21 immutable run manifests) — acceptable? 2. Is `status.phase` of `thinking / searching / working` sufficient? - 3. Keep `messageId` merged with `runId` in v1 and split only when needed? \ No newline at end of file + 3. ~~Keep `messageId` merged with `runId`~~ **resolved**: #34 requires distinct + message / run identities — envelope carries `messageId` only on message-level + events (see Identity model above). \ No newline at end of file From 46dcdd1a4778de552b6dfedfa842c1516efad3b8 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=9F=B3=E6=95=AC=E8=8D=A3?= <15920017+itjingrong@user.noreply.gitee.com> Date: Fri, 11 Sep 2026 16:36:31 +0800 Subject: [PATCH 09/14] =?UTF-8?q?fix(shared):=20=E5=9C=A8=20emit=20?= =?UTF-8?q?=E6=B1=87=E8=81=9A=E7=82=B9=E7=BB=9F=E4=B8=80=E6=8C=89=20run=20?= =?UTF-8?q?=E9=87=8D=E6=8E=92=20sequence?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit run_started(seq 1)与 runtime 自产的首个事件(同样 seq 1)冲突, 导致 replay(runId, 0) 被判为不连续(recoverable: false),renderer 的 幂等去重也会误丢事件。RunManager 现为每个 run 持有 RunProtocol 计数器, 扇出前把所有 AgentEvent 统一重排为 1..N;messageId 随该对象传递, 崩溃/取消兜底合成的 terminal 事件也保住 #34 的身份契约。replay 游标 超出已知最大 sequence 时,由静默视为已同步改为明确不可恢复。 --- packages/shared/src/kernel/run-manager.ts | 87 ++++++++++++-------- packages/shared/src/kernel/stream-history.ts | 8 +- 2 files changed, 61 insertions(+), 34 deletions(-) diff --git a/packages/shared/src/kernel/run-manager.ts b/packages/shared/src/kernel/run-manager.ts index e6a0813..551924b 100644 --- a/packages/shared/src/kernel/run-manager.ts +++ b/packages/shared/src/kernel/run-manager.ts @@ -72,6 +72,18 @@ interface ActiveRun { stop?: RunStop; } +/** + * 单个 run 的协议状态(ADR 0001 / #34)。 + * - sequence:run 内严格单调 +1,由 RunManager 在唯一汇聚点统一重排 + * (runtime 与本管理器都自产事件且各自计数,会冲突在 1 上); + * - messageId:该 run 对应的 assistant message 身份,message 级事件携带。 + * 以局部对象随 run 传递,activeRun 解锁后合成 terminal 事件仍持有原状态。 + */ +interface RunProtocol { + sequence: number; + messageId: string; +} + /** * Starts, observes, persists, and terminates runs. * @@ -92,8 +104,6 @@ export class RunManager { private readonly listeners = new Set<(event: AgentEvent) => void>(); private readonly streamListeners = new Set<(sessionId: string, event: StreamEvent) => void>(); private activeRun: ActiveRun | null = null; - /** 当前 run 对应的 assistant message id(贯穿 message 级事件,见 #34 身份模型)。 */ - private currentMessageId: string | null = null; /** run 级内存事件历史:为 reconnect replay 提供数据源(ADR 0001)。 */ private readonly streamHistory = new StreamEventHistory(); @@ -205,17 +215,23 @@ export class RunManager { usage: createUsage(), runaway: createRunawayState(), }; - this.emit({ - id: randomUUID(), - sessionId, - runId: run.id, - type: 'run_started', - timestamp: now, - sequence: 1, - payload: { run, userMessage }, - }); + // 一次 run = 一次 assistant generation(#34):run 启动时确定稳定的 + // assistant message id 与协议 sequence 计数器,随 run 全程传递。 + const protocol: RunProtocol = { sequence: 0, messageId: randomUUID() }; + this.emit( + { + id: randomUUID(), + sessionId, + runId: run.id, + type: 'run_started', + timestamp: now, + sequence: 1, + payload: { run, userMessage }, + }, + protocol + ); - void this.execute(run, session, workspaceContext, locale); + void this.execute(run, session, workspaceContext, locale, protocol); return run; } @@ -233,18 +249,14 @@ export class RunManager { run: Run, session: SessionMeta, workspaceContext?: WorkspaceContext, - locale?: SupportedLocale + locale?: SupportedLocale, + protocol: RunProtocol = { sequence: 0, messageId: randomUUID() } ): Promise { let failure: ApiError | undefined; let answer = ''; const toolCalls: ToolCall[] = []; let sawTerminal = false; - // 一次 run = 一次 assistant generation(#34):在 run 启动时确定稳定的 - // assistant message id,贯穿 message 级 stream 事件,并作为最终写入的 - // assistantMessage.id(避免流中与落库两套 id)。 - this.currentMessageId = randomUUID(); - try { await this.runtime.ensureSession({ id: run.sessionId, @@ -260,7 +272,7 @@ export class RunManager { workspaceContext, locale, })) { - this.emit(event); + this.emit(event, protocol); if (event.type === 'message_delta' || event.type === 'message_completed') { answer = event.payload.answer; } else if (event.type === 'tool_completed') { @@ -314,7 +326,7 @@ export class RunManager { const isInfraFailure = run.status === 'failed' && isRuntimeInfraCode(run.error?.code); if (!isInfraFailure) { const assistantMessage: Message = { - id: this.currentMessageId ?? randomUUID(), + id: protocol.messageId, role: 'assistant', content: answer || (run.status === 'failed' ? run.error?.message ?? 'Run failed.' : ''), timestamp: now, @@ -334,27 +346,27 @@ export class RunManager { // The run is fully settled (persisted) only now; only then allow the next run. this.activeRun = null; - this.currentMessageId = null; // Adapters emit the terminal event themselves; synthesize it only when the // stream failed before producing one (e.g. runtime spawn failure), so the - // UI always observes a terminal event. + // UI always observes a terminal event. 合成事件仍持有原 run 的 protocol + // (sequence 续排、messageId 不丢),activeRun 解锁不影响。 if (!sawTerminal) { if (stop !== undefined) { - this.emitRunEvent(run, 'run_failed', { error: stopError(stop) }); + this.emitRunEvent(run, protocol, 'run_failed', { error: stopError(stop) }); } else if (cancelled) { - this.emitRunEvent(run, 'run_failed', { + this.emitRunEvent(run, protocol, 'run_failed', { error: { code: 'RUN_CANCELLED', message: 'Run cancelled by user.' }, }); } else if (failure) { - this.emitRunEvent(run, 'run_failed', { error: failure }); + this.emitRunEvent(run, protocol, 'run_failed', { error: failure }); } else { - this.emitRunEvent(run, 'run_completed', { answer, toolCalls }); + this.emitRunEvent(run, protocol, 'run_completed', { answer, toolCalls }); } } } - /** +/** * Account for one runtime event and decide whether the run must stop. A stop * requests cancellation, so the caller stops consuming events and the run * settles as `cancelled` — never as an ordinary success — carrying its partial @@ -422,7 +434,12 @@ export class RunManager { return typeof query === 'string' && query.trim() !== '' ? query : undefined; } - private emitRunEvent(run: Run, type: AgentEvent['type'], payload?: AgentEventPayload): void { + private emitRunEvent( + run: Run, + protocol: RunProtocol, + type: AgentEvent['type'], + payload?: AgentEventPayload + ): void { // Callers pair `type` with the matching payload shape. const event = { id: randomUUID(), @@ -433,14 +450,18 @@ export class RunManager { sequence: 1, payload, } as AgentEvent; - this.emit(event); + this.emit(event, protocol); } - private emit(event: AgentEvent): void { + private emit(event: AgentEvent, protocol: RunProtocol): void { + // ADR 0001 sequence 契约:run 内严格单调 +1。Runtime 与本管理器都会 + // 自产事件且各自计数(run_started 与 runtime 首事件会同时为 1),在 + // 协议唯一汇聚点统一重排,保证 replay 游标与幂等键(runId+sequence)。 + const stamped: AgentEvent = { ...event, sequence: (protocol.sequence += 1) }; for (const listener of this.listeners) { - listener(event); + listener(stamped); } - const mapped = toStreamEvents(event, { messageId: this.currentMessageId ?? undefined }); + const mapped = toStreamEvents(stamped, { messageId: protocol.messageId }); // 无论是否有实时订阅者,都先记录进内存历史,保证 replay 有数据源。 for (const streamEvent of mapped) { this.streamHistory.append(streamEvent); @@ -448,7 +469,7 @@ export class RunManager { if (this.streamListeners.size > 0) { for (const streamEvent of mapped) { for (const listener of this.streamListeners) { - listener(event.sessionId, streamEvent); + listener(stamped.sessionId, streamEvent); } } } diff --git a/packages/shared/src/kernel/stream-history.ts b/packages/shared/src/kernel/stream-history.ts index 559fb5e..508bcb4 100644 --- a/packages/shared/src/kernel/stream-history.ts +++ b/packages/shared/src/kernel/stream-history.ts @@ -53,6 +53,13 @@ export class StreamEventHistory { return { recoverable: false, events: [], atEnd: true }; } + const last = list[list.length - 1]; + // 游标超前(lastSequence 超过已知最大 sequence):客户端状态与历史 + // 分歧,不可能由正常事件流到达 —— 明确不可恢复,不能静默当作已同步。 + if (lastSequence > last.sequence) { + return { recoverable: false, events: [], atEnd: TERMINAL_TYPES.has(last.type) }; + } + const tail = list.filter((e) => e.sequence > lastSequence); const contiguous = tail.length === 0 || @@ -62,7 +69,6 @@ export class StreamEventHistory { return { recoverable: false, events: [], atEnd: false }; } - const last = list[list.length - 1]; return { recoverable: true, events: tail, From d0f160a12f21083927e95f7c2f6298eacbd6e1b7 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=9F=B3=E6=95=AC=E8=8D=A3?= <15920017+itjingrong@user.noreply.gitee.com> Date: Fri, 11 Sep 2026 16:36:40 +0800 Subject: [PATCH 10/14] =?UTF-8?q?test(shared,electron):=20kernel=20?= =?UTF-8?q?=E4=B8=8E=20app=20=E4=B8=A4=E7=BA=A7=20E2E=20=E9=AA=8C=E8=AF=81?= =?UTF-8?q?=E6=B5=81=E5=BC=8F=20replay?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit kernel 级(stream-replay.e2e.test.ts):真实 LocalRuntimeAdapter + RunManager + 持久化。六个用例:完整 run 后全量补发、无实时订阅者时 事件仍入历史、按 lastSequence 中途断线补发、取消路径产生可补发的显式 cancelled 事件、带工具 run 的全量流、未知 run 与伪造游标返回明确的 不可恢复路径。 app 级(e2e/stream-replay.mjs):真实 Electron + preload IPC(CDP)。 四个用例:replay(runId, 0) 与实时投递逐字节一致、断线补发拼接还原 完整流、未知 run 跨 IPC 返回不可恢复、非法 lastSequence 被拒以 INVALID_ARGUMENT。 --- apps/electron/e2e/stream-replay.mjs | 232 +++++++++++++++++ .../src/kernel/stream-replay.e2e.test.ts | 239 ++++++++++++++++++ 2 files changed, 471 insertions(+) create mode 100644 apps/electron/e2e/stream-replay.mjs create mode 100644 packages/shared/src/kernel/stream-replay.e2e.test.ts diff --git a/apps/electron/e2e/stream-replay.mjs b/apps/electron/e2e/stream-replay.mjs new file mode 100644 index 0000000..5bf9453 --- /dev/null +++ b/apps/electron/e2e/stream-replay.mjs @@ -0,0 +1,232 @@ +// Stream Event Protocol v1 — replay 的 app 级 E2E(issue #27)。 +// +// 走真实传输链路:Electron main(AgentKernelHost)→ preload(runs:stream-replay +// / agent:stream)→ renderer。与 kernel 级 E2E(stream-replay.e2e.test.ts)互补, +// 这里验证 IPC 边界上的 reconnect 契约: +// S1. 实时通道收到的事件与 replay(runId, 0) 补发完全一致(transport parity) +// S2. 中途"断线"(只收到前 2 条)→ 按 lastSequence 补发剩余段,拼接无缺口 +// S3. 未知 run → recoverable=false(明确不可恢复路径跨 IPC 成立) +// S4. 非法 lastSequence → INVALID_ARGUMENT(IPC 参数校验) +// +// node e2e/stream-replay.mjs (需先 build:preload / build:main / vite build) + +import { existsSync, mkdirSync, readFileSync, rmSync } from 'node:fs'; +import { dirname, join } from 'node:path'; +import { fileURLToPath } from 'node:url'; +import { createRequire } from 'node:module'; +import { reserveCdpPort, spawnElectron, waitForCdp } from './electron-harness.mjs'; +import { seedLocale } from './seed-locale.mjs'; + +// connectOverCDP 不需要本地浏览器,但 playwright 退出时会清点其浏览器缓存 +// 目录(ms-playwright),在受限环境/CI 里可能被拒。指到临时目录避免副作用。 +process.env.PLAYWRIGHT_BROWSERS_PATH = process.env.PLAYWRIGHT_BROWSERS_PATH ?? '0'; + +const require = createRequire(import.meta.url); +const { chromium } = require('playwright-core'); + +const here = dirname(fileURLToPath(import.meta.url)); +const appRoot = join(here, '..'); +const repoRoot = join(here, '../../..'); +const userDataDir = join(appRoot, 'e2e/.user-data-stream'); +// 日志进 gitignored 的 artifacts 目录(与其他 harness 用例一致)。 +const logPath = join(appRoot, 'e2e/artifacts/stream-replay-electron.log'); +mkdirSync(dirname(logPath), { recursive: true }); + +const KEEP_OPEN = process.env.FINAGENT_E2E_KEEP_OPEN === '1'; + +let failures = 0; +function pass(name) { + console.log(`PASS ${name}`); +} +function fail(name, error) { + failures += 1; + console.error(`FAIL ${name}`); + console.error(String(error?.stack ?? error).slice(0, 2000)); +} + +function assert(condition, message) { + if (!condition) throw new Error(message); +} + +async function main() { + for (const artifact of [ + join(appRoot, 'src/main/index.js'), + join(appRoot, 'src/preload/index.cjs'), + join(appRoot, 'dist/renderer/index.html'), + ]) { + if (!existsSync(artifact)) { + throw new Error(`Missing build artifact: ${artifact} (run build:preload, build:main, vite build)`); + } + } + + // Deterministic: fresh userData, seeded locale. + rmSync(userDataDir, { recursive: true, force: true }); + mkdirSync(userDataDir, { recursive: true }); + seedLocale(userDataDir, 'en-US'); + + const port = await reserveCdpPort(); + const { proc } = spawnElectron({ appRoot, repoRoot, port, userDataDir, logPath }); + // Windows 上 stale 进程不会自动清理;结束时主动 kill。 + const cdpUrl = `http://127.0.0.1:${port}`; + + let browser; + try { + await waitForCdp({ url: cdpUrl, timeoutMs: 60_000, proc, logPath }); + browser = await chromium.connectOverCDP(cdpUrl, { timeout: 30_000 }); + const context = browser.contexts()[0]; + const page = await waitForPage(context, 30_000); + await page.waitForLoadState('domcontentloaded'); + + // Renderer 内安装协议事件采集器(真实 onStreamEvent 通道)。 + await page.evaluate(() => { + window.__streamEvents = []; + window.__unsubStream = window.electronAPI.kernel.onStreamEvent((payload) => { + window.__streamEvents.push(payload); + }); + }); + + // 建 session + 跑一次确定性 run(unsupported 意图:无工具、不触网络)。 + const session = await page.evaluate(async () => { + const result = await window.electronAPI.kernel.createSession('stream replay e2e'); + if (!result.ok) throw new Error(JSON.stringify(result.error)); + return result.data; + }); + assert(session && typeof session.id === 'string', 'createSession did not return an id'); + + const run = await page.evaluate(async (sessionId) => { + const result = await window.electronAPI.kernel.startRun({ + sessionId, + content: '你好,随便聊聊', + }); + if (!result.ok) throw new Error(JSON.stringify(result.error)); + return result.data; + }, session.id); + assert(run && typeof run.id === 'string', 'startRun did not return a run id'); + + // 等 run 终结(run_completed / cancelled / error 任一)。 + await page.waitForFunction( + (runId) => { + const events = window.__streamEvents ?? []; + const mine = events.filter((e) => e.event.runId === runId); + const last = mine[mine.length - 1]; + return Boolean( + last && ['run_completed', 'cancelled', 'error'].includes(last.event.type) + ); + }, + run.id, + { timeout: 30_000 } + ); + + const live = await page.evaluate(() => (window.__streamEvents ?? []).map((e) => e.event)); + assert(live.length >= 5, `expected >= 5 live stream events, got ${live.length}`); + const sequences = live.map((e) => e.sequence); + assert( + sequences.every((seq, i) => seq === i + 1), + `live sequence not strictly 1..N: ${JSON.stringify(sequences)}` + ); + + // S1. 实时事件与 replay(runId, 0) 全量补发一致。 + try { + const replayAll = await page.evaluate(async (input) => { + const result = await window.electronAPI.kernel.streamReplay(input); + if (!result.ok) throw new Error(JSON.stringify(result.error)); + return result.data; + }, { runId: run.id, lastSequence: 0 }); + assert(replayAll.recoverable === true, 'replay(0) not recoverable'); + assert(replayAll.atEnd === true, 'replay(0) atEnd should be true after a completed run'); + assert( + JSON.stringify(replayAll.events) === JSON.stringify(live), + `replayed events differ from live delivery:\nlive=${JSON.stringify(live.map(e => [e.sequence, e.type]))}\nreplay=${JSON.stringify(replayAll.events.map(e => [e.sequence, e.type]))}` + ); + pass('S1: replay(runId, 0) over IPC returns exactly the live-delivered events'); + } catch (error) { + fail('S1: replay(runId, 0) over IPC returns exactly the live-delivered events', error); + } + + // S2. 中途断线:客户端只收到前 2 条,按其 lastSequence 补发剩余段。 + try { + const received = live.slice(0, 2); + const lastSequence = received[received.length - 1].sequence; + const tail = await page.evaluate(async (input) => { + const result = await window.electronAPI.kernel.streamReplay(input); + if (!result.ok) throw new Error(JSON.stringify(result.error)); + return result.data; + }, { runId: run.id, lastSequence }); + assert(tail.recoverable === true, 'tail replay not recoverable'); + assert(tail.events.length === live.length - received.length, 'tail length mismatch'); + assert( + JSON.stringify([...received, ...tail.events]) === JSON.stringify(live), + 'received + replayed tail does not reconstruct the full stream' + ); + assert(tail.atEnd === true, 'tail atEnd should be true'); + pass('S2: mid-run disconnect catch-up (lastSequence) reconstructs the full stream'); + } catch (error) { + fail('S2: mid-run disconnect catch-up (lastSequence) reconstructs the full stream', error); + } + + // S3. 未知 run → 明确不可恢复。 + try { + const unknown = await page.evaluate(async () => { + return window.electronAPI.kernel.streamReplay({ runId: 'never-seen-run', lastSequence: 0 }); + }); + assert(unknown.ok === true, 'streamReplay IPC itself should succeed'); + assert(unknown.data.recoverable === false, 'unknown run must be recoverable=false'); + assert(unknown.data.events.length === 0, 'unknown run must replay nothing'); + pass('S3: unknown run reports the explicit unrecoverable path'); + } catch (error) { + fail('S3: unknown run reports the explicit unrecoverable path', error); + } + + // S4. 非法 lastSequence → INVALID_ARGUMENT。 + try { + const invalid = await page.evaluate(async () => { + return window.electronAPI.kernel.streamReplay({ runId: 'some-run', lastSequence: 'not-a-number' }); + }); + assert(invalid.ok === false, 'invalid lastSequence should fail'); + assert( + invalid.error && invalid.error.code === 'INVALID_ARGUMENT', + `expected INVALID_ARGUMENT, got ${JSON.stringify(invalid.error)}` + ); + pass('S4: invalid lastSequence is rejected with INVALID_ARGUMENT'); + } catch (error) { + fail('S4: invalid lastSequence is rejected with INVALID_ARGUMENT', error); + } + + await page.evaluate(() => window.__unsubStream?.()); + } catch (error) { + // 启动/装配失败:把 electron 日志尾部带出来辅助定位。 + try { + const tail = readFileSync(logPath, 'utf8').slice(-1500); + console.error('--- electron log tail ---\n' + tail); + } catch { + // log 可能不存在。 + } + fail('harness setup', error); + } finally { + await browser?.close().catch(() => undefined); + if (!KEEP_OPEN) { + proc.kill(); + } else { + console.log(`KEEP_OPEN CDP port ${port} — clean up yourself.`); + } + } + + if (failures > 0) { + console.error(`stream-replay E2E failed: ${failures} step(s).`); + process.exit(1); + } + console.log('stream-replay E2E passed.'); + process.exit(0); +} + +async function waitForPage(context, timeoutMs) { + const deadline = Date.now() + timeoutMs; + while (Date.now() < deadline) { + const pages = context.pages(); + if (pages.length > 0) return pages[pages.length - 1]; + await new Promise((resolve) => setTimeout(resolve, 250)); + } + throw new Error('No renderer page appeared in time'); +} + +main(); diff --git a/packages/shared/src/kernel/stream-replay.e2e.test.ts b/packages/shared/src/kernel/stream-replay.e2e.test.ts new file mode 100644 index 0000000..1c8b251 --- /dev/null +++ b/packages/shared/src/kernel/stream-replay.e2e.test.ts @@ -0,0 +1,239 @@ +// Stream Event Protocol v1 — replay 全链路 E2E(issue #27)。 +// +// 与 stream-history.test.ts(纯内存历史单测)不同,这里走真实运行链路: +// RunManager + LocalRuntimeAdapter(FINAGENT_AGENT_PROVIDER=local 时 E2E +// 使用的同一 runtime)+ 真实持久化仓储,验证 ADR 0001 的 replay 契约: +// 1. sequence 每 run 严格单调 +1(所有生产者共同遵守的协议契约); +// 2. 无实时订阅者时事件仍入历史(reconnect 的数据源保证); +// 3. 中途“断线”后按 lastSequence 补发,拼接结果与全量一致; +// 4. 取消路径产生显式 cancelled 事件且可补发(atEnd=true); +// 5. message/run 双身份:message 级事件带 messageId,run 级不带(#34)。 + +import { mkdtemp, rm } from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import { afterEach, beforeEach, describe, expect, it } from 'bun:test'; +import type { StreamEvent, ToolDefinition } from '@finagent/core'; +import { LocalRuntimeAdapter } from '../agent/local-runtime-adapter.ts'; +import { JsonFileStore } from '../storage/json-file-store.ts'; +import { MessageRepository } from '../storage/message-repository.ts'; +import { RunRepository } from '../storage/run-repository.ts'; +import { SessionRepository } from '../storage/session-repository.ts'; +import { SessionManager } from './session-manager.ts'; +import { RunManager } from './run-manager.ts'; + +let dir = ''; +let clock = 1000; + +beforeEach(async () => { + dir = await mkdtemp(join(tmpdir(), 'finagent-stream-e2e-')); + clock = 1000; +}); + +afterEach(async () => { + await rm(dir, { recursive: true, force: true }); +}); + +/** 全链路装配:真实 SessionManager/RunManager + LocalRuntimeAdapter。 */ +function makeStack(runtime = new LocalRuntimeAdapter({ now: () => clock })) { + const store = new JsonFileStore(dir); + const sessions = new SessionManager({ + sessions: new SessionRepository(store), + messages: new MessageRepository(store), + runs: new RunRepository(store), + piSessionDir: join(dir, 'pi-sessions'), + now: () => clock, + }); + const runs = new RunManager({ sessions, runs: new RunRepository(store), runtime, now: () => clock }); + return { sessions, runs }; +} + +/** 可控 registry:execute 由用例注入行为(默认立即成功),不触网络。 */ +function fakeRegistry(behavior: (input: { name: string; args: Record }) => Promise) { + return { + getTools: (): ToolDefinition[] => [], + execute: async (input: { name: string; args: Record }) => { + const details = await behavior(input); + return { + content: [{ type: 'text' as const, text: JSON.stringify(details) }], + details, + provenance: { providerId: 'fake', retrievedAt: clock }, + }; + }, + } as never; +} + +async function waitFor(predicate: () => Promise, timeoutMs = 2000) { + const started = Date.now(); + while (!(await predicate())) { + if (Date.now() - started > timeoutMs) { + throw new Error('waitFor timed out'); + } + await new Promise((resolve) => setTimeout(resolve, 5)); + } +} + +const MESSAGE_LEVEL_TYPES = new Set(['message_started', 'text_delta', 'message_completed', 'cancelled']); + +/** ADR 0001:每 run sequence 从 1 开始严格 +1,所有事件共同遵守。 */ +function expectStrictSequence(events: StreamEvent[]) { + expect(events.map((e) => e.sequence)).toEqual(events.map((_, i) => i + 1)); +} + +/** #34 身份契约:message 级事件必带 messageId,run 级不带。 */ +function expectIdentityContract(events: StreamEvent[]) { + for (const event of events) { + if (MESSAGE_LEVEL_TYPES.has(event.type)) { + expect(typeof event.messageId).toBe('string'); + } else { + expect(event.messageId).toBeUndefined(); + } + } +} + +describe('Stream Event replay E2E(真实 runtime 全链路)', () => { + it('完整 run 后:replay(runId, 0) 补发全量连续事件,atEnd=true,身份契约成立', async () => { + const { sessions, runs } = makeStack(); + const session = await sessions.createSession('E2E replay'); + const run = await runs.startRun(session.id, '你好,随便聊聊'); // unsupported 意图:无 tool,确定性 + await waitFor(async () => !runs.isRunning()); + + const replay = runs.replayStream(run.id, 0); + expect(replay.recoverable).toBe(true); + expect(replay.atEnd).toBe(true); + expectStrictSequence(replay.events); + expectIdentityContract(replay.events); + expect(replay.events.map((e) => e.type)).toEqual([ + 'run_started', + 'message_started', + 'text_delta', + 'message_completed', + 'run_completed', + ]); + // text_delta 载荷为增量文本,拼接后与落库答案一致。 + const text = replay.events + .filter((e): e is Extract => e.type === 'text_delta') + .map((e) => e.payload.text) + .join(''); + const messages = await sessions.listMessages(session.id); + expect(text).toBe(messages[1]?.content); + }); + + it('断线(无实时订阅者):事件仍入历史,重连后全量补发', async () => { + const { sessions, runs } = makeStack(); + const session = await sessions.createSession('offline'); + // 注意:全程不调用 subscribeStream —— 模拟 renderer 掉线期间 run 照常执行。 + const run = await runs.startRun(session.id, '你好,随便聊聊'); + await waitFor(async () => !runs.isRunning()); + + const replay = runs.replayStream(run.id, 0); + expect(replay.recoverable).toBe(true); + expectStrictSequence(replay.events); + expect(replay.events.length).toBeGreaterThanOrEqual(5); + }); + + it('中途断线:按已收 lastSequence 补发剩余段,拼接与全量一致', async () => { + const { sessions, runs } = makeStack(); + const session = await sessions.createSession('mid-run disconnect'); + const received: StreamEvent[] = []; + // 模拟 renderer 订阅后中途掉线:只保留前 2 条。 + const unsubscribe = runs.subscribeStream((_sessionId, event) => { + if (received.length < 2) received.push(event); + }); + + const run = await runs.startRun(session.id, '你好,随便聊聊'); + await waitFor(async () => !runs.isRunning()); + unsubscribe(); + + expect(received.length).toBe(2); + const lastSequence = received[received.length - 1].sequence; + + // “重连”:先全量核对,再按 lastSequence 补发并拼接。 + const full = runs.replayStream(run.id, 0); + const tail = runs.replayStream(run.id, lastSequence); + expect(tail.recoverable).toBe(true); + expect(tail.events[0]?.sequence).toBe(lastSequence + 1); + expect([...received, ...tail.events]).toEqual(full.events); + expect(tail.atEnd).toBe(true); + }); + + it('取消路径:显式 cancelled 事件入历史并可补发,atEnd=true', async () => { + // 慢 registry:execute 挂起直到用例放行,保证 cancel 落在 run 执行窗口内。 + let release: (() => void) | undefined; + const gate = new Promise((resolve) => { + release = resolve; + }); + const runtime = new LocalRuntimeAdapter({ + now: () => clock, + registry: fakeRegistry(async () => { + await gate; + return { lastPrice: 220 }; + }), + }); + const { sessions, runs } = makeStack(runtime); + const session = await sessions.createSession('cancel'); + + const run = await runs.startRun(session.id, 'AAPL.US 行情'); // 路由到 get_quote → 慢 registry + await waitFor(async () => runs.isRunning()); + await runs.cancelRun(session.id, run.id); + release?.(); + await waitFor(async () => !runs.isRunning()); + + const persisted = await sessions.getRun(session.id, run.id); + expect(persisted?.status).toBe('cancelled'); + + const replay = runs.replayStream(run.id, 0); + expect(replay.recoverable).toBe(true); + expect(replay.atEnd).toBe(true); + expectStrictSequence(replay.events); + expectIdentityContract(replay.events); + const terminal = replay.events[replay.events.length - 1]; + expect(terminal?.type).toBe('cancelled'); + if (terminal?.type === 'cancelled') { + expect(terminal.payload.reason).toBe('user'); + } + }); + + it('带工具的 run:tool 事件入历史并可全量补发', async () => { + const runtime = new LocalRuntimeAdapter({ + now: () => clock, + registry: fakeRegistry(async () => ({ lastPrice: 220, changePercent: 1.5 })), + }); + const { sessions, runs } = makeStack(runtime); + const session = await sessions.createSession('with tools'); + + const run = await runs.startRun(session.id, 'AAPL.US 行情'); + await waitFor(async () => !runs.isRunning()); + + const replay = runs.replayStream(run.id, 0); + expect(replay.recoverable).toBe(true); + expectStrictSequence(replay.events); + expectIdentityContract(replay.events); + expect(replay.events.map((e) => e.type)).toEqual([ + 'run_started', + 'tool_started', + 'tool_result', + 'message_started', + 'text_delta', + 'message_completed', + 'run_completed', + ]); + const toolResult = replay.events.find( + (e): e is Extract => e.type === 'tool_result' + ); + expect(toolResult?.payload.name).toBe('get_quote'); + }); + + it('未知 run 与伪造游标:明确不可恢复', async () => { + const { sessions, runs } = makeStack(); + const session = await sessions.createSession('unknown'); + const run = await runs.startRun(session.id, '你好,随便聊聊'); + await waitFor(async () => !runs.isRunning()); + + expect(runs.replayStream('never-seen', 0)).toEqual({ recoverable: false, events: [], atEnd: true }); + // 断在 5,但该 run 实际只有 5 条事件(1..5)——lastSequence 超出末尾不算缺失。 + const beyond = runs.replayStream(run.id, 99); + expect(beyond.recoverable).toBe(false); + expect(beyond.events).toEqual([]); + }); +}); From fd89c5ad2c0b34dab5b4e320e5ef5437973ee0ea Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=9F=B3=E6=95=AC=E8=8D=A3?= <15920017+itjingrong@user.noreply.gitee.com> Date: Fri, 11 Sep 2026 16:57:00 +0800 Subject: [PATCH 11/14] fix(electron): stub subscribeStream/replayStream in kernelHost test fake AgentKernelHost.attach wires the Stream Event v1 channel through kernel.runs.subscribeStream; the fake kernel in kernelHost.test.ts lacked the method, so the transport test crashed with TypeError before asserting. Align the fake with the kernel surface (subscribeStream + replayStream) so focused CI runs green again. --- apps/electron/src/main/kernelHost.test.ts | 2 ++ 1 file changed, 2 insertions(+) diff --git a/apps/electron/src/main/kernelHost.test.ts b/apps/electron/src/main/kernelHost.test.ts index 1123e31..ce9dfe0 100644 --- a/apps/electron/src/main/kernelHost.test.ts +++ b/apps/electron/src/main/kernelHost.test.ts @@ -51,6 +51,7 @@ const fakeSessions = { const fakeRuns = { subscribe: (_listener: (event: AgentEvent) => void) => () => undefined, + subscribeStream: (_listener: (sessionId: string, event: unknown) => void) => () => undefined, startRun: async (sessionId: string, content: string) => ({ id: 'r1', sessionId, @@ -59,6 +60,7 @@ const fakeRuns = { startedAt: 1, }), cancelRun: async () => undefined, + replayStream: (_runId: string, _lastSequence: number) => ({ recoverable: false, events: [], atEnd: true }), }; class FakeAgentKernel { From 8d8ef335e70312853c9ae2c34ead8e5476fdf569 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=9F=B3=E6=95=AC=E8=8D=A3?= <15920017+itjingrong@user.noreply.gitee.com> Date: Fri, 11 Sep 2026 17:00:00 +0800 Subject: [PATCH 12/14] fix(ui): complete kernel channel stub in sessionAtoms test The Stream Event v1 IPC surface (streamReplay / onStreamEvent) landed in the kernel channel types; the test kernel client was not updated, which broke the ui + i18n + electron typecheck gates. Align the stub with the channel so typecheck is green again. --- packages/ui/src/atoms/sessionAtoms.test.ts | 2 ++ 1 file changed, 2 insertions(+) diff --git a/packages/ui/src/atoms/sessionAtoms.test.ts b/packages/ui/src/atoms/sessionAtoms.test.ts index 82b77c3..ff25dcd 100644 --- a/packages/ui/src/atoms/sessionAtoms.test.ts +++ b/packages/ui/src/atoms/sessionAtoms.test.ts @@ -48,7 +48,9 @@ function makeClient(): FinagentClient { listRuns: async () => ({ ok: true as const, data: [] as Run[] }), startRun: async () => ({ ok: false as const, error: { code: 'TEST', message: 'no-op' } }), cancelRun: async () => ({ ok: true as const, data: undefined }), + streamReplay: async () => ({ ok: true as const, data: { recoverable: false, events: [], atEnd: true } }), onAgentEvent: () => () => undefined, + onStreamEvent: () => () => undefined, }, agent: { getTools: async () => ({ ok: true as const, data: [] }), From aca8ba7f1ee204e99f7f48157349b5ce28abcf81 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=9F=B3=E6=95=AC=E8=8D=A3?= <15920017+itjingrong@user.noreply.gitee.com> Date: Fri, 11 Sep 2026 17:05:39 +0800 Subject: [PATCH 13/14] fix(shared,electron): export StreamReplayResult at the package boundary streamReplay returned an inferred StreamReplayResult that TS could not name portably across the @finagent/shared boundary (error TS2742). Export the type from the shared package and annotate the kernelHost surface so every package typecheck gate passes. --- apps/electron/src/main/kernelHost.ts | 3 ++- packages/shared/src/index.ts | 2 ++ packages/shared/src/kernel/index.ts | 1 + 3 files changed, 5 insertions(+), 1 deletion(-) diff --git a/apps/electron/src/main/kernelHost.ts b/apps/electron/src/main/kernelHost.ts index 55235bb..3e3085b 100644 --- a/apps/electron/src/main/kernelHost.ts +++ b/apps/electron/src/main/kernelHost.ts @@ -154,6 +154,7 @@ import { type BriefPortfolioSummary, type MarketPulseSnapshot, type ShareCard, + type StreamReplayResult, type WatchlistQuote, withDemoDataFallback, } from '@finagent/shared'; @@ -573,7 +574,7 @@ export class AgentKernelHost { } /** Stream Event replay(ADR 0001 §Reconnect):按 lastSequence 补发或明确不可恢复。 */ - streamReplay(input: unknown) { + streamReplay(input: unknown): StreamReplayResult { const request = requireObject(input); const lastSequence = request.lastSequence; if (typeof lastSequence !== 'number' || !Number.isFinite(lastSequence)) { diff --git a/packages/shared/src/index.ts b/packages/shared/src/index.ts index 11cb604..f7971e9 100644 --- a/packages/shared/src/index.ts +++ b/packages/shared/src/index.ts @@ -87,6 +87,8 @@ export { AgentKernel, SessionManager, RunManager, + StreamEventHistory, + type StreamReplayResult, BUDGET_KEYS, addUsage, budgetStop, diff --git a/packages/shared/src/kernel/index.ts b/packages/shared/src/kernel/index.ts index ac6f223..4f8bc33 100644 --- a/packages/shared/src/kernel/index.ts +++ b/packages/shared/src/kernel/index.ts @@ -1,6 +1,7 @@ export { AgentKernel, type AgentKernelOptions, type AgentProvider } from './agent-kernel.ts'; export { SessionManager, type SessionManagerOptions } from './session-manager.ts'; export { RunManager, type RunManagerOptions } from './run-manager.ts'; +export { StreamEventHistory, type StreamReplayResult } from './stream-history.ts'; export { BUDGET_KEYS, addUsage, From eb75fc76a2cb567ba7f855b8c359311d1d58d8e4 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=9F=B3=E6=95=AC=E8=8D=A3?= <15920017+itjingrong@user.noreply.gitee.com> Date: Sun, 13 Sep 2026 19:43:31 +0800 Subject: [PATCH 14/14] =?UTF-8?q?fix(shared,electron):=20=E6=94=B6?= =?UTF-8?q?=E6=95=9B=E6=B5=81=E5=8D=8F=E8=AE=AE=E5=8F=96=E6=B6=88=E8=AF=AD?= =?UTF-8?q?=E4=B9=89=E4=B8=8E=E6=B8=B8=E6=A0=87=E6=A0=A1=E9=AA=8C?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - stream-event-adapter:run_failed 按 error.code 归一为 cancelled(user/budget/runtime), 并把已生成的 partial 文本带进 cancelled.partial.text(对齐 ADR 0001) - run-manager:run 全程保留 assistant 文本快照,供 cancelled 事件使用 - kernelHost:streamReplay 游标限定为非负整数;dispose 时一并清理 stream 订阅 - ui/streamAtoms:幂等去重仅在 sequence 回退时扫描,避免长 run 下的 O(n²) - 补齐 adapter / streamAtoms / kernelHost 边界测试,electron E2E 增加非法游标用例 --- apps/electron/e2e/stream-replay.mjs | 23 ++++++----- apps/electron/src/main/kernelHost.test.ts | 23 +++++++++++ apps/electron/src/main/kernelHost.ts | 7 +++- docs/adr/0001-stream-event-protocol.md | 3 +- packages/shared/src/kernel/run-manager.ts | 15 ++++++-- .../src/kernel/stream-event-adapter.test.ts | 38 +++++++++++++++++-- .../shared/src/kernel/stream-event-adapter.ts | 31 ++++++++++++--- packages/ui/src/atoms/streamAtoms.test.ts | 14 +++++++ packages/ui/src/atoms/streamAtoms.ts | 7 ++-- 9 files changed, 133 insertions(+), 28 deletions(-) diff --git a/apps/electron/e2e/stream-replay.mjs b/apps/electron/e2e/stream-replay.mjs index 5bf9453..5c7f304 100644 --- a/apps/electron/e2e/stream-replay.mjs +++ b/apps/electron/e2e/stream-replay.mjs @@ -177,16 +177,21 @@ async function main() { fail('S3: unknown run reports the explicit unrecoverable path', error); } - // S4. 非法 lastSequence → INVALID_ARGUMENT。 + // S4. 非法 lastSequence(非整数 / 负数 / 非数字)→ INVALID_ARGUMENT。 try { - const invalid = await page.evaluate(async () => { - return window.electronAPI.kernel.streamReplay({ runId: 'some-run', lastSequence: 'not-a-number' }); - }); - assert(invalid.ok === false, 'invalid lastSequence should fail'); - assert( - invalid.error && invalid.error.code === 'INVALID_ARGUMENT', - `expected INVALID_ARGUMENT, got ${JSON.stringify(invalid.error)}` - ); + const cursors = ['not-a-number', -1, 1.5]; + for (const lastSequence of cursors) { + const invalid = await page.evaluate( + async (cursor) => + window.electronAPI.kernel.streamReplay({ runId: 'some-run', lastSequence: cursor }), + lastSequence + ); + assert(invalid.ok === false, `invalid lastSequence ${lastSequence} should fail`); + assert( + invalid.error && invalid.error.code === 'INVALID_ARGUMENT', + `expected INVALID_ARGUMENT for ${lastSequence}, got ${JSON.stringify(invalid.error)}` + ); + } pass('S4: invalid lastSequence is rejected with INVALID_ARGUMENT'); } catch (error) { fail('S4: invalid lastSequence is rejected with INVALID_ARGUMENT', error); diff --git a/apps/electron/src/main/kernelHost.test.ts b/apps/electron/src/main/kernelHost.test.ts index ce9dfe0..9967001 100644 --- a/apps/electron/src/main/kernelHost.test.ts +++ b/apps/electron/src/main/kernelHost.test.ts @@ -391,6 +391,29 @@ describe('AgentKernelHost', () => { host.dispose(); }); + it('rejects non-integer or negative stream replay cursors', () => { + const host = new AgentKernelHost(); + const expectInvalid = (input: unknown) => { + try { + host.streamReplay(input); + } catch (error) { + expect(error).toMatchObject({ code: 'INVALID_ARGUMENT' }); + return; + } + throw new Error('expected streamReplay to reject the cursor'); + }; + + for (const lastSequence of [-1, 1.5, Number.NaN, Number.POSITIVE_INFINITY, '3', undefined]) { + expectInvalid({ runId: 'r1', lastSequence }); + } + expect(host.streamReplay({ runId: 'r1', lastSequence: 0 })).toEqual({ + recoverable: false, + events: [], + atEnd: true, + }); + host.dispose(); + }); + it('rejects malformed run payloads', async () => { const host = new AgentKernelHost(); diff --git a/apps/electron/src/main/kernelHost.ts b/apps/electron/src/main/kernelHost.ts index 3e3085b..07006e3 100644 --- a/apps/electron/src/main/kernelHost.ts +++ b/apps/electron/src/main/kernelHost.ts @@ -577,8 +577,9 @@ export class AgentKernelHost { streamReplay(input: unknown): StreamReplayResult { const request = requireObject(input); const lastSequence = request.lastSequence; - if (typeof lastSequence !== 'number' || !Number.isFinite(lastSequence)) { - throw createCodeError('INVALID_ARGUMENT', 'lastSequence must be a finite number.'); + // 游标必须是非负整数:0 表示从头补发,负数/小数/NaN 都是非法客户端状态。 + if (typeof lastSequence !== 'number' || !Number.isInteger(lastSequence) || lastSequence < 0) { + throw createCodeError('INVALID_ARGUMENT', 'lastSequence must be a non-negative integer.'); } return this.kernel.runs.replayStream( requireString(request.runId, 'runId'), @@ -2571,6 +2572,8 @@ export class AgentKernelHost { this.unsubscribe = null; this.unsubscribeEval?.(); this.unsubscribeEval = null; + this.streamUnsubscribe?.(); + this.streamUnsubscribe = null; this.window = null; await this.kernel.dispose(); } diff --git a/docs/adr/0001-stream-event-protocol.md b/docs/adr/0001-stream-event-protocol.md index ac73c18..31286d0 100644 --- a/docs/adr/0001-stream-event-protocol.md +++ b/docs/adr/0001-stream-event-protocol.md @@ -63,7 +63,8 @@ Key behavior change: `message_delta` (full-answer snapshot) becomes `text_delta` | Capability | Contract | |---|---| | Idempotency | Consumers dedupe on `runId + sequence`; replayed events never re-insert text, tool cards, or citations | -| Cancel | renderer `cancelRun` → runtime propagates → runtime **explicitly emits `cancelled`**; partial answer preserved | +| Cancel | renderer `cancelRun` → runtime propagates → runtime **explicitly emits `cancelled`**; partial answer preserved (`payload.partial.text` carries the text produced so far) | +| `run_failed` normalization | `RUN_CANCELLED` → `cancelled{reason:'user'}`; `BUDGET_EXHAUSTED` → `cancelled{reason:'budget'}`; `RETRY_STORM` / `LOOP_DETECTED` → `cancelled{reason:'runtime'}`; everything else → `error` (see `stream-event-adapter`) | | Reconnect | **Implemented (v1, in-memory):** `RunManager.replayStream(runId, lastSequence)` returns the contiguous tail from `StreamEventHistory` (`packages/shared/src/kernel/stream-history.ts`); when the run is unknown or the tail is non-contiguous (buffer eviction), it returns `recoverable: false` — an explicit unrecoverable path the renderer must surface. IPC: `runs:stream-replay` | | Final-state parity | UI `stopReason` and persisted `run.status` come from the same single final state in run-manager (aligns with #18) | | Security (#19) | `status` never exposes chain-of-thought; tool payloads pass redaction before reaching the UI; renderer never executes model-returned code | diff --git a/packages/shared/src/kernel/run-manager.ts b/packages/shared/src/kernel/run-manager.ts index 551924b..cdb4c8b 100644 --- a/packages/shared/src/kernel/run-manager.ts +++ b/packages/shared/src/kernel/run-manager.ts @@ -82,6 +82,8 @@ interface ActiveRun { interface RunProtocol { sequence: number; messageId: string; + /** 已生成的 assistant 文本快照,cancelled.partial.text 的数据源。 */ + text: string; } /** @@ -217,7 +219,7 @@ export class RunManager { }; // 一次 run = 一次 assistant generation(#34):run 启动时确定稳定的 // assistant message id 与协议 sequence 计数器,随 run 全程传递。 - const protocol: RunProtocol = { sequence: 0, messageId: randomUUID() }; + const protocol: RunProtocol = { sequence: 0, messageId: randomUUID(), text: '' }; this.emit( { id: randomUUID(), @@ -250,7 +252,7 @@ export class RunManager { session: SessionMeta, workspaceContext?: WorkspaceContext, locale?: SupportedLocale, - protocol: RunProtocol = { sequence: 0, messageId: randomUUID() } + protocol: RunProtocol = { sequence: 0, messageId: randomUUID(), text: '' } ): Promise { let failure: ApiError | undefined; let answer = ''; @@ -275,6 +277,8 @@ export class RunManager { this.emit(event, protocol); if (event.type === 'message_delta' || event.type === 'message_completed') { answer = event.payload.answer; + // 保留最新文本快照,使 cancelled 事件能带出部分回答(ADR 0001)。 + protocol.text = answer; } else if (event.type === 'tool_completed') { toolCalls.push(event.payload.toolCall); } else if (event.type === 'run_failed') { @@ -366,7 +370,7 @@ export class RunManager { } } -/** + /** * Account for one runtime event and decide whether the run must stop. A stop * requests cancellation, so the caller stops consuming events and the run * settles as `cancelled` — never as an ordinary success — carrying its partial @@ -461,7 +465,10 @@ export class RunManager { for (const listener of this.listeners) { listener(stamped); } - const mapped = toStreamEvents(stamped, { messageId: protocol.messageId }); + const mapped = toStreamEvents(stamped, { + messageId: protocol.messageId, + partialText: protocol.text, + }); // 无论是否有实时订阅者,都先记录进内存历史,保证 replay 有数据源。 for (const streamEvent of mapped) { this.streamHistory.append(streamEvent); diff --git a/packages/shared/src/kernel/stream-event-adapter.test.ts b/packages/shared/src/kernel/stream-event-adapter.test.ts index f497bd2..d8d0e08 100644 --- a/packages/shared/src/kernel/stream-event-adapter.test.ts +++ b/packages/shared/src/kernel/stream-event-adapter.test.ts @@ -105,14 +105,46 @@ describe('toStreamEvents', () => { expect(done.runId).toBe('run-1'); }); - it('用户取消映射为 cancelled(reason=user)', () => { + it('用户取消映射为 cancelled(reason=user,带出部分文本与 messageId)', () => { const [ev] = toStreamEvents( - makeEvent('run_failed', { error: { code: 'RUN_CANCELLED', message: 'Run cancelled by user.' } }) + makeEvent('run_failed', { error: { code: 'RUN_CANCELLED', message: 'Run cancelled by user.' } }), + { messageId: 'msg-9', partialText: 'Apple 已回复一半' } ); expect(ev.type).toBe('cancelled'); if (ev.type === 'cancelled') { expect(ev.payload.reason).toBe('user'); - expect(ev.payload.partial).toEqual({ text: '' }); + expect(ev.payload.partial).toEqual({ text: 'Apple 已回复一半' }); + expect(ev.messageId).toBe('msg-9'); + } + }); + + it('预算/runaway 终止映射为 cancelled(reason=budget / runtime)', () => { + const [budget] = toStreamEvents( + makeEvent('run_failed', { error: { code: 'BUDGET_EXHAUSTED', message: 'Run stopped.' } }) + ); + const [loop] = toStreamEvents( + makeEvent('run_failed', { error: { code: 'LOOP_DETECTED', message: 'Run stopped.' } }) + ); + expect(budget.type).toBe('cancelled'); + if (budget.type === 'cancelled') expect(budget.payload.reason).toBe('budget'); + expect(loop.type).toBe('cancelled'); + if (loop.type === 'cancelled') expect(loop.payload.reason).toBe('runtime'); + }); + + it('run 级事件(tool_*)不携带 messageId', () => { + const toolCall = { + id: 'tc-1', + toolName: 'get_quote', + args: { symbol: 'AAPL.US' }, + startedAt: 1, + status: 'success' as const, + result: { lastPrice: 220 }, + }; + const [result] = toStreamEvents(makeEvent('tool_completed', { toolCall }), { + messageId: 'msg-9', + }); + if (result.type === 'tool_result') { + expect(result.messageId).toBeUndefined(); } }); }); \ No newline at end of file diff --git a/packages/shared/src/kernel/stream-event-adapter.ts b/packages/shared/src/kernel/stream-event-adapter.ts index 30947ad..6fa4713 100644 --- a/packages/shared/src/kernel/stream-event-adapter.ts +++ b/packages/shared/src/kernel/stream-event-adapter.ts @@ -2,20 +2,38 @@ // 供 RunManager 在现有 AgentEvent 广播旁并行产出协议事件(issue #27, // docs/adr/0001-stream-event-protocol.md,Migration step 2)。 -import type { AgentEvent, StreamEvent } from '@finagent/core'; +import type { AgentEvent, StreamCancelReason, StreamEvent } from '@finagent/core'; import { STREAM_EVENT_PROTOCOL_VERSION } from '@finagent/core'; +/** + * run_failed 的 error.code → 取消原因(ADR 0001 §Cross-cutting contracts)。 + * RUN_CANCELLED 为用户主动取消;BUDGET_EXHAUSTED 为预算护栏终止; + * RETRY_STORM / LOOP_DETECTED 为 runaway 检测终止。三者都是"被主动终止", + * 归一为 cancelled;其余失败才归一为 error。 + */ +const CANCEL_REASON_BY_CODE: Record = { + RUN_CANCELLED: 'user', + BUDGET_EXHAUSTED: 'budget', + RETRY_STORM: 'runtime', + LOOP_DETECTED: 'runtime', +}; + /** * 把单个 AgentEvent 映射为一条或多条 StreamEvent。 * - 时间戳:AgentEvent 用 epoch 毫秒 number,协议层用 ISO 8601 UTC string。 * - 身份(对齐 issue #34):runId 恒有;messageId 仅注入到 message 级事件 * (message_started / text_delta / message_completed / cancelled), * 取 run 对应的真实 assistant message id;run 级事件不携带 messageId。 - * - run_failed(code=RUN_CANCELLED) 归一为 cancelled 事件;其余失败归一为 error。 + * - run_failed 按 error.code 归一:主动取消 → cancelled(带 partial 文本), + * 其余失败归一为 error。 */ export function toStreamEvents( event: AgentEvent, - opts?: { messageId?: string } + opts?: { + messageId?: string; + /** 取消时已生成的 assistant 文本(cancelled.partial.text 的数据源)。 */ + partialText?: string; + } ): StreamEvent[] { const base = { protocolVersion: STREAM_EVENT_PROTOCOL_VERSION, @@ -71,13 +89,14 @@ export function toStreamEvents( return [{ ...base, type: 'run_completed', payload: { stopReason: 'completed' } }]; case 'run_failed': { const { error } = event.payload; - if (error.code === 'RUN_CANCELLED') { + const reason = CANCEL_REASON_BY_CODE[error.code]; + if (reason) { return [ { ...base, ...messageish, type: 'cancelled', - payload: { reason: 'user', partial: { text: '' } }, + payload: { reason, partial: { text: opts?.partialText ?? '' } }, }, ]; } @@ -86,4 +105,4 @@ export function toStreamEvents( ]; } } -} \ No newline at end of file +} diff --git a/packages/ui/src/atoms/streamAtoms.test.ts b/packages/ui/src/atoms/streamAtoms.test.ts index 2917d14..f356b35 100644 --- a/packages/ui/src/atoms/streamAtoms.test.ts +++ b/packages/ui/src/atoms/streamAtoms.test.ts @@ -42,6 +42,20 @@ describe('reduceStreamLog', () => { expect(s2.anomalies).toBe(2); }); + it('重复 sequence 落在序列中间也幂等丢弃', () => { + let state = s0; + for (const sequence of [1, 2, 3]) { + state = reduceStreamLog(state, { sessionId: 's1', event: make({ sequence, type: 'text_delta' }) }); + } + const replayed = reduceStreamLog(state, { + sessionId: 's1', + event: make({ sequence: 2, type: 'text_delta' }), + }); + expect(replayed.byRun.get('run-1')?.map((e) => e.sequence)).toEqual([1, 2, 3]); + expect(replayed.drops).toBe(1); + expect(replayed.anomalies).toBe(0); + }); + it('不同 run 互不干扰', () => { const s1 = reduceStreamLog(s0, { sessionId: 's1', event: make({ runId: 'run-a', sequence: 1, type: 'run_started' }) }); const s2 = reduceStreamLog(s1, { sessionId: 's1', event: make({ runId: 'run-b', sequence: 1, type: 'run_started' }) }); diff --git a/packages/ui/src/atoms/streamAtoms.ts b/packages/ui/src/atoms/streamAtoms.ts index 084f32a..291842f 100644 --- a/packages/ui/src/atoms/streamAtoms.ts +++ b/packages/ui/src/atoms/streamAtoms.ts @@ -29,12 +29,13 @@ export function reduceStreamLog(state: StreamLogState, input: StreamLogInput): S const { event } = input; const prev = state.byRun.get(event.runId) ?? []; - // 幂等:同一 run 内相同 sequence 的事件直接丢弃(重放/重复投递)。 - if (prev.some((e) => e.sequence === event.sequence)) { + const lastSeq = prev.length > 0 ? prev[prev.length - 1].sequence : 0; + // 幂等:prev 按 sequence 有序,重复只可能落在已见过的区间内,因此仅在 + // sequence <= lastSeq 时做一次扫描,避免每条事件都全量遍历(长 run 下的 O(n²))。 + if (event.sequence <= lastSeq && prev.some((e) => e.sequence === event.sequence)) { return { ...state, drops: state.drops + 1, last: input }; } - const lastSeq = prev.length > 0 ? prev[prev.length - 1].sequence : 0; const isAnomaly = event.sequence <= lastSeq || event.sequence !== lastSeq + 1; const next = isAnomaly ? [...prev, event].sort((a, b) => a.sequence - b.sequence)