import { test } from "node:test"; import assert from "node:assert/strict"; import http from "node:http"; import type { AddressInfo } from "node:net"; import type { Worker } from "../lib/workers"; import { TranscribeError } from "./transcribeError"; import { runUnitViaRemote } from "./remoteUnit"; // Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test controller/remoteUnit.test.ts // // The unit client against a real node:http stub — no network, no fetch mocks. // What is pinned: the submit → poll → result → DELETE round trip; the failure // CLASSES (5xx/unreachable/silence → transport, so the batch retries the item // elsewhere instead of counting it failed; a 4xx refusal or a "failed" job → // work-class, final); and that the scratch DELETE fires even on success. function worker(baseUrl: string): Worker { return { id: "exec", name: "exec", kind: "remote", enabled: true, priority: 0, tags: ["cpu"], remote: { baseUrl, token: "sekrit" }, }; } type Stub = { baseUrl: string; requests: Array<{ method: string; url: string; body: string }>; close: () => Promise; }; async function startStub(handler: http.RequestListener): Promise { const requests: Stub["requests"] = []; const server = http.createServer((req, res) => { let body = ""; req.on("data", (c) => (body += c)); req.on("end", () => { requests.push({ method: req.method ?? "", url: req.url ?? "", body }); handler(req, res); }); }); await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); const { port } = server.address() as AddressInfo; return { baseUrl: `http://127.0.0.1:${port}`, requests, close: () => new Promise((resolve) => server.close(() => resolve())), }; } function json(res: http.ServerResponse, status: number, body: unknown): void { res.statusCode = status; res.setHeader("content-type", "application/json"); res.end(JSON.stringify(body)); } test("submit → poll → result → DELETE round trip", async () => { let polls = 0; const stub = await startStub((req, res) => { if (req.method === "POST" && req.url === "/api/worker/unit") { return json(res, 202, { remoteJobId: "job1" }); } if (req.url === "/api/worker/unit/job1/events") { polls++; return json(res, 200, { status: polls < 2 ? "running" : "done" }); } if (req.url === "/api/worker/unit/job1/result") { return json(res, 200, { outcome: "done", files: { "attribution.json": "{\"videoId\":\"v\"}" }, }); } if (req.method === "DELETE") return json(res, 200, { ok: true }); return json(res, 404, {}); }); try { const result = await runUnitViaRemote({ worker: worker(stub.baseUrl), op: "attribution-text", channelSlug: "chan", videoId: "vid", files: { "transcript.cues.json": Buffer.from("{}") }, target: {}, pollMs: 10, pollTimeoutMs: 5000, }); assert.equal(result.outcome, "done"); assert.equal(result.files["attribution.json"], '{"videoId":"v"}'); const submit = stub.requests.find((r) => r.method === "POST")!; const body = JSON.parse(submit.body) as { op: string; files: Record; }; assert.equal(body.op, "attribution-text"); assert.equal( Buffer.from(body.files["transcript.cues.json"], "base64").toString(), "{}", ); // The fire-and-forget cleanup reaches the stub shortly after. await new Promise((r) => setTimeout(r, 50)); assert.ok(stub.requests.some((r) => r.method === "DELETE")); } finally { await stub.close(); } }); test("a 5xx submit is a TRANSPORT failure (retry elsewhere)", async () => { const stub = await startStub((_req, res) => json(res, 500, { error: "boom" })); try { await assert.rejects( runUnitViaRemote({ worker: worker(stub.baseUrl), op: "diarization", channelSlug: "chan", videoId: "vid", files: {}, target: {}, pollMs: 10, pollTimeoutMs: 200, }), (err: unknown) => err instanceof TranscribeError && err.failureClass === "transport", ); } finally { await stub.close(); } }); test("a 4xx refusal is a WORK-class failure (no retry)", async () => { const stub = await startStub((_req, res) => json(res, 400, { error: "not a backfill kind" }), ); try { await assert.rejects( runUnitViaRemote({ worker: worker(stub.baseUrl), op: "download", channelSlug: "chan", videoId: "vid", files: {}, target: {}, pollMs: 10, pollTimeoutMs: 200, }), (err: unknown) => err instanceof TranscribeError && err.failureClass === "transcription", ); } finally { await stub.close(); } }); test("a remote job that finishes failed is a WORK-class failure", async () => { const stub = await startStub((req, res) => { if (req.method === "POST" && req.url === "/api/worker/unit") { return json(res, 202, { remoteJobId: "job1" }); } if (req.url?.endsWith("/events")) return json(res, 200, { status: "failed" }); return json(res, 200, { ok: true }); }); try { await assert.rejects( runUnitViaRemote({ worker: worker(stub.baseUrl), op: "attribution-text", channelSlug: "chan", videoId: "vid", files: {}, target: {}, pollMs: 10, pollTimeoutMs: 5000, }), (err: unknown) => err instanceof TranscribeError && err.failureClass === "transcription", ); } finally { await stub.close(); } }); test("a server that goes silent after submit times out as TRANSPORT", async () => { const stub = await startStub((req, res) => { if (req.method === "POST" && req.url === "/api/worker/unit") { return json(res, 202, { remoteJobId: "job1" }); } if (req.method === "DELETE") return json(res, 200, { ok: true }); // Every poll errors — the client keeps polling until its silence timeout. return json(res, 500, {}); }); try { await assert.rejects( runUnitViaRemote({ worker: worker(stub.baseUrl), op: "attribution-text", channelSlug: "chan", videoId: "vid", files: {}, target: {}, pollMs: 10, pollTimeoutMs: 100, }), (err: unknown) => err instanceof TranscribeError && err.failureClass === "transport" && /stopped responding/.test(err.message), ); } finally { await stub.close(); } });