mirror of
https://github.com/heygen-com/hyperframes.git
synced 2026-09-03 04:38:33 +00:00
feat(producer): add request-level render concurrency semaphore (#232)
Add a FIFO semaphore to limit concurrent renders in the producer server, preventing Chrome CPU contention that causes beginFrame failures. - New Semaphore utility class (packages/producer/src/utils/semaphore.ts) - Both blocking render and SSE renderStream handlers acquire/release the semaphore - SSE stream sends a "queued" event when request must wait - New GET /render/queue endpoint exposes active/queued render counts - Configurable via HandlerOptions.maxConcurrentRenders or PRODUCER_MAX_CONCURRENT_RENDERS env var (default: 2) - New --max-concurrent-renders CLI flag (1-10) Co-authored-by: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 4.6
parent
d86e4cb3c3
commit
9115d7364c
@@ -7,6 +7,7 @@
|
||||
* Routes:
|
||||
* POST /render — blocking render, returns JSON
|
||||
* POST /render/stream — SSE streaming render with progress
|
||||
* GET /render/queue — current render queue status
|
||||
* POST /lint — blocking Hyperframe lint
|
||||
* GET /health — health check
|
||||
* GET /outputs/:token — download rendered MP4
|
||||
@@ -36,6 +37,7 @@ import {
|
||||
import { prepareHyperframeLintBody, runHyperframeLint } from "./services/hyperframeLint.js";
|
||||
import { resolveRenderPaths } from "./utils/paths.js";
|
||||
import { defaultLogger, type ProducerLogger } from "./logger.js";
|
||||
import { Semaphore } from "./utils/semaphore.js";
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Types
|
||||
@@ -51,6 +53,8 @@ export interface HandlerOptions {
|
||||
outputUrlPrefix?: string;
|
||||
/** TTL for output artifact download tokens (ms). Default: 15 minutes. */
|
||||
artifactTtlMs?: number;
|
||||
/** Max renders that execute simultaneously. Queued requests wait FIFO. Default: 2. */
|
||||
maxConcurrentRenders?: number;
|
||||
}
|
||||
|
||||
export interface ServerOptions extends HandlerOptions {
|
||||
@@ -232,6 +236,7 @@ export interface RenderHandlers {
|
||||
lint: (c: Context) => Promise<Response>;
|
||||
health: (c: Context) => Response;
|
||||
outputs: (c: Context) => Response;
|
||||
queue: (c: Context) => Response;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -248,6 +253,9 @@ export function createRenderHandlers(options: HandlerOptions = {}): RenderHandle
|
||||
const artifactTtlMs =
|
||||
options.artifactTtlMs ?? Number(process.env.PRODUCER_OUTPUT_ARTIFACT_TTL_MS || 15 * 60 * 1000);
|
||||
const store = createArtifactStore(artifactTtlMs);
|
||||
const maxConcurrentRenders =
|
||||
options.maxConcurrentRenders ?? Number(process.env.PRODUCER_MAX_CONCURRENT_RENDERS || 2);
|
||||
const renderSemaphore = new Semaphore(maxConcurrentRenders);
|
||||
const startTime = Date.now();
|
||||
|
||||
const health = (c: Context): Response =>
|
||||
@@ -316,6 +324,8 @@ export function createRenderHandlers(options: HandlerOptions = {}): RenderHandle
|
||||
const outputDir = dirname(absoluteOutputPath);
|
||||
if (!existsSync(outputDir)) mkdirSync(outputDir, { recursive: true });
|
||||
|
||||
const release = await renderSemaphore.acquire();
|
||||
|
||||
log.info("render started", {
|
||||
requestId,
|
||||
projectDir: input.projectDir,
|
||||
@@ -387,6 +397,7 @@ export function createRenderHandlers(options: HandlerOptions = {}): RenderHandle
|
||||
500,
|
||||
);
|
||||
} finally {
|
||||
release();
|
||||
cleanupTempDir(cleanupProjectDir, log);
|
||||
}
|
||||
};
|
||||
@@ -451,6 +462,17 @@ export function createRenderHandlers(options: HandlerOptions = {}): RenderHandle
|
||||
abortController.abort(new RenderCancelledError("request_aborted"));
|
||||
c.req.raw.signal.addEventListener("abort", onRequestAbort, { once: true });
|
||||
|
||||
if (renderSemaphore.activeCount >= maxConcurrentRenders) {
|
||||
await stream.writeSSE({
|
||||
data: JSON.stringify({
|
||||
type: "queued",
|
||||
requestId,
|
||||
position: renderSemaphore.waitingCount,
|
||||
}),
|
||||
});
|
||||
}
|
||||
const release = await renderSemaphore.acquire();
|
||||
|
||||
try {
|
||||
await executeRenderJob(
|
||||
job,
|
||||
@@ -519,6 +541,7 @@ export function createRenderHandlers(options: HandlerOptions = {}): RenderHandle
|
||||
}),
|
||||
});
|
||||
} finally {
|
||||
release();
|
||||
c.req.raw.signal.removeEventListener("abort", onRequestAbort);
|
||||
cleanupTempDir(cleanupProjectDir, log);
|
||||
}
|
||||
@@ -545,7 +568,14 @@ export function createRenderHandlers(options: HandlerOptions = {}): RenderHandle
|
||||
});
|
||||
};
|
||||
|
||||
return { render, renderStream, lint, health, outputs };
|
||||
const queue = (c: Context): Response =>
|
||||
c.json({
|
||||
maxConcurrentRenders,
|
||||
activeRenders: renderSemaphore.activeCount,
|
||||
queuedRenders: renderSemaphore.waitingCount,
|
||||
});
|
||||
|
||||
return { render, renderStream, lint, health, outputs, queue };
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -562,6 +592,7 @@ export function createProducerApp(options: HandlerOptions = {}): Hono {
|
||||
app.get("/health", handlers.health);
|
||||
app.post("/render", handlers.render);
|
||||
app.post("/render/stream", handlers.renderStream);
|
||||
app.get("/render/queue", handlers.queue);
|
||||
app.post("/lint", handlers.lint);
|
||||
app.get("/outputs/:token", handlers.outputs);
|
||||
|
||||
|
||||
@@ -0,0 +1,39 @@
|
||||
/**
|
||||
* Simple async semaphore for limiting concurrent operations.
|
||||
*/
|
||||
export class Semaphore {
|
||||
private queue: Array<() => void> = [];
|
||||
private active = 0;
|
||||
|
||||
constructor(private readonly maxConcurrent: number) {}
|
||||
|
||||
async acquire(): Promise<() => void> {
|
||||
if (this.active < this.maxConcurrent) {
|
||||
this.active++;
|
||||
return () => this.release();
|
||||
}
|
||||
|
||||
return new Promise<() => void>((resolve) => {
|
||||
this.queue.push(() => {
|
||||
this.active++;
|
||||
resolve(() => this.release());
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
private release(): void {
|
||||
this.active--;
|
||||
const next = this.queue.shift();
|
||||
if (next) next();
|
||||
}
|
||||
|
||||
/** Current number of active slots. */
|
||||
get activeCount(): number {
|
||||
return this.active;
|
||||
}
|
||||
|
||||
/** Number of waiters in the queue. */
|
||||
get waitingCount(): number {
|
||||
return this.queue.length;
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user