| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218 |
- import { createServer, type Server } from 'node:http';
- import { mkdir, mkdtemp, rm } from 'node:fs/promises';
- import { tmpdir } from 'node:os';
- import { join } from 'node:path';
- import { afterAll, beforeAll, describe, expect, it } from 'vitest';
- import { CodexRuntime, type CodexProviderSpec } from './codexRuntime';
- /**
- * 集成冒烟(不进 npm test,跑法:npm run smoke:codex):
- * 真 codex.exe + 只说 Responses 的本地 mock,钉住"Codex 在长时间静默后会把回合判死"这条机制。
- *
- * 病根(源码 rust-v0.155.1):codex-api/src/sse/responses.rs:580-606 在流建立之后,每次轮询都
- * `timeout(stream_idle_timeout_ms, stream.next())`,到期给出可重试的 CodexErr::Stream
- * (protocol/src/error.rs:372-413),默认窗口 300 秒(model-provider-info/src/lib.rs:29)。
- * 端点若在响应头之后长时间静默(自部署单槽推理的 prefill 就是这样),回合就会被判死并重发整轮。
- * 客户端不做协议翻译,所以这条窗口只能靠 provider 参数放宽来治 —— 见 providerService.PROVIDER_DEFAULTS。
- *
- * 这里把静默与空闲窗口都压缩到几十秒内,让同一条机制能在 CI 里被观测,而不是靠推理。
- */
- const SILENCE_MS = 30_000;
- let server: Server | null = null;
- let baseUrl = '';
- let root = '';
- /** 上游被真正问了几次:一次回合只该付一次 prefill,重试就是放大器 */
- let posts = 0;
- let homes = 0;
- let headers: 'late' | 'eager' = 'late';
- /**
- * 静默 SILENCE_MS 后正常吐完一轮 Responses SSE。两种形态的差别是本测试的全部意义:
- * - late:连响应头都不发(流还没建立);
- * - eager:先发响应头 + response.created,再静默(流已建立,长 prefill 的真实形态)。
- * Codex 的 `timeout(stream_idle_timeout_ms, stream.next())` 只在流已建立后武装,
- * 所以只有 eager 形态会在 prefill 静默中被判死 —— 见 rust-v0.155.1
- * codex-api/src/sse/responses.rs:580-606。
- */
- async function listen(): Promise<void> {
- server = createServer((req, res) => {
- let body = '';
- req.on('data', (chunk) => (body += chunk));
- req.on('end', () => {
- const route = req.url?.split('?')[0] ?? '';
- if (route !== '/v1/responses' || req.method !== 'POST') {
- res.writeHead(404, { 'content-type': 'application/json' });
- res.end(JSON.stringify({ error: { message: 'no route' } }));
- return;
- }
- posts += 1;
- const model = (() => {
- try {
- return String(JSON.parse(body).model ?? '');
- } catch {
- return '';
- }
- })();
- const send = (type: string, extra: Record<string, unknown>) =>
- res.write(`event: ${type}\ndata: ${JSON.stringify({ type, ...extra })}\n\n`);
- const created = () =>
- send('response.created', {
- response: { id: 'resp_1', object: 'response', model, status: 'in_progress', output: [] },
- });
- const finish = () => {
- send('response.output_item.added', {
- output_index: 0,
- item: { type: 'message', id: 'msg_1', role: 'assistant', status: 'in_progress', content: [] },
- });
- send('response.output_text.delta', { item_id: 'msg_1', output_index: 0, content_index: 0, delta: '久等了' });
- send('response.output_text.done', { item_id: 'msg_1', output_index: 0, content_index: 0, text: '久等了' });
- send('response.output_item.done', {
- output_index: 0,
- item: {
- type: 'message',
- id: 'msg_1',
- role: 'assistant',
- status: 'completed',
- content: [{ type: 'output_text', text: '久等了' }],
- },
- });
- send('response.completed', {
- response: {
- id: 'resp_1',
- object: 'response',
- model,
- status: 'completed',
- output: [
- {
- type: 'message',
- id: 'msg_1',
- role: 'assistant',
- status: 'completed',
- content: [{ type: 'output_text', text: '久等了' }],
- },
- ],
- usage: { input_tokens: 1000, output_tokens: 3, total_tokens: 1003 },
- },
- });
- res.end();
- };
- const head = () => res.writeHead(200, { 'content-type': 'text/event-stream' });
- if (headers === 'eager') {
- head();
- created();
- setTimeout(finish, SILENCE_MS);
- return;
- }
- setTimeout(() => {
- head();
- created();
- finish();
- }, SILENCE_MS);
- });
- });
- await new Promise<void>((resolve) => server?.listen(0, '127.0.0.1', resolve));
- const address = server.address();
- const port = typeof address === 'object' && address ? address.port : 0;
- baseUrl = `http://127.0.0.1:${port}/v1`;
- }
- function spec(extra: Record<string, number>): CodexProviderSpec {
- return {
- id: 'idletest',
- name: '静默端点',
- baseUrl,
- model: 'mock-responses',
- apiKey: 'sk-idletest-secret-0123456789',
- extra,
- };
- }
- async function runOneTurn(mode: 'late' | 'eager', extra: Record<string, number>) {
- headers = mode;
- const home = join(root, `home-${++homes}`);
- await mkdir(home, { recursive: true });
- const runtime = new CodexRuntime({ codexHome: home, provider: spec(extra) });
- const events: Array<{ method: string; params?: unknown }> = [];
- runtime.on('notification', (n: { method: string; params?: unknown }) => events.push(n));
- try {
- await runtime.start();
- const threadId = await runtime.startThread({ cwd: root, approvalPolicy: 'never', sandbox: 'read-only' });
- const postsBefore = posts;
- const result = await runtime.runTurn({ threadId, prompt: '说句话', approvalPolicy: 'never' });
- return { result, postsInTurn: posts - postsBefore, events };
- } finally {
- await runtime.dispose();
- }
- }
- beforeAll(async () => {
- root = await mkdtemp(join(tmpdir(), 'zsjz-idle-'));
- await listen();
- }, 120_000);
- afterAll(async () => {
- await new Promise<void>((resolve) => (server ? server.close(() => resolve()) : resolve()));
- await rm(root, { recursive: true, force: true }).catch(() => undefined);
- });
- /**
- * 把错误文案取回来:失败的回合由 turn/completed{status:'failed',error} 收尾,
- * 但 'error' 通知先到,而 runTurn 要等带 turnId 的 turn/completed —— 按协议形状从通知里取。
- */
- function firstErrorOf(events: Array<{ method: string; params?: unknown }>): string {
- for (const n of events) {
- if (n.method !== 'error') continue;
- const params = (n.params ?? {}) as { message?: string; error?: { message?: string } };
- const message = params.message ?? params.error?.message ?? '';
- if (message) return message;
- }
- return '';
- }
- describe('Codex 的 SSE 空闲窗口决定长回合的生死', () => {
- it(
- '迟发响应头(流没建立)→ 空闲窗口再小也不会判死:计时器只在流建立后武装',
- async () => {
- const { result, postsInTurn } = await runOneTurn('late', {
- stream_idle_timeout_ms: 5_000,
- request_max_retries: 0,
- stream_max_retries: 0,
- });
- expect(result.status).toBe('completed');
- expect(postsInTurn).toBe(1);
- },
- 120_000,
- );
- it(
- '先建流再静默(长 prefill 的真实形态)→ 空闲窗口到点就把回合判死,且 stream_max_retries=0 时上游只被问一次',
- async () => {
- const { result, postsInTurn, events } = await runOneTurn('eager', {
- stream_idle_timeout_ms: 5_000,
- request_max_retries: 0,
- stream_max_retries: 0,
- });
- expect(result.status).toBe('failed');
- expect(firstErrorOf(events)).toMatch(/idle timeout waiting for SSE/iu);
- expect(postsInTurn).toBe(1);
- },
- 120_000,
- );
- it(
- '同样的静默,把空闲窗口放宽到大于 prefill → 回合正常完成,仍然只问一次',
- async () => {
- const { result, postsInTurn } = await runOneTurn('eager', {
- stream_idle_timeout_ms: 120_000,
- request_max_retries: 0,
- stream_max_retries: 0,
- });
- expect(result.status).toBe('completed');
- expect(result.text).toContain('久等了');
- expect(postsInTurn).toBe(1);
- },
- 120_000,
- );
- });
|