Keyboard shortcuts

Press ← or → to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

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
httpThe universal baseline. Also serves /sutra/health/*.
kafkardkafka.
rabbitmqlapin (AMQP 0.9.1).
aws-sqsAWS SDK.
gcp-pubsubGoogle Cloud client.
amqpfe2o3-amqp (AMQP 1.0).
daprDapr pub/sub: the sidecar pushes to the engine’s own HTTP listener. No vendor client — see Dapr and Knative.
knativeKnative Eventing: a Trigger or Subscription pushes to the engine’s own HTTP listener. No vendor client.
fileAir-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.

daprknative
Inbound routePOST /dapr/{topic}POST /knative/{subscription}
Channel property that binds ittopicsubscription
Pointing the pusher at the engineA declarative Dapr subscription whose route is /dapr/<topic> — the engine does not answer Dapr’s programmatic GET /dapr/subscribeA Trigger or Subscription whose subscriber is /knative/<subscription> on the engine’s HTTP listener
Outbound destinationdapr://<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-completeNot supported, by design — the channel boots with SUTRA.ACK.ON_COMPLETE_UNSUPPORTED and runs on-persistSupported — the push response is held until the instance ends, bounded by on-complete.hold-timeout (default 30 s)
Inbound guardA dapr-topic header that disagrees with the path’s topic is rejectedA 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:

TierExample bindingWhere the decoder comes fromWhat it takes to have one
Built-in formatcodec: urn:sutra:codec:jsonLinked into every engine binaryNothing — it is always there
Package-supplied codeccodec: urn:transferschemas/transfer/*.xsd inside the archive, compiled when the package deploysAuthor the XSDs — no Rust, no engine build
Extension-crate codeccodec: urn:sutra:codec:<name>A crate implementing PayloadCodec, force-linked by a composition rootOne 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’s DeferredAckRegistry until 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