Skip to main content

Stream Functions

A stream is an append-only log of messages. Each message is published to a subject — a dotted string such as orders.ORD-001.created — and a stream stores every message whose subject matches one of the patterns it was created with. Nine functions are available to a Function step:
These functions write and read. Nothing here consumes a stream as it fills. To start a flow when a message arrives, use a stream trigger.

Subjects and patterns

A subject is dot-separated. Two wildcards are available in the patterns you give create-stream and search-stream-items: get-item-from-stream takes a subject, not a pattern. Use search-stream-items when you want to match across many.
Sanitise anything you interpolate into a subject. A subject token built from user input — an email address, a filename — can carry characters that break routing. Map it to an ID first.

create-stream

string
required
Stream name. Must be unique in the account.
array
required
The subject patterns this stream stores. Plural. A payload using subject creates a stream that stores nothing.
string
Free text shown alongside the stream.
string
"limits", "interest" or "workqueue". Default "limits" — messages stay until a size or age limit removes them. Lower case; a capitalised value is rejected.
string
"file" or "memory". Default "file". Lower case.
integer
Maximum age of a retained message, in nanoseconds. One day is 86400000000000.
integer
Maximum total size of the stream in bytes.
integer
Maximum number of messages retained.
integer
Maximum retained per subject. Set to 1 for a stream that only ever holds the latest state per subject.
integer
Maximum size of a single message in bytes.
integer
Copies kept across the cluster. Default 1, maximum 5.
integer
Window in nanoseconds within which duplicate message IDs are rejected. Server default is two minutes.
Payload
Returns
A name that already exists, or subjects that overlap an existing stream, comes back as {"status_code": 409, "body": {"error": "..."}}.
The editor’s form for this function also offers re_publish and a string placement. The configuration the platform reads names those republish and expects placement to be an object, so neither takes effect as the form presents them.

publish-message-to-stream

Publishes one message to a subject. Every stream whose patterns match that subject stores a copy. A subject no stream covers is an error — see Failure.
string
required
The subject to publish to.
object | string
required
The message. An object is stored as JSON; a string is stored verbatim.
Payload
Returns — for an object value:
and for a string value:
Two things happen to an object value on the way in, and neither happens to a string:
  • Any stream or subject key inside the object is removed.
  • A type key is added with the value "created" if the object does not already have one. That key is what aggregation reads, so leave it alone unless you mean to publish a control message.

get-item-from-stream

Reads the messages stored on one subject.
string
required
Stream to read from.
string
required
Exact subject. No wildcards.
integer
Maximum messages to return. Setting this switches to a sequence-bounded read.
integer
Sequence to start from.
integer
Sequence to stop at.
Payload
Returns
Each result is the published message, with subject, sequence and modified added to it. There is no wrapper object to unpick.
A message that is not a JSON object is skipped, silently. Anything published as a plain string is dropped from the results with no error and no gap in results_total. Read those with search-stream-items, which returns raw values.
Only limit reliably bounds the read. The node decides between a plain read and a sequence-bounded read by looking for camel-cased fromSequence / toSequence, but the form and the schema name those fields from_sequence and to_sequence. Passing either on its own therefore does nothing. Send limit as well and both are honoured.
When limit is set and a further page exists, metadata.next_sequence carries the sequence to resume from; it is 0 when you have reached the end.

search-stream-items

Reads messages matching a subject pattern, with their raw values.
string
required
Stream to search.
A subject pattern, wildcards allowed. This matches on the subject, not on message content — there is no full-text search here.
integer
required
Maximum messages to return.
Payload
Returns
value is always the raw stored text. Parse it in an Eval step if you need the structure.

aggregate-stream-items

Reads every message on a subject and folds them, oldest first, into one object. This is how a stream is used to hold state: publish changes as they happen, aggregate to get the current picture.
string
required
Stream to read from.
string
required
A subject pattern, wildcards allowed. Up to 1,000 messages are folded.
Payload
Returns — the folded object, in place of the usual result list:

The fold

Each message is read in order and applied by its type: A merge is a deep merge, so a later message can change one nested field without restating the rest — but it cannot remove a field. That is what unset is for. Aggregation reads the fold markers back out of the messages, so the type key that publish-message-to-stream adds is part of the contract. A message that is not valid JSON breaks the fold.

unset-message-to-stream

Publishes a control message that removes one or more paths from an aggregate. It does not delete anything from the stream — the removal is applied when the subject is next folded.
string
required
Subject to publish to. Use the same subject the values were published on.
string
Dotted path to remove, for example "customer.phone". An array of paths under paths is also accepted.
Payload
The message published is your payload with type: "unset" added.

tombstone-message-to-stream

Publishes a control message that stops the fold. Everything published to the subject afterwards is ignored by aggregate-stream-items, so this is how you mark an aggregate closed.
string
required
Subject to publish to.
The message published is your payload with type: "tombstoned" added.

poison-pill-message-to-stream

Publishes a control message that discards everything folded so far and starts again from empty. Use it to reset an aggregate without deleting its history.
string
required
Subject to publish to.
The message published is your payload with type: "poison-pilled" added.

list-streams

Lists every stream in the account. Takes no parameters — an empty payload {} is correct. The streams behind key-value and object buckets are filtered out of this listing, so a bucket you created will not appear here. Returns

Failure

Every function on this page catches its own errors and returns {"error": "<message>"}. That is a successful step — the run continues and the next step reads an object with an error key where it expected data. A publish to a subject that no stream covers does not quietly vanish — it retries, then comes back as an error in the body. It is still a successful step, so check the result rather than assuming the message landed. See Building Reliable Flows for what a part-completed run leaves behind.

Stream Triggers

Starting a flow from a message

Key-Value Storage

For a current value rather than a history

Eval Step

Parsing raw values out of a search result

Functions Overview

The full catalogue