# HyperFrames distributed render orchestration on Cloud Workflows. # # Plan → BuildChunkList → AssertChunkCount → RenderChunks (parallel) → Assemble # # Mirrors the Step Functions state machine in # `examples/aws-lambda/template.yaml`. Every step POSTs to the same Cloud Run # service URL (passed in as `args.ServiceUrl`) and varies only the body's # `Action`. The service returns the step's small result body on 2xx; on a # non-retryable failure it returns HTTP 400, on a retryable failure HTTP 5xx — # the `retryable` predicate below keys off exactly that split. # # The final returned object accumulates every step's result body so # `getRenderProgress` can read frame totals + per-step durations on success: # { Plan: {...}, Chunks: [{...}, ...], Assemble: {...} } # # Deploy with `gcloud workflows deploy` (the Terraform module / the # `hyperframes cloudrun deploy` command do this for you). main: params: [args] steps: - init: assign: - serviceUrl: ${args.ServiceUrl} - projectGcsUri: ${args.ProjectGcsUri} - planOutputGcsPrefix: ${args.PlanOutputGcsPrefix} - outputGcsUri: ${args.OutputGcsUri} - config: ${args.Config} # Plan v2 is the default. Explicit v1 remains available during the # deprecated monolithic-plan compatibility window. - planProtocol: ${default(map.get(args, "PlanProtocol"), "v2")} # ── Plan (Activity A) ──────────────────────────────────────────────────── - selectPlanProtocol: switch: - condition: ${planProtocol == "v1"} next: planV1 - condition: ${planProtocol == "v2"} next: planV2 next: unsupportedPlanProtocol - unsupportedPlanProtocol: raise: code: PLAN_PROTOCOL_UNSUPPORTED message: ${"PlanProtocol must be v1 or v2; got " + string(planProtocol)} - planV1: try: call: http.post args: url: ${serviceUrl} timeout: 1800 auth: type: OIDC audience: ${serviceUrl} body: Action: plan PlanProtocol: v1 ProjectGcsUri: ${projectGcsUri} PlanOutputGcsPrefix: ${planOutputGcsPrefix} Config: ${config} result: planRespV1 retry: predicate: ${retryable} max_retries: 6 backoff: initial_delay: 2 max_delay: 60 multiplier: 2 next: capturePlanV1 - capturePlanV1: assign: - planResult: ${planRespV1.body} next: validatePlanResult - planV2: try: call: http.post args: url: ${serviceUrl} timeout: 1800 auth: type: OIDC audience: ${serviceUrl} body: Action: plan PlanProtocol: v2 ProjectGcsUri: ${projectGcsUri} PlanOutputGcsPrefix: ${planOutputGcsPrefix} Config: ${config} result: planRespV2 retry: predicate: ${retryable} max_retries: 6 backoff: initial_delay: 2 max_delay: 60 multiplier: 2 next: capturePlanV2 - capturePlanV2: assign: - planResult: ${planRespV2.body} next: validatePlanResult - validatePlanResult: # Fail closed before fan-out. A v2 render may never fall back to a # v1 PlanGcsUri, and a v1 render may never consume v2 locators. switch: - condition: ${planProtocol == "v1" and (not("PlanProtocol" in planResult) or planResult.PlanProtocol == "v1") and ("PlanGcsUri" in planResult) and not("PlanV2ManifestGcsUri" in planResult) and not("PlanV2ArtifactGcsPrefix" in planResult)} next: captureChunkCount - condition: ${planProtocol == "v2" and ("PlanProtocol" in planResult) and planResult.PlanProtocol == "v2" and ("PlanV2ManifestGcsUri" in planResult) and ("PlanV2ArtifactGcsPrefix" in planResult) and not("PlanGcsUri" in planResult)} next: captureChunkCount next: planProtocolLocatorMismatch - planProtocolLocatorMismatch: raise: code: PLAN_PROTOCOL_LOCATOR_MISMATCH message: "Plan response did not match the selected protocol's disjoint locator contract." - captureChunkCount: assign: - chunkCount: ${planResult.ChunkCount} # ── BuildChunkList + AssertChunkCount ────────────────────────────────────── - assertChunkCount: switch: - condition: ${chunkCount > 0} next: buildChunkList next: planProducedZeroChunks - planProducedZeroChunks: raise: code: PLAN_PRODUCED_ZERO_CHUNKS message: "Plan returned ChunkCount=0 — the composition produced no frames. Non-retryable producer-side invariant violation." - buildChunkList: # Pre-size the ordered chunk-URI + per-chunk result lists so the # parallel branches below assign by index (distinct indices, no # read-modify-write race on a shared accumulator). assign: - chunkIndexes: [] - chunkUris: [] - chunkResults: [] - fillLists: for: value: i range: [0, ${chunkCount - 1}] steps: - appendSlots: assign: - chunkIndexes: ${list.concat(chunkIndexes, i)} - chunkUris: ${list.concat(chunkUris, "")} - chunkResults: ${list.concat(chunkResults, "")} # ── RenderChunks (Activity B, fanned out) ────────────────────────────────── - renderChunks: parallel: shared: [chunkUris, chunkResults] # Run up to chunkCount chunks at once, clamped to 20 — Cloud # Workflows hard-caps concurrent branches/iterations per execution # at 20 (https://cloud.google.com/workflows/quotas). Above that, # iterations queue regardless of concurrency_limit, so a config with # maxParallelChunks > 20 still renders correctly; the extra chunks # just wait. All chunkCount iterations always run. concurrency_limit: ${math.min(chunkCount, 20)} for: value: idx in: ${chunkIndexes} steps: - selectChunkProtocol: switch: - condition: ${planProtocol == "v2"} next: renderOneChunkV2 next: renderOneChunkV1 - renderOneChunkV1: try: call: http.post args: url: ${serviceUrl} timeout: 1800 auth: type: OIDC audience: ${serviceUrl} body: Action: renderChunk PlanProtocol: v1 ChunkIndex: ${idx} PlanGcsUri: ${planResult.PlanGcsUri} PlanHash: ${planResult.PlanHash} ChunkOutputGcsPrefix: ${planOutputGcsPrefix} Format: ${planResult.Format} result: chunkResp retry: predicate: ${retryable} max_retries: 4 backoff: initial_delay: 2 max_delay: 60 multiplier: 2 next: storeChunk - renderOneChunkV2: try: call: http.post args: url: ${serviceUrl} timeout: 1800 auth: type: OIDC audience: ${serviceUrl} body: Action: renderChunk PlanProtocol: v2 ChunkIndex: ${idx} PlanV2ManifestGcsUri: ${planResult.PlanV2ManifestGcsUri} PlanV2ArtifactGcsPrefix: ${planResult.PlanV2ArtifactGcsPrefix} PlanHash: ${planResult.PlanHash} ChunkOutputGcsPrefix: ${planOutputGcsPrefix} Format: ${planResult.Format} result: chunkResp retry: predicate: ${retryable} max_retries: 4 backoff: initial_delay: 2 max_delay: 60 multiplier: 2 next: storeChunk - storeChunk: assign: - chunkUris[idx]: ${chunkResp.body.ChunkGcsUri} - chunkResults[idx]: ${chunkResp.body} # ── Assemble (Activity C) ────────────────────────────────────────────────── - selectAssembleProtocol: switch: - condition: ${planProtocol == "v2"} next: assembleV2 next: assembleV1 - assembleV1: try: call: http.post args: url: ${serviceUrl} timeout: 1800 auth: type: OIDC audience: ${serviceUrl} body: Action: assemble PlanProtocol: v1 PlanGcsUri: ${planResult.PlanGcsUri} ChunkGcsUris: ${chunkUris} AudioGcsUri: ${planResult.AudioGcsUri} OutputGcsUri: ${outputGcsUri} Format: ${planResult.Format} # Forward the caller's exact-CFR request (Config.cfr) to assemble. # `"cfr" in config` guards the optional key; when unset this is # false, which the handler reads as the default -c copy path. Cfr: ${("cfr" in config) and config.cfr} result: assembleResp retry: predicate: ${retryable} max_retries: 4 backoff: initial_delay: 2 max_delay: 60 multiplier: 2 next: done - assembleV2: try: call: http.post args: url: ${serviceUrl} timeout: 1800 auth: type: OIDC audience: ${serviceUrl} body: Action: assemble PlanProtocol: v2 PlanV2ManifestGcsUri: ${planResult.PlanV2ManifestGcsUri} PlanV2ArtifactGcsPrefix: ${planResult.PlanV2ArtifactGcsPrefix} PlanHash: ${planResult.PlanHash} ChunkGcsUris: ${chunkUris} # Audio is an assembler-scoped v2 artifact and is materialized # from the manifest, never carried through a v1 AudioGcsUri. AudioGcsUri: null OutputGcsUri: ${outputGcsUri} Format: ${planResult.Format} Cfr: ${("cfr" in config) and config.cfr} result: assembleResp retry: predicate: ${retryable} max_retries: 4 backoff: initial_delay: 2 max_delay: 60 multiplier: 2 next: done - done: return: Plan: ${planResult} Chunks: ${chunkResults} Assemble: ${assembleResp.body} # Retry predicate: retry transient/server failures (403 from Cloud Run IAM # propagation, 429, and 5xx), never the handler's non-retryable 400s (bad # input, plan-hash mismatch, unsupported format, …). The handler does not emit # 403, so that status is always from Cloud Run's authentication edge. # Connection / timeout errors carry no `.code`; retry those too. retryable: params: [e] steps: - classify: switch: - condition: ${not("code" in e)} return: true - condition: ${e.code == 429} return: true - condition: ${e.code == 403} return: true - condition: ${e.code >= 500 and e.code < 600} return: true - nonRetryable: return: false