/** * BUG-1047 (TASK-rectification-latency-20260926): a typed rectification answer * waited 25–90 s behind an unbounded classifier and four thinking steps with * no visible progress, and the run diagnostic could not say where the time * went. * * Covered here (Node 20-runnable, no module mocks): * - D2: each classifier attempt is capped at 10 s and a hang takes the * existing retry → `classifier_unavailable` path; success is timed too. * - D3: stage progress is monotonic, reaches the stream as `turn.progress`, * passes the public allowlist, and never lands in the saved reply or phases. * - D4: step start/end, provider reasoning tokens, classifier timing and * engine call durations reach `RectificationRunDiagnostic` without text. * - D5: `/v5/versions` memo, and the skipped second idle interview call. * The route-level "first byte before the classifier" test needs module mocks * and lives in `rectification-latency-route-20260926.test.ts`. */ import assert from "node:assert/strict"; import { readFileSync } from "node:fs"; import test from "node:test"; import { CASE_ID, RECTIFICATION_SKILL_SHA256, SESSION_ID, TURN_ID, USER_ID, activeFocusFixture, conversationSummaryFixture, dossierFixture, fakeAccounting, receiptHandlers, } from "./rectification-v9-test-support.ts"; import type { ResolvedLanguageModel } from "../src/mastra/model.ts"; import { RECTIFICATION_CLASSIFIER_ATTEMPT_TIMEOUT_MS, classifyTurnIntentWithRetry, } from "../src/lib/rectification-agentic/v9/turn-intent-classifier.ts"; import { createTurnInstrumentation, currentEngineCallTimings, reportTurnProgress, runWithTurnInstrumentation, type RectificationTurnProgressStage, } from "../src/lib/rectification-agentic/v9/turn-instrumentation.ts"; import { RECTIFICATION_ENGINE_VERSIONS_CACHE_TTL_MS, readV9EngineScoringIdentity, } from "../src/lib/rectification-agentic/v9/engine-client.ts"; import { safePublicEvent, turnProgressForChunk } from "../src/lib/rectification-agentic/v9/stream-mapping.ts"; import { diagnosticStepsFromTimings, recordStepChunk, type RectificationStepTiming, } from "../src/lib/rectification-agentic/v9/run-diagnostic.ts"; import { runV9AgentTurn, type V9AgentRunOptions } from "../src/lib/rectification-agentic/v9/agent-run.ts"; import { RECTIFICATION_SKILL_NAME } from "../src/lib/rectification-agentic/v9/case-status.ts"; import { persistNextInterviewIfIdle } from "../src/lib/rectification-agentic/v9/answer-choice.ts"; import { finalizeSuccessfulTurnExit } from "../src/lib/rectification-agentic/v9/turn-exit.ts"; import { RECTIFICATION_TURN_PROGRESS_LABELS } from "../src/lib/rectification-activity-labels.ts"; import { rectificationInitialLiveLabel } from "../src/lib/rectification-surface-state.ts"; const dummyModel = { id: "test-model" } as ResolvedLanguageModel; const PROGRESS_LINES = Object.values(RECTIFICATION_TURN_PROGRESS_LABELS); async function flush(times = 5) { for (let i = 0; i < times; i += 1) await Promise.resolve(); } // ---------------------------------------------------------------- D2 classifier test("D2: a hanging classifier attempt ends at 10 s, retries once, then takes classifier_unavailable", async (t) => { t.mock.timers.enable({ apis: ["setTimeout"] }); assert.equal(RECTIFICATION_CLASSIFIER_ATTEMPT_TIMEOUT_MS, 10_000); const signals: AbortSignal[] = []; const warnings: unknown[][] = []; const originalWarn = console.warn; console.warn = (...args: unknown[]) => { warnings.push(args); }; try { let settled = false; const pending = classifyTurnIntentWithRetry( dummyModel, { userMessage: "2016 年 3 月入学", caseStatus: "collecting_evidence" }, (_model, input) => { if (input.signal) signals.push(input.signal); return new Promise(() => {}); // provider never answers }, ).then((result) => { settled = true; return result; }); await flush(); assert.equal(signals.length, 1); t.mock.timers.tick(9_999); await flush(); assert.equal(settled, false, "no answer before 10 s"); assert.equal(signals[0].aborted, false); t.mock.timers.tick(1); await flush(); assert.equal(signals[0].aborted, true, "the first attempt is aborted at 10 s"); assert.equal(signals.length, 2, "exactly one retry starts"); assert.equal(settled, false); t.mock.timers.tick(10_000); const result = await pending; assert.equal(result.outcome, "classifier_unavailable"); assert.equal(result.classified, null); assert.equal(result.expectedWrite, "unknown"); assert.equal(result.diagnostic?.attempts, 2); assert.equal(result.diagnostic?.timedOutAttempts, 2); assert.equal(signals[1].aborted, true); const line = JSON.stringify(warnings[0]); assert.match(line, /rectification_classifier_unavailable/); assert.match(line, /timedOutAttempts/); assert.doesNotMatch(line, /2016|入学/); } finally { console.warn = originalWarn; } }); test("D2: a first attempt that times out and a second that answers is classified (existing retry)", async (t) => { t.mock.timers.enable({ apis: ["setTimeout"] }); let calls = 0; const pending = classifyTurnIntentWithRetry( dummyModel, { userMessage: "没有", caseStatus: "collecting_evidence" }, async () => { calls += 1; if (calls === 1) return new Promise(() => {}); return { intent: "answer_current_focus", answer_class: "no" } as const; }, ); await flush(); t.mock.timers.tick(10_000); const result = await pending; assert.equal(calls, 2); assert.equal(result.outcome, "classified"); assert.equal(result.diagnostic?.timedOutAttempts, 1); }); test("D2/D4: a successful classification logs its timing without the user text", async () => { const lines: string[] = []; const originalInfo = console.info; console.info = (...args: unknown[]) => { lines.push(args.map(String).join(" ")); }; try { const result = await classifyTurnIntentWithRetry( dummyModel, { userMessage: "2019 年换了工作", caseStatus: "collecting_evidence" }, async () => ({ intent: "provide_new_evidence", answer_class: null, has_new_dated_event: true }), ); assert.equal(result.outcome, "classified"); assert.equal(result.diagnostic?.attempts, 1); assert.equal(result.diagnostic?.timedOutAttempts, 0); assert.equal(typeof result.diagnostic?.elapsedMs, "number"); } finally { console.info = originalInfo; } const line = lines.find((item) => item.includes("RectificationClassifierDiagnostic")); assert.ok(line, "success is logged too"); assert.doesNotMatch(line, /2019|换了工作|provide_new_evidence/); }); test("D2 red line: no model swap, thinking untouched, route uses the default 10 s cap", () => { const classifier = readFileSync(new URL("../src/lib/rectification-agentic/v9/turn-intent-classifier.ts", import.meta.url), "utf8"); const route = readFileSync(new URL("../src/app/api/rectification/agent/route.ts", import.meta.url), "utf8"); // The Agent is still built on the session model with no thinking override. assert.match(classifier, /model: model\.model,/); assert.doesNotMatch(classifier, /thinking:|thinkingTokens|providerOptions|modelSettings|agentGenerationSettings/); assert.doesNotMatch(route, /attemptTimeoutMs/); assert.equal((route.match(/await classifyTurnIntentWithRetry\(/g) ?? []).length, 4); }); // ---------------------------------------------------------------- D3 progress test("D3: stages only move forward, reach the sink, and reset for a retried attempt", async () => { const seen: RectificationTurnProgressStage[] = []; await runWithTurnInstrumentation({ onProgress: (stage) => seen.push(stage) }, async (instrumentation) => { instrumentation.advance("received"); reportTurnProgress("recording"); await Promise.resolve(); reportTurnProgress("received"); // lower: ignored reportTurnProgress("rescoring"); reportTurnProgress("rescoring"); // same: ignored instrumentation.resetStage(); reportTurnProgress("received"); reportTurnProgress("preparing_question"); }); assert.deepEqual(seen, ["received", "recording", "rescoring", "received", "preparing_question"]); // Outside a turn scope the calls are silent no-ops. reportTurnProgress("recording"); assert.deepEqual(currentEngineCallTimings(), []); }); test("D3: the scope follows async work started inside run(), including a ReadableStream start", async () => { const seen: string[] = []; const instrumentation = createTurnInstrumentation(); const stream = instrumentation.run(() => new ReadableStream({ async start(controller) { instrumentation.setProgressSink((stage) => controller.enqueue(stage)); instrumentation.advance("received"); await new Promise((resolve) => setTimeout(resolve, 5)); reportTurnProgress("rescoring"); // e.g. from the engine client, deep in the turn controller.close(); }, })); const reader = stream.getReader(); for (;;) { const { done, value } = await reader.read(); if (done) break; seen.push(value); } assert.deepEqual(seen, ["received", "rescoring"]); }); test("D3: turn.progress and turn.rejected pass the public allowlist, nothing else rides along", () => { assert.deepEqual(safePublicEvent({ type: "turn.progress", stage: "recording", text: "x" }), { type: "turn.progress", stage: "recording", }); assert.equal(safePublicEvent({ type: "turn.progress", stage: "thinking" }), null); assert.deepEqual( safePublicEvent({ type: "turn.rejected", httpStatus: 409, code: "stale_question", message: "这道题已经过期", secret: 1 }), { type: "turn.rejected", httpStatus: 409, code: "stale_question", message: "这道题已经过期" }, ); assert.equal(safePublicEvent({ type: "turn.rejected", httpStatus: 200, message: "ok" }), null); assert.deepEqual( safePublicEvent({ type: "turn.rejected", httpStatus: 500, code: "Bad Code!", message: "m" }), { type: "turn.rejected", httpStatus: 500, message: "m" }, ); }); test("D3: tool events map to the four stages", () => { const call = (toolName: string) => ({ type: "tool-call", payload: { toolName } }) as never; const result = (toolName: string) => ({ type: "tool-result", payload: { toolName, result: { ok: true } } }) as never; assert.equal(turnProgressForChunk(call("rectification-read-case")), null); assert.equal(turnProgressForChunk(call("rectification-record-evidence-batch")), "recording"); assert.equal(turnProgressForChunk(result("rectification-record-evidence-batch")), "preparing_question"); assert.equal(turnProgressForChunk(call("rectification-compare-candidates")), "rescoring"); assert.equal(turnProgressForChunk(call("rectification-set-focus")), "preparing_question"); assert.equal(turnProgressForChunk(call("skill")), null); assert.equal(turnProgressForChunk({ type: "text-delta", payload: { text: "x" } } as never), null); }); test("D3 copy: the four lines, and a typed answer starts on the first one", () => { assert.deepEqual(RECTIFICATION_TURN_PROGRESS_LABELS, { received: "收到,正在对照你的档案…", recording: "正在记下这件事…", rescoring: "正在重新对照盘面…", preparing_question: "正在准备下一个问题…", }); assert.equal(rectificationInitialLiveLabel("message"), RECTIFICATION_TURN_PROGRESS_LABELS.received); const voice = readFileSync(new URL("../docs/VOICE.md", import.meta.url), "utf8"); const design = readFileSync(new URL("../DESIGN.md", import.meta.url), "utf8"); for (const line of PROGRESS_LINES) { assert.ok(voice.includes(line), `VOICE lists ${line}`); assert.ok(design.includes(line), `DESIGN lists ${line}`); } }); // ---------------------------------------------------------------- D4 run diagnostic test("D4: step timeline records start/end, public tools and provider usage per step", () => { const steps: RectificationStepTiming[] = []; const isPublic = (name: string) => name.startsWith("rectification-"); recordStepChunk(steps, { type: "step-start" }, 10, isPublic); recordStepChunk(steps, { type: "tool-call", payload: { toolName: "rectification-read-case", args: { caseId: CASE_ID } } }, 900, isPublic); recordStepChunk(steps, { type: "step-finish", payload: { output: { usage: { inputTokens: 1200, outputTokens: 40, reasoningTokens: 700 } } } }, 1_000, isPublic); recordStepChunk(steps, { type: "step-start" }, 1_050, isPublic); recordStepChunk(steps, { type: "tool-call", payload: { toolName: "skill" } }, 1_060, isPublic); recordStepChunk(steps, { type: "step-finish", payload: { output: { usage: { inputTokens: 1300, outputTokens: 20 } } } }, 4_000, isPublic); const summary = diagnosticStepsFromTimings(steps); assert.equal(summary.stepCount, 2); assert.deepEqual(summary.steps.map((step) => [step.startMs, step.endMs, step.tools]), [ [10, 1_000, ["rectification-read-case"]], [1_050, 4_000, []], ]); assert.equal(summary.reasoningTokens, 700); assert.equal(summary.inputTokens, 2_500); assert.equal(summary.toolCallCount, 1); assert.equal(diagnosticStepsFromTimings([]).reasoningTokens, null, "unreported stays null, not 0"); }); function chunk(type: string, payload?: Record) { return { type, ...(payload ? { payload } : {}) }; } function fakeAgent(chunks: Array<{ type: string; payload?: Record }>) { return { stream: async () => ({ fullStream: (async function* () { for (const item of chunks) yield item; })(), totalUsage: Promise.resolve({ inputTokens: 10, outputTokens: 20 }), }), getSkill: async () => ({ name: RECTIFICATION_SKILL_NAME, instructions: "skill" }), }; } test("D3/D4: an Agent turn reports stages, a full diagnostic, and never persists a progress line", async () => { const accounting = fakeAccounting({ ...receiptHandlers, get_agentic_rectification_case_dossier: () => dossierFixture(), append_agentic_rectification_turn: () => ({ turn_id: TURN_ID }), finalize_agentic_rectification_turn: () => ({ turn_id: TURN_ID, status: "completed", idempotent: false }), }); const emitted: Array<{ type: string }> = []; const stages: RectificationTurnProgressStage[] = []; const infos: string[] = []; const originalInfo = console.info; console.info = (...args: unknown[]) => { infos.push(args.map(String).join(" ")); }; const options: V9AgentRunOptions = { userId: USER_ID, caseId: CASE_ID, sessionId: SESSION_ID, requestId: "aaaaaaaa-bbbb-4ccc-8ddd-eeeeeeeeeeee", action: "read_only", message: null, modelName: "test-model", accounting: accounting.client, billing: { reserve: async () => ({ success: true, status: 200 }), complete: async () => true, release: async () => true, }, emit: (event) => { emitted.push(event); }, classifierDiagnostic: { outcome: "classified", attempts: 1, timedOutAttempts: 0, elapsedMs: 4_200 }, buildAgent: async () => fakeAgent([ chunk("start"), chunk("step-start"), chunk("tool-call", { toolName: "rectification-read-case", args: { caseId: CASE_ID } }), chunk("tool-result", { toolName: "rectification-read-case" }), chunk("step-finish", { output: { usage: { inputTokens: 900, outputTokens: 30, reasoningTokens: 512 } } }), chunk("step-start"), chunk("tool-call", { toolName: "rectification-record-evidence-batch", args: { caseId: CASE_ID } }), chunk("tool-result", { toolName: "rectification-record-evidence-batch", result: { accepted: [] } }), chunk("step-finish", { output: { usage: { inputTokens: 950, outputTokens: 40, reasoningTokens: 256 } } }), chunk("step-start"), chunk("text-delta", { text: "记下了。" }), chunk("step-finish", { output: { usage: { inputTokens: 980, outputTokens: 12 } } }), chunk("finish"), ]) as never, }; let result; try { result = await runWithTurnInstrumentation( { onProgress: (stage) => stages.push(stage) }, () => runV9AgentTurn(options), ); } finally { console.info = originalInfo; } assert.equal(result.ok, true, JSON.stringify(result)); assert.deepEqual(stages, ["recording", "preparing_question"]); // Progress never goes through the persisted phase/answer channel. assert.equal(emitted.some((event) => event.type === "turn.progress"), false); const persistedText = JSON.stringify(accounting.calls.map((call) => call.args)); for (const line of PROGRESS_LINES) assert.ok(!persistedText.includes(line), `not persisted: ${line}`); assert.equal(persistedText.includes("turn.progress"), false); const finalize = accounting.calls.findLast((call) => call.fn === "finalize_agentic_rectification_turn"); assert.equal(finalize?.args.p_assistant_message, "记下了。"); const line = infos.find((item) => item.includes("RectificationRunDiagnostic")); assert.ok(line); const diagnostic = JSON.parse(line); assert.equal(diagnostic.stepCount, 3); assert.equal(diagnostic.reasoningTokens, 768); assert.equal(diagnostic.inputTokens, 2_830); assert.equal(diagnostic.toolCallCount, 2); assert.equal(diagnostic.distinctToolCount, 2); assert.deepEqual(diagnostic.steps.map((step: { tools: string[] }) => step.tools), [ ["rectification-read-case"], ["rectification-record-evidence-batch"], [], ]); for (const step of diagnostic.steps) { assert.equal(typeof step.startMs, "number"); assert.ok(step.endMs >= step.startMs); } assert.deepEqual(diagnostic.classifier, { outcome: "classified", attempts: 1, timedOutAttempts: 0, elapsedMs: 4_200 }); assert.deepEqual(diagnostic.engineCalls, []); assert.equal(diagnostic.attemptNumber, 1); assert.doesNotMatch(line, /记下了|RECTIFICATION_SKILL|北京|1990/); assert.ok(!line.includes(RECTIFICATION_SKILL_SHA256)); }); // ---------------------------------------------------------------- D4 engine timings + D5 versions memo function isolateIdentityEnv(t: { after(fn: () => void): void }) { const saved = { algorithm: process.env.RECTIFICATION_ALGORITHM_VERSION, policy: process.env.RECTIFICATION_DECISION_POLICY_VERSION, engine: process.env.RECTIFICATION_ENGINE_VERSION, }; delete process.env.RECTIFICATION_ALGORITHM_VERSION; delete process.env.RECTIFICATION_DECISION_POLICY_VERSION; delete process.env.RECTIFICATION_ENGINE_VERSION; t.after(() => { for (const [key, value] of [ ["RECTIFICATION_ALGORITHM_VERSION", saved.algorithm], ["RECTIFICATION_DECISION_POLICY_VERSION", saved.policy], ["RECTIFICATION_ENGINE_VERSION", saved.engine], ] as const) { if (value === undefined) delete process.env[key]; else process.env[key] = value; } }); } const VERSIONS = { algorithm_version: "rectification-v5-matrix-scoring-8", decision_policy_version: "rectification-candidate-policy-v3" }; test("D5: /v5/versions is read once per TTL when env does not pin the identity; engine timing is recorded", async (t) => { isolateIdentityEnv(t); let fetches = 0; t.mock.method(globalThis, "fetch", async (url: unknown) => { fetches += 1; assert.ok(String(url).endsWith("/api/rectification/v5/versions")); return Response.json(VERSIONS); }); let now = 1_000_000; t.mock.method(Date, "now", () => now); const timings = await runWithTurnInstrumentation({}, async () => { const first = await readV9EngineScoringIdentity(); const second = await readV9EngineScoringIdentity(); assert.deepEqual(first, { algorithmVersion: VERSIONS.algorithm_version, policyVersion: VERSIONS.decision_policy_version }); assert.deepEqual(second, first); assert.notEqual(second, first, "callers get their own object"); return currentEngineCallTimings(); }); assert.equal(fetches, 1); assert.deepEqual(timings.map((item) => [item.path, item.outcome]), [["/api/rectification/v5/versions", "ok"]]); now += RECTIFICATION_ENGINE_VERSIONS_CACHE_TTL_MS - 1; await readV9EngineScoringIdentity(); assert.equal(fetches, 1, "still inside the TTL"); now += 2; await readV9EngineScoringIdentity(); assert.equal(fetches, 2, "TTL expired: read again"); }); test("D5: failures and incomplete version payloads are never memoised", async (t) => { isolateIdentityEnv(t); let fetches = 0; t.mock.method(globalThis, "fetch", async () => { fetches += 1; return fetches <= 2 ? Response.json({ algorithm_version: VERSIONS.algorithm_version }) : Response.json(VERSIONS); }); assert.equal((await readV9EngineScoringIdentity()).policyVersion, null); assert.equal((await readV9EngineScoringIdentity()).policyVersion, null); assert.equal((await readV9EngineScoringIdentity()).policyVersion, VERSIONS.decision_policy_version); assert.equal(fetches, 3); t.mock.method(globalThis, "fetch", async () => { throw new Error("engine down"); }); assert.deepEqual(await readV9EngineScoringIdentity(), { algorithmVersion: null, policyVersion: null }, "a different transport never sees another transport's memo"); }); // ---------------------------------------------------------------- D5 second idle call function focusDossier() { return dossierFixture({ conversationSummary: conversationSummaryFixture({ activeFocus: activeFocusFixture({ intent: "distinguish_candidates", questionId: "d9:relationship:2023", targetDomain: "relationship", expectedAnswerSchema: { prompt: "2023 年前后,你有没有一段认真开始或结束的关系?", choice: { prompt: "2023 年前后,你有没有一段认真开始或结束的关系?", 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" }, ], }, }, }), }), }); } function focusAccounting() { return fakeAccounting({ ...receiptHandlers, get_agentic_rectification_case_dossier: () => focusDossier(), set_agentic_rectification_conversation_focus: (_fn, args) => ({ focus: { ...activeFocusFixture({ intent: "distinguish_candidates", questionId: String(args.p_question_id) }), asked_turn_id: args.p_asked_turn_id }, idempotent: true, }), }); } test("D5: with the next focus already active, a second idle call is a pure repeat — so skipping it is output-identical", async () => { const accounting = focusAccounting(); const first = await persistNextInterviewIfIdle({ accounting: accounting.client, userId: USER_ID, caseId: CASE_ID, askedTurnId: TURN_ID }); const writesAfterFirst = accounting.calls.filter((call) => call.fn === "set_agentic_rectification_conversation_focus").map((call) => call.args); const second = await persistNextInterviewIfIdle({ accounting: accounting.client, userId: USER_ID, caseId: CASE_ID, askedTurnId: TURN_ID }); const writesAfterSecond = accounting.calls.filter((call) => call.fn === "set_agentic_rectification_conversation_focus").map((call) => call.args); assert.equal(first.focusActive, true); assert.deepEqual(second, first); // The only write the second call makes is the identical idempotent re-link. assert.equal(writesAfterSecond.length, writesAfterFirst.length * 2); assert.deepEqual(writesAfterSecond.slice(writesAfterFirst.length), writesAfterFirst); }); test("D5: the exit gate skips its idle re-read only when told the run settled the interview", async () => { // Both arms first make the Agent run's own post-turn call, as agent-run.ts does. const run = async (interviewSettled: boolean) => { const accounting = focusAccounting(); const idle = await persistNextInterviewIfIdle({ accounting: accounting.client, userId: USER_ID, caseId: CASE_ID, askedTurnId: TURN_ID }); assert.equal(idle.focusActive, true); await finalizeSuccessfulTurnExit({ accounting: accounting.client, userId: USER_ID, caseId: CASE_ID, action: "message", askedTurnId: TURN_ID, interviewSettled, }); return accounting.calls; }; const unsettled = await run(false); const settled = await run(true); const writes = (calls: typeof settled) => calls .filter((call) => call.fn !== "get_agentic_rectification_case_dossier" && call.fn !== "get_agentic_rectification_case_compute") .map((call) => JSON.stringify([call.fn, call.args])); assert.ok(settled.length < unsettled.length, `${settled.length} < ${unsettled.length}`); // Same distinct writes either way: the skipped call only re-linked the same focus. assert.deepEqual([...new Set(writes(settled))].sort(), [...new Set(writes(unsettled))].sort()); }); test("D5 contract: the Agent run reports interviewSettled only from the focus-active exit; the route forwards it", () => { const agentRun = readFileSync(new URL("../src/lib/rectification-agentic/v9/agent-run.ts", import.meta.url), "utf8"); const route = readFileSync(new URL("../src/app/api/rectification/agent/route.ts", import.meta.url), "utf8"); assert.match(agentRun, /interviewSettled = interviewIdle\.focusActive === true;/); assert.match(route, /interviewSettled: result\.interviewSettled === true/); }); // ---------------------------------------------------------------- D3 route shape (source) test("D3 contract: a typed message builds the stream first and classifies inside it", () => { const route = readFileSync(new URL("../src/app/api/rectification/agent/route.ts", import.meta.url), "utf8"); assert.match(route, /const deferToStream = action === "message";/); assert.match(route, /const immediateResponse = deferToStream \? null : await computeImmediateResponse\(\);/); const start = route.indexOf("async start(controller)"); const inStream = route.slice(start); const advance = inStream.indexOf('instrumentation.advance("received")'); const preflight = inStream.indexOf("await computeImmediateResponse()"); const unfocused = inStream.indexOf("await classifyUnfocusedMessage()"); const agent = inStream.indexOf("const result = await runV9AgentTurn"); assert.ok(start > 0 && advance > 0, "first progress line is inside the stream"); assert.ok(advance < preflight && preflight < unfocused && unfocused < agent); // Before the stream exists nothing runs the classifier for a message: the // preflight is only invoked outside the stream when it is not deferred, and // the unfocused classification is only invoked inside the stream. assert.equal((route.match(/await computeImmediateResponse\(\)/g) ?? []).length, 2); assert.equal((route.match(/await classifyUnfocusedMessage\(\)/g) ?? []).length, 1); assert.equal((inStream.match(/await computeImmediateResponse\(\)/g) ?? []).length, 1); // The instrumentation scope wraps the stream construction. assert.match(route, /const body = instrumentation\.run\(\(\) => new ReadableStream\(\{/); });