Stream Functions
A stream is an append-only log of messages. Each message is published to a subject — a dotted string such asorders.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 givecreate-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.
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.
{"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.
- Any
streamorsubjectkey inside the object is removed. - A
typekey 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.
subject, sequence and modified added to it. There is no wrapper object to unpick.
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.
string
required
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.
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.
The fold
Each message is read in order and applied by itstype:
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.type: "unset" added.
tombstone-message-to-stream
Publishes a control message that stops the fold. Everything published to the subject afterwards is ignored byaggregate-stream-items, so this is how you mark an aggregate closed.
string
required
Subject to publish to.
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.
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.
Related
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