From 48d8d9ae7c10596c5be88bcf236c7c6e29dbf86d 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 1/3] 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 4edc88d..3236fc6 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 9aaf6d126df87f94b91d2d4e3ffaeb19efc5a2b9 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 2/3] 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 0bf664c..81edcb8 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 RunManagerTimer { cancel(): void; @@ -107,6 +109,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) { @@ -125,6 +128,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; @@ -477,6 +490,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 86272e127f889834dbf77b439533d1d7ca92a641 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 3/3] 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') {