How to build workflows for AI video ad generation
Generating a video ad is not one AI call. It is a durable workflow involving scenes, providers, money, callbacks, retries and a final render. Here is how I built one for real usage.
TLDR: Model video generation as a versioned graph of small activities. Keep one canonical run, separate decisions from side effects, persist effects before executing them, freeze the plan and its prices, and expose one server-owned summary to every client. Then build a testbed which can exercise prompts and renders without paying an AI provider every time you move a title by twelve pixels.
I have been building the video generation system behind Hubfluencer, a product which can turn a brief into a finished video ad.
At first, the workflow sounds almost suspiciously simple:
- generate a few clips;
- add a voice-over;
- add music;
- render the video.
We could put that in a background job and go home.
Then real usage arrives.
One video provider returns immediately with a request identifier. Another keeps the HTTP request open. The renderer calls us back later. A user changes a scene while the voice is being generated. One scene fails after three others have completed. The process restarts halfway through. Credits were charged, so a retry must not charge them again. A duplicated webhook says that yesterday’s render is complete while a new render is already running.
The four-step job slowly becomes a collection of workers, status columns, polling loops, recovery jobs and comments explaining why one particular if must happen before another one.
This is the point where a video pipeline stops being a function and becomes a workflow engine, whether we admit it or not.
I eventually admitted it.
This article is about the useful parts of that implementation: how to represent the workflow, where to put provider calls, how to survive retries and callbacks, how to give the interface one truthful status, and how to test rendering without spending money on every iteration.
Start with the product, not the graph
A workflow diagram can become architecture fan fiction very quickly.
Before choosing a DAG library, list what the user is actually trying to obtain. In Hubfluencer’s editor, a typical ad has:
- a brief and a scenario;
- several generated scenes;
- optional narration;
- a generated voice;
- music;
- a composed MP4;
- a verified delivery artifact.
The last item is easy to miss. A render provider returning 200 OK does not mean the user has a downloadable video. The file may not have reached storage yet. It may be empty. Its callback may refer to an older attempt. Completion should mean that the product’s delivery contract is satisfied, not that the last provider call looked optimistic.
The product also supports several ways of getting there. A user can generate one scene manually, generate all missing scenes, run the complete autopilot, regenerate a scene, create several variants, or rerender existing footage after changing the copy.
These are related workflows, but they are not one enormous universal workflow with forty optional branches.
In the implementation, they have separate identities:
editor.autopilot@1
editor.batch@1
editor.scene.generate@1
editor.scene.regenerate@1
editor.scene.variants@1
editor.voice.generate@1
editor.music.generate@1
editor.render@1
short.generation@1
short.rerender@1
A workflow run is one user-visible attempt. It is not the project itself.
The same project can have a completed autopilot run, two later scene-regeneration runs and a final rerender run. Keeping those attempts separate gives each one its own inputs, costs, failures and history. It also prevents a project-level status field from becoming a compressed autobiography of everything that ever happened to the video.
Describe the work as data
The workflow definition is ordinary application code which builds an immutable plan.
Here is a reduced version of a short-video workflow:
Plan.new(scopes: ["full_generation", "music", "render"])
|> Plan.step("segment.generate:1",
activity: "short.segment.generate",
phase: "footage",
retry: media_retry(),
timeout: segment_timeout()
)
|> Plan.step("segment.generate:2",
activity: "short.segment.generate",
phase: "footage",
depends_on: ["segment.generate:1"],
retry: media_retry(),
timeout: segment_timeout()
)
|> Plan.step("music.generate",
activity: "short.music.generate",
phase: "music",
depends_on: ["segment.generate:2"]
)
|> Plan.step("render.prepare",
activity: "short.render.prepare",
phase: "render",
depends_on: [
"segment.generate:1",
"segment.generate:2",
"music.generate"
]
)
|> Plan.step("render.submit",
activity: "short.render.submit",
phase: "render",
depends_on: ["render.prepare"]
)
|> Plan.step("delivery.verify",
activity: "short.delivery.verify",
phase: "delivery",
depends_on: ["render.submit"]
)
Nothing executes while this plan is being built. It is data: step keys, dependencies, phases, retry policies, timeouts, charges and expected artifacts.
That boring property is extremely useful.
We can validate the graph before saving it. We can reject missing dependencies, duplicate steps and cycles. We can calculate progress. We can display the future work. We can test the plan without calling a video model. We can freeze a fingerprint of the complete execution contract.
For editor ads, scenes are a dynamic fan-out. The scenario determines how many scene activities need to exist, then voice and rendering depend on the whole scene group:
plan
|> Plan.fan_out("editor.scenes", scenes, fn scene, index ->
{
"scene.generate:#{index}",
[
activity: "editor.scene.generate",
phase: "scenes",
input: %{"segment_id" => scene.id},
retry: scene_retry()
]
}
end)
|> Plan.step("voice.generate",
activity: "editor.voice.generate",
phase: "voice",
depends_on: [Plan.group("editor.scenes")]
)
The group dependency is the fan-in barrier. Voice cannot begin until every required scene has reached an acceptable terminal state.
I prefer ordinary data builders to an elaborate macro language here. Video products keep inventing new exceptions. A plain builder can compose functions, inspect values and produce a structure which remains easy to print in a failing test.
Freeze the workflow which the user bought
Imagine that a user starts an eight-scene generation on Monday. On Tuesday we deploy a new retry policy, change a model and add a quality-review step. Their run is still waiting for a provider callback.
Which workflow should resume?
The Monday one.
Every run stores its workflow name and version, frozen input, policy and plan fingerprint. On startup, the application rebuilds the plan for every active run and compares the fingerprint:
with {:ok, definition} <- Registry.fetch(run.workflow_name, run.workflow_version),
{:ok, plan} <- build_plan(definition, run.input, run.policy),
:ok <- validate_activities(plan),
^run.plan_fingerprint <- Plan.fingerprint(plan) do
:ok
else
_ -> {:error, :plan_drift}
end
If the plan changed without a version change, the application fails closed instead of quietly applying new semantics to work already in flight.
This may sound strict. It is strict because workflow changes can alter money, retry behaviour and which artifact gets published. Silently accepting drift is convenient right until an old run executes a newly added paid step.
A model alias also belongs in the frozen policy. latest is not a reproducible video model. Store the exact model, duration, audio mode, price basis and any provider options which affect the result.
One reducer decides what happens next
The most important boundary in the system is between deciding and doing.
A command or callback becomes a typed signal:
start_run
activity_succeeded
activity_failed
external_callback_received
retry_run
provide_input
cancel_run
lease_expired
The dispatcher loads the run and its steps under a database lock, gives the signal to a reducer, then receives a decision. The reducer does not call Veo, ElevenLabs, Forge or a push service. It returns data describing:
- state patches;
- append-only events;
- local database operations;
- external effects to execute later.
This gives the workflow one transition authority. A worker cannot decide that the run is complete. A webhook cannot independently refund credits. A stale cleanup job cannot invent a recovery path. They can only report facts back to the reducer as signals.
The pattern looks roughly like this:
Repo.transaction(fn ->
snapshot = Runs.lock_snapshot!(run_id)
decision = Reducer.decide(snapshot, signal)
apply_state_patches(decision)
append_events(decision.events)
apply_local_operations(decision.local_ops)
persist_effects(decision.effects)
end)
This is not full event sourcing. Current state lives in materialized run and step rows, because rebuilding a large workflow from its complete history on every poll would be theatrical. The event log exists for audit and debugging. Tests can replay it when useful.
Keep local operations local
Not every consequence of a transition should become a background job.
If an operation only touches our own database, it can commit inside the same transaction as the workflow state:
- charge or refund credits;
- register an artifact;
- append an in-app notification;
- update compatibility fields;
- mark the run terminal.
External work cannot join that transaction:
- call a video provider;
- synthesize speech;
- upload or render media;
- send a push notification;
- cancel a provider request.
Those operations become durable effects. The effect row is persisted in the decision transaction, then a generic background worker executes it.
Decision transaction
├── update run and steps
├── append events
├── charge credits
├── insert effect: activity.execute
└── enqueue one job containing effect_id
The queue job is deliberately boring:
- claim the effect with a lease;
- resolve the versioned activity;
- execute it;
- dispatch a success or failure signal.
It never chooses the next step.
This is an outbox pattern adapted to a workflow engine. If the process crashes after the transaction, the effect still exists. A reconciler can enqueue it again. If the worker runs twice, the effect and activity idempotency fences make the second execution harmless.
For video generation, this separation also makes timeouts understandable. A step can be running while a local worker owns a leased effect, or waiting_external while a provider owns a request with a deadline. There is no vague third state where an old background job might still be doing something somewhere.
Retries need identities, not hope
Queues are usually at-least-once systems. Webhooks are also at-least-once in practice, even when their documentation is feeling confident.
Every meaningful boundary needs an identity:
- a stable run ID;
- a stable logical step ID;
- an attempt number;
- a unique effect idempotency key;
- a provider submission token;
- a callback deduplication key;
- a monotonic run revision for commands.
A retry increments the attempt of the same logical step. It creates a new effect row with a new key. It does not pretend the previous effect never existed.
run: 7a9...
step: scene.generate:2
attempt: 1
provider token: 7a9.../scene.generate:2/1
If attempt 1 times out and attempt 2 starts, a late callback from attempt 1 can be authenticated, recorded and ignored. It cannot publish over attempt 2 because it carries the wrong attempt fence.
The same principle applies to user commands. A cancel command includes the run revision the client observed. If the workflow has moved to a successor run, the stale cancellation fails rather than cancelling the new work.
Idempotency is not a utility function added near the controller. It is part of the domain model.
Callbacks need an inbox
External callbacks create an annoying race.
A provider may call back before the submit activity has committed the provider’s request identifier on the step. If the webhook handler tries to find the step immediately, it finds nothing. Returning an error asks the provider to retry, perhaps. Accepting and discarding the callback loses the result.
The implementation stores authenticated callbacks in an inbox before interpreting them:
generation_inbox
├── provider
├── external_ref
├── callback type
├── dedupe key
├── redacted payload
├── run_id / step_id, initially nullable
└── status: pending | dispatched | ignored | failed
Both sides try to attach the message. The submit transaction checks for an early callback after saving its external reference. An inbox resolver checks pending messages after the fact. Whichever side sees both halves first dispatches one typed callback signal.
A duplicate callback converges on the same inbox row. A stale callback becomes ignored, not deleted. That small distinction is useful when somebody asks why a provider says it delivered a video which the product did not publish.
Authenticate before inserting. Redact before persisting. Provider payloads have a habit of containing signed URLs, user prompts and implementation details which do not belong in logs forever.
Billing is part of the transition
Video generation costs real money. That changes the design.
The run freezes a quote, an approved cap and a billing policy. Each paid step declares its price and refund semantics. The transition which authorizes paid work also writes the credit ledger entry in the same database transaction.
This avoids the worst possible asynchronous effect:
run says "charged"
charge job never ran
provider generation did run
Refunds use stable ledger keys derived from the run, step and attempt. Repeating the same settlement cannot refund twice. The invariants are simple enough to state and important enough to audit continuously:
charged >= refunded >= 0
net = charged - refunded
net <= approved cap
ledger user = run user
There are two useful payment modes in the product.
A short generation is prepaid. The complete 15-credit run is charged at admission and fully refunded if the delivery contract fails.
An editor workflow is pay as you go. Each scene, voice and music activity carries its own frozen charge. Completed scenes remain useful if a later scene fails, so refunding the entire workflow would be wrong.
Do not bolt billing onto a generic workflow after it works. Paid side effects change retries, cancellation, idempotency and what “completed” means.
Concurrency should describe resources
A user can reasonably regenerate scene four while editing the caption. They should not be able to start two generators which both believe they own scene four.
Each run declares the resources it owns:
scene:<segment_id>
voice
music
render
full_generation
The start transaction checks active runs for intersecting scopes on the same project. The check is serialized with a database advisory lock.
This is more precise than one project_locked boolean. A scene run can lock one scene. A render can lock the render and the media versions it snapshots. Autopilot can hold the union of everything it intends to touch.
Be literal here. A magical full_generation token does not conflict with scene:42 unless the conflict algorithm explicitly says it does. In my implementation, conflict is set intersection, so broad workflows include every concrete scope they need to exclude.
Simple rules are easier to audit than clever lock hierarchies.
Give clients one truthful summary
Before the workflow engine, status was spread across project rows, scenes, audio, renders and background jobs. The API derived one interpretation, the web app another and the agent integration a third.
That is how a product becomes simultaneously “generating”, “failed” and “ready”, depending on which screen you opened.
The server now projects a generation_summary from the canonical run:
{
"run_id": "7a9...",
"workflow": "editor.autopilot",
"workflow_version": 1,
"revision": 12,
"status": "running",
"phase": "scenes",
"label": "Generating scene 2 of 4",
"progress": { "completed": 7, "total": 18, "percent": 39 },
"current_step": {
"key": "scene.generate:2",
"attempt": 1,
"max_attempts": 3
},
"blockers": [],
"available_actions": ["cancel_generation"],
"cost": {
"quoted_credits": 28,
"charged_credits": 10,
"refunded_credits": 0,
"net_credits": 10
},
"timing": { "poll_after_ms": 5000 }
}
The client does not infer whether it may retry. It invokes retry only when available_actions advertises it. It does not invent a polling interval. It polls only while poll_after_ms is a positive integer. It does not assume that a render row means the complete workflow is ready.
This contract is particularly helpful for agents. An agent does not need a two-page prompt explaining seventeen combinations of internal status fields. It reads the same summary as the interface and acts on the actions the server currently allows.
The workflow engine owns lifecycle. The product still owns creative state. A video can have a completed historical render which became stale after a scene edit. That is not a second workflow status; it is a delivery-freshness fact. Keep those concepts separate instead of forcing every product concern into the run state machine.
Human decisions are workflow states
Some video decisions should pause automation.
A variant workflow generates three candidate takes for one scene, then waits for the user to choose. A voice-first project may discover that one narration line cannot fit any supported scene duration. A workflow may need approval for a larger spend cap.
These are not failures. They are needs_input states with typed blockers:
{
"status": "needs_input",
"blockers": [
{
"type": "variant_select",
"step_key": "scene.variants",
"details": {
"candidate_segment_ids": [41, 42, 43]
}
}
],
"available_actions": ["provide_generation_input", "cancel_generation"]
}
The input command carries the run ID, revision and blocking step. A delayed answer to an old selection screen cannot alter a newer run.
This made the engine more useful than a pipeline abstraction. It can represent automation and deliberate human pauses with the same durable semantics.
Rendering needs its own testbed
A workflow engine test suite can prove that callbacks deduplicate and refunds settle. It cannot tell us whether a logo covers the caption or a closing card lasts one frame too long.
AI video also makes visual testing expensive. Generating fresh Veo footage and ElevenLabs audio every time I change typography would be a fairly direct way to turn CSS iteration into a budget line.
The Hubfluencer testbed has two legs.
The local renderer
The fastest path never starts the API.
Checked-in JSON fixtures contain the exact props consumed by the Remotion compositions. Synthetic videos, images and audio live under the renderer’s public directory. A fixture can be opened in Remotion Studio, rendered as one still frame, or exported as a complete local MP4.
bun run testbed:assets
bun run dev
bun run testbed:still fixtures/short-vibe-cinematic.json --frame 30
bun run testbed:render fixtures/editor-logo-closing.json
The still command is the workhorse. A frame renders in seconds, which is fast enough to adjust an overlay, run it again and keep thinking about the design instead of the infrastructure.
Fixtures cover product situations, not random property combinations:
- a neutral short;
- a visually styled short;
- a conversion short with poster and CTA;
- a three-scene editor ad;
- an editor ad with a logo and closing card;
- a landscape editor export;
- voice-first scenes with authored durations;
- several candidate takes for one scene.
Tests validate every fixture against the real composition schema. If the render contract changes, stale fixtures fail in CI instead of becoming archaeological JSON.
The API-attached testbed
Some problems live before the renderer. Prompt construction, pacing, storage keys and payload assembly need the application.
A development-only module seeds real project rows pointing at sample media. It can preview the exact prompts, build the exact wire payload and send it through the local render server. It does not call an AI provider or charge credits.
alias Hubfluencer.Dev.Testbed
short = Testbed.seed_short(
headline: "Espresso anywhere",
short_visual_language: "premium_editorial"
)
Testbed.short_prompts(short)
{:ok, payload} = Testbed.payload(short)
Testbed.validate_payload(payload, "short")
This works because the architecture separates asset generation from payload assembly. AI workers create media and store stable object keys upstream. The render request only needs project rows which point at valid media. Testbed rows can point at synthetic local assets or previously generated samples and exercise the rest of the path unchanged.
The testbed also captures real authenticated render requests with secrets and callback tokens removed. A strange production-shaped payload can be replayed locally, then converted into a permanent synthetic fixture after customer copy and signed URLs are removed.
Captures are temporary and sensitive. They stay git-ignored and local. “Useful for debugging” is not a data-retention policy.
Test contracts at their natural boundary
The complete system has several kinds of tests because one giant end-to-end test would be slow, expensive and vague when it fails.
Plan tests verify topology:
- every dependency exists;
- the graph is acyclic;
- fan-out order is deterministic;
- workflow and activity versions are registered;
- the frozen fingerprint is stable;
- required artifact declarations exist.
Reducer tests verify transitions without providers:
- duplicate signals are no-ops;
- stale attempts cannot publish;
- a retry creates one new effect;
- cancellation is monotonic;
- a user-input step parks and resumes;
- terminal runs cannot return to running unless their policy allows it.
Persistence tests attack the transaction boundaries:
- two starts with one idempotency key create one run and one charge;
- concurrent runs cannot acquire the same scope;
- a callback arriving before submission attaches later;
- duplicate callbacks dispatch once;
- refunds never exceed charges;
- a crash cannot commit state without its effect.
Activity contract tests use fake providers but real schemas. They verify request construction, result normalization, redaction and error classification.
Renderer tests use local media and the real composition. One-frame renders cover layout cheaply. Full MP4 renders cover timing, audio and asset loading.
Finally, a small number of paid smoke tests call the real providers. They prove that credentials, provider APIs and current model contracts still work. They should not be the first place where we discover a cycle in our graph.
What I would build first
You do not need seven tables and an admin dashboard on day one.
For a new video product, I would build this vertical slice:
- Define one workflow, such as
ad.generate@1, with three small steps: generate one scene, submit one render and verify delivery. - Store one canonical run and its steps in the database.
- Make the plan immutable and fingerprint it.
- Process all commands under a run lock through one reducer.
- Persist external effects before a generic worker executes them.
- Give each effect, provider submission and callback an idempotency key.
- Freeze the approved price on the run and charge in the transition transaction.
- Return one
generation_summaryfrom the server. - Add a local renderer fixture using synthetic media.
- Write the duplicate-callback and process-crash tests before adding a second provider.
Then grow from observed pressure.
Add fan-out when the product has several scenes. Add needs_input when there is a real review decision. Add scopes when concurrent edits become useful. Add an inbox when callbacks exist. Add compensation when the product has something concrete to compensate.
Do not begin by implementing every box from a distributed-systems book. Also do not keep a fifteen-minute paid workflow in one background function because the first demo passed.
Mistakes which were expensive enough to remember
Letting workers orchestrate
A worker which calls a provider and then decides the next step is convenient. Soon every worker contains a partial copy of the state machine. Workers should report outcomes. The reducer should decide.
Treating the queue as the source of truth
A queue knows whether a job is scheduled or executing. It does not know whether a provider accepted the request, whether a callback is current, or whether the user cancelled the logical attempt. Persist workflow liveness in your own model.
Deriving status in every client
Raw fields feel flexible. They export your orchestration bugs to every interface. Project one server-owned lifecycle and keep the clients boring.
Retrying without attempt fences
A retry without an attempt identity allows yesterday’s success to overwrite today’s result. Put the attempt on provider requests, callbacks and publications.
Charging in a background effect
Money and the state which authorizes it should commit together. External work is uncertain enough already.
Testing every visual change through AI generation
Rendering and generation are different systems. Give the renderer local fixtures and a one-frame fast path. Your design work should not depend on a remote model queue.
Storing signed URLs as artifact identity
Signed URLs expire. Store stable object keys and create URLs at the transport boundary. A finished workflow should not become unreadable four hours later.
Editing a workflow in place
An active run bought a specific graph, model policy and price. Change the version when those semantics change.
The useful abstraction
An AI video workflow is not a clever chain of prompts.
It is a durable agreement between the user, our application and several unreliable external systems. The agreement says which work was approved, which attempt owns each result, how much may be spent, what can be retried, when a human must decide, and what counts as a delivered video.
Once that agreement is explicit, the AI calls become surprisingly ordinary activities. They can be slow, expensive and occasionally strange, but they no longer own the product state.
Build the workflow around the real ad lifecycle. Keep decisions pure. Persist external intentions. Freeze what the user approved. Let the server tell every client the same story.
And keep a synthetic espresso video around.
You will render it far more often than you expect.
Michaël Mazurczak
Fullstack developer, Lyon, France