@powerhousedao/reactor-workflow
v6.2.3
Published
The Powerhouse workflow engine: the reactor-side runtime that runs workflow documents, and the piece loader, worker pool and executor beneath it.
Readme
@powerhousedao/reactor-workflow
The workflow engine that runs on a reactor. It owns the runtime that turns a
powerhouse/workflow document into runs — triggers, the run journal, managed
secrets, connections — and, beneath it, the machinery that runs an Activepieces
piece in a child process.
The document models, their editors and Workflow Studio live in
@powerhousedao/workflow; the GraphQL subgraph that serves this
runtime lives in apps/switchboard, the host that
composes the runtime.
The seam
src/pieces/ loader, descriptor, worker pool, egress, executor, expressions
src/reactor/ trigger supervisor, coordinator, run journal, secret store, portssrc/pieces runs a piece. It knows nothing about reactors, documents or
Powerhouse packages: give it a piece name, a config and a connection value and
it returns an output. The boundary is enforced by lint — nothing under
src/pieces may import src/reactor, nor any @powerhousedao/* package other
than @powerhousedao/pieces-framework, whose contract it implements — so the
piece layer stays something you can reason about on its own.
src/reactor is everything that only makes sense on a reactor: which workflow a
trigger belongs to, where a run is journaled, whose credentials a step may
resolve. It depends on the piece layer, never the other way round.
What comes from the piece framework
The contract this engine implements is declared once, in
@powerhousedao/pieces-framework, and taken from there
rather than restated here.
- The types.
ApAction,ApTrigger,ApPieceandApPropertyderive fromActionBase,TriggerBase,PieceBaseand the property schemas; the contexts fromStore,ServerContext,FilesService,ConnectionsManager,FlowsContext,RunContext,TriggerHookContextandSetScheduleRequest; the connection shapes fromAppConnectionTypeandAppConnectionValue.PackagePiece,ReactorServiceandDEDUPE_KEY_PROPERTYare the framework's own Powerhouse half. - The enums stay strings here. A piece bundle inlines its own copy of the
framework, so a
PropertyTypeorTriggerStrategyread off one shares no identity with ours. Every such value is compared as a string; nothing insrc/piecesusesinstanceofor enum identity across that boundary. - Prop coercion, from
@powerhousedao/pieces-framework/host, which carries the Activepieces engine's own property processors.context/normalize.tsdispatches to them; only file props stay ours, because attachment andapfile://refs, the size ceiling and a host-injected fetcher have no upstream equivalent. AnApFilea processor builds is flattened to a plain object at that boundary: a class instance does not survive the worker IPC. - Prop validation, from the same place. Before an action's
run()or any trigger hook butonDisable, a prop left unset takes itsdefaultValueand the coerced values go through the engine'svalidateProperty. A failure is aPropsValidationErrornaming each field, as inTitle (title): Expected string, received: undefined; the step fails and piece code never runs. Two departures: a JSON or OBJECT prop whose text does not parse reaches the piece as that text, as coercion already hands it on; and every declared prop is checked, so a required one absent from the input fails, where upstream's processor checks only the keys it was given. A DYNAMIC prop's children are checked in the host, before the worker is asked: the editor writes them as the step'spropertySettings[].schemain the SET_STEP_CONFIG (or SET_TRIGGER) of the edit they belong to, and the run reads them from the published step. A value that is not an object, or lacks a required child, fails the step (or parks the trigger) naming each missing field. Nothing stores whether a step is complete: the editor computes it where it shows it, and gates Publish on it. - The SSRF table, likewise from
./host.worker/egress.tsclassifies an address withssrfIpClassifier.isBlockedIp; the connect-time socket and DNS hooks, the per-request policy and the allow-lists are ours. The one range the classifier reads as unicast and we still refuse is the deprecated IPv4-compatible::/96block, which carries the metadata endpoint. - Error formatting, again from
./host. A thrown piece error passes throughformatPieceErrorbefore redaction, so the HTTP status, request, response and the text of an HTML error page reach the run journal. Redaction runs last, over the formatter's output as well.
How the host composes it
The engine names no host type. WorkflowRuntimeHostDeps (src/reactor/host.ts)
is what the runtime reads — a relational db, a reactor client, the read and
write checks, and optionally the HTTP scope its webhook endpoints live under.
createWorkflowRuntime(deps) returns a configured runtime; nothing here
constructs one by itself.
Switchboard is that host (apps/switchboard/src/workflow-runtime.mts). It
resolves the workflows flag, builds the runtime from what startAPI hands
back, and owns the GraphQL face in apps/switchboard/src/workflow/, registered
live with the GraphQL manager the way a package subgraph is. reactor-api knows
nothing about workflows.
Document operations reach the runtime through a read model, not a
processor: WorkflowTriggersReadModel (WORKFLOW_TRIGGERS_READ_MODEL) is
registered on the reactor's read-model coordinator at the
WORKFLOW_TRIGGERS_READ_MODEL_STAGE (post_ready) stage, and its
indexOperations is runtime.onOperations. The durable cursor BaseReadModel
gives it means an operation written while the runtime is down catches up on the
next boot instead of vanishing; a fresh registration starts at head, so history
is never replayed. onOperations journals a matched fire before it returns, so
the cursor never passes an event that is not yet durable.
A workflow document's DELETE_DOCUMENT disarms it as disabling does: its
deliveries stop once the operation is indexed, a piece trigger's onDisable
runs, and then its trigger row, its FLOW ctx.store partition, its webhook
token and its dedupe keys are deleted. A restart finds nothing to re-arm.
runsPage pages newest first on when a run was journaled (enqueued_at),
which starting a PENDING run leaves alone, so a run keeps its place between
pages. startedAt is when it began executing.
The webhook endpoints are registered under the @powerhousedao/workflow
namespace, which this package exports as WORKFLOW_PACKAGE_NAME.
The worker entry
A piece runs in a forked node child, never in the reactor's process. The
transport finds that child's code by walking up to this package's own
package.json and reading dist/worker-entry.js, so pnpm build must have run
before anything executes a piece — including the suites here.
Blocks
A step names its block the way Activepieces does, with three fields:
pieceName, pieceVersion and actionName. The trigger has pieceName,
pieceVersion and triggerName. pieceVersion is always an exact semver; the
workflow model refuses anything else. A block's identity is the piece and the
name (blockKey in @powerhousedao/pieces-framework/block-type), so the same
action at another version is the same block.
The core piece
The engine's own blocks are the built-in piece @powerhousedao/piece-core, in
src/pieces/core/. Its actions are branch and assert; its triggers are
manual, schedule and webhook. It is written with createPiece,
createAction and createTrigger, and its version is this package's version.
- The runtime registers it itself (
src/pieces/builtin.ts); it is always installed and never comes from a package list. - It is host-bound, like
@powerhousedao/piece-reactor: it always runs the installed copy. - It is this package's own code, so it is described and run in the reactor
process, not in the worker. Its descriptor comes from the same
buildDescriptoras any piece, which also reads its output ports and form hints (showWhen,emptyChoice, the trigger'sdisplay: "schedule"). branchandassertrun throughActivepiecesBlockExecutorlike any piece action, with the same resolution and journaling.branchleaves ontrueorfalseby its result.- The triggers are fed by the host:
reactor/schedule.ts,reactor/webhook.tsand thefiremutation. Their hooks are never called.
Block resolution
One policy (reactor/block-resolver.ts) resolves every block, for steps,
trigger arming, design-time descriptors, output trees, connection checks and
step tests. A block whose version is not an exact semver resolves to missing.
- Candidates. The installed package piece (
local). For a name outside@activepieces/, the configured registry'sGET /pieces/<name>/versions(registry); npm's packument is read only when there is no registry or it answers 404 (npm), so a name the registry owns never comes from npm. For@activepieces/*, the npm packument (activepieces: fetched from their CDN, then npm). Listings are cached for five minutes; a source that does not answer within two seconds contributes nothing. - Exact wins,
localfirst on a tie, thenregistry,activepieces,npm. - Otherwise the closest (
rankClosestVersions): the highest of the same major (the same minor for0.x) at or above the pin iscompatible; anything else isfallback. A candidate without the block's action or trigger is skipped for the next in the same order (Skipped 0.1.0: it has no action "send_request"in the note). At most five candidates are described per resolution. - Host-bound pieces (
@powerhousedao/piece-core,@powerhousedao/piece-reactor) always run the installed copy, with matchinstalled. - Missing is the only failure: no source has the piece, or none of the candidates described has an action or trigger of that name.
A version mismatch never blocks. The resolution is journaled on the step
(piece_version, piece_source, version_match, version_note) and on the
trigger state; a run's warnings counts its fallback steps and its edges on
ports their source never takes. blockResolutions(workflowId) reports the
resolution for every draft block before anything runs. The bundle is
fetched from the chosen source only, and cached under <cache>/<source>/.
Execution order
The coordinator (pieces/engine/coordinator.ts) runs one step at a time. It
passes over the steps array in order until a pass changes nothing:
- Edges are decided by their source. The trigger's edges are decided when
the run starts, on its
nextport. A step's edges are all decided once the step runs or is skipped. An edge is taken when its port is the one the step took and its condition, if any, holds. - A step waits for every inbound edge. Once all are decided, the step runs if any of them was taken, and is skipped otherwise. A join is an OR.
- Siblings run in array order. Two steps that become ready in the same pass
run in the order the
stepsarray lists them, whatever their ports. A step listed before the step that feeds it waits for the next pass. - Entries. A step with no inbound edges runs only when the workflow has no trigger. With a trigger, it never runs.
- Never reached. A step in a cycle, or fed by an edge from a step that does
not exist, is skipped when the run ends. So is every step not yet reached
when a step fails with no
erroredge taken.
Which port a step takes changes which steps run, never the order steps are
reached in. The studio's step outline (stepOutline in the workflow editor)
lists steps in this order.
Expressions
Every string in a step's config is a template, nested strings in objects and
arrays included. A field's propertySettings[].mode is editor-only: it picks
the control (the typed one, or the free expression box) and has no effect at
run time.
- Literal braces.
\{{is a literal{{:"Dear \{{name}}"reaches the piece asDear {{name}}. - Paths.
trigger.payload.x,steps.<key>.output.y,steps.<key>.errorandvariables.<key>. Brackets read keys that are not plain names, and array indexes:steps.fetch.output.headers["content.type"],output.items[0]. - Raw or text. When the trimmed field is exactly one
{{…}}, the field takes the raw value (a number stays a number, an object an object). Any other string is text, and each expression is interpolated:nullas the empty string, objects and arrays as JSON. - Unresolved references fail the step. A path that names nothing fails with
Unresolved reference {{steps.x.output.y}}. The optional form{{steps.x.output.y?}}resolves tonullinstead. A path whose value isnullis not missing. - Fallbacks.
a || b || 'default'takes the first term whose value is notnullor"". A missing path falls through to the next term; the last term fails when missing, unless it is a literal or ends in?. Literals are single- or double-quoted, with\escapes.
Edge conditions are templates too. A condition with an unresolved reference fails the run.
./testing
@powerhousedao/reactor-workflow/testing re-exports the piece layer: the
loader, the worker and its pool, the executor, the coordinator, and the
in-memory secret and connection resolvers. It is what a piece author's own
package uses to run its piece the way this reactor will, without standing up a
reactor to do it.
Configuration
Everything the engine itself reads from the environment is prefixed
PH_WORKFLOWS_, and every one of these is declared under config in
packages/workflow/powerhouse.manifest.json
— that manifest is the published surface a host reads, this table is the
explanation behind it.
| Variable | Default | What it sets |
| ------------------------------------- | ------------------ | ----------------------------------------------------------------------------------------- |
| PH_WORKFLOWS_SECRETS_MASTER_KEY | generated key file | 64 hex chars (32 bytes) encrypting connection secrets at rest (reactor/secret-store.ts) |
| PH_WORKFLOWS_EGRESS_ALLOW_ADDRESSES | unset | Addresses or CIDRs a piece may reach, widening the default policy (reactor/lib.ts) |
| PH_WORKFLOWS_RUN_CONCURRENCY | 4 | Runs executing at once; one forked node child each (worker/pool.ts) |
| PH_WORKFLOWS_RUN_QUEUE_DEPTH | 0 | Runs that may wait for a slot before new ones are refused; 0 waits without limit |
| PH_WORKFLOWS_POLL_INTERVAL_MS | 60000 | Cadence for a polling trigger that names none of its own |
| PH_WORKFLOWS_WEBHOOK_RECONCILE_MS | 900000 | How often a webhook trigger re-registers with its provider |
| PH_WORKFLOWS_WEBHOOK_TIMEOUT_MS | 30000 | How long a sync-mode delivery holds the provider's socket |
| PH_WORKFLOWS_PIECE_MAX_FILE_BYTES | 8388608 | File-size ceiling for FILE-property hydration and ctx.files.write |
| PH_WORKFLOWS_RUN_RETENTION_DAYS | unset (off) | Deletes finished runs older than this many days (reactor/run-retention.ts) |
Each numeric one parses as Number(raw) || default: a value that is not a
positive number falls back silently rather than failing at boot.
Run retention is off by default: the run journal keeps every run. With
PH_WORKFLOWS_RUN_RETENTION_DAYS set, a sweep runs when the journal opens and
hourly after, deleting runs that finished before the window together with their
step executions and run documents, 500 runs per transaction. Unfinished runs
are never pruned. The same sweep drops trigger dedupe keys older than the
longest dedupe TTL (24h); a deleted workflow's keys go when it is deleted.
The secrets key is not optional in production. Unset, loadKey generates
./.ph/secrets.key — relative to the working directory, like the bundle cache
and the attachment staging dir. A host whose working directory does not survive
a restart would come back with a new key, so the store guards against it:
- A host passes
secretsKeyFile: falsewhen its database outlives the working directory, and the store then requiresPH_WORKFLOWS_SECRETS_MASTER_KEY(MasterKeyRequiredError). Switchboard does this when its read model is on Postgres. - The first key a namespace is used with is fingerprinted into
secret_key_check. Any other key is refused (MasterKeyMismatchError) instead of failing later as an undecryptable secret.
Only the host process reads any of these. The worker child is forked with an
empty environment, so the two settings it enforces travel on the wire instead:
the egress policy is compiled per request in the child, and the file ceiling is
stamped onto every request in PieceWorker.execute and installed by the child
before it dispatches.
From the host
Workflows are turned on by switchboard, not here: PH_WORKFLOWS_ENABLED
("1"/"0"/"true"/"false"), which loses to the host's own workflows
option and wins over workflows.enabled in powerhouse.config.json. Two more
switchboard-side settings shape what the runtime can do, and keep their own
names because they are not workflow settings:
PH_REGISTRY_URL— orpackageRegistryUrlinpowerhouse.config.json. See the paragraph below.PUBLIC_URL(thenRENDER_EXTERNAL_URL, thenHEROKU_APP_DEFAULT_DOMAIN_NAME) — the origin a minted webhook URL carries. Unset, endpoints are advertised ashttp://localhost:<port>, which is not something a provider can call. This is notPH_SWITCHBOARD_PUBLIC_URL, which sets the attachment service base URL instead.
Pieces are also read from the registry the host installs packages from —
packageRegistryUrl in powerhouse.config.json, or PH_REGISTRY_URL. It
serves the same list, detail and bundle endpoints cloud.activepieces.com and
its CDN do, and is read ahead of both: its listing merges into the catalog and
into block search, and its tarball is the first download source tried. A host
that installs packages from no registry reads the Activepieces CDN and npm
only. There is no second setting: a package installed from that registry
already ships pieces that run in the worker, so a bundle fetched from it is no
more trusted than one that arrived inside a package.
OAuth2 connections
An OAUTH2 connection brings its own app. startOAuth (reactor/oauth.ts)
builds the provider's authorize URL from the piece's PieceAuth.OAuth2
({prop} placeholders filled from the connection's config, PKCE when the
piece asks for it) and records a pending attempt in the oauth namespace,
good for one exchange within 10 minutes. The host serves the redirect —
Switchboard at <workflow package base>/oauth/callback — and passes what the
provider sent to completeOAuth, which exchanges the code, stores the token
set as a managed secret named token on the connection and runs the
connection check. Resolution refreshes the token 15 minutes before it
expires and rotates that secret in place.
Token requests leave from the reactor process, so they are held to the egress
policy pieces run under: https to a public address, or an address named in
PH_WORKFLOWS_EGRESS_ALLOW_ADDRESSES, over http too.
Known missing features
This is what a piece can declare or call that this engine does not run. It is tracked in #3081, #3090, #3091 and #3095.
Some are rejected rather than run wrongly. A rejected feature is refused
wherever a user meets it: the catalog, pieceActions, pieceTriggers and
block search carry the reason as unsupported, and the editor lists the
block disabled; blockDescriptor throws Piece "<name>": <reason> (or
Trigger "<name>" of "<piece>": <reason>); enabling a trigger parks it in
ERROR with that message and no retry; a step or hook that reaches the
worker anyway fails before piece code runs, except onDisable, which still
releases what an earlier enable registered. The reason reads
<feature> is not supported yet (<issue URL>).
Triggers
TriggerStrategy.APP_WEBHOOKandcontext.app.createListeners: rejected, asTriggerStrategy APP_WEBHOOK(#3081). There is no app-level endpoint or listener table to route a delivery by. A strategy the engine does not know, or none at all, is rejected the same way. One helper decides the strategy (triggerDeliveryin@powerhousedao/pieces-framework/workflow), and a trigger whose descriptor cannot be read isERRORwith a retry, never polled.renewConfiguration/onRenew: aWEBHOOKtrigger'sCRONstrategy runsonRenewon that cron, in UTC. The next renewal time is stored on the trigger row, so it survives a restart. A failedonRenewsetsrenew_error(renewErrorintriggerStates), apart from the poll'slast_error, and retries with backoff capped at the next cron slot. The trigger staysENABLED. Any other strategy butNONE, or a cron that does not parse, is rejected asrenewConfiguration(#3090).TriggerStrategy.MANUALon a piece trigger: rejected, asTriggerStrategy.MANUAL(#3091). The core piece'smanualtrigger is fed by the host'sfiremutation and is unaffected.- Every
WEBHOOKtrigger'srun()is called every 15 minutes without apayload, as a reconciliation sweep. Arunthat only maps the delivery either fails or fires a spurious run (#3090). setSchedule({ cronExpression })is run as a fixed interval; wall-clock time and timezone are lost.onStartis never called. The trigger context'sserveris a throwing stub.- Outside a delivery a hook's
payloadisundefined; upstream passes{}(#3090). - A webhook payload carries no raw body, and its signature headers
(
x-signature,x-hub-signature-256,stripe-signature,authorization) arrive redacted, sorun()cannot verify the sender's signature (#3090).
Auth
- CustomAuth
refresh: rejected, asCustomAuth refresh(#3091). authas an array runs through the method whose type matches the connection's. The piece is refused only when none of its methods can run, with the first method's reason.- OAuth2 runs with the connection's own app only: the authorization-code
grant, with the
client_idin the connection's config and theclient_secretin its secret refs. Theclient_credentialsgrant is rejected, asOAuth2 client credentials(#3091). There are no operator-configured or Powerhouse-hosted apps, and no ActivepiecesCLOUD_OAUTH2/PLATFORM_OAUTH2connections. - A token is refreshed at most once at a time per process. Replicas are not coordinated, so two can refresh one token at once; a provider that rotates refresh tokens may then revoke one of them.
- OIDC: rejected, as
OIDC auth(#3091). Its connections are refused at check and run too. serverinvalidateandgetConnectionIdentifieris a throwing stub.- A CUSTOM_AUTH value's props reach the piece as stored, not coerced: a
Property.Numberprop arrives as the string it was entered as.
Props
refreshOnSearch: a dropdown'ssearchValueis never sent.- An optional prop set to
nullreachesrun()asnull; upstream passesundefined. - A DYNAMIC prop's value is not coerced against the props it resolved to, so an ARRAY inside it arrives as parallel arrays rather than rows.
- Dynamic resolvers nested in ARRAY items or DYNAMIC output can't be called.
- CUSTOM props carry only their type.
- DYNAMIC prop keys are not escaped, and an
options()that throws gets no disabled-dropdown fallback (#3091).
Actions
errorHandlingOptions(retry, continue on failure) is ignored.testis never called.requireAuthdefaults tofalsein the descriptor; upstream defaults totrue.run.stop,run.respond,run.pause, waitpoints andgenerateResumeUrlthrow.connections.get,tags,server,agentandflows.listthrow.ctx.storechecks the key length onputonly;getanddeleteof an over-long key answer as if it were absent.flows.current.version.idis a constant.project.idisreactoron every reactor: the reactor is the project, as it is forctx.store's PROJECT scope.
Piece
deprecatedis not in the descriptor.- Every piece gets the current context shape, whatever its
getContextInfosays. - Reads of context members outside the documented surface are tracked but not reported.
- A piece's setup markdown can describe Activepieces features this engine
does not serve, such as the webhook URL's
/syncand/testforms (#3095).
Running the tests
pnpm build # dist/worker-entry.js, which the piece suites fork
pnpm testBundles fetched from npm are cached in node_modules/.cache/ap-bundles; the
suites that need one skip when it cannot be fetched, so the offline run is
smaller but green. Two suites need the docling piece's own package checked out
beside this repo, and skip otherwise.
The suites set the variables above themselves where they need to — a mock
service on loopback is reached by widening
PH_WORKFLOWS_EGRESS_ALLOW_ADDRESSES, and the secret-store suites set a master
key so no run picks up a developer's key file. The two live piece suites read
DOCLING_E2E_URL / DOCLING_E2E_API_KEY and PAPERLESS_E2E_URL /
PAPERLESS_E2E_USER / PAPERLESS_E2E_PASSWORD, and skip when unset.
test/upstream/ holds Activepieces' own engine tests, generated by
pieces-framework's sync (see its
UPSTREAM.md)
and run against this engine through the adapters in test/upstream-adapters/.
A case this engine is known to fail runs as it.fails and names the issue that
fixes it; never edit those files by hand. To regenerate them:
pnpm --filter @powerhousedao/pieces-framework sync-upstream -- --tag 0.91.0 --from <activepieces checkout>Design documents
docs/plan holds the specifications this engine was built from —
the automation spec (08), the piece-loading architecture (06), the secrets
service (09) and the HTTP routes (10) are the ones worth reading first.
