// Regression: when the consumer disconnects from a streaming server action // (e.g. user navigates away mid-job), the producer kept calling // controller.enqueue on the closed stream and uncaughtException leaked out // of the RSC encoder with "Controller is already closed". streamCommand.ts // now nullifies the controller and aborts the producer on cancel; this spec // guards against the regression. import { test, expect } from "@playwright/test"; import { channelStage, generateReport, resetData, } from "./helpers"; import { baseUrl } from "./baseUrl"; async function readUncaughtCount(): Promise<{ uncaught: number; unhandled: number; messages: string[]; }> { const res = await fetch(`${baseUrl}/api/test/uncaught-count`); return res.json() as Promise<{ uncaught: number; unhandled: number; messages: string[]; }>; } async function clearUncaughtCount(): Promise { await fetch(`${baseUrl}/api/test/uncaught-count`, { method: "DELETE" }); } test("disconnecting from a running job stream does not crash with 'Controller is already closed'", async ({ page, }) => { test.setTimeout(60_000); await resetData("slow-pipeline-channel"); await clearUncaughtCount(); // Start the slow sync — fake-ytdlp sleeps 30s before producing output, so // the producer is alive and pushing into the stream for the duration. await generateReport(page, "slow-channel"); await page.goto(channelStage("slow-channel", "playlist")); await page.getByRole("button", { name: "Sync" }).click(); await expect(page.getByLabel("Sync output")).toContainText("test-slow", { timeout: 15_000, }); // Simulate consumer disconnect by navigating to a different page. The // RSC response stream attached to the previous page is torn down, which // cancels our underlying ReadableStream source. The job itself keeps // running on the server (writes to disk) — that's intentional, since // another client could reconnect. await page.goto("/jobs"); // Cancel the still-running job via the jobs page so the test doesn't // wait the full 30s for fake-ytdlp to wake up. The actual regression // signal is the uncaught-counter check after this — which catches the // "Controller is already closed" error that the prior implementation // produced whenever the source kept enqueueing post-disconnect. const jobRow = page .getByRole("row") .filter({ hasText: "slow-channel" }) .first(); await expect(jobRow).toContainText("running", { timeout: 10_000 }); await jobRow.getByRole("button", { name: /^Cancel$/ }).click(); await expect(jobRow).toContainText("cancelled", { timeout: 15_000 }); // Settle so any post-cancel async cleanup has a chance to fire its error // path before we check the counter. await page.waitForTimeout(1_000); const counts = await readUncaughtCount(); expect( counts.messages.filter((m) => m.includes("Controller is already closed")), ).toEqual([]); expect(counts.uncaught).toBe(0); });