Repository navigation
Support per-node-type activity worker routing (dedicated Docker images/tools per node type) #141
Description
Activity
Open Question 1 — Worker secrets/config
Verdict: Yes — each specialized worker gets its own
.env/.env.exampleand its ownenv.ts, following the pattern already established byapps/backendandapps/execution-workertoday.Justification:
This repo's existing convention is one deployable unit = one env contract, not a shared/global config:
apps/backend/.env.exampleandapps/execution-worker/.env.exampleare separate files with overlapping-but-distinct keys (OPENROUTER_API_KEY,AI_MODELappear in both, independently), and the comment inapps/backend/.env.example:16is explicit: "The execution worker keeps its own key for running workflows" — this is the precedent question already answered once, for a different key.apps/execution-worker/src/env.tsimplementsrequireEnv/envOras a local, fail-fast module scoped to that package — not imported fromapps/backendor any shared@workflow-builder/*package. There's no shared env module in the monorepo today; each app owns its full contract.deploy/ai-studio/docker-compose.ymlreinforces this at the ops layer: each service gets its ownenvironment:block wired at the compose level, with unit-specific required-vs-optional treatment (e.g.OPENROUTER_API_KEY: ${OPENROUTER_API_KEY:?set ...}forbackend).
A specialized worker (e.g.
apps/execution-worker-coding-agent) will need its own credentials (its own model/API key for its dedicated image, possibly a different provider or scoped token) that the main worker has no reason to hold or fail-fast-check for. Making it depend on or extend the main worker'senv.tswould either force irrelevant vars into the main worker's contract or create an implicit coupling between two independently-deployed images — contradicting the "each Docker image is a separate deployable unit" premise of this issue itself.Recommendation: Repeat the pattern exactly — new
apps/<name>/directory, ownpackage.json, own.env.example, ownenv.tswithrequireEnv/envOr. Do not add conditional/optional secrets to the sharedexecution-workerenv module.taskQueuerouting determines which worker picks up an activity; it says nothing about how that worker is configured, and this repo's existing convention already treats those as orthogonal.Open Question 2 — Heartbeats vs. timeout
Verdict: Yes — heartbeats, and they should be a first-class
ActivityProfilefield, not bolted on separately.Why timeout alone is insufficient here:
DEFAULT_NODE_ACTIVITY_PROFILEalready picks a generousstartToCloseTimeoutwith a low retry cap specifically to bound cost on partial failures (activity-profiles.ts:25-30). That reasoning holds for an LLM call that either finishes or errors quickly. A CLI-agent subprocess is a different failure mode: it can hang — wedged on stdin, a stuck child process, a network partition mid-stream — without erroring and without finishing. With onlystartToCloseTimeout, Temporal has no signal that the activity is dead until the full timeout elapses, so a 30-minute profile means up to 30 minutes wasted per attempt before a single retry fires, multiplied byretry.maximumAttempts. AheartbeatTimeoutlets the activity self-report liveness on a much shorter cadence (e.g. every 30s) so a stalled worker slot is reclaimed and retried in seconds, not the entire timeout window — whilestartToCloseTimeoutstill bounds the case where it's alive but genuinely slow.Fits the existing pattern:
node-activity-options.ts:37-53andprofile-validation.tsalready treat a profile as a closed, whole-object shape (PROFILE_KEYS,RETRY_KEYS) explicitly to avoid silently inheriting Temporal defaults (README:96).heartbeatTimeoutshould follow the same rule — optional, but validated the same way asstartToCloseTimeout(sameDurationString+ protobuf-Durationbounds check) — rather than becoming an implicit default that only long-running node types remember to set.Concrete change:
// activity-profiles.ts export type ActivityProfile = { startToCloseTimeout: DurationString; heartbeatTimeout?: DurationString; // required for activities that call Context.current().heartbeat() retry: { maximumAttempts: number }; };
Add
heartbeatTimeouttoPROFILE_KEYSinprofile-validation.ts, pass it through inresolveFromValidatedProfiles, and document in the README that any node executor doing long-running subprocess/stream work must callheartbeat()on an interval shorter than its profile'sheartbeatTimeout— otherwise Temporal cancels it as stalled.Open Question 3 — Event shape for CLI output
Verdict: reuse the existing event pipeline as-is (option a) — no new event type needed for the MVP.
Why:
execution_events.payload_jsonis already an untypedjsonbcolumn (apps/backend/src/db/schema.ts:55) andemitEvent(executionId, type, payload?: unknown, nodeId?)(apps/execution-worker/src/database.ts:15) accepts any JSON-serializable payload. A captured stdout/stderr blob fits without a migration.- The graph runner's event model is strictly binary per node today:
node_started→ one ofnode_completed/node_failed/node_skipped(packages/execution-core/src/graph-runner.ts:312-321). There is no existing notion of a partial/intermediate event for a single node — adding one would be a new concept threaded through the runner, the SSEformatEvent(apps/backend/src/routes/executions.ts:237), and the frontend execution-log renderer, not just the worker. - Precedent:
ai-agent.tsalready returns the entire result as one{ output: { response: ... } }payload on completion (apps/execution-worker/src/activities/ai-agent.ts:59), and errors are logged in one shot vialogger?.error('llm call failed', {...})(:63) rather than streamed. A CLI-agent activity should follow the same convention: buffer stdout/stderr for the run and attach it to the singlenode_completed/node_failedpayload.
Recommendation:
Ship the CLI-agent node with buffered output onnode_completed/node_failed(zero schema change, zero new SSE/UI work) as the MVP for this issue's routing work. Track true incremental streaming (anode_progress/partial event type, new SSE handling, log-panel append-as-you-go UI) as an explicit, separate follow-up once per-node-type worker routing has landed — same incremental-scope pattern the replay-test backlog already uses for "scenarios still worth adding." Don't couple the routing change to a streaming-UX change; they're independently valuable and independently risky.Open Question 4 — Replica configurability
Verdict: No — don't add a static replica-count config field; use
docker compose up --scale.Why:
deploy/ai-studio/docker-compose.ymlhas zerodeploy:/replicasblocks today —backendandworkerare both defined as strictly single-instance services (restart: unless-stoppedonly, no scaling primitive at all). There's no existing pattern to extend, only one to introduce.- The
workerservice (docker-compose.yml:99-115) already differs frombackendpurely bycommand:, sharing oneai-studio-runtimeimage built from a singleruntimetarget (Dockerfile:19-30, comment: "backend + worker, command chosen per compose service"). That's this repo's established precedent for "N behaviors, one image" — it argues for adding new specialized worker services (newcommand:/task-queue per node-type group), not a new config schema for how many of each to run. - Docker Compose already ships the mechanism for this:
docker compose up --scale worker-coding-agent=3 --scale worker-default=1. That's free, needs zero code or schema changes, and works today with plaindocker compose(no Swarm/deploy:block required for scaling itself, just for resource limits if ever needed). - No Kubernetes/Nomad/Helm manifests exist anywhere in the repo — introducing a bespoke replica-count config or orchestration layer would go against this project's stated KISS/YAGNI stance ("don't overengineer," "justify every line of code you write").
Recommendation: When this issue lands, give each specialized worker its own compose service (own
command:/task-queue env, same sharedai-studio-runtimeimage), and documentdocker compose up --scale <service>=Nper service as the scaling story. Revisit only if/when the project actually adopts a real orchestrator beyond docker-compose — it doesn't today.Replied across all three in #148. Short version: a direction we want and it's on our roadmap. Keeping this open as the design record.
Reacted by Tom Brandenburg
Problem / Motivation
Today, every node type (
trigger,decision,ai-agent,visualize, and any future type) executes as an Activity inside the sameexecution-workerprocess, built from the same Docker image. This means:ffmpeg, etc.) would bloat the shared image for every node type, or conflict with existing dependencies.Goal
Allow specific node types to be executed by a dedicated worker process, running from its own Docker image with its own tools installed, while all other node types keep running on the existing general-purpose worker — with no changes to the workflow/interpreter code.
Non-goals
runGraph) or the deterministic workflow contract at all — this is purely a worker-topology/deployment change.Proposed Design: Temporal Task-Queue Routing
Temporal already supports this natively: any
proxyActivities(...)call can target a specific task queue name. Workers subscribe to (poll) task queues; Temporal server only ever dispatches a task to a worker polling the matching queue. So:taskQueuefield on the per-node-type activity profile, so specific node types can be pinned to a non-default queue.No change needed to:
runGraph, the graph JSON model, the UI, the DB schema, or the backend'ssubmit()call — the workflow code doesn't know or care which container ran an activity.Required Changes
1.
packages/temporal— extend the activity-profile typepackages/temporal/src/workflow/activity-profiles.ts:Keep
DEFAULT_NODE_ACTIVITY_PROFILEunchanged (notaskQueue= current behavior, zero risk to existing node types).2.
packages/temporal— wire it intoproxyActivitiespackages/temporal/src/workflow/run-workflow.ts, insiderunner.executeNode, pass the resolved profile'staskQueuestraight through toproxyActivities(...)(structurally compatible with Temporal's realActivityOptionsalready).resolveFromValidatedProfiles/profile-validation.tsneed to accept and pass through the optionaltaskQueue.3.
apps/execution-worker— a second, minimal entrypointA new activity-only worker entrypoint (no
workflowsPathneeded — Temporal supports activity-only workers) that registers only the specialized executor(s) and polls the new task queue name.4. New Docker image for the specialized worker
docker-compose.ymlalongside the existingworkerservice, connecting to the same Temporal server/namespace on a different task queue..envif needed) to run it standalone.5. Validation
profile-validation.ts(already validated atWorker.createtime and inside the sandbox) should also guard against a node type routed to ataskQueuewith no worker polling it — otherwise the activity silently sits pending forever, which is an easy and hard-to-diagnose misconfiguration.6. Testing
taskQueueactually produces aScheduleActivityTaskcommand targeting that queue (inspect the recorded Event History, akin to the existing replay tests).taskQueueis pure routing metadata and doesn't leak into the sandbox's decision-making.Rollout / Compatibility
Fully additive and backward compatible: existing node types with no
taskQueuein their profile behave exactly as today. Can be introduced node-type by node-type.Open Questions
.env/credentials, separate from the main worker?startToCloseTimeout?emitEvent/execution-log pipeline as-is, or does it need a richer event shape (streaming partial output)?Summary
Add an optional
taskQueueto the existing per-node-typeActivityProfile, thread it through toproxyActivities, and support standing up one additional worker process/Docker image per specialized node type that polls that queue — no changes to the graph interpreter, the workflow contract, or any other node type's behavior.