Rfcs

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.heartbeat event every 10 s carrying {instance_id, hostname, role, version, started_at, healthy, checks} on the shared EventBus. Role comes from the SERVICE_ROLE env var (defaults api, overridden to admin-api in 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 exposes GET /admin/workers.
  • FlowDeployment / FlowDeploymentVariant — named module→flow binding with weighted + sticky-session variants (backend/modules/assistant uses the "assistant" deployment).

Goals

  1. Operators can pin a deployment variant to one or more roles — e.g. ["gpu-worker"] or ["api", "worker"].
  2. 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).
  3. Variants without target roles keep executing on whatever already handled them — fully backwards compatible.
  4. Admin UI lets operators pick target roles from a live menu of currently-announced worker roles.
  5. 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 RunAsync to 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:

go
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:

go
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 → before PickDeploymentVariant, filter out variants whose TargetRoles have no currently-healthy match (OR across roles — a variant with two pinned roles stays eligible if either has a healthy worker). Empty TargetRoles always passes through.
  • If all variants are filtered, ResolveDeployment errors 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:

go
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):

  1. New cmd/worker binary, or SERVICE_ROLE=worker on the generic api binary. It blank-imports flowexec and workers but no HTTP surface; it subscribes to a flow.run topic on the EventBus.
  2. Executor.RunAsync on the admin-api publishes a flow.run EventMessage with headers {target_role, run_id, flow_id, version_id, input} rather than calling engine.Start in-process.
  3. Each worker filters by headers[target_role] against its own SERVICE_ROLE — matching workers consume, others ignore.
  4. Workers emit FlowEvents onto run.events.<run_id>; the admin-api SSE endpoint subscribes to that topic and forwards to clients.
  5. Admin UI gets no new concepts — TargetRoles already 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 TargetRoles keep working (treated as "any role").
  • Existing cmd/api binaries that don't load workers-admin pass nil as RoleChecker → variant filtering off → no behaviour change.
  • Existing admin-api with phase 3a loaded: the assistant deployment still routes identically because its lone variant has no TargetRoles pinned.

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 RoleChecker preserves 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).