From 74d7bfde4870bbc1c6c4471cfc004964807df3e7 Mon Sep 17 00:00:00 2001 From: James Date: Sun, 26 Jul 2026 05:19:55 +0000 Subject: [PATCH] feat(gcp-cloud-run): publish plan v2 directly to GCS --- .../src/gcsPlanV2Publisher.test.ts | 167 ++++++++++++++++++ .../gcp-cloud-run/src/gcsPlanV2Publisher.ts | 136 ++++++++++++++ packages/gcp-cloud-run/src/index.ts | 4 + packages/gcp-cloud-run/src/server.test.ts | 23 ++- packages/gcp-cloud-run/src/server.ts | 83 ++++----- 5 files changed, 356 insertions(+), 57 deletions(-) create mode 100644 packages/gcp-cloud-run/src/gcsPlanV2Publisher.test.ts create mode 100644 packages/gcp-cloud-run/src/gcsPlanV2Publisher.ts diff --git a/packages/gcp-cloud-run/src/gcsPlanV2Publisher.test.ts b/packages/gcp-cloud-run/src/gcsPlanV2Publisher.test.ts new file mode 100644 index 000000000..6d2c4a3a7 --- /dev/null +++ b/packages/gcp-cloud-run/src/gcsPlanV2Publisher.test.ts @@ -0,0 +1,167 @@ +// fallow-ignore-file code-duplication complexity +import { afterEach, describe, expect, it } from "bun:test"; +import { createHash } from "node:crypto"; +import { mkdtempSync, rmSync, statSync, writeFileSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { asStorage, FakeGcs } from "./__fixtures__/fakeGcs.js"; +import { GcsPlanV2ArtifactPublisher } from "./gcsPlanV2Publisher.js"; + +const roots: string[] = []; + +afterEach(() => { + for (const root of roots) rmSync(root, { recursive: true, force: true }); + roots.length = 0; +}); + +function makeSource(contents: string): { + readonly root: string; + readonly path: string; + readonly digest: string; + readonly sizeBytes: number; +} { + const root = mkdtempSync(join(tmpdir(), "hf-gcs-plan-v2-publisher-")); + roots.push(root); + const path = join(root, "artifact.bin"); + writeFileSync(path, contents); + return { + root, + path, + digest: createHash("sha256").update(contents).digest("hex"), + sizeBytes: statSync(path).size, + }; +} + +function manifestFor(digest: string, marker = "one"): string { + return JSON.stringify({ + planHash: marker, + artifacts: [{ path: "compiled/index.html", sha256: digest, sizeBytes: 5 }], + }); +} + +describe("GcsPlanV2ArtifactPublisher", () => { + it("trims an arbitrary trailing-slash run in linear time", () => { + const publisher = new GcsPlanV2ArtifactPublisher({ + storage: asStorage(new FakeGcs()), + planOutputGcsPrefix: `gs://bucket/render${"/".repeat(10_000)}`, + }); + + expect(publisher.artifactPrefix).toBe("gs://bucket/render/v2/artifacts/sha256"); + expect(publisher.manifestUri).toBe("gs://bucket/render/v2/manifest.json"); + }); + + it("publishes immutable blobs before the fixed-key manifest", async () => { + const source = makeSource("hello"); + const gcs = new FakeGcs(); + const artifactPrefix = "gs://bucket/render/v2/artifacts/sha256"; + const manifestUri = "gs://bucket/render/v2/manifest.json"; + const publisher = new GcsPlanV2ArtifactPublisher({ + storage: asStorage(gcs), + planOutputGcsPrefix: "gs://bucket/render", + temporaryRoot: source.root, + }); + + await publisher.putBlob({ + sourcePath: source.path, + sha256: source.digest, + sizeBytes: source.sizeBytes, + }); + const manifest = manifestFor(source.digest); + await publisher.commitManifest(manifest); + + const blobUri = `${artifactPrefix}/${source.digest.slice(0, 2)}/${source.digest}`; + expect(gcs.ops.filter((operation) => operation.kind === "upload").map((op) => op.uri)).toEqual([ + blobUri, + manifestUri, + ]); + expect(gcs.objects.get(manifestUri)?.toString("utf8")).toBe(manifest); + }); + + it("refuses to expose a manifest that references an unpublished digest", async () => { + const source = makeSource("hello"); + const gcs = new FakeGcs(); + const manifestUri = "gs://bucket/render/v2/manifest.json"; + const publisher = new GcsPlanV2ArtifactPublisher({ + storage: asStorage(gcs), + planOutputGcsPrefix: "gs://bucket/render", + temporaryRoot: source.root, + }); + + await expect(publisher.commitManifest(manifestFor(source.digest))).rejects.toMatchObject({ + name: "PlanV2IntegrityError", + }); + expect(gcs.objects.has(manifestUri)).toBe(false); + }); + + it("rejects malformed digests before constructing a GCS object key", async () => { + const source = makeSource("hello"); + const gcs = new FakeGcs(); + const publisher = new GcsPlanV2ArtifactPublisher({ + storage: asStorage(gcs), + planOutputGcsPrefix: "gs://bucket/render", + temporaryRoot: source.root, + }); + + await expect( + publisher.putBlob({ + sourcePath: source.path, + sha256: "../outside-prefix", + sizeBytes: source.sizeBytes, + }), + ).rejects.toMatchObject({ name: "PlanV2IntegrityError" }); + expect(gcs.ops.filter((operation) => operation.kind === "upload")).toHaveLength(0); + }); + + it("reuses matching objects and rejects a conflicting fixed-key manifest", async () => { + const source = makeSource("hello"); + const gcs = new FakeGcs(); + const options = { + storage: asStorage(gcs), + planOutputGcsPrefix: "gs://bucket/render", + temporaryRoot: source.root, + }; + const blob = { + sourcePath: source.path, + sha256: source.digest, + sizeBytes: source.sizeBytes, + }; + const first = new GcsPlanV2ArtifactPublisher(options); + await first.putBlob(blob); + await first.commitManifest(manifestFor(source.digest, "one")); + + const retry = new GcsPlanV2ArtifactPublisher(options); + await retry.putBlob(blob); + await retry.commitManifest(manifestFor(source.digest, "one")); + expect(gcs.ops.filter((operation) => operation.kind === "upload")).toHaveLength(2); + + const conflict = new GcsPlanV2ArtifactPublisher(options); + await conflict.putBlob(blob); + await expect(conflict.commitManifest(manifestFor(source.digest, "two"))).rejects.toMatchObject({ + name: "PLAN_ARTIFACT_DIGEST_MISMATCH", + }); + expect(gcs.ops.filter((operation) => operation.kind === "upload")).toHaveLength(2); + }); + + it("leaves durable remote CAS blobs intact when publication aborts", async () => { + const source = makeSource("hello"); + const gcs = new FakeGcs(); + const publisher = new GcsPlanV2ArtifactPublisher({ + storage: asStorage(gcs), + planOutputGcsPrefix: "gs://bucket/render", + temporaryRoot: source.root, + }); + const blob = { + sourcePath: source.path, + sha256: source.digest, + sizeBytes: source.sizeBytes, + }; + + await publisher.putBlob(blob); + await publisher.abort(); + await publisher.abort(); + expect(gcs.objects.size).toBe(1); + await expect(publisher.putBlob(blob)).rejects.toMatchObject({ + name: "PlanV2IntegrityError", + }); + }); +}); diff --git a/packages/gcp-cloud-run/src/gcsPlanV2Publisher.ts b/packages/gcp-cloud-run/src/gcsPlanV2Publisher.ts new file mode 100644 index 000000000..8ca05c028 --- /dev/null +++ b/packages/gcp-cloud-run/src/gcsPlanV2Publisher.ts @@ -0,0 +1,136 @@ +// fallow-ignore-file code-duplication +import { createHash } from "node:crypto"; +import { mkdirSync, mkdtempSync, rmSync, statSync, writeFileSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import type { Storage } from "@google-cloud/storage"; +import { + PlanV2IntegrityError, + type PlanV2ArtifactPublisher, + type PlanV2PublishBlob, +} from "@hyperframes/producer/distributed"; +import { parseGcsUri, uploadContentAddressedFileToGcs } from "./gcsTransport.js"; + +export interface GcsPlanV2ArtifactPublisherOptions { + readonly storage: Storage; + /** Validated render output prefix from which all v2 object keys are derived. */ + readonly planOutputGcsPrefix: string; + /** Planner-local scratch parent for the small manifest upload file. */ + readonly temporaryRoot?: string; +} + +function isRecord(value: unknown): value is Record { + return value !== null && typeof value === "object" && !Array.isArray(value); +} + +function assertSha256(value: unknown, label: string): string { + if (typeof value !== "string" || !/^[a-f0-9]{64}$/.test(value)) { + throw new PlanV2IntegrityError(`${label} must be a lowercase SHA-256 digest`); + } + return value; +} + +function manifestDigests(manifestBytes: string): ReadonlySet { + let value: unknown; + try { + value = JSON.parse(manifestBytes); + } catch { + throw new PlanV2IntegrityError("GCS publisher received invalid manifest JSON"); + } + if (!isRecord(value) || !Array.isArray(value.artifacts)) { + throw new PlanV2IntegrityError("GCS publisher manifest requires an artifacts array"); + } + return new Set( + value.artifacts.map((artifact, index) => { + if (!isRecord(artifact)) { + throw new PlanV2IntegrityError(`GCS publisher artifacts[${index}] must be an object`); + } + return assertSha256(artifact.sha256, `GCS publisher artifacts[${index}].sha256`); + }), + ); +} + +function trimTrailingSlash(value: string): string { + let end = value.length; + while (end > 0 && value.charCodeAt(end - 1) === 47) end -= 1; + return value.slice(0, end); +} + +/** + * Manifest-last GCS implementation of the producer's plan-v2 publication seam. + * + * Every path remains private to the planner container. Remote workers receive + * only the manifest URI and artifact prefix and materialize their own target. + */ +export class GcsPlanV2ArtifactPublisher implements PlanV2ArtifactPublisher { + readonly artifactPrefix: string; + readonly manifestUri: string; + readonly #storage: Storage; + readonly #temporaryRoot: string; + readonly #publishedDigests = new Set(); + #state: "open" | "committed" | "aborted" = "open"; + + constructor(options: Readonly) { + const outputPrefix = `${trimTrailingSlash(options.planOutputGcsPrefix)}/v2`; + parseGcsUri(outputPrefix); + this.#storage = options.storage; + this.artifactPrefix = `${outputPrefix}/artifacts/sha256`; + this.manifestUri = `${outputPrefix}/manifest.json`; + this.#temporaryRoot = options.temporaryRoot ?? tmpdir(); + mkdirSync(this.#temporaryRoot, { recursive: true }); + } + + async putBlob(blob: Readonly): Promise { + this.#assertOpen("publish a blob"); + const digest = assertSha256(blob.sha256, "GCS published blob sha256"); + const sourceSize = statSync(blob.sourcePath).size; + if (sourceSize !== blob.sizeBytes) { + throw new PlanV2IntegrityError( + `GCS published blob size changed for ${digest}: expected ${blob.sizeBytes}, got ${sourceSize}`, + ); + } + const uri = `${this.artifactPrefix}/${digest.slice(0, 2)}/${digest}`; + await uploadContentAddressedFileToGcs(this.#storage, blob.sourcePath, uri, digest); + this.#publishedDigests.add(digest); + } + + async commitManifest(manifestBytes: string): Promise { + this.#assertOpen("commit a manifest"); + for (const digest of manifestDigests(manifestBytes)) { + if (!this.#publishedDigests.has(digest)) { + throw new PlanV2IntegrityError( + `cannot commit GCS manifest before referenced blob is durable: ${digest}`, + ); + } + } + + const manifestDigest = createHash("sha256").update(manifestBytes, "utf8").digest("hex"); + const stagingDir = mkdtempSync(join(this.#temporaryRoot, "hf-plan-v2-manifest-")); + const manifestPath = join(stagingDir, "manifest.json"); + try { + writeFileSync(manifestPath, manifestBytes, "utf8"); + await uploadContentAddressedFileToGcs( + this.#storage, + manifestPath, + this.manifestUri, + manifestDigest, + "application/json", + ); + this.#state = "committed"; + } finally { + rmSync(stagingDir, { recursive: true, force: true }); + } + } + + async abort(): Promise { + if (this.#state === "open") this.#state = "aborted"; + // Immutable CAS blobs may be shared with or reused by another retry. + // Unreferenced blobs expire under the bucket's intermediate lifecycle. + } + + #assertOpen(operation: string): void { + if (this.#state !== "open") { + throw new PlanV2IntegrityError(`cannot ${operation} after publisher is ${this.#state}`); + } + } +} diff --git a/packages/gcp-cloud-run/src/index.ts b/packages/gcp-cloud-run/src/index.ts index 4967d5a02..be61a2b6d 100644 --- a/packages/gcp-cloud-run/src/index.ts +++ b/packages/gcp-cloud-run/src/index.ts @@ -53,6 +53,10 @@ export { uploadContentAddressedFileToGcs, uploadFileToGcs, } from "./gcsTransport.js"; +export { + GcsPlanV2ArtifactPublisher, + type GcsPlanV2ArtifactPublisherOptions, +} from "./gcsPlanV2Publisher.js"; // ── Client-side SDK ───────────────────────────────────────────────────────── export { deploySite, type DeploySiteOptions, type SiteHandle } from "./sdk/deploySite.js"; diff --git a/packages/gcp-cloud-run/src/server.test.ts b/packages/gcp-cloud-run/src/server.test.ts index 7b982c437..0028e9d2d 100644 --- a/packages/gcp-cloud-run/src/server.test.ts +++ b/packages/gcp-cloud-run/src/server.test.ts @@ -18,17 +18,18 @@ import { afterEach, describe, expect, it } from "bun:test"; import { createHash } from "node:crypto"; import { existsSync, mkdirSync, mkdtempSync, readFileSync, rmSync, writeFileSync } from "node:fs"; import { tmpdir } from "node:os"; -import { join } from "node:path"; +import { dirname, join } from "node:path"; import { CURRENT_PLAN_PROTOCOL, - createPlanV2FromV1, PLAN_V2_INTEGRITY_UNRECOVERABLE, PlanV2IntegrityError, PlanProtocolUnsupportedError, type AssembleResult, type ChunkResult, type PlanResult, - type PlanV2Result, + type PlanV2ArtifactPublisher, + type PlanV2Manifest, + publishPlanV2FromV1, } from "@hyperframes/producer/distributed"; import { recomputePlanHashFromPlanDir } from "../../producer/src/services/render/stages/freezePlan.js"; import { asStorage, FakeGcs } from "./__fixtures__/fakeGcs.js"; @@ -244,14 +245,18 @@ describe("dispatch", () => { const gcs = new FakeGcs(); await seedProjectTar(gcs, "gs://b/sites/v2/project.tar.gz"); const root = mkTmp("hf-v2-e2e-"); - const planV2 = async ( - _projectDir: string, + const planV2WithPublisher = async ( + projectDir: string, _config: unknown, - planV2Dir: string, - ): Promise => { + publisher: PlanV2ArtifactPublisher, + options: Readonly<{ stagingParentDir?: string }>, + ): Promise => { const v1Dir = join(root, "v1"); makeMinimalV1PlanDir(v1Dir, true); - return createPlanV2FromV1(v1Dir, planV2Dir); + const manifest = await publishPlanV2FromV1(v1Dir, publisher); + expect(options.stagingParentDir).toBe(dirname(projectDir)); + expect(existsSync(join(dirname(projectDir), "plan-v2"))).toBe(false); + return manifest; }; const renderChunk = async ( planDir: string, @@ -278,7 +283,7 @@ describe("dispatch", () => { writeFileSync(finalOutput, "v2-output"); return { framesEncoded: 30, fileSize: 9 }; }; - const deps = depsWith(gcs, { planV2, renderChunk, assemble }); + const deps = depsWith(gcs, { planV2WithPublisher, renderChunk, assemble }); const planned = await dispatch( { diff --git a/packages/gcp-cloud-run/src/server.ts b/packages/gcp-cloud-run/src/server.ts index 98120a869..02959b7c2 100644 --- a/packages/gcp-cloud-run/src/server.ts +++ b/packages/gcp-cloud-run/src/server.ts @@ -31,11 +31,11 @@ import { listPlanV2ArtifactsForTarget, materializePlanV2Target, plan, - planV2, + planV2WithPublisher, type PlanResult, type PlanV2Artifact, + type PlanV2Manifest, type PlanV2MaterializationTarget, - type PlanV2Result, readPlanV2Manifest, renderChunk, } from "@hyperframes/producer/distributed"; @@ -56,12 +56,11 @@ import { downloadGcsObjectToFile, downloadGcsObjectToFileVerified, parseGcsUri, - sha256File, tarDirectory, untarDirectory, - uploadContentAddressedFileToGcs, uploadFileToGcs, } from "./gcsTransport.js"; +import { GcsPlanV2ArtifactPublisher } from "./gcsPlanV2Publisher.js"; /** * Lazily-constructed Storage client. Cached at module scope so warm @@ -84,7 +83,7 @@ export interface HandlerDeps { storage?: Storage; primitives?: { plan: typeof plan; - planV2?: typeof planV2; + planV2WithPublisher?: typeof planV2WithPublisher; renderChunk: typeof renderChunk; assemble: typeof assemble; }; @@ -352,8 +351,8 @@ async function handlePlan(event: PlanEvent, deps?: HandlerDeps): Promise> { const started = Date.now(); const storage = deps?.storage ?? getStorage(); - const primitive = deps?.primitives?.planV2 ?? planV2; + const primitive = deps?.primitives?.planV2WithPublisher ?? planV2WithPublisher; primeChrome(deps); const work = mkdtempSync(join(deps?.tmpRoot ?? tmpdir(), "hf-cr-plan-v2-")); const projectArchive = join(work, "project.tar.gz"); const projectDir = join(work, "project"); - const planV2Dir = join(work, "plan-v2"); try { await downloadGcsObjectToFile(storage, event.ProjectGcsUri, 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.PlanOutputGcsPrefix)}/v2`; - const artifactPrefix = `${outputPrefix}/artifacts/sha256`; - const uniqueArtifacts = [ - ...new Map(manifest.artifacts.map((artifact) => [artifact.sha256, artifact])).values(), - ]; - await mapConcurrent(uniqueArtifacts, 16, async (artifact) => { - await uploadContentAddressedFileToGcs( - storage, - planV2BlobPath(planV2Dir, artifact.sha256), - planV2BlobUri(artifactPrefix, artifact.sha256), - artifact.sha256, - ); - }); - - const manifestUri = `${outputPrefix}/manifest.json`; - await uploadContentAddressedFileToGcs( + const publisher = new GcsPlanV2ArtifactPublisher({ storage, - result.manifestPath, - manifestUri, - await sha256File(result.manifestPath), - "application/json", - ); + planOutputGcsPrefix: event.PlanOutputGcsPrefix, + temporaryRoot: work, + }); + const manifest: PlanV2Manifest = await primitive(projectDir, { ...event.Config }, publisher, { + stagingParentDir: work, + }); return { Action: "plan", PlanProtocol: "v2", - PlanV2ManifestGcsUri: manifestUri, - PlanV2ArtifactGcsPrefix: artifactPrefix, - PlanHash: result.planHash, - ChunkCount: result.chunkCount, - TotalFrames: result.totalFrames, - Fps: result.fps, - Width: result.width, - Height: result.height, - Format: result.format, + PlanV2ManifestGcsUri: publisher.manifestUri, + PlanV2ArtifactGcsPrefix: publisher.artifactPrefix, + PlanHash: manifest.planHash, + ChunkCount: manifest.chunkCount, + TotalFrames: manifest.totalFrames, + Fps: manifest.fps, + Width: manifest.width, + Height: manifest.height, + Format: manifest.format, HasAudio: manifest.artifacts.some((artifact) => artifact.path === "audio.aac"), AudioGcsUri: null, - FfmpegVersion: result.ffmpegVersion, - ProducerVersion: result.producerVersion, + FfmpegVersion: manifest.ffmpegVersion, + ProducerVersion: manifest.producerVersion, DurationMs: Date.now() - started, }; } finally { @@ -736,7 +714,16 @@ async function mapConcurrent( await fn(values[index]!); } } - await Promise.all(Array.from({ length: Math.min(concurrency, values.length) }, () => worker())); + const results = await Promise.allSettled( + Array.from({ length: Math.min(concurrency, values.length) }, () => worker()), + ); + const failure = results.find( + (result): result is PromiseRejectedResult => result.status === "rejected", + ); + // Invocation cleanup removes the work directory in `finally`. Drain all + // sibling downloads before surfacing an error so a late GCS stream cannot + // keep writing into scratch after another artifact fails verification. + if (failure) throw failure.reason; } async function downloadChunkObjects(