feat(aws-lambda): support plan protocol v2 (#2789)

* feat(aws-lambda): support plan protocol v2

* fix(aws-lambda): align SAM v2 terminal errors
This commit is contained in:
James Russo
2026-07-25 23:42:51 -04:00
committed by GitHub
parent c6fdd9c015
commit 5bf61d6df0
46 changed files with 4687 additions and 114 deletions
+304 -6
View File
@@ -20,8 +20,15 @@ import {
type AssembleResult,
type ChunkResult,
type DistributedRenderConfig,
listPlanV2ArtifactsForTarget,
materializePlanV2Target,
plan,
planV2,
type PlanResult,
type PlanV2Artifact,
type PlanV2MaterializationTarget,
type PlanV2Result,
readPlanV2Manifest,
renderChunk,
} from "@hyperframes/producer/distributed";
import { resolveChromeExecutablePath } from "./chromium.js";
@@ -39,9 +46,12 @@ import type {
} from "./events.js";
import {
downloadS3ObjectToFile,
downloadS3ObjectToFileVerified,
parseS3Uri,
sha256File,
tarDirectory,
untarDirectory,
uploadContentAddressedFileToS3,
uploadFileToS3,
} from "./s3Transport.js";
@@ -67,6 +77,7 @@ export interface HandlerDeps {
s3?: S3Client;
primitives?: {
plan: typeof plan;
planV2?: typeof planV2;
renderChunk: typeof renderChunk;
assemble: typeof assemble;
};
@@ -109,6 +120,7 @@ export async function handler(event: LambdaEvent, deps?: HandlerDeps): Promise<L
}
}
} catch (err) {
normalizeTerminalErrorName(err);
// Log before re-throwing so CloudWatch captures the structured
// error context alongside Lambda's default stack trace. Otherwise
// ops only sees the trace and has to correlate with execution
@@ -116,6 +128,7 @@ export async function handler(event: LambdaEvent, deps?: HandlerDeps): Promise<L
logEvent({
event: "handler_error",
action: unwrapped.Action,
input: summarizeEvent(unwrapped),
message: err instanceof Error ? err.message : String(err),
name: err instanceof Error ? err.name : undefined,
});
@@ -123,6 +136,25 @@ export async function handler(event: LambdaEvent, deps?: HandlerDeps): Promise<L
}
}
/**
* AWS Lambda reports `Error.name` to Step Functions, while producer errors
* expose stable machine codes separately. Normalize the terminal codes
* whose historical class names differ from their orchestration contracts.
*/
// The explicit error-name mapping is the public Step Functions failure contract.
// fallow-ignore-next-line complexity
function normalizeTerminalErrorName(error: unknown): void {
if (!error || typeof error !== "object") return;
const candidate = error as { code?: unknown; name?: string };
if (
candidate.code === "PLAN_PROTOCOL_UNSUPPORTED" ||
candidate.code === "PLAN_TOO_LARGE" ||
candidate.code === "PLAN_V2_INTEGRITY_UNRECOVERABLE"
) {
candidate.name = candidate.code;
}
}
/**
* Walk through Step Functions' Map-state and Task-state envelopes until
* the discriminated event is found.
@@ -178,6 +210,8 @@ function logEvent(payload: Record<string, unknown>): void {
* the routable fields (S3 URIs, chunk index, format) needed to triage
* a failure from CloudWatch.
*/
// Keep event variants together so logs share one redaction and summarization boundary.
// fallow-ignore-next-line complexity
function summarizeEvent(
event: PlanEvent | RenderChunkEvent | AssembleEvent,
): Record<string, unknown> {
@@ -186,18 +220,25 @@ function summarizeEvent(
return {
projectS3Uri: event.ProjectS3Uri,
planOutputS3Prefix: event.PlanOutputS3Prefix,
planProtocol: event.PlanProtocol ?? "v1",
format: event.Config.format,
fps: event.Config.fps,
};
case "renderChunk":
return {
planS3Uri: event.PlanS3Uri,
planProtocol: event.PlanProtocol ?? "v1",
...(event.PlanProtocol === "v2"
? { planV2ManifestS3Uri: event.PlanV2ManifestS3Uri }
: { planS3Uri: event.PlanS3Uri }),
chunkIndex: event.ChunkIndex,
format: event.Format,
};
case "assemble":
return {
planS3Uri: event.PlanS3Uri,
planProtocol: event.PlanProtocol ?? "v1",
...(event.PlanProtocol === "v2"
? { planV2ManifestS3Uri: event.PlanV2ManifestS3Uri }
: { planS3Uri: event.PlanS3Uri }),
chunkCount: event.ChunkS3Uris.length,
hasAudio: event.AudioS3Uri !== null,
outputS3Uri: event.OutputS3Uri,
@@ -226,6 +267,9 @@ function primeRuntimeEnv(): void {
// ── Plan ────────────────────────────────────────────────────────────────────
async function handlePlan(event: PlanEvent, deps?: HandlerDeps): Promise<PlanLambdaResult> {
if (event.PlanProtocol === "v2") {
return handlePlanV2(event, deps);
}
const started = Date.now();
const s3 = deps?.s3 ?? getS3Client();
const primitive = deps?.primitives?.plan ?? plan;
@@ -297,12 +341,90 @@ async function handlePlan(event: PlanEvent, deps?: HandlerDeps): Promise<PlanLam
}
}
// Plan v2 orchestration is kept as one transactional boundary: stage, CAS upload,
// and manifest-last publication must stay ordered and fail together.
// fallow-ignore-next-line complexity
async function handlePlanV2(
event: Extract<PlanEvent, { PlanProtocol: "v2" }>,
deps?: HandlerDeps,
): Promise<Extract<PlanLambdaResult, { PlanProtocol: "v2" }>> {
const started = Date.now();
const s3 = deps?.s3 ?? getS3Client();
const primitive = deps?.primitives?.planV2 ?? planV2;
if (!deps?.skipChromeResolution && !process.env.PRODUCER_HEADLESS_SHELL_PATH) {
process.env.PRODUCER_HEADLESS_SHELL_PATH = await resolveChromeExecutablePath();
}
const work = mkdtempSync(join(deps?.tmpRoot ?? tmpdir(), "hf-lambda-plan-v2-"));
const projectArchive = join(work, "project.tar.gz");
const projectDir = join(work, "project");
const planV2Dir = join(work, "plan-v2");
try {
await downloadS3ObjectToFile(s3, event.ProjectS3Uri, projectArchive);
await untarDirectory(projectArchive, projectDir);
const result: PlanV2Result = await primitive(projectDir, { ...event.Config }, planV2Dir);
const manifest = readPlanV2Manifest(planV2Dir);
if (manifest.planHash !== result.planHash) {
throwPlanHashMismatch(result.planHash, manifest.planHash);
}
const outputPrefix = `${trimTrailingSlash(event.PlanOutputS3Prefix)}/v2`;
const artifactPrefix = `${outputPrefix}/artifacts/sha256`;
const uniqueArtifacts = [
...new Map(manifest.artifacts.map((artifact) => [artifact.sha256, artifact])).values(),
];
await mapConcurrent(uniqueArtifacts, 16, async (artifact) => {
const localPath = planV2BlobPath(planV2Dir, artifact.sha256);
await uploadContentAddressedFileToS3(
s3,
localPath,
planV2BlobUri(artifactPrefix, artifact.sha256),
artifact.sha256,
);
});
// Publish the manifest only after every referenced blob is durable.
const manifestUri = `${outputPrefix}/manifest.json`;
await uploadContentAddressedFileToS3(
s3,
result.manifestPath,
manifestUri,
await sha256File(result.manifestPath),
"application/json",
);
return {
Action: "plan",
PlanProtocol: "v2",
PlanV2ManifestS3Uri: manifestUri,
PlanV2ArtifactS3Prefix: artifactPrefix,
PlanHash: result.planHash,
ChunkCount: result.chunkCount,
TotalFrames: result.totalFrames,
Fps: result.fps,
Width: result.width,
Height: result.height,
Format: result.format,
HasAudio: manifest.artifacts.some((artifact) => artifact.path === "audio.aac"),
AudioS3Uri: null,
FfmpegVersion: result.ffmpegVersion,
ProducerVersion: result.producerVersion,
DurationMs: Date.now() - started,
};
} finally {
cleanupDir(work);
}
}
// ── RenderChunk ─────────────────────────────────────────────────────────────
async function handleRenderChunk(
event: RenderChunkEvent,
deps?: HandlerDeps,
): Promise<RenderChunkLambdaResult> {
if (event.PlanProtocol === "v2") {
return handleRenderChunkV2(event, deps);
}
const started = Date.now();
const s3 = deps?.s3 ?? getS3Client();
const primitive = deps?.primitives?.renderChunk ?? renderChunk;
@@ -364,6 +486,53 @@ async function handleRenderChunk(
}
}
// The v2 chunk handler deliberately keeps download, verified materialization,
// render, and upload in one lifecycle so cleanup and errors remain atomic.
// fallow-ignore-next-line complexity
async function handleRenderChunkV2(
event: Extract<RenderChunkEvent, { PlanProtocol: "v2" }>,
deps?: HandlerDeps,
): Promise<RenderChunkLambdaResult> {
const started = Date.now();
const s3 = deps?.s3 ?? getS3Client();
const primitive = deps?.primitives?.renderChunk ?? renderChunk;
if (!deps?.skipChromeResolution && !process.env.PRODUCER_HEADLESS_SHELL_PATH) {
process.env.PRODUCER_HEADLESS_SHELL_PATH = await resolveChromeExecutablePath();
}
const work = mkdtempSync(join(deps?.tmpRoot ?? tmpdir(), "hf-lambda-chunk-v2-"));
try {
const planDir = await downloadAndMaterializePlanV2(
s3,
event,
{ role: "chunk", chunkIndex: event.ChunkIndex },
work,
);
const chunkOutputBase = join(
work,
event.Format === "png-sequence"
? `chunk-${pad(event.ChunkIndex)}`
: `chunk-${pad(event.ChunkIndex)}${formatExtension(event.Format)}`,
);
const result = await primitive(planDir, event.ChunkIndex, chunkOutputBase);
const chunkUri = await uploadChunkOutput(
s3,
result,
event.ChunkOutputS3Prefix,
event.ChunkIndex,
);
return {
Action: "renderChunk",
ChunkS3Uri: chunkUri,
ChunkIndex: event.ChunkIndex,
Sha256: result.sha256,
FramesEncoded: result.framesEncoded,
DurationMs: Date.now() - started,
};
} finally {
cleanupDir(work);
}
}
async function uploadChunkOutput(
s3: S3Client,
result: ChunkResult,
@@ -393,6 +562,9 @@ async function handleAssemble(
event: AssembleEvent,
deps?: HandlerDeps,
): Promise<AssembleLambdaResult> {
if (event.PlanProtocol === "v2") {
return handleAssembleV2(event, deps);
}
const started = Date.now();
const s3 = deps?.s3 ?? getS3Client();
const primitive = deps?.primitives?.assemble ?? assemble;
@@ -442,6 +614,123 @@ async function handleAssemble(
}
}
// Assembly mirrors the chunk lifecycle while adding assembler-only artifacts;
// keeping the steps local makes its temporary-storage ownership explicit.
// fallow-ignore-next-line complexity
async function handleAssembleV2(
event: Extract<AssembleEvent, { PlanProtocol: "v2" }>,
deps?: HandlerDeps,
): Promise<AssembleLambdaResult> {
const started = Date.now();
const s3 = deps?.s3 ?? getS3Client();
const primitive = deps?.primitives?.assemble ?? assemble;
const work = mkdtempSync(join(deps?.tmpRoot ?? tmpdir(), "hf-lambda-assemble-v2-"));
try {
const planDir = await downloadAndMaterializePlanV2(s3, event, { role: "assembler" }, work);
// `downloadAndMaterializePlanV2` materializes atomically. Audio is
// assembler-only and lives at the familiar v1-compatible location.
const audioPath = existsSync(join(planDir, "audio.aac")) ? join(planDir, "audio.aac") : null;
const chunkPaths = await downloadChunkObjects(s3, event.ChunkS3Uris, work, event.Format);
const finalOutput =
event.Format === "png-sequence"
? join(work, "output-frames")
: join(work, `output${formatExtension(event.Format)}`);
const result = await primitive(planDir, chunkPaths, audioPath, finalOutput, {
cfr: event.Cfr === true,
});
if (event.Format === "png-sequence") {
const tarball = `${finalOutput}.tar.gz`;
await tarDirectory(finalOutput, tarball);
await uploadFileToS3(s3, tarball, event.OutputS3Uri, "application/gzip");
} else {
await uploadFileToS3(s3, finalOutput, event.OutputS3Uri);
}
return {
Action: "assemble",
OutputS3Uri: event.OutputS3Uri,
FramesEncoded: result.framesEncoded,
FileSize: result.fileSize,
DurationMs: Date.now() - started,
};
} finally {
cleanupDir(work);
}
}
async function downloadAndMaterializePlanV2(
s3: S3Client,
event: {
PlanV2ManifestS3Uri: string;
PlanV2ArtifactS3Prefix: string;
PlanHash: string;
},
target: PlanV2MaterializationTarget,
work: string,
): Promise<string> {
const transportDir = join(work, "plan-v2");
mkdirSync(transportDir, { recursive: true });
await downloadS3ObjectToFile(s3, event.PlanV2ManifestS3Uri, join(transportDir, "plan.json"));
const manifest = readPlanV2Manifest(transportDir);
if (manifest.planHash !== event.PlanHash) {
throwPlanHashMismatch(event.PlanHash, manifest.planHash);
}
const artifacts = listPlanV2ArtifactsForTarget(manifest, target);
const uniqueArtifacts = [
...new Map(artifacts.map((artifact) => [artifact.sha256, artifact])).values(),
];
await mapConcurrent(uniqueArtifacts, 16, async (artifact) => {
await downloadPlanV2Artifact(s3, event.PlanV2ArtifactS3Prefix, transportDir, artifact);
});
const planDir = join(work, "plan");
materializePlanV2Target(transportDir, target, planDir);
return planDir;
}
async function downloadPlanV2Artifact(
s3: S3Client,
artifactPrefix: string,
planV2Dir: string,
artifact: Readonly<PlanV2Artifact>,
): Promise<void> {
await downloadS3ObjectToFileVerified(
s3,
planV2BlobUri(artifactPrefix, artifact.sha256),
planV2BlobPath(planV2Dir, artifact.sha256),
artifact.sha256,
);
}
function planV2BlobPath(planV2Dir: string, digest: string): string {
return join(planV2Dir, "artifacts", "sha256", digest.slice(0, 2), digest);
}
function planV2BlobUri(prefix: string, digest: string): string {
return `${trimTrailingSlash(prefix)}/${digest.slice(0, 2)}/${digest}`;
}
function throwPlanHashMismatch(expected: string, actual: string): never {
const error = new Error(
`PLAN_HASH_MISMATCH: event PlanHash=${expected} did not match v2 manifest planHash=${actual}`,
);
error.name = "PLAN_HASH_MISMATCH";
throw error;
}
async function mapConcurrent<T>(
values: readonly T[],
concurrency: number,
fn: (value: T) => Promise<void>,
): Promise<void> {
let cursor = 0;
async function worker(): Promise<void> {
while (cursor < values.length) {
const index = cursor++;
await fn(values[index]!);
}
}
await Promise.all(Array.from({ length: Math.min(concurrency, values.length) }, () => worker()));
}
async function downloadChunkObjects(
s3: S3Client,
uris: string[],
@@ -479,16 +768,25 @@ async function downloadChunkObjects(
// ── Helpers ─────────────────────────────────────────────────────────────────
/** Collect every S3 URI that the handler will touch for a given event. */
// This is an exhaustive event-union projection used only for safe log summaries.
// fallow-ignore-next-line complexity
function getEventS3Uris(event: PlanEvent | RenderChunkEvent | AssembleEvent): string[] {
switch (event.Action) {
case "plan":
return [event.ProjectS3Uri, event.PlanOutputS3Prefix];
case "renderChunk":
return [event.PlanS3Uri, event.ChunkOutputS3Prefix];
return event.PlanProtocol === "v2"
? [event.PlanV2ManifestS3Uri, event.PlanV2ArtifactS3Prefix, event.ChunkOutputS3Prefix]
: [event.PlanS3Uri, event.ChunkOutputS3Prefix];
case "assemble":
return [event.PlanS3Uri, ...event.ChunkS3Uris, event.OutputS3Uri, event.AudioS3Uri].filter(
(u): u is string => u != null,
);
return [
...(event.PlanProtocol === "v2"
? [event.PlanV2ManifestS3Uri, event.PlanV2ArtifactS3Prefix]
: [event.PlanS3Uri]),
...event.ChunkS3Uris,
event.OutputS3Uri,
event.AudioS3Uri,
].filter((u): u is string => u != null);
}
}