mirror of
https://github.com/heygen-com/hyperframes.git
synced 2026-09-10 22:20:14 +00:00
refactor(render): centralize execution ownership (#2158)
This commit is contained in:
@@ -1,144 +1,7 @@
|
|||||||
/**
|
/** Tests for render failure-detail construction. */
|
||||||
* Tests for the cancel/error-path helpers in `./cleanup.ts`.
|
|
||||||
*/
|
|
||||||
|
|
||||||
import { describe, expect, it, vi } from "vitest";
|
import { describe, expect, it } from "vitest";
|
||||||
import { mkdtempSync, mkdirSync, rmSync, writeFileSync, existsSync } from "node:fs";
|
import { buildRenderErrorDetails } from "./cleanup.js";
|
||||||
import { tmpdir } from "node:os";
|
|
||||||
import { join } from "node:path";
|
|
||||||
import type { CaptureSession } from "@hyperframes/engine";
|
|
||||||
import type { FileServerHandle } from "../fileServer.js";
|
|
||||||
import { buildRenderErrorDetails, cleanupRenderResources, safeCleanup } from "./cleanup.js";
|
|
||||||
|
|
||||||
function makeLog() {
|
|
||||||
return { error: vi.fn(), warn: vi.fn(), info: vi.fn(), debug: vi.fn() };
|
|
||||||
}
|
|
||||||
|
|
||||||
describe("safeCleanup", () => {
|
|
||||||
it("returns normally when the operation succeeds", async () => {
|
|
||||||
const log = makeLog();
|
|
||||||
const op = vi.fn().mockResolvedValue(undefined);
|
|
||||||
await safeCleanup("close x", op, log);
|
|
||||||
expect(op).toHaveBeenCalledOnce();
|
|
||||||
expect(log.debug).not.toHaveBeenCalled();
|
|
||||||
});
|
|
||||||
|
|
||||||
it("swallows thrown errors and logs them at debug", async () => {
|
|
||||||
const log = makeLog();
|
|
||||||
await safeCleanup(
|
|
||||||
"close x",
|
|
||||||
() => {
|
|
||||||
throw new Error("boom");
|
|
||||||
},
|
|
||||||
log,
|
|
||||||
);
|
|
||||||
expect(log.debug).toHaveBeenCalledWith("Cleanup failed (close x)", { error: "boom" });
|
|
||||||
});
|
|
||||||
|
|
||||||
it("swallows async rejections", async () => {
|
|
||||||
const log = makeLog();
|
|
||||||
await safeCleanup("close x", async () => Promise.reject(new Error("async boom")), log);
|
|
||||||
expect(log.debug).toHaveBeenCalledWith("Cleanup failed (close x)", { error: "async boom" });
|
|
||||||
});
|
|
||||||
});
|
|
||||||
|
|
||||||
describe("cleanupRenderResources", () => {
|
|
||||||
it("closes fileServer, probeSession, then removes workDir (non-debug)", async () => {
|
|
||||||
const log = makeLog();
|
|
||||||
const workDir = mkdtempSync(join(tmpdir(), "cleanup-test-"));
|
|
||||||
writeFileSync(join(workDir, "marker.txt"), "x");
|
|
||||||
|
|
||||||
const order: string[] = [];
|
|
||||||
const fileServer = {
|
|
||||||
close: () => {
|
|
||||||
order.push("fileServer.close");
|
|
||||||
},
|
|
||||||
} as unknown as FileServerHandle;
|
|
||||||
const probeSession = {
|
|
||||||
_markClosed: () => {
|
|
||||||
order.push("probeSession.close");
|
|
||||||
},
|
|
||||||
} as unknown as CaptureSession;
|
|
||||||
|
|
||||||
// closeCaptureSession is the engine helper; the helper itself isn't
|
|
||||||
// mockable per-call without intercepting the module import. Instead
|
|
||||||
// we verify the higher-level invariants: fileServer.close was called,
|
|
||||||
// and the workDir was rmSync'd.
|
|
||||||
await cleanupRenderResources({
|
|
||||||
fileServer,
|
|
||||||
probeSession: null, // skip probe to keep this test focused on the workDir invariant
|
|
||||||
workDir,
|
|
||||||
debug: false,
|
|
||||||
log,
|
|
||||||
label: "cancel",
|
|
||||||
});
|
|
||||||
|
|
||||||
expect(order).toEqual(["fileServer.close"]);
|
|
||||||
expect(existsSync(workDir)).toBe(false);
|
|
||||||
void probeSession; // suppress unused-var (kept as a doc of the surface)
|
|
||||||
});
|
|
||||||
|
|
||||||
it("keeps workDir when debug=true", async () => {
|
|
||||||
const log = makeLog();
|
|
||||||
const workDir = mkdtempSync(join(tmpdir(), "cleanup-debug-"));
|
|
||||||
writeFileSync(join(workDir, "marker.txt"), "x");
|
|
||||||
|
|
||||||
await cleanupRenderResources({
|
|
||||||
fileServer: null,
|
|
||||||
probeSession: null,
|
|
||||||
workDir,
|
|
||||||
debug: true,
|
|
||||||
log,
|
|
||||||
label: "error",
|
|
||||||
});
|
|
||||||
|
|
||||||
expect(existsSync(workDir)).toBe(true);
|
|
||||||
rmSync(workDir, { recursive: true, force: true });
|
|
||||||
});
|
|
||||||
|
|
||||||
it("is a no-op for missing workDir thanks to rmSync force:true", async () => {
|
|
||||||
const log = makeLog();
|
|
||||||
const workDir = join(tmpdir(), `cleanup-missing-${Date.now()}`);
|
|
||||||
expect(existsSync(workDir)).toBe(false);
|
|
||||||
|
|
||||||
await cleanupRenderResources({
|
|
||||||
fileServer: null,
|
|
||||||
probeSession: null,
|
|
||||||
workDir,
|
|
||||||
debug: false,
|
|
||||||
log,
|
|
||||||
label: "error",
|
|
||||||
});
|
|
||||||
|
|
||||||
// No throw; nothing logged at debug for the rmSync step.
|
|
||||||
expect(log.debug).not.toHaveBeenCalled();
|
|
||||||
});
|
|
||||||
|
|
||||||
it("logs (and continues past) a fileServer.close that throws", async () => {
|
|
||||||
const log = makeLog();
|
|
||||||
const workDir = mkdtempSync(join(tmpdir(), "cleanup-throw-"));
|
|
||||||
|
|
||||||
const fileServer = {
|
|
||||||
close: () => {
|
|
||||||
throw new Error("server stuck");
|
|
||||||
},
|
|
||||||
} as unknown as FileServerHandle;
|
|
||||||
|
|
||||||
await cleanupRenderResources({
|
|
||||||
fileServer,
|
|
||||||
probeSession: null,
|
|
||||||
workDir,
|
|
||||||
debug: false,
|
|
||||||
log,
|
|
||||||
label: "error",
|
|
||||||
});
|
|
||||||
|
|
||||||
expect(log.debug).toHaveBeenCalledWith("Cleanup failed (close file server (error))", {
|
|
||||||
error: "server stuck",
|
|
||||||
});
|
|
||||||
expect(existsSync(workDir)).toBe(false);
|
|
||||||
});
|
|
||||||
});
|
|
||||||
|
|
||||||
describe("buildRenderErrorDetails", () => {
|
describe("buildRenderErrorDetails", () => {
|
||||||
const baseDiagnostics = { videoExtractionFailures: 0, imageDecodeFailures: 0 };
|
const baseDiagnostics = { videoExtractionFailures: 0, imageDecodeFailures: 0 };
|
||||||
@@ -235,6 +98,3 @@ describe("buildRenderErrorDetails", () => {
|
|||||||
expect(failed.hdrDiagnostics).toEqual({ videoExtractionFailures: 2, imageDecodeFailures: 0 });
|
expect(failed.hdrDiagnostics).toEqual({ videoExtractionFailures: 2, imageDecodeFailures: 0 });
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
// Quiet unused-import warning — these are referenced via type-only paths.
|
|
||||||
void mkdirSync;
|
|
||||||
|
|||||||
@@ -1,80 +1,11 @@
|
|||||||
/**
|
/** Error-detail construction shared by render failure paths. */
|
||||||
* Sequencer cleanup + error-details helpers shared by the cancel and
|
|
||||||
* error paths in `executeRenderJob`.
|
|
||||||
*/
|
|
||||||
|
|
||||||
import { rmSync } from "node:fs";
|
|
||||||
import { freemem } from "node:os";
|
import { freemem } from "node:os";
|
||||||
import {
|
import { type SubTimelineWaitOutcome } from "@hyperframes/engine";
|
||||||
type CaptureSession,
|
|
||||||
type SubTimelineWaitOutcome,
|
|
||||||
closeCaptureSession,
|
|
||||||
} from "@hyperframes/engine";
|
|
||||||
import type { FileServerHandle } from "../fileServer.js";
|
|
||||||
import { defaultLogger, type ProducerLogger } from "../../logger.js";
|
|
||||||
import type { HdrDiagnostics, RenderJob } from "../renderOrchestrator.js";
|
import type { HdrDiagnostics, RenderJob } from "../renderOrchestrator.js";
|
||||||
import { normalizeErrorMessage } from "../../utils/errorMessage.js";
|
import { normalizeErrorMessage } from "../../utils/errorMessage.js";
|
||||||
import type { RenderObservabilitySummary } from "./observability.js";
|
import type { RenderObservabilitySummary } from "./observability.js";
|
||||||
|
|
||||||
/**
|
|
||||||
* Wrap a cleanup operation so it never throws, but logs any failure.
|
|
||||||
* The sequencer needs to keep tearing down resources even when one of
|
|
||||||
* them is stuck (e.g. a `fileServer.close()` hitting a TCP race); a
|
|
||||||
* thrown cleanup error would mask the original render failure.
|
|
||||||
*/
|
|
||||||
export async function safeCleanup(
|
|
||||||
label: string,
|
|
||||||
fn: () => Promise<void> | void,
|
|
||||||
log: ProducerLogger = defaultLogger,
|
|
||||||
): Promise<void> {
|
|
||||||
try {
|
|
||||||
await fn();
|
|
||||||
} catch (err) {
|
|
||||||
log.debug(`Cleanup failed (${label})`, {
|
|
||||||
error: err instanceof Error ? err.message : String(err),
|
|
||||||
});
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Close the file server, close the probe session, and remove the
|
|
||||||
* working directory. Each step runs through `safeCleanup` so a stuck
|
|
||||||
* resource doesn't mask the original render error.
|
|
||||||
*/
|
|
||||||
export async function cleanupRenderResources(input: {
|
|
||||||
fileServer: FileServerHandle | null;
|
|
||||||
probeSession: CaptureSession | null;
|
|
||||||
workDir: string;
|
|
||||||
debug: boolean;
|
|
||||||
log: ProducerLogger;
|
|
||||||
/** Suffix appended to safeCleanup labels. Pinned to the existing diagnostic payloads. */
|
|
||||||
label: "cancel" | "error";
|
|
||||||
}): Promise<void> {
|
|
||||||
const { fileServer, probeSession, workDir, debug, log, label } = input;
|
|
||||||
if (fileServer) {
|
|
||||||
const fs = fileServer;
|
|
||||||
await safeCleanup(
|
|
||||||
`close file server (${label})`,
|
|
||||||
() => {
|
|
||||||
fs.close();
|
|
||||||
},
|
|
||||||
log,
|
|
||||||
);
|
|
||||||
}
|
|
||||||
if (probeSession) {
|
|
||||||
const session = probeSession;
|
|
||||||
await safeCleanup(`close probe session (${label})`, () => closeCaptureSession(session), log);
|
|
||||||
}
|
|
||||||
if (!debug) {
|
|
||||||
// `force: true` swallows ENOENT, so no need to existsSync first.
|
|
||||||
await safeCleanup(
|
|
||||||
`remove workDir (${label})`,
|
|
||||||
() => rmSync(workDir, { recursive: true, force: true, maxRetries: 3, retryDelay: 100 }),
|
|
||||||
log,
|
|
||||||
);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Build the `RenderJob.errorDetails` shape downstream consumers (SSE,
|
* Build the `RenderJob.errorDetails` shape downstream consumers (SSE,
|
||||||
* sync `/render` response, queue introspection) read on failure.
|
* sync `/render` response, queue introspection) read on failure.
|
||||||
|
|||||||
@@ -0,0 +1,137 @@
|
|||||||
|
import { describe, expect, it, vi } from "vitest";
|
||||||
|
import type { ProducerLogger } from "../../logger.js";
|
||||||
|
import type { RenderJob } from "../renderOrchestrator.js";
|
||||||
|
import { RenderExecutionContext } from "./renderExecutionContext.js";
|
||||||
|
|
||||||
|
function logger(): ProducerLogger {
|
||||||
|
return { error: vi.fn(), warn: vi.fn(), info: vi.fn(), debug: vi.fn() };
|
||||||
|
}
|
||||||
|
|
||||||
|
function job(): RenderJob {
|
||||||
|
return {
|
||||||
|
id: "job-1",
|
||||||
|
status: "queued",
|
||||||
|
outcome: null,
|
||||||
|
warnings: [],
|
||||||
|
progress: 0,
|
||||||
|
currentStage: "queued",
|
||||||
|
config: { fps: { num: 30, den: 1 }, quality: "standard" },
|
||||||
|
createdAt: new Date(),
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
describe("RenderExecutionContext", () => {
|
||||||
|
it("freezes request identity and scopes every log record", () => {
|
||||||
|
const base = logger();
|
||||||
|
const context = new RenderExecutionContext({
|
||||||
|
request: { renderJobId: "job-1", projectDir: "/project", outputPath: "/out.mp4" },
|
||||||
|
logger: base,
|
||||||
|
});
|
||||||
|
|
||||||
|
expect(Object.isFrozen(context.request)).toBe(true);
|
||||||
|
context.logger.info("started", { phase: "compile" });
|
||||||
|
expect(base.info).toHaveBeenCalledWith("started", {
|
||||||
|
renderJobId: "job-1",
|
||||||
|
phase: "compile",
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
it("runs disposers once in reverse acquisition order and contains failures", async () => {
|
||||||
|
const base = logger();
|
||||||
|
const calls: string[] = [];
|
||||||
|
const context = new RenderExecutionContext({
|
||||||
|
request: { renderJobId: "job-1", projectDir: "/project", outputPath: "/out.mp4" },
|
||||||
|
logger: base,
|
||||||
|
});
|
||||||
|
context.defer("first", () => calls.push("first"));
|
||||||
|
context.defer("broken", () => {
|
||||||
|
calls.push("broken");
|
||||||
|
throw new Error("close failed");
|
||||||
|
});
|
||||||
|
context.defer("last", async () => calls.push("last"));
|
||||||
|
|
||||||
|
await Promise.all([context.dispose(), context.dispose()]);
|
||||||
|
|
||||||
|
expect(calls).toEqual(["last", "broken", "first"]);
|
||||||
|
expect(base.debug).toHaveBeenCalledWith(
|
||||||
|
"Cleanup failed (broken)",
|
||||||
|
expect.objectContaining({ renderJobId: "job-1", error: "close failed" }),
|
||||||
|
);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("can release ownership after a resource is closed early", async () => {
|
||||||
|
const dispose = vi.fn();
|
||||||
|
const context = new RenderExecutionContext({
|
||||||
|
request: { renderJobId: "job-1", projectDir: "/project", outputPath: "/out.mp4" },
|
||||||
|
logger: logger(),
|
||||||
|
});
|
||||||
|
const release = context.defer("resource", dispose);
|
||||||
|
release();
|
||||||
|
|
||||||
|
await context.dispose();
|
||||||
|
|
||||||
|
expect(dispose).not.toHaveBeenCalled();
|
||||||
|
});
|
||||||
|
|
||||||
|
it("serializes progress before disposal completes", async () => {
|
||||||
|
const deliveries: string[] = [];
|
||||||
|
const context = new RenderExecutionContext({
|
||||||
|
request: { renderJobId: "job-1", projectDir: "/project", outputPath: "/out.mp4" },
|
||||||
|
logger: logger(),
|
||||||
|
progressSink: async (_job, message) => {
|
||||||
|
await Promise.resolve();
|
||||||
|
deliveries.push(message);
|
||||||
|
},
|
||||||
|
});
|
||||||
|
context.defer("resource", () => deliveries.push("disposed"));
|
||||||
|
const renderJob = job();
|
||||||
|
context.onProgress?.(renderJob, "first");
|
||||||
|
context.onProgress?.(renderJob, "second");
|
||||||
|
|
||||||
|
await context.dispose();
|
||||||
|
|
||||||
|
expect(deliveries).toEqual(["first", "second", "disposed"]);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("does not let a broken cleanup logger reject disposal", async () => {
|
||||||
|
const base = logger();
|
||||||
|
vi.mocked(base.debug).mockImplementation(() => {
|
||||||
|
throw new Error("logger unavailable");
|
||||||
|
});
|
||||||
|
const context = new RenderExecutionContext({
|
||||||
|
request: { renderJobId: "job-1", projectDir: "/project", outputPath: "/out.mp4" },
|
||||||
|
logger: base,
|
||||||
|
});
|
||||||
|
context.defer("broken", () => {
|
||||||
|
throw new Error("close failed");
|
||||||
|
});
|
||||||
|
|
||||||
|
await expect(context.dispose()).resolves.toBeUndefined();
|
||||||
|
});
|
||||||
|
|
||||||
|
it("composes caller cancellation with a deadline", async () => {
|
||||||
|
const controller = new AbortController();
|
||||||
|
const context = new RenderExecutionContext({
|
||||||
|
request: { renderJobId: "job-1", projectDir: "/project", outputPath: "/out.mp4" },
|
||||||
|
logger: logger(),
|
||||||
|
signal: controller.signal,
|
||||||
|
deadlineAtMs: Date.now() + 60_000,
|
||||||
|
});
|
||||||
|
controller.abort();
|
||||||
|
|
||||||
|
expect(() => context.assertActive()).toThrow("render_cancelled");
|
||||||
|
await context.dispose();
|
||||||
|
});
|
||||||
|
|
||||||
|
it("is already aborted when its deadline has elapsed", async () => {
|
||||||
|
const context = new RenderExecutionContext({
|
||||||
|
request: { renderJobId: "job-1", projectDir: "/project", outputPath: "/out.mp4" },
|
||||||
|
logger: logger(),
|
||||||
|
deadlineAtMs: Date.now() - 1,
|
||||||
|
});
|
||||||
|
|
||||||
|
expect(context.signal?.aborted).toBe(true);
|
||||||
|
expect(() => context.assertActive()).toThrow("render_cancelled");
|
||||||
|
await context.dispose();
|
||||||
|
});
|
||||||
|
});
|
||||||
@@ -0,0 +1,111 @@
|
|||||||
|
import type { ProducerLogger } from "../../logger.js";
|
||||||
|
import type { ProgressCallback } from "../renderOrchestrator.js";
|
||||||
|
import { OrderedRenderEventPublisher } from "./renderEventPublisher.js";
|
||||||
|
|
||||||
|
export interface RenderExecutionRequestIdentity {
|
||||||
|
renderJobId: string;
|
||||||
|
projectDir: string;
|
||||||
|
outputPath: string;
|
||||||
|
}
|
||||||
|
|
||||||
|
export type RenderDisposer = () => Promise<void> | void;
|
||||||
|
|
||||||
|
interface RegisteredDisposer {
|
||||||
|
active: boolean;
|
||||||
|
label: string;
|
||||||
|
dispose: RenderDisposer;
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface RenderExecutionContextOptions {
|
||||||
|
request: RenderExecutionRequestIdentity;
|
||||||
|
logger: ProducerLogger;
|
||||||
|
progressSink?: ProgressCallback;
|
||||||
|
signal?: AbortSignal;
|
||||||
|
deadlineAtMs?: number;
|
||||||
|
}
|
||||||
|
|
||||||
|
function scopedLogger(
|
||||||
|
logger: ProducerLogger,
|
||||||
|
fields: Readonly<Record<string, unknown>>,
|
||||||
|
): ProducerLogger {
|
||||||
|
const merge = (meta?: Record<string, unknown>) => ({ ...fields, ...meta });
|
||||||
|
return {
|
||||||
|
error: (message, meta) => logger.error(message, merge(meta)),
|
||||||
|
warn: (message, meta) => logger.warn(message, merge(meta)),
|
||||||
|
info: (message, meta) => logger.info(message, merge(meta)),
|
||||||
|
debug: (message, meta) => logger.debug(message, merge(meta)),
|
||||||
|
isLevelEnabled: logger.isLevelEnabled?.bind(logger),
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
function composeSignal(signal?: AbortSignal, deadlineAtMs?: number): AbortSignal | undefined {
|
||||||
|
if (deadlineAtMs === undefined) return signal;
|
||||||
|
const remainingMs = Math.max(0, deadlineAtMs - Date.now());
|
||||||
|
const deadlineSignal =
|
||||||
|
remainingMs === 0
|
||||||
|
? AbortSignal.abort(new Error("render_deadline_exceeded"))
|
||||||
|
: AbortSignal.timeout(Math.min(remainingMs, 2_147_483_647));
|
||||||
|
return signal ? AbortSignal.any([signal, deadlineSignal]) : deadlineSignal;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Owns one render execution's control-plane resources and terminal cleanup. */
|
||||||
|
export class RenderExecutionContext {
|
||||||
|
readonly request: Readonly<RenderExecutionRequestIdentity>;
|
||||||
|
readonly signal: AbortSignal | undefined;
|
||||||
|
readonly logger: ProducerLogger;
|
||||||
|
readonly events: OrderedRenderEventPublisher;
|
||||||
|
readonly onProgress: ProgressCallback | undefined;
|
||||||
|
|
||||||
|
private readonly disposers: RegisteredDisposer[] = [];
|
||||||
|
private disposePromise: Promise<void> | null = null;
|
||||||
|
|
||||||
|
constructor(options: RenderExecutionContextOptions) {
|
||||||
|
this.request = Object.freeze({ ...options.request });
|
||||||
|
this.signal = composeSignal(options.signal, options.deadlineAtMs);
|
||||||
|
this.logger = scopedLogger(options.logger, { renderJobId: options.request.renderJobId });
|
||||||
|
this.events = new OrderedRenderEventPublisher(options.progressSink, this.logger);
|
||||||
|
this.onProgress = options.progressSink
|
||||||
|
? (job, message) => this.events.publish(job, message)
|
||||||
|
: undefined;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Register cleanup in acquisition order; execution is LIFO and exactly once. */
|
||||||
|
defer(label: string, dispose: RenderDisposer): () => void {
|
||||||
|
if (this.disposePromise) throw new Error("Cannot register a disposer after context disposal");
|
||||||
|
const entry: RegisteredDisposer = { active: true, label, dispose };
|
||||||
|
this.disposers.push(entry);
|
||||||
|
return () => {
|
||||||
|
entry.active = false;
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
assertActive(createError: () => Error = () => new Error("render_cancelled")): void {
|
||||||
|
if (this.signal?.aborted) throw createError();
|
||||||
|
}
|
||||||
|
|
||||||
|
dispose(): Promise<void> {
|
||||||
|
if (this.disposePromise) return this.disposePromise;
|
||||||
|
this.disposePromise = this.disposeOnce();
|
||||||
|
return this.disposePromise;
|
||||||
|
}
|
||||||
|
|
||||||
|
private async disposeOnce(): Promise<void> {
|
||||||
|
await this.events.flush();
|
||||||
|
for (const entry of this.disposers.reverse()) {
|
||||||
|
if (!entry.active) continue;
|
||||||
|
entry.active = false;
|
||||||
|
try {
|
||||||
|
await entry.dispose();
|
||||||
|
} catch (error) {
|
||||||
|
try {
|
||||||
|
this.logger.debug(`Cleanup failed (${entry.label})`, {
|
||||||
|
error: error instanceof Error ? error.message : String(error),
|
||||||
|
});
|
||||||
|
} catch {
|
||||||
|
// Cleanup and its diagnostics are both best-effort: neither may mask
|
||||||
|
// the render outcome that caused disposal.
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -100,11 +100,9 @@ import {
|
|||||||
} from "./fileServer.js";
|
} from "./fileServer.js";
|
||||||
import { defaultLogger, type ProducerLogger } from "../logger.js";
|
import { defaultLogger, type ProducerLogger } from "../logger.js";
|
||||||
import { createMemorySampler, type MemorySampler, updateJobStatus } from "./render/shared.js";
|
import { createMemorySampler, type MemorySampler, updateJobStatus } from "./render/shared.js";
|
||||||
import { buildRenderErrorDetails, cleanupRenderResources, safeCleanup } from "./render/cleanup.js";
|
import { buildRenderErrorDetails } from "./render/cleanup.js";
|
||||||
import {
|
import { publishRenderFailure } from "./render/renderEventPublisher.js";
|
||||||
OrderedRenderEventPublisher,
|
import { RenderExecutionContext } from "./render/renderExecutionContext.js";
|
||||||
publishRenderFailure,
|
|
||||||
} from "./render/renderEventPublisher.js";
|
|
||||||
import { ArtifactTransaction } from "./render/artifactTransaction.js";
|
import { ArtifactTransaction } from "./render/artifactTransaction.js";
|
||||||
import { normalizeErrorMessage } from "../utils/errorMessage.js";
|
import { normalizeErrorMessage } from "../utils/errorMessage.js";
|
||||||
import { formatCaptureFrameName } from "../utils/paths.js";
|
import { formatCaptureFrameName } from "../utils/paths.js";
|
||||||
@@ -1576,11 +1574,51 @@ export async function executeRenderJob(
|
|||||||
const pipelineStart = Date.now();
|
const pipelineStart = Date.now();
|
||||||
const baseLog = job.config.logger ?? defaultLogger;
|
const baseLog = job.config.logger ?? defaultLogger;
|
||||||
const logPath = job.config.debug ? join(workDir, "render.log") : null;
|
const logPath = job.config.debug ? join(workDir, "render.log") : null;
|
||||||
const log = logPath ? createRenderFileLogger(logPath, baseLog) : baseLog;
|
const execution = new RenderExecutionContext({
|
||||||
const eventPublisher = new OrderedRenderEventPublisher(progressSink, log);
|
request: { renderJobId: job.id, projectDir, outputPath },
|
||||||
const onProgress: ProgressCallback | undefined = progressSink
|
logger: logPath ? createRenderFileLogger(logPath, baseLog) : baseLog,
|
||||||
? (progressJob, message) => eventPublisher.publish(progressJob, message)
|
progressSink,
|
||||||
: undefined;
|
signal: abortSignal,
|
||||||
|
});
|
||||||
|
const log = execution.logger;
|
||||||
|
execution.defer("remove workDir", () => {
|
||||||
|
if (job.config.debug) return;
|
||||||
|
if (job.status === "complete" && process.env.KEEP_TEMP === "1") {
|
||||||
|
log.info("KEEP_TEMP=1 — leaving workDir on disk for inspection", { workDir });
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
rmSync(workDir, { recursive: true, force: true, maxRetries: 3, retryDelay: 100 });
|
||||||
|
});
|
||||||
|
|
||||||
|
try {
|
||||||
|
await executeRenderPipeline({
|
||||||
|
job,
|
||||||
|
projectDir,
|
||||||
|
outputPath,
|
||||||
|
workDir,
|
||||||
|
logPath,
|
||||||
|
pipelineStart,
|
||||||
|
execution,
|
||||||
|
});
|
||||||
|
} finally {
|
||||||
|
await execution.dispose();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
async function executeRenderPipeline(input: {
|
||||||
|
job: RenderJob;
|
||||||
|
projectDir: string;
|
||||||
|
outputPath: string;
|
||||||
|
workDir: string;
|
||||||
|
logPath: string | null;
|
||||||
|
pipelineStart: number;
|
||||||
|
execution: RenderExecutionContext;
|
||||||
|
}): Promise<void> {
|
||||||
|
const { job, projectDir, outputPath, workDir, logPath, pipelineStart, execution } = input;
|
||||||
|
const log = execution.logger;
|
||||||
|
const eventPublisher = execution.events;
|
||||||
|
const onProgress = execution.onProgress;
|
||||||
|
const executionSignal = execution.signal;
|
||||||
let fileServer: FileServerHandle | null = null;
|
let fileServer: FileServerHandle | null = null;
|
||||||
let probeSession: CaptureSession | null = null;
|
let probeSession: CaptureSession | null = null;
|
||||||
let lastBrowserConsole: string[] = [];
|
let lastBrowserConsole: string[] = [];
|
||||||
@@ -1652,21 +1690,35 @@ export async function executeRenderJob(
|
|||||||
const count = captureAttempts.filter((a) => a.reason === "transient-retry").length;
|
const count = captureAttempts.filter((a) => a.reason === "transient-retry").length;
|
||||||
if (count > 0) updateCaptureObservability({ transientRetries: count });
|
if (count > 0) updateCaptureObservability({ transientRetries: count });
|
||||||
};
|
};
|
||||||
// Declared outside the try so `finally` can stop the interval, but
|
// The execution context's dynamic disposer reads this binding, so any
|
||||||
// the sampler is created INSIDE the try so a synchronous throw
|
// sampler acquired by the pipeline is stopped by the unconditional outer
|
||||||
// between declaration and the try-block (currently impossible, but
|
// finally even when setup or terminal reporting throws.
|
||||||
// defensible if more setup ever lands here) can't leak the interval.
|
|
||||||
let memSampler: MemorySampler | null = null;
|
let memSampler: MemorySampler | null = null;
|
||||||
// "routed" = the parallel router fired and held; "reverted" = fired but
|
// "routed" = the parallel router fired and held; "reverted" = fired but
|
||||||
// the self-verify retry rolled back; undefined = never fired.
|
// the self-verify retry rolled back; undefined = never fired.
|
||||||
let deParallelRouter: "routed" | "reverted" | undefined;
|
let deParallelRouter: "routed" | "reverted" | undefined;
|
||||||
|
|
||||||
|
execution.defer("rollback staged artifact", () => artifactTransaction.rollback());
|
||||||
|
execution.defer("close file server", () => {
|
||||||
|
if (!fileServer) return;
|
||||||
|
closeFileServerSafely(fileServer, "renderExecutionContext", log);
|
||||||
|
fileServer = null;
|
||||||
|
});
|
||||||
|
execution.defer("close probe session", async () => {
|
||||||
|
if (!probeSession) return;
|
||||||
|
const session = probeSession;
|
||||||
|
probeSession = null;
|
||||||
|
await closeCaptureSession(session);
|
||||||
|
});
|
||||||
|
execution.defer("stop memory sampler", () => {
|
||||||
|
memSampler?.stop();
|
||||||
|
memSampler = null;
|
||||||
|
});
|
||||||
|
|
||||||
try {
|
try {
|
||||||
memSampler = createMemorySampler();
|
memSampler = createMemorySampler();
|
||||||
const assertNotAborted = () => {
|
const assertNotAborted = () => {
|
||||||
if (abortSignal?.aborted) {
|
execution.assertActive(() => new RenderCancelledError("render_cancelled"));
|
||||||
throw new RenderCancelledError("render_cancelled");
|
|
||||||
}
|
|
||||||
};
|
};
|
||||||
|
|
||||||
job.startedAt = new Date();
|
job.startedAt = new Date();
|
||||||
@@ -1949,7 +2001,7 @@ export async function executeRenderJob(
|
|||||||
cfg,
|
cfg,
|
||||||
log,
|
log,
|
||||||
composition,
|
composition,
|
||||||
abortSignal,
|
abortSignal: executionSignal,
|
||||||
assertNotAborted,
|
assertNotAborted,
|
||||||
// Copy (don't symlink) extracted frames on Windows — symlinkSync throws
|
// Copy (don't symlink) extracted frames on Windows — symlinkSync throws
|
||||||
// EPERM there without Developer Mode/admin, which failed local renders.
|
// EPERM there without Developer Mode/admin, which failed local renders.
|
||||||
@@ -2041,7 +2093,7 @@ export async function executeRenderJob(
|
|||||||
compiledDir,
|
compiledDir,
|
||||||
duration: probeResult.duration,
|
duration: probeResult.duration,
|
||||||
audios: composition.audios,
|
audios: composition.audios,
|
||||||
abortSignal,
|
abortSignal: executionSignal,
|
||||||
assertNotAborted,
|
assertNotAborted,
|
||||||
}),
|
}),
|
||||||
);
|
);
|
||||||
@@ -2739,7 +2791,7 @@ export async function executeRenderJob(
|
|||||||
buildCaptureOptions,
|
buildCaptureOptions,
|
||||||
createRenderVideoFrameInjector,
|
createRenderVideoFrameInjector,
|
||||||
hdrDiagnostics,
|
hdrDiagnostics,
|
||||||
abortSignal,
|
abortSignal: executionSignal,
|
||||||
assertNotAborted,
|
assertNotAborted,
|
||||||
onProgress,
|
onProgress,
|
||||||
}),
|
}),
|
||||||
@@ -2797,7 +2849,7 @@ export async function executeRenderJob(
|
|||||||
},
|
},
|
||||||
buildCaptureOptions,
|
buildCaptureOptions,
|
||||||
createRenderVideoFrameInjector,
|
createRenderVideoFrameInjector,
|
||||||
abortSignal,
|
abortSignal: executionSignal,
|
||||||
assertNotAborted,
|
assertNotAborted,
|
||||||
onProgress,
|
onProgress,
|
||||||
dedupPerfs,
|
dedupPerfs,
|
||||||
@@ -2820,7 +2872,7 @@ export async function executeRenderJob(
|
|||||||
// which errors qualify.
|
// which errors qualify.
|
||||||
const isVerifyError = isDrawElementVerificationError(err);
|
const isVerifyError = isDrawElementVerificationError(err);
|
||||||
const isCancellation =
|
const isCancellation =
|
||||||
err instanceof RenderCancelledError || abortSignal?.aborted === true;
|
err instanceof RenderCancelledError || executionSignal?.aborted === true;
|
||||||
if (
|
if (
|
||||||
!shouldRetryViaPinnedFallback({
|
!shouldRetryViaPinnedFallback({
|
||||||
isVerifyError,
|
isVerifyError,
|
||||||
@@ -3009,7 +3061,7 @@ export async function executeRenderJob(
|
|||||||
dedupPerfs,
|
dedupPerfs,
|
||||||
buildCaptureOptions,
|
buildCaptureOptions,
|
||||||
createRenderVideoFrameInjector,
|
createRenderVideoFrameInjector,
|
||||||
abortSignal,
|
abortSignal: executionSignal,
|
||||||
assertNotAborted,
|
assertNotAborted,
|
||||||
onProgress,
|
onProgress,
|
||||||
}),
|
}),
|
||||||
@@ -3058,7 +3110,7 @@ export async function executeRenderJob(
|
|||||||
enableChunkedEncode,
|
enableChunkedEncode,
|
||||||
chunkedEncodeSize,
|
chunkedEncodeSize,
|
||||||
engineConfig: cfg,
|
engineConfig: cfg,
|
||||||
abortSignal,
|
abortSignal: executionSignal,
|
||||||
assertNotAborted,
|
assertNotAborted,
|
||||||
onProgress,
|
onProgress,
|
||||||
}),
|
}),
|
||||||
@@ -3113,7 +3165,7 @@ export async function executeRenderJob(
|
|||||||
audioOutputPath,
|
audioOutputPath,
|
||||||
outputPath: stagedOutputPath,
|
outputPath: stagedOutputPath,
|
||||||
hasAudio,
|
hasAudio,
|
||||||
abortSignal,
|
abortSignal: executionSignal,
|
||||||
assertNotAborted,
|
assertNotAborted,
|
||||||
onProgress,
|
onProgress,
|
||||||
}),
|
}),
|
||||||
@@ -3189,7 +3241,6 @@ export async function executeRenderJob(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// ── Cleanup ─────────────────────────────────────────────────────────
|
|
||||||
if (job.config.debug) {
|
if (job.config.debug) {
|
||||||
// Copy output MP4 (or single-file alpha output) into the debug dir for
|
// Copy output MP4 (or single-file alpha output) into the debug dir for
|
||||||
// easy access. Skipped for png-sequence: outputPath is a directory, not
|
// easy access. Skipped for png-sequence: outputPath is a directory, not
|
||||||
@@ -3199,16 +3250,6 @@ export async function executeRenderJob(
|
|||||||
const debugOutput = join(workDir, `output${videoExt}`);
|
const debugOutput = join(workDir, `output${videoExt}`);
|
||||||
copyFileSync(stagedOutputPath, debugOutput);
|
copyFileSync(stagedOutputPath, debugOutput);
|
||||||
}
|
}
|
||||||
} else if (process.env.KEEP_TEMP === "1") {
|
|
||||||
log.info("KEEP_TEMP=1 — leaving workDir on disk for inspection", { workDir });
|
|
||||||
} else {
|
|
||||||
await safeCleanup(
|
|
||||||
"remove workDir",
|
|
||||||
() => {
|
|
||||||
rmSync(workDir, { recursive: true, force: true, maxRetries: 3, retryDelay: 100 });
|
|
||||||
},
|
|
||||||
log,
|
|
||||||
);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
artifactTransaction.commit();
|
artifactTransaction.commit();
|
||||||
@@ -3216,19 +3257,10 @@ export async function executeRenderJob(
|
|||||||
updateJobStatus(job, "complete", "Render complete", 100, onProgress);
|
updateJobStatus(job, "complete", "Render complete", 100, onProgress);
|
||||||
await eventPublisher.flush();
|
await eventPublisher.flush();
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
await safeCleanup("rollback staged artifact", () => artifactTransaction.rollback(), log);
|
if (error instanceof RenderCancelledError || executionSignal?.aborted) {
|
||||||
if (error instanceof RenderCancelledError || abortSignal?.aborted) {
|
|
||||||
job.error = error instanceof Error ? error.message : "render_cancelled";
|
job.error = error instanceof Error ? error.message : "render_cancelled";
|
||||||
updateJobStatus(job, "cancelled", "Render cancelled", job.progress, onProgress);
|
updateJobStatus(job, "cancelled", "Render cancelled", job.progress, onProgress);
|
||||||
await eventPublisher.flush();
|
await eventPublisher.flush();
|
||||||
await cleanupRenderResources({
|
|
||||||
fileServer,
|
|
||||||
probeSession,
|
|
||||||
workDir,
|
|
||||||
debug: Boolean(job.config.debug),
|
|
||||||
log,
|
|
||||||
label: "cancel",
|
|
||||||
});
|
|
||||||
throw error instanceof RenderCancelledError
|
throw error instanceof RenderCancelledError
|
||||||
? error
|
? error
|
||||||
: new RenderCancelledError("render_cancelled");
|
: new RenderCancelledError("render_cancelled");
|
||||||
@@ -3361,17 +3393,6 @@ export async function executeRenderJob(
|
|||||||
.slice(-5),
|
.slice(-5),
|
||||||
});
|
});
|
||||||
|
|
||||||
await cleanupRenderResources({
|
|
||||||
fileServer,
|
|
||||||
probeSession,
|
|
||||||
workDir,
|
|
||||||
debug: Boolean(job.config.debug),
|
|
||||||
log,
|
|
||||||
label: "error",
|
|
||||||
});
|
|
||||||
|
|
||||||
throw error;
|
throw error;
|
||||||
} finally {
|
|
||||||
memSampler?.stop();
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user