Statestore Eventing

Publish and subscribe with durable topics on the built-in statestore — message queue triggers that need no external broker.

Statestore eventing gives you durable publish/subscribe topics on the built-in statestore — no Kafka, no external broker, no extra infrastructure.

Statestore eventing is available starting with Fission v1.27.0. A topic is a durable, replayable stream in the statestore event log. Any producer appends events to it: the fission topic publish command, or an async invocation result destination. A message queue trigger with --mqtype statestore subscribes a function to the topic. Fission delivers each event to the function at least once, with retries and an error topic for events that keep failing.

flowchart TB
  pub["Publisher<br/>(fission topic publish, async destination)"]:::user
  pub -->|"<b>1.</b> append"| stream["Topic Stream"]:::store
  stream -->|"<b>2.</b> read from cursor"| mqt["Statestore MQ Consumer"]:::fission
  mqt -->|"<b>3.</b> invoke via router"| pod["Function Pod"]:::pod
  mqt -.->|"retries exhausted"| err["Error Topic"]:::store
  mqt -.->|"response body"| resp["Response Topic"]:::store

  classDef user fill:#ffffff,stroke:#94a3b8,color:#1f2a43
  classDef fission fill:#e8f0fe,stroke:#2d70de,color:#1f2a43
  classDef pod fill:#e6f7f1,stroke:#11a37f,color:#1f2a43,stroke-dasharray:5 3
  classDef store fill:#fff7e0,stroke:#dba514,color:#1f2a43,stroke-dasharray:5 3

When to use it

You needUse
Function-to-function events inside the cluster, zero extra infrastructureStatestore eventing (this page)
Events from an external broker you already run (Kafka, SQS, RabbitMQ, …)KEDA message queue triggers
High throughput, partitioned ordering, or consumer groupsAn external broker

The statestore provider targets small and medium event volumes. When you outgrow it, change --mqtype to kafka and point at a broker — the trigger fields stay the same.

Prerequisites

Eventing needs the statestore; the embedded mode is enough:

helm upgrade --install fission fission-charts/fission-all \
  --namespace fission \
  --set statestore.enabled=true

The chart value eventing.enabled defaults to true, so a statestore-enabled install already runs the eventing consumer. Without the statestore, fission topic commands fail with eventing is not enabled on this cluster (requires the statestore).

When internal service authentication is enabled, set FISSION_INTERNAL_AUTH_SECRET so fission topic publish and fission topic peek can sign their requests.

Worked example

Wire a function to a topic named orders, publish an event, and watch it flow.

Create the consumer function:

// process-order.js
module.exports = async function (context) {
    console.log("processing order:", JSON.stringify(context.request.body));
    return { status: 200, body: "ok" };
}
fission env create --name node --image ghcr.io/fission/node-env
fission fn create --name process-order --env node --code process-order.js

Create the trigger before you publish. A new trigger starts at the head of the stream: it sees only events published after it starts.

fission mqtrigger create --name order-consumer \
  --mqtkind fission --mqtype statestore \
  --topic orders --function process-order \
  --pollinginterval 1 \
  --maxretries 3 --errortopic orders-errors
FlagMeaning
--mqtkind fissionThe statestore provider runs in the classic head; the default keda kind rejects it.
--topicTopic to consume; 1–249 characters of [a-zA-Z0-9._-].
--pollingintervalSeconds between reads on an idle topic; the CLI default of 30 adds up to 30 s of delivery latency.
--maxretriesDelivery retries per event before the event goes to the error topic.
--errortopicTopic that receives events that exhaust their retries.
--resptopicTopic that receives the function’s response body (optional).

Publish an event:

$ fission topic publish --topic orders --data '{"orderId":"A-1042"}'
published to topic "orders" in namespace "default" (statestore)

Peek at the topic to confirm the event is durably stored:

$ fission topic peek --topic orders
head: 1
SEQ TYPE             AGE PAYLOAD
1   application/json 10s {"orderId":"A-1042"}

Verify delivery in the function’s log:

$ fission fn log --name process-order
... processing order: {"orderId":"A-1042"}

fission topic publish also accepts --mqtype kafka to publish to a broker topic through the egress queue; this page covers the default statestore type.

Publish from a function

Async invocations can fan their results out to a topic instead of a single destination function:

fission fn update --name resize-image \
  --async-on-success-topic image-resized \
  --async-on-failure-topic image-failures

Every statestore trigger on image-resized then receives the result envelope of each successful delivery. See Asynchronous Invocation for the async delivery pipeline itself.

Delivery semantics

BehaviorAs implemented
GuaranteeAt-least-once per trigger; make consumers idempotent.
Start positionStream head at first subscribe; no backlog replay.
OrderingSingle stream, roughly FIFO; no partitions.
Fan-outEach trigger on a topic keeps its own durable cursor and receives every event.
SuccessAny 2xx response from the function.
RetryUp to --maxretries retries per event, 500 ms apart.
ExhaustedEvent is published to --errortopic; without one it is dropped with a log line.
Response topicThe 2xx response body is published to --resptopic, best effort.

The consumer advances its cursor only after an event reaches terminal handling: delivered, or routed to the error topic. One failing event therefore cannot wedge the topic, and a crash mid-batch redelivers the tail rather than skipping it.

Retention

A reaper in the statestore MQ consumer trims each subscribed topic once per minute:

  • Events that every trigger on the topic has consumed are trimmed — no live subscriber loses an unconsumed event.
  • Two backstops trim past a stalled subscriber: events older than 7 days, and any backlog beyond 100,000 events per topic. A subscriber that resumes after a backstop trim logs the gap and counts it in the fission_eventing_gap_events_total metric.

A topic with no statestore trigger is not trimmed at all — the orphan-stream age sweep is not implemented yet. Instead, a per-topic backlog cap of 10,000 events bounds the growth. Publishes to a capped topic fail with topic backlog cap reached instead of dropping silently. To recover a capped orphan topic, create a statestore trigger on it. The trigger’s cursor starts at the head, so the reaper trims the old backlog within about a minute and publish flow resumes. The trimmed backlog is not delivered.

Limits

  • One durable cursor per trigger; there are no consumer groups, so you cannot parallelize one trigger across replicas.
  • Throughput is bounded by the statestore, not by Fission.
  • Topics are namespace-scoped; a trigger can only consume topics in its own namespace.
  • Exactly-once delivery is out of scope; design consumers to tolerate duplicates.