Fundamentals

Operations

The steps that transform, filter, and route data inside a capability.

What are operations?

Operations are the verbs of the DSL. They run in the order you write them, and the exchange passes through each one in sequence.

craft()
  .id('process-order')
  .from(timer({ interval: 60_000 }))
  .transform((body) => normalise(body))
  .filter((ex) => ex.body.amount > 0)
  .enrich(http({ url: '/inventory' }))
  .tap(log())
  .to(destination)

Operation categories

Capability-level (before .from())

Capability-level operations configure the capability itself. They go before .from() and apply to the whole capability, not to one step: .id(), .title(), .description(), .input(), .output(), .tag(), .authorize(), .enabled() and .batch(), plus any wrapper that should cover the whole capability rather than one step.

.from() names the source and opens the pipeline. Everything before it is configuration; everything after it operates on exchanges. The checks and wrappers placed before .from() run in a fixed order whatever order you write them in; the filter chain has that order.

Transform

Transform operations reshape the data as it flows through the pipeline. They receive the current exchange and return a new version of it.

The distinction between them is how much of the exchange they expose. .transform() receives the body only and returns the new body, the right choice for most data reshaping. .process() receives the full exchange, giving access to headers and context. .map() projects fields into a new typed shape. .enrich() pulls data in through an adapter's fetch; the result replaces the body unless you pass an aggregator such as only() to merge it in. .header() sets metadata without touching the body at all. .authenticate() establishes the caller's identity: it mints a trusted Principal from claims you have verified (an e-mail sender, a webhook signature) so a later .authorize() can gate on it.

Flow control

Flow control operations decide which exchanges continue and how they are split or merged.

.filter() drops exchanges that do not match a predicate: the exchange does not continue downstream. Return { reason: "..." } instead of false to record why in telemetry. .validate() checks the body using a Validator adapter or callable function; on failure it throws (hitting the error handler if configured). .schema() is sugar for .validate(schema(...)) and validates against a Standard Schema (Zod, Valibot, ArkType), throwing RC5002 with formatted issue details on failure. .split() fans an array body out into one exchange per item, so each can be processed independently. .aggregate() collects those back into a single exchange. .choice() routes exchanges through different sub-pipelines based on predicates, taking variadic when(...) branches and an optional otherwise(...) fallback. .multicast() fans the exchange out to several paths in parallel, waits for all of them to settle, then continues the original downstream.

.defer() is the one flow-control operation that stops an exchange without ending it. It defers the exchange durably and answers the caller straight away with a Deferred value, because the answer it waits for (a human approval, an out-of-band decision) arrives in hours or days and no transport can be held open that long. .resume() is the other half: it addresses that deferred exchange by its signed token, from any route on any transport, and re-enters the pipeline at the step after the defer with the answer on ex.deferral.result. A route with a reachable defer therefore has two shapes of output, its own and Deferred, and each source renders the second one its own way (202 on HTTP, the value itself on direct(), an ack on a queue). See defer and resume.

Wrappers

Wrappers change how the next operation runs. They do not stand alone: place one immediately before the operation it wraps. Chained after .from(), a wrapper covers that one step. Every wrapper except .delay() can also go before .from(), where it covers the whole capability instead (see the filter chain).

.retry() re-runs the next operation on failure, with configurable backoff (factor growth, a maxBackoff ceiling, and optional jitter). .timeout() throws RC5011 when it takes too long (the abandoned work is not cancelled; the pipeline just stops waiting). .delay() adds a pause before it runs (step scope only). .error() catches any error and lets you provide a fallback body. .cache() skips re-running it when the same caller has sent the same input before (it caches what the operation returned, so a failure it returns instead of throwing is cached too; see cache). .throttle() rate-limits it to a maximum number of calls per time window, pacing exchanges that exceed the rate (or rejecting them with mode: 'reject'). .concurrency() bounds how many exchanges run an operation at the same instant (a bulkhead): where .throttle() caps calls per time window, .concurrency() caps simultaneity, protecting connection pools or memory-bound steps from overload, queueing the overflow or rejecting it with RC5026. .circuitBreaker() counts the operation's failures and, past a threshold, stops calling it for a cooldown, failing fast with RC5025 (or a fallback you supply) instead.

Multiple wrappers can be stacked. They apply in outside-in order, so the first listed is the outermost. This means the order changes the semantics:

// Each retry attempt gets a fresh 5s timeout
.retry({ maxAttempts: 3 })
.timeout(5000)
.process(slowOp)

// Total 30s budget shared across all retry attempts
.timeout(30000)
.retry({ maxAttempts: 3 })
.process(flakyOp)

A wrapper on a step behaves the same way whichever wrapper it is:

  • It covers exactly one step. In .retry().transform(a).transform(b) only a is retried. Put a wrapper before each step you want covered, or wrap a step that nests others, such as a .choice(), to cover the whole block.
  • It needs a step after it. A wrapper left with nothing to wrap, at the end of a capability or just before the next one starts, is refused with RC2001 when the capability is built, rather than attaching to whatever comes next.
  • When it gives up, the error moves outward. A retry that runs out of attempts, a timeout that fires, or a step-scope .error() handler that itself throws passes the error to the next wrapper in the stack, then to a capability-level .error() if there is one, and otherwise to the default error path (route:error, context:error, route:exchange:failed). The capability keeps running, and the next exchange is processed normally.
  • A step-scope .error() recovers in place. Its return value replaces the body and the pipeline continues with the next step, as if the wrapped step had succeeded. Error handling covers both scopes.

Side effects

.to() hands the exchange to a destination adapter (a push-out: the body flows through unchanged, receipts land on headers) or an enricher (a pull-in: the fetched value replaces the body). The pipeline then continues with that body, so steps can follow a .to(). Prefer one .to() per capability for its primary output and use .tap() for the rest; the single-to-per-route lint rule warns on a second one.

.tap() is fire-and-forget. It gets a deep copy of the exchange with the correlation ID preserved and runs in the background while the main pipeline continues immediately. Use .tap() for logging, metrics, and auditing that should never slow down the critical path.


Operations

Every operation, each with its full signature, options and examples.

Error handling

Recover from a failing step or a failing capability with .error().

Filter chain

The fixed order the checks and wrappers before .from() run in.

Previous
The Exchange