Channels and transports
Processes don’t get called — they’re triggered by messages arriving on channels.
channels.yaml (one per package, see Deployment packages) declares each
channel: which transport it rides, which codec decodes it, how it acknowledges receipt, and who
may send to it.
Transports
One neutral transport SPI, nine implementations, all self-registering behind the same lifecycle
trait — the engine binds, activates, and drains every one of them through a single generic path
with no if transport == "..." branching anywhere:
transport: | Notes |
|---|---|
http | The universal baseline. Also serves /sutra/health/*. |
kafka | rdkafka. |
rabbitmq | lapin (AMQP 0.9.1). |
aws-sqs | AWS SDK. |
gcp-pubsub | Google Cloud client. |
amqp | fe2o3-amqp (AMQP 1.0). |
dapr | Dapr pub/sub: the sidecar pushes to the engine’s own HTTP listener. No vendor client — see Dapr and Knative. |
knative | Knative Eventing: a Trigger or Subscription pushes to the engine’s own HTTP listener. No vendor client. |
file | Air-gapped: file-spool inbound + file:// outbound sink, no network dependency. |
Two further transport: values are engine-internal rather than vendor clients — they have no wire
protocol and no listener of their own. local delivers in-process to another channel; pull parks
the delivery as a task a worker fetches instead of dialing anything, which is the
external-task surface.
flowchart LR
H["http"] --> SPI
B["kafka · rabbitmq · aws-sqs<br/>gcp-pubsub · amqp"] --> SPI
PU["dapr · knative<br/>pushed to the HTTP listener"] --> SPI
F["file<br/>air-gapped spool"] --> SPI
LP["local · pull<br/>engine-internal, no wire protocol"] -.-> SPI
SPI["One transport SPI<br/>bind · activate · drain, generically"] --> CD["The channel's codec<br/>decode + schema-validate"]
CD --> SUB["The processes that subscribed<br/>q:source channel + messageType"]
Which transport a channel rides changes only how the bytes arrive: bind, activate and drain are one generic path, and everything past the doorway — decode, validation, subscription — is identical for all of them.
A hardened or air-gapped build selects a subset of transports at compile time via Cargo features on
the distribution crate (cargo build -p sutra-dist --no-default-features --features file), so the
unlinked vendor clients (rdkafka, the AWS/GCP SDKs, lapin, fe2o3-amqp) are not compiled in at
all. dapr and knative are features too, though neither links a vendor client — leaving them out
only removes their routes and outbound sinks. An operator can additionally restrict which
transports a running binary accepts via SUTRA_ALLOWED_TRANSPORTS — a channel declaring a
disallowed transport fails the deployment with a clear diagnostic, not a silent no-op.
Dapr and Knative
Dapr pub/sub and Knative Eventing each have a transport crate of their own (sutra-transport-dapr,
sutra-transport-knative), but neither links a vendor client or opens a connection. Both are
push transports: the Dapr sidecar, or a Knative Trigger or Subscription, delivers each event as
a CloudEvents HTTP POST to a route on the engine’s own HTTP listener. There is no long-lived
consumer to leader-elect — the pusher’s at-least-once delivery is the guarantee, and inbox dedup
absorbs its redeliveries. CloudEvents extraction follows the channel’s cloudevents-mode, as on
http. Outbound, each rewrites its destination to an HTTP URL and sends it through the same HTTP
sink the http transport uses.
dapr | knative | |
|---|---|---|
| Inbound route | POST /dapr/{topic} | POST /knative/{subscription} |
| Channel property that binds it | topic | subscription |
| Pointing the pusher at the engine | A declarative Dapr subscription whose route is /dapr/<topic> — the engine does not answer Dapr’s programmatic GET /dapr/subscribe | A Trigger or Subscription whose subscriber is /knative/<subscription> on the engine’s HTTP listener |
| Outbound destination | dapr://<pubsub>/<topic>, published through the local sidecar’s /v1.0/publish/<pubsub>/<topic> | knative://<namespace>/<broker>, posted to the Broker ingress |
| Outbound config (process-wide) | SUTRA_SINK_DAPR_SIDECAR_PORT (default 3500) | SUTRA_SINK_KNATIVE_BROKER_INGRESS, else K_SINK, else a built-in in-cluster default |
ack-mode: on-complete | Not supported, by design — the channel boots with SUTRA.ACK.ON_COMPLETE_UNSUPPORTED and runs on-persist | Supported — the push response is held until the instance ends, bounded by on-complete.hold-timeout (default 30 s) |
| Inbound guard | A dapr-topic header that disagrees with the path’s topic is rejected | A request carrying some but not all of ce-id, ce-source, ce-type is rejected (400) |
The sidecar port and the Broker ingress are process-wide because an engine process has one sidecar
and one ingress: a channel’s sidecar.port or broker.url is validated but has no effect.
channels:
- name: orders-in # a Dapr subscription routes topic orders.created here
transport: dapr
codec: urn:sutra:codec:json
topic: orders.created
- name: shipments-in # a Knative Trigger delivers to /knative/shipments
transport: knative
codec: urn:sutra:codec:json
subscription: shipments
ack-mode: on-complete
on-complete.hold-timeout: PT20S # ISO-8601 or bare seconds; keep it below the Trigger's timeout
Why the two differ on on-complete. On both, the HTTP response to the push is the
acknowledgement; there is no separate ack to defer. Knative bounds that response with a timeout the
operator sets on the Trigger or Subscription (delivery.timeout), so the engine can hold it until
the instance completes (202), fails (422, which Knative sends to the deadLetterSink rather than
retrying) or outlives on-complete.hold-timeout (202 with a warning — that one delivery degrades
to on-persist). Keep the hold timeout below the sender’s timeout, or every held delivery turns
into a redelivery. Dapr’s equivalent bound belongs to the pub/sub component — Redis Streams’
processingTimeout, Service Bus’s handlerTimeoutInSec — which the engine cannot see, and holding a
response past it makes the broker deliver the same message again, concurrently. So Dapr runs
on-persist, which answers once the dispatch has run to its first wait state or to completion. The
full per-transport matrix is in Acknowledgement modes.
Binding a channel
# examples/money-transfer/.../channels.yaml (abridged)
channels:
- name: transfer-request
transport: http
bind: "POST /channels/transfer-request"
codec: urn:transfer
cloudevents-mode: none
ack-mode: on-complete
auth:
scheme: apikey
apikey:
value: transfer-demo-key
header: X-Api-Key
A channel does not name a process. Processes subscribe to channels: a start event’s
<q:source channel="transfer-request" messageTypeValue="TransferRequest"/> is what routes a
decoded inbound to it (see The q: namespace). This is what lets one channel feed
several processes and one process listen on several channels — the money-transfer example runs
the same transfer flow off three different intake channels (http, rabbitmq, kafka), each a
separate <q:source> on the same underlying <bpmn:transaction>.
Fan-out. By default a channel is point-to-point: for a given (channel, messageType), exactly
one subscribing process is allowed (enforced at deploy time). Setting broadcast: true fans a
decoded message out to every subscribing process, one instance each — genuine pub/sub.
flowchart LR
H["intake channel<br/>http"] --> P["one process<br/>one q:source per channel"]
R["intake channel<br/>rabbitmq"] --> P
K["intake channel<br/>kafka"] --> P
N["one channel<br/>+ messageType"] --> S1["exactly one subscribing process<br/>enforced at deploy time"]
N -.->|"broadcast: true"| S2["every subscribing process<br/>one instance each"]
The channel never names a process; the process names the channel. That is what lets one flow be fed
by three transports at once, and why a second subscriber to the same (channel, messageType) is a
deploy-time error unless the channel opts into broadcast.
Concurrency admission. A channel may optionally declare maxConcurrentInstances — an
admission cap on simultaneously active instances from that channel (a suspended instance still
holds its slot). Absent, the channel is unbounded.
Codecs — format × schema
A codec is a channel-facing named decoder: a parser (json / xml / yaml / csv / a
domain wire format) plus an optional schema. The schema is the load-bearing part — it’s what
yields a messageType and structural validation. A codec with no schema (a media codec, or a
channel with no codec at all) still decodes, but only a catch-all <q:source> (no
messageTypeValue/messageTypePattern) can subscribe to it, and sutra lint warns.
Codec names share one urn:sutra:codec:<name> namespace, and codec: in channels.yaml looks
identical whichever tier the decoder came from. There are three:
| Tier | Example binding | Where the decoder comes from | What it takes to have one |
|---|---|---|---|
| Built-in format | codec: urn:sutra:codec:json | Linked into every engine binary | Nothing — it is always there |
| Package-supplied codec | codec: urn:transfer | schemas/transfer/*.xsd inside the archive, compiled when the package deploys | Author the XSDs — no Rust, no engine build |
| Extension-crate codec | codec: urn:sutra:codec:<name> | A crate implementing PayloadCodec, force-linked by a composition root | One Cargo dependency and one line in that composition root |
A built-in format is a pure parser — it decodes, but carries no schema, so it lands in the
catch-all-subscription case above. This distribution bundles six of them — json, xml, yaml,
csv, raw-text, raw-bytes — and no domain codec at all. Bind one when the payload genuinely
has no contract to check, or when that contract is enforced somewhere else.
A package-supplied codec is the zero-install path to a typed contract, and the one both
example apps take. schemas/transfer/ holds the XSDs plus a codec-manifest.yaml declaring
schemaKind: xsd and the formats the codec accepts; the engine compiles them at deploy time and
names the codec after the folder path (urn:transfer — see
Deployment packages for the exact folding rule). Nothing about it is
second-class: it yields message types, structural validation, and the shape every FEEL path is
checked against at load time, exactly as a codec written in Rust does. The engine cannot tell that
the schema arrived in an archive rather than a crate.
An extension-crate codec is what a wire format needs when a folder of schemas cannot express
it — a grammar that isn’t XML or JSON at all (a fixed-width block structure, or a delimited segment
stream), or a whole profile: an envelope grammar, a mapping from wire-level message names to
schemas, and versioned editions revved on the standard’s own release cadence. It implements
PayloadCodec from sutra-codec-spi, inventory::submit!s a BuiltinCodec next to the impl, and
claims its name in the same namespace; a distribution that wants it adds the dependency and
force-links it from its own composition root. Every message standard is served this way, by
proprietary extension crates built outside this repository. An engine binary that links one
resolves it exactly like a built-in; see
Domain neutrality and the SPI model for why none of them
lives here.
What an extension codec can express — an enveloped profile, generically
The clearest illustration of what the codec SPI has to be able to carry is a market venue’s own profile of a standard rather than the bare standard itself. A generic-format codec is the easy case: one instance document validated against a schema, with the message type read straight off the document’s own root element or namespace. A profile codec decodes a family of messages for a venue that doesn’t fit that shape — every message travels inside a venue-specific envelope, wrapping a header plus a body, and the venue pins its own schema versions on its own release cadence. Worse for identification purposes, the same underlying message can sometimes back several different wrappers, so a bare schema namespace can’t tell you which one you’re looking at.
Such a codec’s message type is therefore the wrapper element’s local name, not a schema
namespace slug — which lets declared_message_types() return the venue’s own closed set of
wrapper names, rather than the open type set a bare-standard codec declares, so a message-type
applicability check (a rules-manifest.yaml entry, a q:source messageTypeValue pin) actually
fires for a profile-bound module rather than silently no-op.
Decoding validates the envelope grammar, the header, and the body against the venue-pinned base
schemas the codec crate carries, out of the box. The payload view projects the body exactly as a
bare-standard document would be — a compatibility guarantee: every alias, DMN input, or fixture
written against the underlying message shape reads identically whether the underlying codec is the
profile variant or the plain standard. Outbound can be a template-rendered passthrough — encode()
returning the rendered envelope bytes verbatim rather than assembling one — so template drift
against the venue’s schemas has to be caught by validating rendered output in tests, not by a
runtime encode path. A second venue following the identical envelope/wrapper/edition pattern
reuses the same machinery rather than re-implementing it.
None of that requires a change to the engine, the channel layer, or this repository — which is the
point of the codec SPI. A deployment can go one step further still and override a codec’s own
schemas: see
Deployment packages for
the schemaKind bundle mechanism that lets an archive map its own schema editions per wrapper.
Acknowledgement modes
ack-mode decides when the engine acknowledges an inbound relative to processing:
on-persist— ack as soon as the message is durably captured, before the process runs. Broker default. On HTTP this is what makes a channel asynchronous (202 Accepted, no business body).on-complete— ack only once the instance reaches a terminal state. HTTP default (classic synchronous request/reply — hold the connection, return the reply body). On a broker, this defers the ack via the engine’sDeferredAckRegistryuntil the instance completes or fails.
The full per-transport wiring matrix, the bounded-registry knobs, and when to pick which mode live in Acknowledgement modes — that’s the operating-chapter deep dive this summary points at.
Secrets on a channel
A channel’s credentials (broker username/password, an API key) are never literal values in
channels.yaml — they’re scheme references (env:NAME, secret:…, vault:…,
aws-secrets:…), resolved at channel startup through one vendor-neutral resolver SPI. Package-time
validation rejects a literal secret outright.
Next
- The q: namespace — how a BPMN process subscribes to a channel.
- External tasks: the pull worker surface — the
pulltransport in full. - Acknowledgement modes — the operator-facing deep dive.