mirror of
https://github.com/THU-MAIC/OpenMAIC.git
synced 2026-10-02 09:24:43 +08:00
refactor(video-export): extract stable RenderExecutor seam (#1104)
* refactor(video-export): extract render executor seam * fix(video-export): preserve executor failure semantics * test(video-export): lock the render HTTP contract --------- Co-authored-by: wyuc <wang-yc24@mails.tsinghua.edu.cn>
This commit is contained in:
@@ -182,10 +182,14 @@ completion time is independently bounded.
|
||||
|
||||
## Scalability
|
||||
|
||||
The service is built with two swap points so it can move from a single OSS host
|
||||
The service is built with three swap points so it can move from a single OSS host
|
||||
to a horizontally-scaled demo deployment without changing the HTTP contract or
|
||||
the app:
|
||||
|
||||
- **`RenderExecutor`** (`src/render-executor.ts`) — `InProcessExecutor` adapts
|
||||
the current HyperFrames producer to stable progress, cancellation, deadline,
|
||||
failure, and performance types. A bounded local or remote executor can replace
|
||||
it without changing `RenderCoordinator` or the routes.
|
||||
- **`JobStore`** (`src/job-store.ts`) — Part A ships `InMemoryJobStore`. A
|
||||
`RedisJobStore` implementing the same interface lets any replica serve poll /
|
||||
download requests.
|
||||
@@ -194,6 +198,9 @@ the app:
|
||||
whose `locate` returns a presigned URL makes the download route `302` the
|
||||
browser straight to object storage, bypassing the proxy.
|
||||
|
||||
`RenderCoordinator` owns admission, queueing, job state, artifact registration,
|
||||
and cleanup while depending only on those three interfaces.
|
||||
|
||||
Chunked distributed rendering (`@hyperframes/producer/distributed`) to cut
|
||||
single-job latency is a further, separate follow-up.
|
||||
|
||||
|
||||
+16
-14
@@ -30,10 +30,11 @@ import { config } from './config.js';
|
||||
import { InMemoryJobStore } from './job-store.js';
|
||||
import { LocalDiskArtifactStore } from './artifact-store.js';
|
||||
import {
|
||||
RenderManager,
|
||||
RenderCoordinator,
|
||||
RenderRejectedError,
|
||||
makeProjectDir as defaultMakeProjectDir,
|
||||
} from './render-manager.js';
|
||||
} from './render-coordinator.js';
|
||||
import { InProcessExecutor } from './render-executor.js';
|
||||
import { InvalidProjectError, unzipProject as defaultUnzipProject } from './unzip.js';
|
||||
import { capBodyStream } from './capped-stream.js';
|
||||
import { Semaphore } from './semaphore.js';
|
||||
@@ -50,7 +51,7 @@ class BadRequestError extends Error {}
|
||||
export interface AppDeps {
|
||||
jobs: JobStore;
|
||||
artifacts: ArtifactStore;
|
||||
manager: RenderManager;
|
||||
coordinator: RenderCoordinator;
|
||||
/** Bounds concurrent *buffering + extraction* (the whole RAM-heavy section). */
|
||||
extractionGate: Semaphore;
|
||||
/** Extract a validated archive into a dir. Overridable in tests. */
|
||||
@@ -89,7 +90,7 @@ function parseOptions(form: FormData): RenderOptions | string {
|
||||
* not held in RAM. This is what stops a near-cap burst from OOMing the box.
|
||||
*/
|
||||
export function createApp(deps: AppDeps): Hono {
|
||||
const { jobs, artifacts, manager, extractionGate } = deps;
|
||||
const { jobs, artifacts, coordinator, extractionGate } = deps;
|
||||
const unzipProject = deps.unzipProject ?? defaultUnzipProject;
|
||||
const makeProjectDir = deps.makeProjectDir ?? defaultMakeProjectDir;
|
||||
|
||||
@@ -115,7 +116,7 @@ export function createApp(deps: AppDeps): Hono {
|
||||
// (queue full / per-identity limit) never enters buffering or extraction.
|
||||
let reservation;
|
||||
try {
|
||||
reservation = manager.reserve(identity);
|
||||
reservation = coordinator.reserve(identity);
|
||||
} catch (error) {
|
||||
if (error instanceof RenderRejectedError) return c.json({ error: error.message }, 429);
|
||||
throw error;
|
||||
@@ -164,12 +165,12 @@ export function createApp(deps: AppDeps): Hono {
|
||||
projectDir = await makeProjectDir();
|
||||
const bytes = new Uint8Array(await file.arrayBuffer());
|
||||
await unzipProject(bytes, projectDir);
|
||||
return manager.submit(reservation, projectDir, options);
|
||||
return coordinator.submit(reservation, projectDir, options);
|
||||
});
|
||||
return c.json({ jobId }, 202);
|
||||
} catch (error) {
|
||||
manager.release(reservation);
|
||||
if (projectDir) await manager.cleanupProject(projectDir);
|
||||
coordinator.release(reservation);
|
||||
if (projectDir) await coordinator.cleanupProject(projectDir);
|
||||
if (error instanceof UploadTooLargeError) return c.json({ error: error.message }, 413);
|
||||
if (error instanceof BadRequestError) return c.json({ error: error.message }, 400);
|
||||
if (error instanceof InvalidProjectError) return c.json({ error: error.message }, 400);
|
||||
@@ -194,7 +195,7 @@ export function createApp(deps: AppDeps): Hono {
|
||||
});
|
||||
|
||||
app.delete('/render/:jobId', async (c) => {
|
||||
const ok = await manager.cancel(c.req.param('jobId'));
|
||||
const ok = await coordinator.cancel(c.req.param('jobId'));
|
||||
if (!ok) return c.json({ error: 'Job not found' }, 404);
|
||||
return c.json({ cancelled: true });
|
||||
});
|
||||
@@ -232,20 +233,21 @@ export function createApp(deps: AppDeps): Hono {
|
||||
/** Wire the production collaborators and start the server (skipped under tests). */
|
||||
async function main(): Promise<void> {
|
||||
const artifacts = new LocalDiskArtifactStore();
|
||||
// Assigned after `jobs` so its reap callback can close over the manager.
|
||||
const executor = new InProcessExecutor();
|
||||
// Assigned after `jobs` so its reap callback can close over the coordinator.
|
||||
// eslint-disable-next-line prefer-const
|
||||
let manager: RenderManager;
|
||||
let coordinator: RenderCoordinator;
|
||||
const jobs = new InMemoryJobStore(config.jobTtlMs, (record) => {
|
||||
// A reaped job's artifact + project dir go with it.
|
||||
void artifacts.remove(record.id);
|
||||
void manager.cleanupProject(record.projectDir);
|
||||
void coordinator.cleanupProject(record.projectDir);
|
||||
});
|
||||
manager = new RenderManager(jobs, artifacts);
|
||||
coordinator = new RenderCoordinator(executor, jobs, artifacts);
|
||||
|
||||
const app = createApp({
|
||||
jobs,
|
||||
artifacts,
|
||||
manager,
|
||||
coordinator,
|
||||
// Bounds concurrent buffering + extraction so the per-archive RAM ceiling
|
||||
// can't stack across a burst of admitted requests.
|
||||
extractionGate: new Semaphore(config.maxConcurrentExtractions),
|
||||
|
||||
@@ -0,0 +1,264 @@
|
||||
/**
|
||||
* RenderCoordinator — owns admission, queueing, job state, artifacts, and
|
||||
* cleanup while delegating rendering policy to the RenderExecutor seam.
|
||||
*
|
||||
* Admission is split from enqueue so a caller is bounded before archive
|
||||
* extraction: reserve(identity) atomically claims a slot, submit() consumes it,
|
||||
* and release() undoes it when extraction fails.
|
||||
*/
|
||||
import { randomUUID } from 'node:crypto';
|
||||
import { mkdtemp, rm } from 'node:fs/promises';
|
||||
import { join } from 'node:path';
|
||||
import type { ArtifactStore } from './artifact-store.js';
|
||||
import { config } from './config.js';
|
||||
import type { JobStore } from './job-store.js';
|
||||
import type { RenderExecutor } from './render-executor.js';
|
||||
import type {
|
||||
RenderCancelledFailure,
|
||||
RenderExecutionResult,
|
||||
RenderFailedFailure,
|
||||
RenderJobRecord,
|
||||
RenderOptions,
|
||||
} from './types.js';
|
||||
|
||||
/** Thrown when admission control rejects a submission (mapped to HTTP 429). */
|
||||
export class RenderRejectedError extends Error {}
|
||||
|
||||
/** An accepted admission slot, returned by RenderCoordinator.reserve. */
|
||||
export interface Reservation {
|
||||
identity: string;
|
||||
consumed: boolean;
|
||||
}
|
||||
|
||||
interface QueuedJob {
|
||||
record: RenderJobRecord;
|
||||
options: RenderOptions;
|
||||
abort: AbortController;
|
||||
}
|
||||
|
||||
export interface RenderCoordinatorOptions {
|
||||
maxConcurrency?: number;
|
||||
maxQueue?: number;
|
||||
maxJobsPerUser?: number;
|
||||
jobDeadlineMs?: number;
|
||||
}
|
||||
|
||||
export class RenderCoordinator {
|
||||
private running = 0;
|
||||
private readonly queue: QueuedJob[] = [];
|
||||
private readonly controllers = new Map<string, AbortController>();
|
||||
private readonly activeByIdentity = new Map<string, number>();
|
||||
private pending = 0;
|
||||
private readonly maxConcurrency: number;
|
||||
private readonly maxQueue: number;
|
||||
private readonly maxJobsPerUser: number;
|
||||
private readonly jobDeadlineMs: number;
|
||||
|
||||
constructor(
|
||||
private readonly executor: RenderExecutor,
|
||||
private readonly jobs: JobStore,
|
||||
private readonly artifacts: ArtifactStore,
|
||||
options: RenderCoordinatorOptions = {},
|
||||
) {
|
||||
this.maxConcurrency = options.maxConcurrency ?? config.maxConcurrency;
|
||||
this.maxQueue = options.maxQueue ?? config.maxQueue;
|
||||
this.maxJobsPerUser = options.maxJobsPerUser ?? config.maxJobsPerUser;
|
||||
this.jobDeadlineMs = options.jobDeadlineMs ?? config.jobDeadlineMs;
|
||||
}
|
||||
|
||||
/** Total jobs occupying the system: reserved + queued + running. */
|
||||
private get inSystem(): number {
|
||||
return this.pending + this.queue.length + this.running;
|
||||
}
|
||||
|
||||
/** Claim an admission slot before buffering or extracting the archive. */
|
||||
reserve(identity: string): Reservation {
|
||||
if (this.inSystem >= this.maxQueue) {
|
||||
throw new RenderRejectedError('The render queue is full; try again shortly.');
|
||||
}
|
||||
if (this.maxJobsPerUser > 0) {
|
||||
const active = this.activeByIdentity.get(identity) ?? 0;
|
||||
if (active >= this.maxJobsPerUser) {
|
||||
throw new RenderRejectedError(
|
||||
`A render is already in progress (limit ${this.maxJobsPerUser}).`,
|
||||
);
|
||||
}
|
||||
}
|
||||
this.activeByIdentity.set(identity, (this.activeByIdentity.get(identity) ?? 0) + 1);
|
||||
this.pending += 1;
|
||||
return { identity, consumed: false };
|
||||
}
|
||||
|
||||
/** Release a reservation that will not become a job. */
|
||||
release(reservation: Reservation): void {
|
||||
if (reservation.consumed) return;
|
||||
reservation.consumed = true;
|
||||
this.pending = Math.max(0, this.pending - 1);
|
||||
this.decrementIdentity(reservation.identity);
|
||||
}
|
||||
|
||||
private decrementIdentity(identity: string): void {
|
||||
const next = (this.activeByIdentity.get(identity) ?? 0) - 1;
|
||||
if (next <= 0) this.activeByIdentity.delete(identity);
|
||||
else this.activeByIdentity.set(identity, next);
|
||||
}
|
||||
|
||||
/** Enqueue a render against a held reservation and return its stable job id. */
|
||||
async submit(
|
||||
reservation: Reservation,
|
||||
projectDir: string,
|
||||
options: RenderOptions,
|
||||
): Promise<string> {
|
||||
if (reservation.consumed) throw new RenderRejectedError('Reservation already used');
|
||||
|
||||
reservation.consumed = true;
|
||||
this.pending = Math.max(0, this.pending - 1);
|
||||
|
||||
const id = randomUUID();
|
||||
const now = Date.now();
|
||||
const record: RenderJobRecord = {
|
||||
id,
|
||||
userId: reservation.identity,
|
||||
status: 'queued',
|
||||
progress: 0,
|
||||
currentStage: 'queued',
|
||||
createdAtMs: now,
|
||||
updatedAtMs: now,
|
||||
projectDir,
|
||||
};
|
||||
try {
|
||||
await this.jobs.create(record);
|
||||
} catch (error) {
|
||||
this.decrementIdentity(reservation.identity);
|
||||
throw error;
|
||||
}
|
||||
|
||||
const abort = new AbortController();
|
||||
this.controllers.set(id, abort);
|
||||
this.queue.push({ record, options, abort });
|
||||
this.pump();
|
||||
return id;
|
||||
}
|
||||
|
||||
/** Cancel a queued or running job through the same AbortSignal executor seam. */
|
||||
async cancel(id: string): Promise<boolean> {
|
||||
const controller = this.controllers.get(id);
|
||||
if (!controller) return false;
|
||||
controller.abort();
|
||||
|
||||
const queuedIdx = this.queue.findIndex((queued) => queued.record.id === id);
|
||||
if (queuedIdx >= 0) {
|
||||
const [queued] = this.queue.splice(queuedIdx, 1);
|
||||
this.controllers.delete(id);
|
||||
if (queued.record.userId) this.decrementIdentity(queued.record.userId);
|
||||
const failure: RenderCancelledFailure = {
|
||||
code: 'cancelled',
|
||||
message: 'Render cancelled',
|
||||
};
|
||||
await this.jobs.update(id, {
|
||||
status: 'cancelled',
|
||||
currentStage: 'cancelled',
|
||||
failure,
|
||||
});
|
||||
await this.cleanupProject(queued.record.projectDir);
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
private pump(): void {
|
||||
while (this.running < this.maxConcurrency && this.queue.length > 0) {
|
||||
const next = this.queue.shift()!;
|
||||
this.running += 1;
|
||||
void this.run(next);
|
||||
}
|
||||
}
|
||||
|
||||
private async finishNonSuccess(
|
||||
id: string,
|
||||
projectDir: string,
|
||||
result: Exclude<RenderExecutionResult, { status: 'succeeded' }>,
|
||||
): Promise<void> {
|
||||
try {
|
||||
await this.jobs.update(id, {
|
||||
status: result.status,
|
||||
currentStage: result.status,
|
||||
failure: result.failure,
|
||||
error: result.failure.message,
|
||||
...(result.performance ? { performance: result.performance } : {}),
|
||||
});
|
||||
} finally {
|
||||
await this.cleanupProject(projectDir);
|
||||
}
|
||||
}
|
||||
|
||||
private async run({ record, options, abort }: QueuedJob): Promise<void> {
|
||||
const { id, projectDir } = record;
|
||||
const outputPath = join(projectDir, 'output.mp4');
|
||||
try {
|
||||
await this.jobs.update(id, { status: 'running', currentStage: 'preparing' });
|
||||
const result = await this.executor.execute({
|
||||
projectDir,
|
||||
outputPath,
|
||||
options,
|
||||
signal: abort.signal,
|
||||
deadlineMs: this.jobDeadlineMs,
|
||||
onProgress: async (progress) => {
|
||||
await this.jobs.update(id, {
|
||||
status: 'running',
|
||||
progress: progress.progress,
|
||||
currentStage: progress.stage,
|
||||
...(progress.framesRendered !== undefined
|
||||
? { framesRendered: progress.framesRendered }
|
||||
: {}),
|
||||
...(progress.totalFrames !== undefined ? { totalFrames: progress.totalFrames } : {}),
|
||||
});
|
||||
},
|
||||
});
|
||||
|
||||
if (result.status !== 'succeeded') {
|
||||
await this.finishNonSuccess(id, projectDir, result);
|
||||
return;
|
||||
}
|
||||
|
||||
if (abort.signal.aborted) {
|
||||
await this.finishNonSuccess(id, projectDir, {
|
||||
status: 'cancelled',
|
||||
failure: { code: 'cancelled', message: 'Render cancelled' },
|
||||
...(result.performance ? { performance: result.performance } : {}),
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
await this.artifacts.put(id, outputPath);
|
||||
await this.jobs.update(id, {
|
||||
status: 'succeeded',
|
||||
progress: 1,
|
||||
currentStage: 'complete',
|
||||
outputPath,
|
||||
...(result.performance ? { performance: result.performance } : {}),
|
||||
});
|
||||
} catch (error) {
|
||||
await this.artifacts.remove(id).catch(() => {});
|
||||
const failure: RenderFailedFailure = {
|
||||
code: 'execution_failed',
|
||||
message: error instanceof Error ? error.message : String(error),
|
||||
};
|
||||
await this.finishNonSuccess(id, projectDir, { status: 'failed', failure });
|
||||
} finally {
|
||||
this.controllers.delete(id);
|
||||
if (record.userId) this.decrementIdentity(record.userId);
|
||||
this.running -= 1;
|
||||
this.pump();
|
||||
}
|
||||
}
|
||||
|
||||
/** Best-effort recursive delete of a job's unzipped project dir. */
|
||||
async cleanupProject(dir: string): Promise<void> {
|
||||
await rm(dir, { recursive: true, force: true }).catch(() => {});
|
||||
}
|
||||
}
|
||||
|
||||
/** Create a fresh per-render project directory under the configured tmp root. */
|
||||
export async function makeProjectDir(): Promise<string> {
|
||||
return mkdtemp(join(config.tmpDir, 'render-'));
|
||||
}
|
||||
@@ -0,0 +1,207 @@
|
||||
/**
|
||||
* RenderExecutor — the stable seam between render lifecycle policy and a
|
||||
* concrete rendering engine. Callers provide cancellation, a deadline, and a
|
||||
* progress sink; adapters return domain results and performance data without
|
||||
* leaking engine-specific job or error types.
|
||||
*/
|
||||
import {
|
||||
createRenderJob,
|
||||
executeRenderJob,
|
||||
RenderCancelledError,
|
||||
type RenderConfigInput,
|
||||
type RenderJob,
|
||||
type RenderPerfSummary,
|
||||
} from '@hyperframes/producer';
|
||||
import { config } from './config.js';
|
||||
import type {
|
||||
RenderExecutionRequest,
|
||||
RenderExecutionResult,
|
||||
RenderOptions,
|
||||
RenderPerformanceSummary,
|
||||
} from './types.js';
|
||||
|
||||
export interface RenderExecutor {
|
||||
execute(request: RenderExecutionRequest): Promise<RenderExecutionResult>;
|
||||
}
|
||||
|
||||
interface ProducerBridge {
|
||||
createJob(options: RenderConfigInput): RenderJob;
|
||||
executeJob(
|
||||
job: RenderJob,
|
||||
projectDir: string,
|
||||
outputPath: string,
|
||||
onProgress: (job: RenderJob) => void | Promise<void>,
|
||||
signal: AbortSignal,
|
||||
): Promise<void>;
|
||||
}
|
||||
|
||||
const producerBridge: ProducerBridge = {
|
||||
createJob: createRenderJob,
|
||||
executeJob: executeRenderJob,
|
||||
};
|
||||
|
||||
export interface InProcessExecutorOptions {
|
||||
workers?: number;
|
||||
requireBeginFrame?: boolean;
|
||||
}
|
||||
|
||||
/** Build the engine-specific config entirely inside the production adapter. */
|
||||
export function buildProducerJobConfig(
|
||||
options: RenderOptions,
|
||||
workers = config.producerWorkers,
|
||||
): RenderConfigInput {
|
||||
const producerOptions: RenderConfigInput = {
|
||||
fps: options.fps,
|
||||
quality: options.quality,
|
||||
format: options.format,
|
||||
};
|
||||
if (workers !== undefined) producerOptions.workers = workers;
|
||||
return producerOptions;
|
||||
}
|
||||
|
||||
function performanceSummary(
|
||||
summary: RenderPerfSummary | undefined,
|
||||
): RenderPerformanceSummary | undefined {
|
||||
if (!summary) return undefined;
|
||||
return {
|
||||
totalElapsedMs: summary.totalElapsedMs,
|
||||
stages: { ...summary.stages },
|
||||
workers: summary.workers,
|
||||
totalFrames: summary.totalFrames,
|
||||
...(summary.drawElement?.mode ? { captureMode: summary.drawElement.mode } : {}),
|
||||
...(summary.peakRssMb !== undefined ? { peakRssMb: summary.peakRssMb } : {}),
|
||||
...(summary.tmpPeakBytes !== undefined ? { tmpPeakBytes: summary.tmpPeakBytes } : {}),
|
||||
};
|
||||
}
|
||||
|
||||
function message(error: unknown): string {
|
||||
return error instanceof Error ? error.message : String(error);
|
||||
}
|
||||
|
||||
/** In-process adapter around the current HyperFrames producer. */
|
||||
export class InProcessExecutor implements RenderExecutor {
|
||||
private readonly workers: number | undefined;
|
||||
private readonly requireBeginFrame: boolean;
|
||||
|
||||
constructor(
|
||||
options: InProcessExecutorOptions = {},
|
||||
private readonly producer: ProducerBridge = producerBridge,
|
||||
) {
|
||||
this.workers = options.workers ?? config.producerWorkers;
|
||||
this.requireBeginFrame = options.requireBeginFrame ?? config.requireBeginFrame;
|
||||
}
|
||||
|
||||
async execute(request: RenderExecutionRequest): Promise<RenderExecutionResult> {
|
||||
if (request.signal.aborted) {
|
||||
return {
|
||||
status: 'cancelled',
|
||||
failure: { code: 'cancelled', message: 'Render cancelled' },
|
||||
};
|
||||
}
|
||||
|
||||
const abort = new AbortController();
|
||||
let abortCause: 'cancelled' | 'deadline' | null = null;
|
||||
const cancel = () => {
|
||||
if (abortCause !== null) return;
|
||||
abortCause = 'cancelled';
|
||||
abort.abort();
|
||||
};
|
||||
request.signal.addEventListener('abort', cancel, { once: true });
|
||||
|
||||
const deadline = setTimeout(
|
||||
() => {
|
||||
if (abortCause !== null) return;
|
||||
abortCause = 'deadline';
|
||||
abort.abort();
|
||||
},
|
||||
Math.max(0, request.deadlineMs),
|
||||
);
|
||||
deadline.unref?.();
|
||||
|
||||
let job: RenderJob | undefined;
|
||||
try {
|
||||
job = this.producer.createJob(buildProducerJobConfig(request.options, this.workers));
|
||||
await this.producer.executeJob(
|
||||
job,
|
||||
request.projectDir,
|
||||
request.outputPath,
|
||||
async (current) => {
|
||||
const progress =
|
||||
typeof current.progress === 'number'
|
||||
? Math.max(0, Math.min(1, current.progress / 100))
|
||||
: 0;
|
||||
await request.onProgress({
|
||||
progress,
|
||||
stage: current.currentStage || current.status,
|
||||
...(typeof current.framesRendered === 'number'
|
||||
? { framesRendered: current.framesRendered }
|
||||
: {}),
|
||||
...(typeof current.totalFrames === 'number'
|
||||
? { totalFrames: current.totalFrames }
|
||||
: {}),
|
||||
});
|
||||
},
|
||||
abort.signal,
|
||||
);
|
||||
|
||||
const performance = performanceSummary(job.perfSummary);
|
||||
if (abortCause === 'deadline') {
|
||||
return {
|
||||
status: 'failed',
|
||||
failure: { code: 'deadline_exceeded', message: 'Render exceeded the deadline' },
|
||||
...(performance ? { performance } : {}),
|
||||
};
|
||||
}
|
||||
if (abortCause === 'cancelled') {
|
||||
return {
|
||||
status: 'cancelled',
|
||||
failure: { code: 'cancelled', message: 'Render cancelled' },
|
||||
...(performance ? { performance } : {}),
|
||||
};
|
||||
}
|
||||
|
||||
const captureMode = performance?.captureMode;
|
||||
if (this.requireBeginFrame && captureMode !== 'beginframe') {
|
||||
return {
|
||||
status: 'failed',
|
||||
failure: {
|
||||
code: 'unsupported_capture_mode',
|
||||
message:
|
||||
`Producer did not resolve beginFrame capture (actual=${captureMode ?? 'unknown'}). ` +
|
||||
'Check PRODUCER_HEADLESS_SHELL_PATH and Chromium compatibility.',
|
||||
},
|
||||
...(performance ? { performance } : {}),
|
||||
};
|
||||
}
|
||||
|
||||
return { status: 'succeeded', ...(performance ? { performance } : {}) };
|
||||
} catch (error) {
|
||||
const performance = performanceSummary(job?.perfSummary);
|
||||
if (
|
||||
abortCause === 'deadline' ||
|
||||
(abortCause === null && error instanceof RenderCancelledError && error.reason === 'timeout')
|
||||
) {
|
||||
return {
|
||||
status: 'failed',
|
||||
failure: { code: 'deadline_exceeded', message: 'Render exceeded the deadline' },
|
||||
...(performance ? { performance } : {}),
|
||||
};
|
||||
}
|
||||
if (abortCause === 'cancelled') {
|
||||
return {
|
||||
status: 'cancelled',
|
||||
failure: { code: 'cancelled', message: message(error) },
|
||||
...(performance ? { performance } : {}),
|
||||
};
|
||||
}
|
||||
return {
|
||||
status: 'failed',
|
||||
failure: { code: 'execution_failed', message: message(error) },
|
||||
...(performance ? { performance } : {}),
|
||||
};
|
||||
} finally {
|
||||
clearTimeout(deadline);
|
||||
request.signal.removeEventListener('abort', cancel);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,292 +0,0 @@
|
||||
/**
|
||||
* RenderManager — owns a render job's whole lifecycle: admission (concurrency,
|
||||
* per-identity, and global-queue guards), the FIFO queue, driving
|
||||
* `@hyperframes/producer`, feeding progress into the JobStore, registering the
|
||||
* artifact, a per-job wall-clock deadline, and cleanup.
|
||||
*
|
||||
* Admission is split from enqueue so a caller is bounded *before* the expensive
|
||||
* archive extraction: `reserve(identity)` atomically claims a slot (or throws),
|
||||
* the route extracts, then `submit()` consumes the reservation. `release()`
|
||||
* undoes a reservation if extraction/submit fails. All counters are plain
|
||||
* fields mutated synchronously on the single-threaded event loop, so the
|
||||
* check-then-increment is atomic; a Redis-backed store would make it
|
||||
* distributed for the demo layer.
|
||||
*/
|
||||
import { randomUUID } from 'node:crypto';
|
||||
import { mkdtemp, rm } from 'node:fs/promises';
|
||||
import { join } from 'node:path';
|
||||
import { createRenderJob, executeRenderJob, type RenderConfigInput } from '@hyperframes/producer';
|
||||
import type { JobStore } from './job-store.js';
|
||||
import type { ArtifactStore } from './artifact-store.js';
|
||||
import type { RenderJobRecord, RenderOptions } from './types.js';
|
||||
import { config } from './config.js';
|
||||
|
||||
/** Thrown when admission control rejects a submission (mapped to HTTP 429). */
|
||||
export class RenderRejectedError extends Error {}
|
||||
|
||||
/** An accepted admission slot, returned by {@link RenderManager.reserve}. */
|
||||
export interface Reservation {
|
||||
identity: string;
|
||||
consumed: boolean;
|
||||
}
|
||||
|
||||
interface QueuedJob {
|
||||
record: RenderJobRecord;
|
||||
options: RenderOptions;
|
||||
abort: AbortController;
|
||||
}
|
||||
|
||||
/**
|
||||
* Build the producer job config with an explicit worker count. Producer's env
|
||||
* `concurrency` is only an auto-sizing hint and can be raised by its minimum
|
||||
* parallel-frame rule; `job.config.workers` is the authoritative user choice.
|
||||
*/
|
||||
export function buildProducerJobConfig(
|
||||
options: RenderOptions,
|
||||
workers = config.producerWorkers,
|
||||
): RenderConfigInput {
|
||||
const producerOptions: RenderConfigInput = {
|
||||
fps: options.fps,
|
||||
quality: options.quality,
|
||||
format: options.format,
|
||||
};
|
||||
if (workers !== undefined) producerOptions.workers = workers;
|
||||
return producerOptions;
|
||||
}
|
||||
|
||||
/** Assert the worker-reported capture mode when the deployment requires beginFrame. */
|
||||
export function assertRequiredCaptureMode(
|
||||
captureMode: string | undefined,
|
||||
requireBeginFrame: boolean,
|
||||
): void {
|
||||
if (!requireBeginFrame) return;
|
||||
if (captureMode !== 'beginframe') {
|
||||
throw new Error(
|
||||
`Producer did not resolve beginFrame capture (actual=${captureMode ?? 'unknown'}). ` +
|
||||
'Check PRODUCER_HEADLESS_SHELL_PATH and Chromium compatibility.',
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
export class RenderManager {
|
||||
private running = 0;
|
||||
private readonly queue: QueuedJob[] = [];
|
||||
/** Live AbortControllers for queued/running jobs, keyed by jobId (for cancel). */
|
||||
private readonly controllers = new Map<string, AbortController>();
|
||||
/** Active (reserved + queued + running) count per identity, for the per-user guard. */
|
||||
private readonly activeByIdentity = new Map<string, number>();
|
||||
/** Reserved-but-not-yet-submitted count, included in the global queue-depth cap. */
|
||||
private pending = 0;
|
||||
|
||||
constructor(
|
||||
private readonly jobs: JobStore,
|
||||
private readonly artifacts: ArtifactStore,
|
||||
) {}
|
||||
|
||||
/** Total jobs occupying the system: reserved + queued + running. */
|
||||
private get inSystem(): number {
|
||||
return this.pending + this.queue.length + this.running;
|
||||
}
|
||||
|
||||
/**
|
||||
* Atomically claim an admission slot for `identity`, or throw
|
||||
* {@link RenderRejectedError}. Must be paired with {@link consume} (on
|
||||
* success) or {@link release} (on failure). Reserve before extracting the
|
||||
* archive so a rejected caller never triggers a decompression.
|
||||
*/
|
||||
reserve(identity: string): Reservation {
|
||||
if (this.inSystem >= config.maxQueue) {
|
||||
throw new RenderRejectedError('The render queue is full; try again shortly.');
|
||||
}
|
||||
if (config.maxJobsPerUser > 0) {
|
||||
const active = this.activeByIdentity.get(identity) ?? 0;
|
||||
if (active >= config.maxJobsPerUser) {
|
||||
throw new RenderRejectedError(
|
||||
`A render is already in progress (limit ${config.maxJobsPerUser}).`,
|
||||
);
|
||||
}
|
||||
}
|
||||
this.activeByIdentity.set(identity, (this.activeByIdentity.get(identity) ?? 0) + 1);
|
||||
this.pending += 1;
|
||||
return { identity, consumed: false };
|
||||
}
|
||||
|
||||
/** Release a reservation that will not become a job (extraction/submit failed). */
|
||||
release(reservation: Reservation): void {
|
||||
if (reservation.consumed) return;
|
||||
reservation.consumed = true;
|
||||
this.pending = Math.max(0, this.pending - 1);
|
||||
this.decrementIdentity(reservation.identity);
|
||||
}
|
||||
|
||||
private decrementIdentity(identity: string): void {
|
||||
const next = (this.activeByIdentity.get(identity) ?? 0) - 1;
|
||||
if (next <= 0) this.activeByIdentity.delete(identity);
|
||||
else this.activeByIdentity.set(identity, next);
|
||||
}
|
||||
|
||||
/**
|
||||
* Enqueue a render against a held reservation. `projectDir` already contains
|
||||
* the unzipped project (with index.html). Returns the new jobId.
|
||||
*/
|
||||
async submit(
|
||||
reservation: Reservation,
|
||||
projectDir: string,
|
||||
options: RenderOptions,
|
||||
): Promise<string> {
|
||||
if (reservation.consumed) {
|
||||
throw new RenderRejectedError('Reservation already used');
|
||||
}
|
||||
// Convert the reservation into a real queued job: the identity count stays,
|
||||
// but it's no longer "pending".
|
||||
reservation.consumed = true;
|
||||
this.pending = Math.max(0, this.pending - 1);
|
||||
|
||||
const id = randomUUID();
|
||||
const now = Date.now();
|
||||
const record: RenderJobRecord = {
|
||||
id,
|
||||
userId: reservation.identity,
|
||||
status: 'queued',
|
||||
progress: 0,
|
||||
currentStage: 'queued',
|
||||
createdAtMs: now,
|
||||
updatedAtMs: now,
|
||||
projectDir,
|
||||
};
|
||||
// If persisting the job fails, the identity slot claimed at reserve() would
|
||||
// otherwise leak (run() never runs to decrement it). Release it here.
|
||||
try {
|
||||
await this.jobs.create(record);
|
||||
} catch (error) {
|
||||
this.decrementIdentity(reservation.identity);
|
||||
throw error;
|
||||
}
|
||||
|
||||
const abort = new AbortController();
|
||||
this.controllers.set(id, abort);
|
||||
this.queue.push({ record, options, abort });
|
||||
this.pump();
|
||||
return id;
|
||||
}
|
||||
|
||||
/** Cancel a queued or running job. */
|
||||
async cancel(id: string): Promise<boolean> {
|
||||
const controller = this.controllers.get(id);
|
||||
if (!controller) return false;
|
||||
controller.abort();
|
||||
// If still queued (not yet running), drop it and finalize now.
|
||||
const queuedIdx = this.queue.findIndex((q) => q.record.id === id);
|
||||
if (queuedIdx >= 0) {
|
||||
const [q] = this.queue.splice(queuedIdx, 1);
|
||||
this.controllers.delete(id);
|
||||
if (q.record.userId) this.decrementIdentity(q.record.userId);
|
||||
await this.jobs.update(id, { status: 'cancelled', currentStage: 'cancelled' });
|
||||
await this.cleanupProject(q.record.projectDir);
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
/** Start as many queued jobs as the concurrency budget allows. */
|
||||
private pump(): void {
|
||||
while (this.running < config.maxConcurrency && this.queue.length > 0) {
|
||||
const next = this.queue.shift()!;
|
||||
this.running++;
|
||||
// Fire-and-forget: run() owns its own error handling and always decrements.
|
||||
void this.run(next);
|
||||
}
|
||||
}
|
||||
|
||||
private async run({ record, options, abort }: QueuedJob): Promise<void> {
|
||||
const { id, projectDir } = record;
|
||||
const outputPath = join(projectDir, 'output.mp4');
|
||||
// Wall-clock watchdog: abort a render that overruns the deadline so it can't
|
||||
// hold a concurrency slot + scratch dir indefinitely. executeRenderJob
|
||||
// honors the same AbortSignal we pass for user cancellation. `timedOut`
|
||||
// distinguishes a deadline abort (→ failed) from a user cancel (→ cancelled).
|
||||
let timedOut = false;
|
||||
const deadline = setTimeout(() => {
|
||||
timedOut = true;
|
||||
abort.abort();
|
||||
}, config.jobDeadlineMs);
|
||||
if (typeof deadline.unref === 'function') deadline.unref();
|
||||
try {
|
||||
await this.jobs.update(id, { status: 'running', currentStage: 'preparing' });
|
||||
|
||||
const job = createRenderJob(buildProducerJobConfig(options));
|
||||
|
||||
await executeRenderJob(
|
||||
job,
|
||||
projectDir,
|
||||
outputPath,
|
||||
async (j) => {
|
||||
// Producer mutates the same `job` object; mirror the fields we expose.
|
||||
// Producer's `progress` is 0..100; our HTTP contract is 0..1, so
|
||||
// normalize here (clamped) — success is set to 1 below.
|
||||
const progress =
|
||||
typeof j.progress === 'number' ? Math.max(0, Math.min(1, j.progress / 100)) : 0;
|
||||
await this.jobs.update(id, {
|
||||
status: 'running',
|
||||
progress,
|
||||
currentStage: j.currentStage || j.status,
|
||||
...(typeof j.framesRendered === 'number' ? { framesRendered: j.framesRendered } : {}),
|
||||
...(typeof j.totalFrames === 'number' ? { totalFrames: j.totalFrames } : {}),
|
||||
});
|
||||
},
|
||||
abort.signal,
|
||||
);
|
||||
|
||||
assertRequiredCaptureMode(job.perfSummary?.drawElement?.mode, config.requireBeginFrame);
|
||||
|
||||
if (abort.signal.aborted) {
|
||||
// Deadline overrun is a failure, not a user cancellation.
|
||||
await this.jobs.update(id, {
|
||||
status: timedOut ? 'failed' : 'cancelled',
|
||||
currentStage: timedOut ? 'failed' : 'cancelled',
|
||||
...(timedOut ? { error: 'Render exceeded the deadline' } : {}),
|
||||
});
|
||||
await this.cleanupProject(projectDir);
|
||||
return;
|
||||
}
|
||||
|
||||
await this.artifacts.put(id, outputPath);
|
||||
await this.jobs.update(id, {
|
||||
status: 'succeeded',
|
||||
progress: 1,
|
||||
currentStage: 'complete',
|
||||
outputPath,
|
||||
});
|
||||
} catch (error) {
|
||||
// A deadline abort surfaces as a thrown RenderCancelledError; report it as
|
||||
// failed (with a clear reason), reserving `cancelled` for user cancels.
|
||||
const cancelledByUser = abort.signal.aborted && !timedOut;
|
||||
await this.jobs.update(id, {
|
||||
status: cancelledByUser ? 'cancelled' : 'failed',
|
||||
currentStage: cancelledByUser ? 'cancelled' : 'failed',
|
||||
error: timedOut
|
||||
? 'Render exceeded the deadline'
|
||||
: error instanceof Error
|
||||
? error.message
|
||||
: String(error),
|
||||
});
|
||||
// On failure/cancel the artifact is worthless — reclaim the project dir now.
|
||||
await this.cleanupProject(projectDir);
|
||||
} finally {
|
||||
clearTimeout(deadline);
|
||||
this.controllers.delete(id);
|
||||
if (record.userId) this.decrementIdentity(record.userId);
|
||||
this.running--;
|
||||
this.pump();
|
||||
}
|
||||
}
|
||||
|
||||
/** Best-effort recursive delete of a job's unzipped project dir. */
|
||||
async cleanupProject(dir: string): Promise<void> {
|
||||
await rm(dir, { recursive: true, force: true }).catch(() => {});
|
||||
}
|
||||
}
|
||||
|
||||
/** Create a fresh, empty per-render project directory under the configured tmp root. */
|
||||
export async function makeProjectDir(): Promise<string> {
|
||||
return mkdtemp(join(config.tmpDir, 'render-'));
|
||||
}
|
||||
@@ -18,9 +18,76 @@ export interface RenderOptions {
|
||||
format: 'mp4';
|
||||
}
|
||||
|
||||
/** Progress emitted by any render executor, normalized for service callers. */
|
||||
export interface RenderProgress {
|
||||
/** Fraction complete in the stable 0..1 service range. */
|
||||
progress: number;
|
||||
stage: string;
|
||||
framesRendered?: number;
|
||||
totalFrames?: number;
|
||||
}
|
||||
|
||||
/** Executor-independent performance data retained for diagnostics and benchmarks. */
|
||||
export interface RenderPerformanceSummary {
|
||||
totalElapsedMs: number;
|
||||
stages: Record<string, number>;
|
||||
workers: number;
|
||||
totalFrames: number;
|
||||
captureMode?: string;
|
||||
peakRssMb?: number;
|
||||
tmpPeakBytes?: number;
|
||||
}
|
||||
|
||||
/** Stable failure classification produced at the RenderExecutor seam. */
|
||||
export type RenderFailureCode =
|
||||
| 'cancelled'
|
||||
| 'deadline_exceeded'
|
||||
| 'unsupported_capture_mode'
|
||||
| 'execution_failed';
|
||||
|
||||
export interface RenderCancelledFailure {
|
||||
code: 'cancelled';
|
||||
message: string;
|
||||
}
|
||||
|
||||
export interface RenderFailedFailure {
|
||||
code: Exclude<RenderFailureCode, 'cancelled'>;
|
||||
message: string;
|
||||
}
|
||||
|
||||
export type RenderFailure = RenderCancelledFailure | RenderFailedFailure;
|
||||
|
||||
/** Everything an executor needs to run one render without knowing about HTTP jobs. */
|
||||
export interface RenderExecutionRequest {
|
||||
projectDir: string;
|
||||
outputPath: string;
|
||||
options: RenderOptions;
|
||||
/** User or coordinator cancellation. */
|
||||
signal: AbortSignal;
|
||||
/** Wall-clock budget starting when execution begins. */
|
||||
deadlineMs: number;
|
||||
onProgress: (progress: RenderProgress) => void | Promise<void>;
|
||||
}
|
||||
|
||||
export type RenderExecutionResult =
|
||||
| {
|
||||
status: 'succeeded';
|
||||
performance?: RenderPerformanceSummary;
|
||||
}
|
||||
| {
|
||||
status: 'cancelled';
|
||||
failure: RenderCancelledFailure;
|
||||
performance?: RenderPerformanceSummary;
|
||||
}
|
||||
| {
|
||||
status: 'failed';
|
||||
failure: RenderFailedFailure;
|
||||
performance?: RenderPerformanceSummary;
|
||||
};
|
||||
|
||||
/**
|
||||
* A render job's observable state. `progress` is 0..1 (producer's native
|
||||
* range); the HTTP layer surfaces it as-is and the client scales to a percent.
|
||||
* A render job's observable state. `progress` is normalized to 0..1; the HTTP
|
||||
* layer surfaces it as-is and the client scales it to a percent.
|
||||
*/
|
||||
export interface RenderJobRecord {
|
||||
id: string;
|
||||
@@ -38,6 +105,10 @@ export interface RenderJobRecord {
|
||||
projectDir: string;
|
||||
/** Absolute path to the rendered MP4 once `succeeded`. */
|
||||
outputPath?: string;
|
||||
/** Domain failure retained independently from the HTTP-compatible error string. */
|
||||
failure?: RenderFailure;
|
||||
/** Executor-independent diagnostics for completed or failed attempts. */
|
||||
performance?: RenderPerformanceSummary;
|
||||
}
|
||||
|
||||
export function isTerminal(status: RenderJobStatus): boolean {
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
import { describe, expect, it } from 'vitest';
|
||||
import { assertRequiredCaptureMode, buildProducerJobConfig } from '../src/render-manager.js';
|
||||
import { buildProducerJobConfig } from '../src/render-executor.js';
|
||||
|
||||
describe('buildProducerJobConfig', () => {
|
||||
const options = { fps: 30, quality: 'standard', format: 'mp4' } as const;
|
||||
@@ -15,15 +15,4 @@ describe('buildProducerJobConfig', () => {
|
||||
it('leaves workers unset when no explicit override is supplied', () => {
|
||||
expect(buildProducerJobConfig(options, undefined)).toEqual(options);
|
||||
});
|
||||
|
||||
it('rejects a required beginFrame profile when workers report screenshot mode', () => {
|
||||
expect(() => assertRequiredCaptureMode('screenshot', true)).toThrow(/beginFrame/i);
|
||||
expect(() => assertRequiredCaptureMode('beginframe|screenshot', true)).toThrow(/beginFrame/i);
|
||||
expect(() => assertRequiredCaptureMode(undefined, true)).toThrow(/beginFrame/i);
|
||||
});
|
||||
|
||||
it('accepts the resolved beginFrame mode and does nothing when not required', () => {
|
||||
expect(() => assertRequiredCaptureMode('beginframe', true)).not.toThrow();
|
||||
expect(() => assertRequiredCaptureMode('screenshot', false)).not.toThrow();
|
||||
});
|
||||
});
|
||||
|
||||
+21
-48
@@ -1,5 +1,5 @@
|
||||
/**
|
||||
* Admission-control arithmetic for {@link RenderManager}. These guard the
|
||||
* Admission-control arithmetic for {@link RenderCoordinator}. These guard the
|
||||
* reserve → submit/release lifecycle that bounds a caller *before* the archive
|
||||
* is extracted:
|
||||
* - the global queue-depth cap (`RENDER_MAX_QUEUE`) counts reserved slots;
|
||||
@@ -7,55 +7,28 @@
|
||||
* - `release()` fully undoes a reservation (the leak the route fix depends on:
|
||||
* if a post-reserve step like makeProjectDir throws, the slot must come back).
|
||||
*
|
||||
* We drive the manager directly with in-memory stores so nothing invokes the
|
||||
* We drive the coordinator directly with in-memory stores so nothing invokes the
|
||||
* real Chromium/FFmpeg producer — this is pure counter arithmetic.
|
||||
*/
|
||||
import { describe, it, expect } from 'vitest';
|
||||
import { RenderManager, RenderRejectedError } from '../src/render-manager.js';
|
||||
import type { JobStore } from '../src/job-store.js';
|
||||
import type { ArtifactStore, ArtifactLocation } from '../src/artifact-store.js';
|
||||
import type { RenderJobRecord } from '../src/types.js';
|
||||
import { RenderCoordinator, RenderRejectedError } from '../src/render-coordinator.js';
|
||||
import {
|
||||
createMemoryArtifactStore,
|
||||
createMemoryJobStore,
|
||||
succeedingExecutor,
|
||||
} from './support/fakes.js';
|
||||
|
||||
function fakeJobStore(): JobStore {
|
||||
const jobs = new Map<string, RenderJobRecord>();
|
||||
return {
|
||||
async create(record) {
|
||||
jobs.set(record.id, record);
|
||||
},
|
||||
async get(id) {
|
||||
return jobs.get(id) ?? null;
|
||||
},
|
||||
async update(id, patch) {
|
||||
const existing = jobs.get(id);
|
||||
if (existing) jobs.set(id, { ...existing, ...patch });
|
||||
},
|
||||
async remove(id) {
|
||||
jobs.delete(id);
|
||||
},
|
||||
async list() {
|
||||
return [...jobs.values()];
|
||||
},
|
||||
async countActiveForUser() {
|
||||
return 0;
|
||||
},
|
||||
};
|
||||
function newCoordinator(): RenderCoordinator {
|
||||
return new RenderCoordinator(
|
||||
succeedingExecutor,
|
||||
createMemoryJobStore(),
|
||||
createMemoryArtifactStore().store,
|
||||
);
|
||||
}
|
||||
|
||||
const fakeArtifacts: ArtifactStore = {
|
||||
async put() {},
|
||||
async locate(): Promise<ArtifactLocation | null> {
|
||||
return null;
|
||||
},
|
||||
async remove() {},
|
||||
};
|
||||
|
||||
function newManager(): RenderManager {
|
||||
return new RenderManager(fakeJobStore(), fakeArtifacts);
|
||||
}
|
||||
|
||||
describe('RenderManager admission control', () => {
|
||||
describe('RenderCoordinator admission control', () => {
|
||||
it('reserve then release fully restores the per-identity slot', () => {
|
||||
const m = newManager();
|
||||
const m = newCoordinator();
|
||||
// Default RENDER_MAX_JOBS_PER_USER is 1.
|
||||
const r = m.reserve('alice');
|
||||
// A second reserve for the same identity is now rejected...
|
||||
@@ -66,7 +39,7 @@ describe('RenderManager admission control', () => {
|
||||
});
|
||||
|
||||
it('release is idempotent and does not double-decrement', () => {
|
||||
const m = newManager();
|
||||
const m = newCoordinator();
|
||||
const r = m.reserve('bob');
|
||||
m.release(r);
|
||||
m.release(r); // no-op, must not free a slot that isn't held
|
||||
@@ -78,7 +51,7 @@ describe('RenderManager admission control', () => {
|
||||
});
|
||||
|
||||
it('enforces the per-identity cap across distinct identities independently', () => {
|
||||
const m = newManager();
|
||||
const m = newCoordinator();
|
||||
const a = m.reserve('alice');
|
||||
const b = m.reserve('bob'); // different identity: allowed
|
||||
expect(a.identity).toBe('alice');
|
||||
@@ -89,7 +62,7 @@ describe('RenderManager admission control', () => {
|
||||
});
|
||||
|
||||
it('rejects reservations once the global queue is full', () => {
|
||||
const m = newManager();
|
||||
const m = newCoordinator();
|
||||
// Reserve up to RENDER_MAX_QUEUE (default 20) with unique identities so the
|
||||
// per-user guard never fires first, then the next reserve trips the queue cap.
|
||||
const held = [];
|
||||
@@ -104,11 +77,11 @@ describe('RenderManager admission control', () => {
|
||||
// submit() consumes the reservation and persists the job; if create() throws
|
||||
// (a fallible JobStore, e.g. a future Redis backend), run() never runs to
|
||||
// decrement the identity — so submit() must decrement it itself.
|
||||
const store = fakeJobStore();
|
||||
const store = createMemoryJobStore();
|
||||
store.create = async () => {
|
||||
throw new Error('store down');
|
||||
};
|
||||
const m = new RenderManager(store, fakeArtifacts);
|
||||
const m = new RenderCoordinator(succeedingExecutor, store, createMemoryArtifactStore().store);
|
||||
const r = m.reserve('carol');
|
||||
await expect(
|
||||
m.submit(r, '/tmp/whatever', { fps: 30, quality: 'draft', format: 'mp4' }),
|
||||
@@ -0,0 +1,161 @@
|
||||
import { access, mkdtemp, rm } from 'node:fs/promises';
|
||||
import { tmpdir } from 'node:os';
|
||||
import { join } from 'node:path';
|
||||
import { afterEach, describe, expect, it } from 'vitest';
|
||||
import type { JobStore } from '../src/job-store.js';
|
||||
import { RenderCoordinator } from '../src/render-coordinator.js';
|
||||
import type { RenderExecutor } from '../src/render-executor.js';
|
||||
import type {
|
||||
RenderExecutionRequest,
|
||||
RenderExecutionResult,
|
||||
RenderJobRecord,
|
||||
} from '../src/types.js';
|
||||
import { createMemoryArtifactStore, createMemoryJobStore } from './support/fakes.js';
|
||||
|
||||
const scratch: string[] = [];
|
||||
|
||||
afterEach(async () => {
|
||||
await Promise.all(scratch.splice(0).map((path) => rm(path, { recursive: true, force: true })));
|
||||
});
|
||||
|
||||
class FakeExecutor implements RenderExecutor {
|
||||
readonly requests: RenderExecutionRequest[] = [];
|
||||
|
||||
constructor(
|
||||
private readonly handler: (request: RenderExecutionRequest) => Promise<RenderExecutionResult>,
|
||||
) {}
|
||||
|
||||
async execute(request: RenderExecutionRequest): Promise<RenderExecutionResult> {
|
||||
this.requests.push(request);
|
||||
return this.handler(request);
|
||||
}
|
||||
}
|
||||
|
||||
async function projectDir(): Promise<string> {
|
||||
const path = await mkdtemp(join(tmpdir(), 'render-coordinator-'));
|
||||
scratch.push(path);
|
||||
return path;
|
||||
}
|
||||
|
||||
async function waitForJob(
|
||||
jobs: JobStore,
|
||||
id: string,
|
||||
predicate: (job: RenderJobRecord) => boolean,
|
||||
): Promise<RenderJobRecord> {
|
||||
for (let attempt = 0; attempt < 100; attempt += 1) {
|
||||
const job = await jobs.get(id);
|
||||
if (job && predicate(job)) return job;
|
||||
await new Promise((resolve) => setTimeout(resolve, 5));
|
||||
}
|
||||
throw new Error(`Timed out waiting for job ${id}`);
|
||||
}
|
||||
|
||||
const renderOptions = { fps: 30, quality: 'standard', format: 'mp4' } as const;
|
||||
|
||||
describe('RenderCoordinator through the RenderExecutor seam', () => {
|
||||
it('persists normalized progress, performance, and the artifact on success', async () => {
|
||||
const jobs = createMemoryJobStore();
|
||||
const artifacts = createMemoryArtifactStore();
|
||||
const performance = {
|
||||
totalElapsedMs: 800,
|
||||
stages: { captureMs: 600 },
|
||||
workers: 1,
|
||||
totalFrames: 30,
|
||||
captureMode: 'beginframe',
|
||||
};
|
||||
const executor = new FakeExecutor(async (request) => {
|
||||
await request.onProgress({
|
||||
progress: 0.5,
|
||||
stage: 'capturing',
|
||||
framesRendered: 15,
|
||||
totalFrames: 30,
|
||||
});
|
||||
return { status: 'succeeded', performance };
|
||||
});
|
||||
const coordinator = new RenderCoordinator(executor, jobs, artifacts.store, {
|
||||
jobDeadlineMs: 12_345,
|
||||
});
|
||||
const dir = await projectDir();
|
||||
const id = await coordinator.submit(coordinator.reserve('alice'), dir, renderOptions);
|
||||
|
||||
const job = await waitForJob(jobs, id, (current) => current.status === 'succeeded');
|
||||
expect(executor.requests).toHaveLength(1);
|
||||
expect(executor.requests[0].deadlineMs).toBe(12_345);
|
||||
expect(job).toMatchObject({
|
||||
status: 'succeeded',
|
||||
progress: 1,
|
||||
currentStage: 'complete',
|
||||
framesRendered: 15,
|
||||
totalFrames: 30,
|
||||
performance,
|
||||
});
|
||||
expect(artifacts.paths.get(id)).toBe(join(dir, 'output.mp4'));
|
||||
});
|
||||
|
||||
it('routes running-job cancellation through the executor signal and cleans up', async () => {
|
||||
const jobs = createMemoryJobStore();
|
||||
const artifacts = createMemoryArtifactStore();
|
||||
const executor = new FakeExecutor(
|
||||
(request) =>
|
||||
new Promise((resolve) => {
|
||||
request.signal.addEventListener('abort', () => {
|
||||
resolve({
|
||||
status: 'cancelled',
|
||||
failure: { code: 'cancelled', message: 'Render cancelled' },
|
||||
});
|
||||
});
|
||||
}),
|
||||
);
|
||||
const coordinator = new RenderCoordinator(executor, jobs, artifacts.store);
|
||||
const dir = await projectDir();
|
||||
const id = await coordinator.submit(coordinator.reserve('bob'), dir, renderOptions);
|
||||
await waitForJob(jobs, id, () => executor.requests.length === 1);
|
||||
|
||||
expect(await coordinator.cancel(id)).toBe(true);
|
||||
const job = await waitForJob(jobs, id, (current) => current.status === 'cancelled');
|
||||
expect(job.failure).toEqual({ code: 'cancelled', message: 'Render cancelled' });
|
||||
await expect(access(dir)).rejects.toThrow();
|
||||
});
|
||||
|
||||
it('keeps deadline failure classification from a replaceable executor', async () => {
|
||||
const jobs = createMemoryJobStore();
|
||||
const artifacts = createMemoryArtifactStore();
|
||||
const executor = new FakeExecutor(async () => ({
|
||||
status: 'failed',
|
||||
failure: { code: 'deadline_exceeded', message: 'Render exceeded the deadline' },
|
||||
}));
|
||||
const coordinator = new RenderCoordinator(executor, jobs, artifacts.store, {
|
||||
jobDeadlineMs: 42,
|
||||
});
|
||||
const dir = await projectDir();
|
||||
const id = await coordinator.submit(coordinator.reserve('carol'), dir, renderOptions);
|
||||
|
||||
const job = await waitForJob(jobs, id, (current) => current.status === 'failed');
|
||||
expect(executor.requests[0].deadlineMs).toBe(42);
|
||||
expect(job).toMatchObject({
|
||||
status: 'failed',
|
||||
error: 'Render exceeded the deadline',
|
||||
failure: { code: 'deadline_exceeded' },
|
||||
});
|
||||
expect(artifacts.paths.has(id)).toBe(false);
|
||||
await expect(access(dir)).rejects.toThrow();
|
||||
});
|
||||
|
||||
it('classifies unexpected executor errors and still performs cleanup', async () => {
|
||||
const jobs = createMemoryJobStore();
|
||||
const artifacts = createMemoryArtifactStore();
|
||||
const executor = new FakeExecutor(async () => {
|
||||
throw new Error('executor unavailable');
|
||||
});
|
||||
const coordinator = new RenderCoordinator(executor, jobs, artifacts.store);
|
||||
const dir = await projectDir();
|
||||
const id = await coordinator.submit(coordinator.reserve('dana'), dir, renderOptions);
|
||||
|
||||
const job = await waitForJob(jobs, id, (current) => current.status === 'failed');
|
||||
expect(job).toMatchObject({
|
||||
error: 'executor unavailable',
|
||||
failure: { code: 'execution_failed', message: 'executor unavailable' },
|
||||
});
|
||||
await expect(access(dir)).rejects.toThrow();
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,183 @@
|
||||
import { createRenderJob, RenderCancelledError } from '@hyperframes/producer';
|
||||
import { describe, expect, it } from 'vitest';
|
||||
import { InProcessExecutor } from '../src/render-executor.js';
|
||||
import type { RenderExecutionRequest } from '../src/types.js';
|
||||
|
||||
const options = { fps: 30, quality: 'standard', format: 'mp4' } as const;
|
||||
|
||||
function request(overrides: Partial<RenderExecutionRequest> = {}): RenderExecutionRequest {
|
||||
return {
|
||||
projectDir: '/tmp/project',
|
||||
outputPath: '/tmp/project/output.mp4',
|
||||
options,
|
||||
signal: new AbortController().signal,
|
||||
deadlineMs: 1_000,
|
||||
onProgress() {},
|
||||
...overrides,
|
||||
};
|
||||
}
|
||||
|
||||
function setPerformance(job: ReturnType<typeof createRenderJob>, captureMode = 'beginframe'): void {
|
||||
job.perfSummary = {
|
||||
renderId: job.id,
|
||||
totalElapsedMs: 1_200,
|
||||
fps: 30,
|
||||
quality: 'standard',
|
||||
workers: 2,
|
||||
chunkedEncode: false,
|
||||
chunkSizeFrames: null,
|
||||
compositionDurationSeconds: 2,
|
||||
totalFrames: 60,
|
||||
resolution: { width: 1920, height: 1080 },
|
||||
videoCount: 0,
|
||||
audioCount: 0,
|
||||
stages: { compileMs: 100, captureMs: 900, encodeMs: 200 },
|
||||
drawElement: {
|
||||
mode: captureMode,
|
||||
workerEncode: false,
|
||||
verifyArmed: 0,
|
||||
verifyChecked: 0,
|
||||
verifyInitMs: 0,
|
||||
selfVerifyFallback: false,
|
||||
blankSuspects: 0,
|
||||
blankDeterministicAccepts: 0,
|
||||
blankRecaptures: 0,
|
||||
boundaryFrames: 0,
|
||||
ncprFallbacks: 0,
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
describe('InProcessExecutor', () => {
|
||||
it('normalizes progress and maps producer performance into domain data', async () => {
|
||||
const progress = [];
|
||||
const executor = new InProcessExecutor(
|
||||
{ workers: 2, requireBeginFrame: true },
|
||||
{
|
||||
createJob(options) {
|
||||
return createRenderJob(options);
|
||||
},
|
||||
async executeJob(job, _projectDir, _outputPath, onProgress) {
|
||||
job.progress = 150;
|
||||
job.currentStage = 'capturing';
|
||||
job.framesRendered = 60;
|
||||
job.totalFrames = 60;
|
||||
await onProgress(job);
|
||||
setPerformance(job);
|
||||
},
|
||||
},
|
||||
);
|
||||
|
||||
const result = await executor.execute(
|
||||
request({
|
||||
onProgress(update) {
|
||||
progress.push(update);
|
||||
},
|
||||
}),
|
||||
);
|
||||
|
||||
expect(progress).toEqual([
|
||||
{ progress: 1, stage: 'capturing', framesRendered: 60, totalFrames: 60 },
|
||||
]);
|
||||
expect(result).toEqual({
|
||||
status: 'succeeded',
|
||||
performance: {
|
||||
totalElapsedMs: 1_200,
|
||||
stages: { compileMs: 100, captureMs: 900, encodeMs: 200 },
|
||||
workers: 2,
|
||||
totalFrames: 60,
|
||||
captureMode: 'beginframe',
|
||||
},
|
||||
});
|
||||
});
|
||||
|
||||
it('classifies a user abort as cancellation', async () => {
|
||||
const abort = new AbortController();
|
||||
const executor = new InProcessExecutor(
|
||||
{},
|
||||
{
|
||||
createJob(options) {
|
||||
return createRenderJob(options);
|
||||
},
|
||||
async executeJob(_job, _projectDir, _outputPath, _onProgress, signal) {
|
||||
await new Promise<void>((resolve) => signal.addEventListener('abort', () => resolve()));
|
||||
throw new RenderCancelledError('cancelled', 'aborted');
|
||||
},
|
||||
},
|
||||
);
|
||||
|
||||
const execution = executor.execute(request({ signal: abort.signal }));
|
||||
abort.abort();
|
||||
|
||||
await expect(execution).resolves.toEqual({
|
||||
status: 'cancelled',
|
||||
failure: { code: 'cancelled', message: 'cancelled' },
|
||||
});
|
||||
});
|
||||
|
||||
it('preserves the first user-cancellation cause when producer shutdown crosses the deadline', async () => {
|
||||
const abort = new AbortController();
|
||||
const executor = new InProcessExecutor(
|
||||
{},
|
||||
{
|
||||
createJob(options) {
|
||||
return createRenderJob(options);
|
||||
},
|
||||
async executeJob(_job, _projectDir, _outputPath, _onProgress, signal) {
|
||||
await new Promise<void>((resolve) => signal.addEventListener('abort', () => resolve()));
|
||||
await new Promise((resolve) => setTimeout(resolve, 20));
|
||||
throw new RenderCancelledError('cancelled after cleanup', 'aborted');
|
||||
},
|
||||
},
|
||||
);
|
||||
|
||||
const execution = executor.execute(request({ signal: abort.signal, deadlineMs: 5 }));
|
||||
abort.abort();
|
||||
|
||||
await expect(execution).resolves.toEqual({
|
||||
status: 'cancelled',
|
||||
failure: { code: 'cancelled', message: 'cancelled after cleanup' },
|
||||
});
|
||||
});
|
||||
|
||||
it('enforces the deadline and classifies it independently from cancellation', async () => {
|
||||
const executor = new InProcessExecutor(
|
||||
{},
|
||||
{
|
||||
createJob(options) {
|
||||
return createRenderJob(options);
|
||||
},
|
||||
async executeJob(_job, _projectDir, _outputPath, _onProgress, signal) {
|
||||
await new Promise<void>((resolve) => signal.addEventListener('abort', () => resolve()));
|
||||
throw new RenderCancelledError('deadline', 'aborted');
|
||||
},
|
||||
},
|
||||
);
|
||||
|
||||
await expect(executor.execute(request({ deadlineMs: 1 }))).resolves.toEqual({
|
||||
status: 'failed',
|
||||
failure: { code: 'deadline_exceeded', message: 'Render exceeded the deadline' },
|
||||
});
|
||||
});
|
||||
|
||||
it('classifies a required capture-mode mismatch without leaking producer status', async () => {
|
||||
const executor = new InProcessExecutor(
|
||||
{ requireBeginFrame: true },
|
||||
{
|
||||
createJob(options) {
|
||||
return createRenderJob(options);
|
||||
},
|
||||
async executeJob(job) {
|
||||
setPerformance(job, 'screenshot');
|
||||
},
|
||||
},
|
||||
);
|
||||
|
||||
const result = await executor.execute(request());
|
||||
expect(result.status).toBe('failed');
|
||||
if (result.status === 'failed') {
|
||||
expect(result.failure.code).toBe('unsupported_capture_mode');
|
||||
expect(result.failure.message).toMatch(/beginFrame/i);
|
||||
}
|
||||
});
|
||||
});
|
||||
@@ -6,60 +6,42 @@
|
||||
* (multipart buffering → file read → extraction) at once. Everything else waits
|
||||
* with its body unconsumed, so a burst of near-cap uploads can't stack in memory.
|
||||
*
|
||||
* We drive the real Hono app (`createApp`) with a fake manager/stores and a
|
||||
* We drive the real Hono app (`createApp`) with a fake coordinator/stores and a
|
||||
* `unzipProject` stub that parks — recording how many calls are simultaneously
|
||||
* "inside" — so we can assert the peak never exceeds the gate's permit count.
|
||||
*/
|
||||
import { describe, it, expect, beforeAll } from 'vitest';
|
||||
import { Semaphore } from '../src/semaphore.js';
|
||||
import { mkdtemp, rm, writeFile } from 'node:fs/promises';
|
||||
import { tmpdir } from 'node:os';
|
||||
import { join } from 'node:path';
|
||||
import { afterEach, beforeAll, describe, expect, it } from 'vitest';
|
||||
import type { ArtifactStore } from '../src/artifact-store.js';
|
||||
import type { JobStore } from '../src/job-store.js';
|
||||
import type { ArtifactStore, ArtifactLocation } from '../src/artifact-store.js';
|
||||
import type { RenderCoordinatorOptions } from '../src/render-coordinator.js';
|
||||
import type { RenderExecutor } from '../src/render-executor.js';
|
||||
import { Semaphore } from '../src/semaphore.js';
|
||||
import type { RenderJobRecord } from '../src/types.js';
|
||||
import {
|
||||
createMemoryArtifactStore,
|
||||
createMemoryJobStore,
|
||||
succeedingExecutor,
|
||||
} from './support/fakes.js';
|
||||
|
||||
// Prevent main.ts from binding a port when we import it.
|
||||
process.env.RENDER_SERVICE_NO_LISTEN = 'true';
|
||||
|
||||
// Loaded in beforeAll after the env guard above is set.
|
||||
let createApp: typeof import('../src/main.js').createApp;
|
||||
let RenderManager: typeof import('../src/render-manager.js').RenderManager;
|
||||
let RenderCoordinator: typeof import('../src/render-coordinator.js').RenderCoordinator;
|
||||
const scratch: string[] = [];
|
||||
|
||||
beforeAll(async () => {
|
||||
({ createApp } = await import('../src/main.js'));
|
||||
({ RenderManager } = await import('../src/render-manager.js'));
|
||||
({ RenderCoordinator } = await import('../src/render-coordinator.js'));
|
||||
});
|
||||
|
||||
function fakeJobStore(): JobStore {
|
||||
const jobs = new Map<string, RenderJobRecord>();
|
||||
return {
|
||||
async create(r) {
|
||||
jobs.set(r.id, r);
|
||||
},
|
||||
async get(id) {
|
||||
return jobs.get(id) ?? null;
|
||||
},
|
||||
async update(id, patch) {
|
||||
const e = jobs.get(id);
|
||||
if (e) jobs.set(id, { ...e, ...patch });
|
||||
},
|
||||
async remove(id) {
|
||||
jobs.delete(id);
|
||||
},
|
||||
async list() {
|
||||
return [...jobs.values()];
|
||||
},
|
||||
async countActiveForUser() {
|
||||
return 0;
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
const fakeArtifacts: ArtifactStore = {
|
||||
async put() {},
|
||||
async locate(): Promise<ArtifactLocation | null> {
|
||||
return null;
|
||||
},
|
||||
async remove() {},
|
||||
};
|
||||
afterEach(async () => {
|
||||
await Promise.all(scratch.splice(0).map((path) => rm(path, { recursive: true, force: true })));
|
||||
});
|
||||
|
||||
/** Build a valid-looking multipart body for `POST /render`. */
|
||||
function renderRequest(sizeBytes = 4096, identity = 'anon'): Request {
|
||||
@@ -75,6 +57,43 @@ function renderRequest(sizeBytes = 4096, identity = 'anon'): Request {
|
||||
});
|
||||
}
|
||||
|
||||
function testApp(
|
||||
executor: RenderExecutor,
|
||||
options: {
|
||||
jobs?: JobStore;
|
||||
artifacts?: ArtifactStore;
|
||||
coordinator?: RenderCoordinatorOptions;
|
||||
makeProjectDir?: () => Promise<string>;
|
||||
} = {},
|
||||
) {
|
||||
const jobs = options.jobs ?? createMemoryJobStore();
|
||||
const artifacts = options.artifacts ?? createMemoryArtifactStore().store;
|
||||
const coordinator = new RenderCoordinator(executor, jobs, artifacts, options.coordinator);
|
||||
const app = createApp({
|
||||
jobs,
|
||||
artifacts,
|
||||
coordinator,
|
||||
extractionGate: new Semaphore(1),
|
||||
unzipProject: async () => {},
|
||||
...(options.makeProjectDir ? { makeProjectDir: options.makeProjectDir } : {}),
|
||||
});
|
||||
return { app, artifacts, coordinator, jobs };
|
||||
}
|
||||
|
||||
async function waitForPoll(
|
||||
app: ReturnType<typeof createApp>,
|
||||
jobId: string,
|
||||
status: RenderJobRecord['status'],
|
||||
): Promise<Record<string, unknown>> {
|
||||
for (let attempt = 0; attempt < 100; attempt += 1) {
|
||||
const response = await app.fetch(new Request(`http://test/render/${jobId}`));
|
||||
const body = (await response.json()) as Record<string, unknown>;
|
||||
if (body.status === status) return body;
|
||||
await new Promise((resolve) => setTimeout(resolve, 5));
|
||||
}
|
||||
throw new Error(`Timed out waiting for ${jobId} to reach ${status}`);
|
||||
}
|
||||
|
||||
describe('POST /render buffering/extraction bound', () => {
|
||||
it('never lets more than the permit count into the buffering+extraction section', async () => {
|
||||
const PERMITS = 2;
|
||||
@@ -99,14 +118,15 @@ describe('POST /render buffering/extraction bound', () => {
|
||||
let n = 0;
|
||||
const makeProjectDir = async () => `/tmp/fake-${n++}`;
|
||||
|
||||
const jobs = fakeJobStore();
|
||||
const jobs = createMemoryJobStore();
|
||||
const artifacts = createMemoryArtifactStore().store;
|
||||
// A big per-user cap so all REQUESTS are admitted (we're testing the gate,
|
||||
// not the per-identity guard); unique identities would also work.
|
||||
const manager = new RenderManager(jobs, fakeArtifacts);
|
||||
const coordinator = new RenderCoordinator(succeedingExecutor, jobs, artifacts);
|
||||
const app = createApp({
|
||||
jobs,
|
||||
artifacts: fakeArtifacts,
|
||||
manager,
|
||||
artifacts,
|
||||
coordinator,
|
||||
extractionGate: new Semaphore(PERMITS),
|
||||
unzipProject,
|
||||
makeProjectDir,
|
||||
@@ -137,3 +157,150 @@ describe('POST /render buffering/extraction bound', () => {
|
||||
expect(peak).toBe(PERMITS);
|
||||
});
|
||||
});
|
||||
|
||||
describe('render HTTP contract through a replaceable executor', () => {
|
||||
it('preserves submit, polling, and file download behavior', async () => {
|
||||
const dir = await mkdtemp(join(tmpdir(), 'render-route-success-'));
|
||||
scratch.push(dir);
|
||||
const executor: RenderExecutor = {
|
||||
async execute(request) {
|
||||
await request.onProgress({
|
||||
progress: 0.5,
|
||||
stage: 'capturing',
|
||||
framesRendered: 12,
|
||||
totalFrames: 24,
|
||||
});
|
||||
await writeFile(request.outputPath, Buffer.from('fake-mp4'));
|
||||
return { status: 'succeeded' };
|
||||
},
|
||||
};
|
||||
const { app } = testApp(executor, { makeProjectDir: async () => dir });
|
||||
|
||||
const submit = await app.fetch(renderRequest());
|
||||
expect(submit.status).toBe(202);
|
||||
const { jobId } = (await submit.json()) as { jobId: string };
|
||||
|
||||
await expect(waitForPoll(app, jobId, 'succeeded')).resolves.toEqual({
|
||||
jobId,
|
||||
status: 'succeeded',
|
||||
progress: 1,
|
||||
currentStage: 'complete',
|
||||
framesRendered: 12,
|
||||
totalFrames: 24,
|
||||
done: true,
|
||||
});
|
||||
|
||||
const download = await app.fetch(new Request(`http://test/render/${jobId}/download`));
|
||||
expect(download.status).toBe(200);
|
||||
expect(download.headers.get('content-type')).toBe('video/mp4');
|
||||
expect(download.headers.get('content-disposition')).toBe(`attachment; filename="${jobId}.mp4"`);
|
||||
expect(Buffer.from(await download.arrayBuffer()).toString()).toBe('fake-mp4');
|
||||
});
|
||||
|
||||
it('preserves queued cancellation and its polling shape', async () => {
|
||||
const dir = await mkdtemp(join(tmpdir(), 'render-route-queued-cancel-'));
|
||||
scratch.push(dir);
|
||||
const { app } = testApp(succeedingExecutor, {
|
||||
coordinator: { maxConcurrency: 0 },
|
||||
makeProjectDir: async () => dir,
|
||||
});
|
||||
const submit = await app.fetch(renderRequest());
|
||||
const { jobId } = (await submit.json()) as { jobId: string };
|
||||
|
||||
const cancel = await app.fetch(
|
||||
new Request(`http://test/render/${jobId}`, { method: 'DELETE' }),
|
||||
);
|
||||
expect(cancel.status).toBe(200);
|
||||
await expect(cancel.json()).resolves.toEqual({ cancelled: true });
|
||||
await expect(waitForPoll(app, jobId, 'cancelled')).resolves.toEqual({
|
||||
jobId,
|
||||
status: 'cancelled',
|
||||
progress: 0,
|
||||
currentStage: 'cancelled',
|
||||
done: true,
|
||||
});
|
||||
});
|
||||
|
||||
it('preserves the running-cancellation error in polling', async () => {
|
||||
const dir = await mkdtemp(join(tmpdir(), 'render-route-running-cancel-'));
|
||||
scratch.push(dir);
|
||||
let started!: () => void;
|
||||
const executionStarted = new Promise<void>((resolve) => {
|
||||
started = resolve;
|
||||
});
|
||||
const executor: RenderExecutor = {
|
||||
async execute(request) {
|
||||
started();
|
||||
await new Promise<void>((resolve) =>
|
||||
request.signal.addEventListener('abort', () => resolve(), { once: true }),
|
||||
);
|
||||
return {
|
||||
status: 'cancelled',
|
||||
failure: { code: 'cancelled', message: 'render_cancelled' },
|
||||
};
|
||||
},
|
||||
};
|
||||
const { app } = testApp(executor, { makeProjectDir: async () => dir });
|
||||
const submit = await app.fetch(renderRequest());
|
||||
const { jobId } = (await submit.json()) as { jobId: string };
|
||||
await executionStarted;
|
||||
|
||||
const cancel = await app.fetch(
|
||||
new Request(`http://test/render/${jobId}`, { method: 'DELETE' }),
|
||||
);
|
||||
expect(cancel.status).toBe(200);
|
||||
await expect(waitForPoll(app, jobId, 'cancelled')).resolves.toEqual({
|
||||
jobId,
|
||||
status: 'cancelled',
|
||||
progress: 0,
|
||||
currentStage: 'cancelled',
|
||||
error: 'render_cancelled',
|
||||
done: true,
|
||||
});
|
||||
});
|
||||
|
||||
it('preserves redirect downloads and terminal download errors', async () => {
|
||||
const jobs = createMemoryJobStore();
|
||||
const now = Date.now();
|
||||
await jobs.create({
|
||||
id: 'ready',
|
||||
status: 'succeeded',
|
||||
progress: 1,
|
||||
currentStage: 'complete',
|
||||
createdAtMs: now,
|
||||
updatedAtMs: now,
|
||||
projectDir: '/tmp/ready',
|
||||
outputPath: '/tmp/ready/output.mp4',
|
||||
});
|
||||
await jobs.create({
|
||||
id: 'failed',
|
||||
status: 'failed',
|
||||
progress: 0.25,
|
||||
currentStage: 'failed',
|
||||
error: 'boom',
|
||||
createdAtMs: now,
|
||||
updatedAtMs: now,
|
||||
projectDir: '/tmp/failed',
|
||||
});
|
||||
const artifacts: ArtifactStore = {
|
||||
async put() {},
|
||||
async locate(id) {
|
||||
return id === 'ready' ? { kind: 'url', href: 'https://example.test/video.mp4' } : null;
|
||||
},
|
||||
async remove() {},
|
||||
};
|
||||
const { app } = testApp(succeedingExecutor, { artifacts, jobs });
|
||||
|
||||
const redirect = await app.fetch(new Request('http://test/render/ready/download'));
|
||||
expect(redirect.status).toBe(302);
|
||||
expect(redirect.headers.get('location')).toBe('https://example.test/video.mp4');
|
||||
|
||||
const notReady = await app.fetch(new Request('http://test/render/failed/download'));
|
||||
expect(notReady.status).toBe(409);
|
||||
await expect(notReady.json()).resolves.toEqual({ error: 'Job not ready (status: failed)' });
|
||||
|
||||
const missing = await app.fetch(new Request('http://test/render/missing/download'));
|
||||
expect(missing.status).toBe(404);
|
||||
await expect(missing.json()).resolves.toEqual({ error: 'Job not found' });
|
||||
});
|
||||
});
|
||||
|
||||
@@ -0,0 +1,57 @@
|
||||
import type { ArtifactLocation, ArtifactStore } from '../../src/artifact-store.js';
|
||||
import type { JobStore } from '../../src/job-store.js';
|
||||
import type { RenderExecutor } from '../../src/render-executor.js';
|
||||
import type { RenderJobRecord } from '../../src/types.js';
|
||||
|
||||
export function createMemoryJobStore(): JobStore {
|
||||
const jobs = new Map<string, RenderJobRecord>();
|
||||
return {
|
||||
async create(record) {
|
||||
jobs.set(record.id, record);
|
||||
},
|
||||
async get(id) {
|
||||
return jobs.get(id) ?? null;
|
||||
},
|
||||
async update(id, patch) {
|
||||
const current = jobs.get(id);
|
||||
if (current) jobs.set(id, { ...current, ...patch, updatedAtMs: Date.now() });
|
||||
},
|
||||
async remove(id) {
|
||||
jobs.delete(id);
|
||||
},
|
||||
async list() {
|
||||
return [...jobs.values()];
|
||||
},
|
||||
async countActiveForUser(userId) {
|
||||
return [...jobs.values()].filter(
|
||||
(job) => job.userId === userId && (job.status === 'queued' || job.status === 'running'),
|
||||
).length;
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
export function createMemoryArtifactStore(): {
|
||||
paths: Map<string, string>;
|
||||
store: ArtifactStore;
|
||||
} {
|
||||
const paths = new Map<string, string>();
|
||||
const store: ArtifactStore = {
|
||||
async put(id, sourcePath) {
|
||||
paths.set(id, sourcePath);
|
||||
},
|
||||
async locate(id): Promise<ArtifactLocation | null> {
|
||||
const path = paths.get(id);
|
||||
return path ? { kind: 'file', path } : null;
|
||||
},
|
||||
async remove(id) {
|
||||
paths.delete(id);
|
||||
},
|
||||
};
|
||||
return { paths, store };
|
||||
}
|
||||
|
||||
export const succeedingExecutor: RenderExecutor = {
|
||||
async execute() {
|
||||
return { status: 'succeeded' };
|
||||
},
|
||||
};
|
||||
Reference in New Issue
Block a user