// Characterization of assembleStreamingPipeline — the streaming transform-chain assembly extracted // from handleChatCore's streaming success path (chatCore god-file decomposition, #3501). All // transform factories are injected; a fake "stream" records each pipeThrough so the exact chain // order and the branch conditions (PII explicit vs feature-flag, progress, echo) are observable // without real ReadableStreams. Locks: transform order, the progress header side-effect, and the // PII branch precedence (explicit createPiiTransform wins over the feature flag). import { test } from "node:test"; import assert from "node:assert/strict"; const { assembleStreamingPipeline } = await import("../../open-sse/handlers/chatCore/streamingPipeline.ts"); type PipeWithDisconnectParameters = Parameters< typeof import("../../open-sse/utils/streamHandler.ts").pipeWithDisconnect >; type ShapeForClientFormatParameters = Parameters< typeof import("../../open-sse/utils/sseHeartbeat.ts").shapeForClientFormat >; function assertPipelineInputContracts(args: Parameters[0]) { const providerResponse: PipeWithDisconnectParameters[0] = args.providerResponse; const transformStream: PipeWithDisconnectParameters[1] = args.transformStream; const streamController: PipeWithDisconnectParameters[2] = args.streamController; const clientResponseFormat: ShapeForClientFormatParameters[0] = args.clientResponseFormat; return { providerResponse, transformStream, streamController, clientResponseFormat }; } // A fake stream: each pipeThrough appends the transform's tag and returns a new fake stream. function fakeStream(tag: string, log: string[]) { return { tag, pipeThrough(t: { __tag: string }) { log.push(t.__tag); return fakeStream(t.__tag, log); }, }; } function makeDeps(over: Record = {}) { const log: string[] = []; const deps = { wantsProgress: () => false, pipeWithDisconnect: (..._a: unknown[]) => fakeStream("pii-base", log), isFeatureFlagEnabled: () => false, createPiiSseTransform: () => ({ __tag: "pii-flag" }), createProgressTransform: () => ({ __tag: "progress" }), createSseHeartbeatTransform: () => ({ __tag: "heartbeat" }), shapeForClientFormat: (f: unknown) => f, createModelEchoTransform: () => ({ __tag: "echo" }), ...over, } as Parameters[1]; return { deps, log }; } function baseArgs(over: Record = {}) { return { providerResponse: {}, transformStream: {}, streamController: { signal: {} as AbortSignal }, createPiiTransform: undefined, clientRawRequestHeaders: null, clientResponseFormat: "openai", echoModel: null, responseHeaders: {} as Record, ...over, } as Parameters[0]; } test("pipeline arguments preserve the dependency parameter contracts", () => { const args = baseArgs({ providerResponse: new Response(), transformStream: new TransformStream(), streamController: { signal: new AbortController().signal, startTime: Date.now(), isConnected: () => true, handleDisconnect: () => {}, handleComplete: () => {}, markClientTerminalSeen: () => {}, handleError: () => {}, abort: () => {}, clientResponseFormat: "openai", }, clientResponseFormat: "openai", }); const contracts = assertPipelineInputContracts(args); assert.ok(contracts.providerResponse instanceof Response); assert.ok(contracts.transformStream instanceof TransformStream); assert.equal(contracts.streamController.clientResponseFormat, "openai"); assert.equal(contracts.clientResponseFormat, "openai"); }); test("baseline (no pii, no progress, no echo) → only heartbeat in the chain", () => { const { deps, log } = makeDeps(); assembleStreamingPipeline(baseArgs(), deps); assert.deepEqual(log, ["heartbeat"]); }); test("feature-flag PII → pii-flag transform applied before heartbeat", () => { const { deps, log } = makeDeps({ isFeatureFlagEnabled: () => true }); assembleStreamingPipeline(baseArgs(), deps); assert.deepEqual(log, ["pii-flag", "heartbeat"]); }); test("explicit createPiiTransform wins over the feature flag", () => { const { deps, log } = makeDeps({ isFeatureFlagEnabled: () => true }); const explicit = () => ({ __tag: "pii-explicit" }); assembleStreamingPipeline(baseArgs({ createPiiTransform: explicit }), deps); assert.deepEqual(log, ["pii-explicit", "heartbeat"]); }); test("progress enabled → progress transform + progress header set", () => { const { deps, log } = makeDeps({ wantsProgress: () => true }); const args = baseArgs(); assembleStreamingPipeline(args, deps); assert.deepEqual(log, ["progress", "heartbeat"]); assert.ok(Object.values(args.responseHeaders).includes("enabled")); }); test("progress disabled → no progress header", () => { const { deps } = makeDeps({ wantsProgress: () => false }); const args = baseArgs(); assembleStreamingPipeline(args, deps); assert.deepEqual(args.responseHeaders, {}); }); test("echoModel set → echo transform applied last", () => { const { deps, log } = makeDeps(); assembleStreamingPipeline(baseArgs({ echoModel: "alias-x" }), deps); assert.deepEqual(log, ["heartbeat", "echo"]); }); test("full chain order: pii → progress → heartbeat → echo", () => { const { deps, log } = makeDeps({ isFeatureFlagEnabled: () => true, wantsProgress: () => true, }); assembleStreamingPipeline(baseArgs({ echoModel: "alias-x" }), deps); assert.deepEqual(log, ["pii-flag", "progress", "heartbeat", "echo"]); }); test("pipeline assembly creates performance mark and measure entries", () => { performance.clearMarks(); performance.clearMeasures(); const { deps } = makeDeps(); assembleStreamingPipeline(baseArgs(), deps); assert.equal(performance.getEntriesByName("omni-pipeline-start").length, 1, "start mark"); assert.equal(performance.getEntriesByName("omni-pipeline-end").length, 1, "end mark"); assert.equal(performance.getEntriesByName("omni-pipeline").length, 1, "measure"); }); test("re-entering pipeline clears previous timeline entries", () => { performance.clearMarks(); performance.clearMeasures(); const { deps } = makeDeps(); // First call creates entries assembleStreamingPipeline(baseArgs(), deps); // Second call — clearMarks/clearMeasures runs before new marks, so count stays 1 assembleStreamingPipeline(baseArgs(), deps); assert.equal(performance.getEntriesByName("omni-pipeline-start").length, 1, "start cleared"); assert.equal(performance.getEntriesByName("omni-pipeline-end").length, 1, "end cleared"); assert.equal(performance.getEntriesByName("omni-pipeline").length, 1, "measure cleared"); }); test("performance measure has positive duration", () => { performance.clearMarks(); performance.clearMeasures(); const { deps } = makeDeps(); assembleStreamingPipeline(baseArgs(), deps); const [entry] = performance.getEntriesByName("omni-pipeline"); assert.ok(entry, "measure exists"); assert.equal(entry.entryType, "measure"); assert.ok(entry.duration >= 0, `duration ${entry.duration} >= 0`); });