RFC 0001 — Worker Assignment for Flow Deployments
Worker Assignment for Flow Deployments
Problem
Today flows execute in-process on whichever backend instance handles the
triggering request (RunAsync → in-proc Engine.Start). Both cmd/api
and cmd/admin-api blank-import flowexec, so either binary can run a
flow. There is no way to say "this flow variant only runs on
instances with this capability" — no GPU-pinning, no region pinning,
no worker-only execution path, and no preparation for the upcoming
bundler work that will produce per-flow worker images.
Context primitives already in place
- go-modules/workers — every process publishes a
worker.heartbeatevent every 10 s carrying{instance_id, hostname, role, version, started_at, healthy, checks}on the shared EventBus. Role comes from theSERVICE_ROLEenv var (defaultsapi, overridden toadmin-apiin compose; arbitrary labels for future workers). - go-modules/workers/admin — admin-api aggregator holds the live
registry keyed by
instance_id, prunes stale entries after 30 s, and exposesGET /admin/workers. - FlowDeployment / FlowDeploymentVariant — named module→flow
binding with weighted + sticky-session variants
(
backend/modules/assistantuses the"assistant"deployment).
Goals
- Operators can pin a deployment variant to one or more roles —
e.g.
["gpu-worker"]or["api", "worker"]. - If no healthy worker reports a pinned role, the variant is skipped in weighted selection and traffic falls through to another variant (rather than error every Nth request).
- Variants without target roles keep executing on whatever already handled them — fully backwards compatible.
- Admin UI lets operators pick target roles from a live menu of currently-announced worker roles.
- Future: split flow execution off onto dedicated worker processes addressed by role, without redesigning the routing table.
Non-goals
- Introducing a separate worker binary in this phase (Phase 3b).
- Changing
RunAsyncto publish events instead of running in-process (Phase 3b). - Replicated/quorum execution, leader election per variant, or any sharding inside a single run (future, not scoped).
Phase 3a — annotation + role-aware resolver (shipped)
Wire format
FlowDeploymentVariant gains one field:
type FlowDeploymentVariant struct {
Label string `bson:"label" json:"label"`
FlowID string `bson:"flowId" json:"flowId"`
VersionID string `bson:"versionId,omitempty" json:"versionId,omitempty"`
Weight int `bson:"weight" json:"weight"`
TargetRoles []string `bson:"targetRoles,omitempty" json:"targetRoles,omitempty"` // NEW
Labels map[string]string `bson:"labels,omitempty" json:"labels,omitempty"`
}
TargetRoles is advisory for Phase 3a: execution still happens
in-process on the calling binary. The annotation is enforced at the
resolver level — by Phase 3b it also drives event routing.
Resolver
store.ResolveDeployment gains an optional RoleChecker parameter:
type RoleChecker func(role string) bool // nil = accept every variant
func ResolveDeployment(ctx, s, deploymentID, meta, prng, roles RoleChecker)
(*ResolvedDeployment, error)
Semantics:
roles == nil→ pre-3a behaviour, every variant considered.roles != nil→ beforePickDeploymentVariant, filter out variants whoseTargetRoleshave no currently-healthy match (ORacross roles — a variant with two pinned roles stays eligible if either has a healthy worker). EmptyTargetRolesalways passes through.- If all variants are filtered,
ResolveDeploymenterrors explicitly so callers can 503 instead of silently running.
The filtered variants are then weight-renormalised by
PickDeploymentVariant — a canary pinned to a dead role simply loses
its weight share, which flows entirely to stable.
Registry accessor
go-modules/workers/admin exposes:
func (m *Module) IsRoleHealthy(role string) bool // matches RoleChecker
func (m *Module) HealthyRoles() map[string]bool // snapshot for UIs
"Healthy" = at least one registered worker reports role AND its last
heartbeat carries healthy == true (the publisher aggregates
HealthRegistry.All() per heartbeat).
Wiring
Consumer modules that route via deployments discover workers-admin in
Configure(registry) and wire its IsRoleHealthy in as RoleChecker.
When workers-admin isn't loaded (e.g. on cmd/api without the admin
submodule), the checker stays nil → pre-3a behaviour.
Applied in this phase to backend/modules/assistant. Any future module
routing via deployments follows the same pattern.
UI
/flows/deployments/[id] variant editor gains a Target roles
picker populated from GET /admin/workers, grouped/badged by health.
A variant whose roles are all dead shows a warning banner; the runtime
behaviour matches (the variant is skipped).
The deployment list shows the pinned roles inline as a chip per variant.
Phase 3b — dedicated workers (deferred)
Sketch for when worker-only processes become useful (expected when the flow bundler lands):
- New
cmd/workerbinary, orSERVICE_ROLE=workeron the generic api binary. It blank-importsflowexecandworkersbut no HTTP surface; it subscribes to aflow.runtopic on the EventBus. Executor.RunAsyncon the admin-api publishes aflow.runEventMessage with headers{target_role, run_id, flow_id, version_id, input}rather than callingengine.Startin-process.- Each worker filters by
headers[target_role]against its ownSERVICE_ROLE— matching workers consume, others ignore. - Workers emit
FlowEvents ontorun.events.<run_id>; the admin-api SSE endpoint subscribes to that topic and forwards to clients. - Admin UI gets no new concepts —
TargetRolesalready drives the new dispatch path.
Trade-offs to decide in 3b:
- SSE aggregation: admin-api-side fan-out vs workers publishing directly and clients subscribing to coordination streams.
- Ack / redelivery semantics for failed worker runs (timeout → re- publish to a different worker? give up?).
- Runtime handlers that register per-node (assistant-handoff-request) need a corresponding registration path on workers.
Migration
- Existing deployments without
TargetRoleskeep working (treated as "any role"). - Existing
cmd/apibinaries that don't load workers-admin passnilasRoleChecker→ variant filtering off → no behaviour change. - Existing admin-api with phase 3a loaded: the
assistantdeployment still routes identically because its lone variant has noTargetRolespinned.
Testing
Unit test in go-flowdsl/flowexec/store/memstore/deployments_test.go: TestResolveDeployment_TargetRolesFilter covers:
- Canary pinned to dead role is filtered across 200 varied sessions; every resolve lands on stable.
- Nil
RoleCheckerpreserves weighted random behaviour (some sessions do hit canary). - When every variant's role is dead, resolution errors with "no variants with healthy target roles".
End-to-end: rebuild admin-api, set a TargetRole the cluster doesn't
report → variant skip observable in the live registry (workers-admin
reports api/admin-api only; pin a variant to worker and confirm
the resolver filters it).