/** * v0.8.0 P0-5 回归测试(第二批):流可靠性四件套的剩余项。 * * 1. parseSSEStream 集成级空闲超时:挂死 ReadableStream → 60s 无数据 → * SseUpstreamError(504)(P1-2 防线的集成真实性,此前仅纯函数层覆盖); * 2. SseUpstreamError(504) 经引擎 isRetryableError 归类 → 退避重试后成功; * 3. REASONING_DELTA 持续输出中 abort → USER_INTERRUPT("思考中点停止"场景); * 4. 引擎 P4-2:上一 run 卡死(不响应 abort)→ 新 run 等待 30s → 强制 abort * → 再等 5s → 抛错(防"会话永久排队")。 */ import { describe, it, expect, vi, afterEach } from 'vitest'; import { parseSSEStream, SseUpstreamError } from '../../adapters/shared/sse-stream'; import { AgentLoopEngine } from '../engine'; import type { MetonaRequest, MetonaStreamEvent, IMetonaProviderAdapter } from '../../types'; import { MetonaStreamEventType } from '../../types'; const SYSTEM_PROMPT = { roleDefinition: 'test', outputConstraints: '', safetyGuidelines: '', }; afterEach(() => { vi.useRealTimers(); }); // ===== 1. 集成级空闲超时 ===== describe('parseSSEStream integrated idle timeout (P0-5 #2)', () => { function makeHungBody(): ReadableStream { return { getReader: () => ({ // 永不 resolve 的 read —— 模拟"服务器保活但不再推数据" read: () => new Promise(() => {}), releaseLock: () => {}, }), } as unknown as ReadableStream; } it('挂死流 60s 无数据 → SseUpstreamError(504)', async () => { vi.useFakeTimers(); const iterator = parseSSEStream(makeHungBody(), 'r', 's', 1); let thrown: unknown = null; const firstRead = iterator.next().catch((e: unknown) => { thrown = e; return { done: true, value: undefined } as const; }); // 推进 61s:IDLE_TIMEOUT_MS=60_000 触发 idle controller → read() 竞速失败 await vi.advanceTimersByTimeAsync(61_000); await firstRead; expect(thrown).toBeInstanceOf(SseUpstreamError); expect((thrown as SseUpstreamError).status).toBe(504); expect((thrown as Error).message).toContain('idle timeout'); }); }); // ===== 2. SseUpstreamError(504) → 引擎重试 ===== describe('engine: SseUpstreamError(504) retry classification (P0-5 #3)', () => { it('504 首败 → 退避重试 → 第二次流成功 → COMPLETED', async () => { let call = 0; const adapter = { providerId: 'mock' as const, supportedModels: ['m'], supportsToolCalling: true, supportsThinking: true, getContextWindow: () => 128_000, async send(): Promise { throw new Error('not used'); }, async *sendStream(_r: MetonaRequest): AsyncIterable { call++; if (call === 1) { throw new SseUpstreamError('Stream idle timeout after 60000ms (no data received)', { status: 504, }); } yield { type: MetonaStreamEventType.TEXT_DELTA, requestId: 'r', sessionId: 's', iteration: 1, seq: 0, timestamp: Date.now(), delta: 'recovered', }; yield { type: MetonaStreamEventType.DONE, requestId: 'r', sessionId: 's', iteration: 1, seq: 1, timestamp: Date.now(), finishReason: 'stop', }; }, } as unknown as IMetonaProviderAdapter; const engine = new AgentLoopEngine( { maxIterations: 3, totalTimeoutMs: 15_000, retryCount: 2, contextWindow: 128_000 }, adapter, ); const events: MetonaStreamEvent[] = []; engine.on('streamEvent', (e: MetonaStreamEvent) => events.push(e)); const output = await engine.runStream( { role: 'user', content: 'hi', timestamp: Date.now() }, 's1', [], SYSTEM_PROMPT, ); expect(output.terminationReason).toBe('completed'); expect(output.finalAnswer).toBe('recovered'); expect(call).toBe(2); // 重试经过 RETRY 事件(引擎消费后向渲染层转 STREAM_RESET) expect(events.some((e) => e.type === MetonaStreamEventType.STREAM_RESET)).toBe(true); }); }); // ===== 3. REASONING_DELTA 中途 abort ===== describe('engine: abort during reasoning stream (P0-5 #4)', () => { it('思考增量持续输出中 abort → USER_INTERRUPT,迟到 reasoning 不再产出事件', async () => { const adapter = { providerId: 'mock' as const, supportedModels: ['m'], supportsToolCalling: true, supportsThinking: true, getContextWindow: () => 128_000, async send(): Promise { throw new Error('not used'); }, // 持续每 10ms 推送一个 reasoning delta,直到 abort 或 500 个(挂起) async *sendStream(_r: MetonaRequest): AsyncIterable { for (let i = 0; i < 500; i++) { if (_r.params.thinkingEnabled === false) return; yield { type: MetonaStreamEventType.REASONING_DELTA, requestId: 'r', sessionId: 's', iteration: 1, seq: i, timestamp: Date.now(), delta: `t${i};`, }; await new Promise((r) => setTimeout(r, 10)); } yield { type: MetonaStreamEventType.DONE, requestId: 'r', sessionId: 's', iteration: 1, seq: 999, timestamp: Date.now(), }; }, } as unknown as IMetonaProviderAdapter; const engine = new AgentLoopEngine( { maxIterations: 3, totalTimeoutMs: 30_000, retryCount: 0, contextWindow: 128_000 }, adapter, ); const events: MetonaStreamEvent[] = []; engine.on('streamEvent', (e: MetonaStreamEvent) => events.push(e)); const run = engine.runStream( { role: 'user', content: 'hi', timestamp: Date.now() }, 's1', [], SYSTEM_PROMPT, ); // 等待至少 3 个 reasoning delta 到达(确认流确实在思考态输出中) await new Promise((r) => setTimeout(r, 80)); expect( events.filter((e) => e.type === MetonaStreamEventType.REASONING_DELTA).length, ).toBeGreaterThanOrEqual(3); engine.abort(); const output = await run; expect(output.terminationReason).toBe('user_interrupt'); const lastReasoningIdx = events .map((e) => e.type) .lastIndexOf(MetonaStreamEventType.REASONING_DELTA); const doneIdx = events.findIndex((e) => e.type === MetonaStreamEventType.DONE); // abort 生效后不再有迟到 reasoning(abort 尾巴过滤在引擎侧) expect(lastReasoningIdx).toBeLessThan(doneIdx); }, 15_000); }); // ===== 4. P4-2: 上一 run 卡死 → 30s 强制 abort + 5s 再等待 → 抛错 ===== describe('engine: stuck previous run force-abort path (P4-2)', () => { it('上一 run 不响应 abort → 新 run 30s 超时强制 abort → 5s 后仍卡 → 抛错', async () => { vi.useFakeTimers(); let released = false; const stuckAdapter = { providerId: 'mock' as const, supportedModels: ['m'], supportsToolCalling: true, supportsThinking: true, getContextWindow: () => 128_000, async send(): Promise { throw new Error('not used'); }, // 第一次 run:永远挂起且不响应 abort(模拟卡死的底层 fetch); // 挂起生成器无 yield 属场景定义(require-yield 由下行指令豁免) // eslint-disable-next-line require-yield async *sendStream(): AsyncIterable { while (!released) { await new Promise(() => {}); } return; }, } as unknown as IMetonaProviderAdapter; const engine = new AgentLoopEngine( { maxIterations: 3, totalTimeoutMs: 600_000, retryCount: 0, contextWindow: 128_000 }, stuckAdapter, ); const firstRun = engine.runStream( { role: 'user', content: 'first', timestamp: Date.now() }, 's1', [], SYSTEM_PROMPT, ); firstRun.catch(() => {}); // 新消息进入 → 等待上一 run 结束(30s 超时)→ 强制 abort → 5s 再等待 → 抛错 const secondRun = engine.runStream( { role: 'user', content: 'second', timestamp: Date.now() }, 's1', [], SYSTEM_PROMPT, ); const expectation = expect(secondRun).rejects.toThrow(/Previous run did not finish/i); // 驱动 30s 等待超时 + 5s 强制 abort 等待 await vi.advanceTimersByTimeAsync(30_500); await vi.advanceTimersByTimeAsync(5_500); // 清理:放行挂死流,让 firstRun 收敛 released = true; await expectation; }, 20_000); });