- D2: each intent-classifier attempt is capped at 10 s (same session model, thinking untouched); a hang takes the existing retry -> classifier_unavailable path. Success is timed too (RectificationClassifierDiagnostic). - D3: a typed message builds the NDJSON stream first; the first line is turn.progress "received", then the classifier and deterministic replies run inside the stream. Preflight rejections become turn.rejected (old status, code, message) and the client handles them like the old HTTP rejection. Stage lines 收到,正在对照你的档案… / 正在记下这件事… / 正在重新对照盘面… / 正在准备下一个问题… are driven by existing tool events and engine calls, are transient (live row only) and never persisted. VOICE / DESIGN updated. - D4: RectificationRunDiagnostic records per-step start/end, provider token usage incl. reasoning tokens, classifier timing and per-engine-call durations (AsyncLocalStorage scope per turn); RectificationTurnDiagnostic for deterministic turns. No user text, birth data or model text. - D5: /v5/versions memo (30 s, complete identities only, per transport); the exit gate skips its second persistNextInterviewIfIdle when the run's own call found the next focus already active (provably identical). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_017eEAG8HD3mm8gsKXgk8uU8
183 lines
11 KiB
TypeScript
183 lines
11 KiB
TypeScript
import assert from "node:assert/strict";
|
|
import { spawnSync } from "node:child_process";
|
|
import { fileURLToPath } from "node:url";
|
|
import test from "node:test";
|
|
|
|
// BUG-1047 D3, route level (TASK-rectification-latency-20260926). The real
|
|
// POST /api/rectification/agent handler with fake persistence: a typed answer
|
|
// to a pending choice gets its first stream byte — the `turn.progress`
|
|
// received line — while the intent classifier is still running, and the
|
|
// deterministic reply follows on the same stream after the awaited exit gate.
|
|
// The progress line never reaches the persisted assistant message.
|
|
// Uses node:test module mocks (Node >= 22.3, same as the other route tests).
|
|
test("a typed message streams turn.progress before the classifier answers, then the reply after the exit gate", () => {
|
|
const script = String.raw`
|
|
import assert from 'node:assert/strict';
|
|
import { mock } from 'node:test';
|
|
import { pathToFileURL } from 'node:url';
|
|
import { CASE_ID, SESSION_ID, USER_ID, TURN_ID, FOCUS_ID, dossierFixture,
|
|
conversationSummaryFixture, activeFocusFixture, receiptHandlers } from './tests/rectification-v9-test-support.ts';
|
|
import * as realClassifier from './src/lib/rectification-agentic/v9/turn-intent-classifier.ts';
|
|
import * as realTurnExit from './src/lib/rectification-agentic/v9/turn-exit.ts';
|
|
import { RECTIFICATION_USER_COPY } from './src/lib/rectification-agentic/user-copy.ts';
|
|
import { RECTIFICATION_TURN_PROGRESS_LABELS } from './src/lib/rectification-activity-labels.ts';
|
|
|
|
const prompt = '2023 年前后,有没有换过工作或职责明显变化?';
|
|
const focus = activeFocusFixture({ intent: 'distinguish_candidates', questionId: 'd10:career:2023',
|
|
expectedAnswerSchema: { prompt, probe_id: 'probe:career:2023', choice: { prompt,
|
|
option_a: '明确发生且时间吻合', option_b: '发生过但程度较弱', option_c: '没有这回事', option_d: '这段记不清楚',
|
|
options: [
|
|
{ key: 'A', label: '明确发生且时间吻合', answer_class: 'yes' },
|
|
{ key: 'B', label: '发生过但程度较弱', answer_class: 'weak_yes' },
|
|
{ key: 'C', label: '没有这回事', answer_class: 'no' },
|
|
{ key: 'D', label: '这段记不清楚', answer_class: 'unsure' },
|
|
] } } });
|
|
const calls = [];
|
|
const order = [];
|
|
const accounting = { rpc: async (fn, args = {}) => {
|
|
calls.push({ fn, args });
|
|
let data;
|
|
if (fn === 'get_agentic_rectification_case') data = { ...dossierFixture().case, session_id: SESSION_ID, status: 'collecting_evidence' };
|
|
else if (fn === 'get_agentic_rectification_case_dossier') data = dossierFixture({
|
|
conversationSummary: conversationSummaryFixture({ activeFocus: focus }) });
|
|
else if (fn === 'append_agentic_rectification_turn') { order.push('append'); data = { turn_id: TURN_ID, idempotent: false }; }
|
|
else if (receiptHandlers[fn]) data = await receiptHandlers[fn](fn, args);
|
|
else throw new Error('Unexpected persistence RPC: ' + fn);
|
|
return { data: structuredClone(data), error: null };
|
|
}, from: () => { throw new Error('no table access expected'); } };
|
|
|
|
let releaseClassifier;
|
|
const classifierGate = new Promise((resolve) => { releaseClassifier = resolve; });
|
|
let classifierStarted = 0;
|
|
mock.module('server-only', { namedExports: {} });
|
|
mock.module('@/lib/supabase/server', { namedExports: { createServerSupabaseClient: async () => ({
|
|
auth: { getUser: async () => ({ data: { user: { id: USER_ID } }, error: null }) },
|
|
from: () => ({ select() { return this; }, eq() { return this; },
|
|
maybeSingle: async () => ({ data: { id: SESSION_ID, session_type: 'birth_time_rectification',
|
|
agentic_rectification_case_id: CASE_ID, model_id: 'synthetic', model_config_version: 1 }, error: null }) })
|
|
}) } });
|
|
mock.module('@/lib/supabase/admin', { namedExports: { createAdminSupabaseClient: () => accounting } });
|
|
mock.module('@/lib/product-access', { namedExports: { isProductEnabled: async () => true } });
|
|
mock.module('@/lib/feature-flags', { namedExports: { loadRuntimeFeatureFlags: async () => new Map([
|
|
['rectification_runtime_version', { enabled: true }]
|
|
]) } });
|
|
mock.module('@/lib/model-catalog', { namedExports: { resolveSessionLanguageModel: async () => ({
|
|
id: 'synthetic', configVersion: 1, model: {}
|
|
}) } });
|
|
mock.module('@/mastra/agentic-rectification', { namedExports: { getRectificationV9Agent: async () => {
|
|
throw new Error('the deterministic path must not build an Agent');
|
|
} } });
|
|
mock.module('@/lib/rectification-agentic/v9/turn-intent-classifier', { namedExports: {
|
|
...realClassifier,
|
|
classifyTurnIntentWithRetry: async () => {
|
|
classifierStarted += 1;
|
|
await classifierGate;
|
|
order.push('classified');
|
|
return { classified: null, expectedWrite: 'unknown', outcome: 'classifier_unavailable',
|
|
diagnostic: { outcome: 'classifier_unavailable', attempts: 2, timedOutAttempts: 2, elapsedMs: 20000 } };
|
|
},
|
|
} });
|
|
mock.module('@/lib/rectification-agentic/v9/turn-exit', { namedExports: {
|
|
...realTurnExit,
|
|
finalizeSuccessfulTurnExit: async (input) => {
|
|
order.push('exit:' + input.action + ':' + input.askedTurnId);
|
|
},
|
|
} });
|
|
|
|
const { POST } = await import(pathToFileURL(process.cwd() + '/src/app/api/rectification/agent/route.ts').href);
|
|
const sentAt = Date.now();
|
|
const response = await POST(new Request('https://example.invalid/api/rectification/agent', {
|
|
method: 'POST', headers: { 'content-type': 'application/json' },
|
|
body: JSON.stringify({ caseId: CASE_ID, sessionId: SESSION_ID, requestId: TURN_ID, action: 'message',
|
|
message: '换过,2023 年 3 月' })
|
|
}));
|
|
assert.equal(response.status, 200);
|
|
const reader = response.body.getReader();
|
|
const decoder = new TextDecoder();
|
|
const first = await reader.read();
|
|
const firstMs = Date.now() - sentAt;
|
|
const firstLine = decoder.decode(first.value).trim().split('\n')[0];
|
|
assert.deepEqual(JSON.parse(firstLine), { type: 'turn.progress', stage: 'received' });
|
|
assert.ok(firstMs <= 300, 'first byte within 300 ms, was ' + firstMs);
|
|
// The classifier is running (or about to) and has not answered: nothing else streamed yet.
|
|
await new Promise((resolve) => setTimeout(resolve, 50));
|
|
assert.equal(classifierStarted, 1);
|
|
assert.deepEqual(order, []);
|
|
releaseClassifier();
|
|
let rest = decoder.decode(first.value).trim().split('\n').slice(1).join('\n');
|
|
for (;;) {
|
|
const { done, value } = await reader.read();
|
|
if (done) break;
|
|
rest += decoder.decode(value, { stream: true });
|
|
}
|
|
const events = rest.trim().split('\n').filter(Boolean).map((line) => JSON.parse(line));
|
|
const types = events.map((event) => event.type);
|
|
assert.deepEqual(types.filter((type) => type !== 'turn.progress'), ['answer.delta', 'run.completed'], JSON.stringify(events));
|
|
assert.equal(events.find((event) => event.type === 'answer.delta').text, RECTIFICATION_USER_COPY.classifierUnavailableReply);
|
|
assert.equal(events.find((event) => event.type === 'run.completed').turnId, TURN_ID);
|
|
// Classified → persisted → awaited exit gate → only then the reply bytes.
|
|
assert.deepEqual(order, ['classified', 'append', 'exit:message:' + TURN_ID]);
|
|
const append = calls.find((call) => call.fn === 'append_agentic_rectification_turn');
|
|
assert.equal(append.args.p_assistant_message, RECTIFICATION_USER_COPY.classifierUnavailableReply);
|
|
for (const line of Object.values(RECTIFICATION_TURN_PROGRESS_LABELS)) {
|
|
assert.equal(JSON.stringify(calls).includes(line), false, 'progress line persisted: ' + line);
|
|
}
|
|
console.log(JSON.stringify({ firstMs, types: ['turn.progress', ...types] }));
|
|
`;
|
|
const result = spawnSync(process.execPath, ["--experimental-test-module-mocks", "--import", "tsx", "--input-type=module", "--eval", script], {
|
|
cwd: fileURLToPath(new URL("../", import.meta.url)), encoding: "utf8", timeout: 60_000,
|
|
});
|
|
assert.equal(result.status, 0, result.stderr + result.stdout);
|
|
console.log(result.stdout.trim());
|
|
});
|
|
|
|
test("a rejected typed message arrives as turn.rejected on the stream with the old status, code and message", () => {
|
|
const script = String.raw`
|
|
import assert from 'node:assert/strict';
|
|
import { mock } from 'node:test';
|
|
import { pathToFileURL } from 'node:url';
|
|
import { CASE_ID, SESSION_ID, USER_ID, TURN_ID, dossierFixture, receiptHandlers } from './tests/rectification-v9-test-support.ts';
|
|
const accounting = { rpc: async (fn, args = {}) => {
|
|
if (fn === 'get_agentic_rectification_case') return { data: { ...dossierFixture().case, session_id: SESSION_ID, status: 'collecting_evidence' }, error: null };
|
|
if (fn === 'get_agentic_rectification_case_dossier') return { data: null, error: { message: 'agentic_rectification_case_terminal' } };
|
|
if (receiptHandlers[fn]) return { data: await receiptHandlers[fn](fn, args), error: null };
|
|
throw new Error('Unexpected persistence RPC: ' + fn);
|
|
} };
|
|
mock.module('server-only', { namedExports: {} });
|
|
mock.module('@/lib/supabase/server', { namedExports: { createServerSupabaseClient: async () => ({
|
|
auth: { getUser: async () => ({ data: { user: { id: USER_ID } }, error: null }) },
|
|
from: () => ({ select() { return this; }, eq() { return this; },
|
|
maybeSingle: async () => ({ data: { id: SESSION_ID, session_type: 'birth_time_rectification',
|
|
agentic_rectification_case_id: CASE_ID, model_id: 'synthetic', model_config_version: 1 }, error: null }) })
|
|
}) } });
|
|
mock.module('@/lib/supabase/admin', { namedExports: { createAdminSupabaseClient: () => accounting } });
|
|
mock.module('@/lib/product-access', { namedExports: { isProductEnabled: async () => true } });
|
|
mock.module('@/lib/feature-flags', { namedExports: { loadRuntimeFeatureFlags: async () => new Map([
|
|
['rectification_runtime_version', { enabled: true }]
|
|
]) } });
|
|
mock.module('@/lib/model-catalog', { namedExports: { resolveSessionLanguageModel: async () => ({
|
|
id: 'synthetic', configVersion: 1, model: {}
|
|
}) } });
|
|
mock.module('@/mastra/agentic-rectification', { namedExports: { getRectificationV9Agent: async () => {
|
|
throw new Error('no Agent expected');
|
|
} } });
|
|
const { POST } = await import(pathToFileURL(process.cwd() + '/src/app/api/rectification/agent/route.ts').href);
|
|
const response = await POST(new Request('https://example.invalid/api/rectification/agent', {
|
|
method: 'POST', headers: { 'content-type': 'application/json' },
|
|
body: JSON.stringify({ caseId: CASE_ID, sessionId: SESSION_ID, requestId: TURN_ID, action: 'message', message: '换过' })
|
|
}));
|
|
assert.equal(response.status, 200);
|
|
const events = (await response.text()).trim().split('\n').map((line) => JSON.parse(line));
|
|
assert.deepEqual(events, [
|
|
{ type: 'turn.progress', stage: 'received' },
|
|
{ type: 'turn.rejected', httpStatus: 409, code: 'case_terminal', message: '该校正已结束,不能继续修改' },
|
|
]);
|
|
console.log(JSON.stringify(events.map((event) => event.type)));
|
|
`;
|
|
const result = spawnSync(process.execPath, ["--experimental-test-module-mocks", "--import", "tsx", "--input-type=module", "--eval", script], {
|
|
cwd: fileURLToPath(new URL("../", import.meta.url)), encoding: "utf8", timeout: 60_000,
|
|
});
|
|
assert.equal(result.status, 0, result.stderr + result.stdout);
|
|
console.log(result.stdout.trim());
|
|
});
|