# Subscriptions

[Wiki](../README.md) / [Operations](README.md)

**On this page**

- [Creating a subscription](#creating-a-subscription)
- [Selectors: trigger queues](#selectors-trigger-queues)
- [Payloads: full rows and keys](#payloads-full-rows-and-keys)
- [Consuming](#consuming)
- [A consumer program](#a-consumer-program)
- [Delivery guarantees](#delivery-guarantees)
- [Starting points and replay](#starting-points-and-replay)
- [Managing subscriptions](#managing-subscriptions)
- [Where state lives](#where-state-lives)
- [Retention and expiry](#retention-and-expiry)
- [Errors](#errors)

A subscription is a named, server-tracked position in a database's
[change stream](streams.md). Each of its **consumers** receives every message,
with its own position, acknowledgements and redeliveries, so several services
can follow the same changes independently and resume after restarts without
storing anything themselves.

Subscriptions need `GDB_CHANGE_STREAMS=on` and a source database with
`CHANGE STREAM ON`. Without selectors a subscription delivers every change.
With selectors it is a **trigger queue** that delivers only the label,
label-property and relationship changes it names. From version 1.4.0.

## Creating a subscription

```cypher
ALTER DATABASE galactus SET CHANGE STREAM ON;
CREATE SUBSCRIPTION audit ON DATABASE galactus;
CREATE SUBSCRIPTION search ON DATABASE galactus
  OPTIONS {start: 'snapshot', ackTimeout: '30s', maxDeliveries: 10, maxLag: '10GiB', consumerIdleTimeout: '7d'};
```

| Option | Default | Meaning |
|---|---|---|
| `start` | `'now'` | Where a new consumer starts: `'now'`, `'earliest'`, `'snapshot'` or a change id |
| `payload` | `'changes'` | `'changes'`, `'keys'` or `'full'`: what each message's `row` carries |
| `ackTimeout` | `'30s'` | An unacknowledged message is delivered again after this |
| `maxDeliveries` | `10` | After this many deliveries a message is dead-lettered |
| `maxInFlight` | `1000` | Unacknowledged messages per consumer before `read` returns fewer |
| `maxLag` | unset | Expire a consumer whose unread stream exceeds this size |
| `consumerIdleTimeout` | unset | Drop a consumer that has not read for this long |

These statements run auto-commit and need `ALTER DATABASE` on the source
database.

## Selectors: trigger queues

```cypher
CREATE SUBSCRIPTION fulfilment ON DATABASE orders
  FOR (n:Order) ON CREATE, DELETE, PROPERTY status UPDATE
  FOR ()-[r:PAID]->() ON CREATE, PROPERTY amount
  OPTIONS {payload: 'full'};
```

Several `FOR` clauses combine (any of them matches). Each one names a node
label `(n:Label)` or a relationship type `()-[r:TYPE]->()` and the operations
that fire it:

| Operation | Fires when |
|---|---|
| `CREATE` | A node with the label, or a relationship of the type, is created |
| `DELETE` | It is deleted, including by `DETACH DELETE` |
| `UPDATE` | Any property, or (nodes) any label, of an existing one changes |
| `LABEL ADDED` / `LABEL REMOVED` | That label is added to or removed from an existing node |
| `PROPERTY p` | Property `p` is created, updated or deleted (including when the entity holding it is created or deleted) |
| `PROPERTY p CREATE` / `UPDATE` / `DELETE` | Only that kind of change to `p` |

- **Labels at commit:** labels are evaluated as they were at commit, so
  removing `:Order` still fires `FOR (n:Order) ON LABEL REMOVED`.
- **One message per entity per transaction:** a queue sends one message per
  entity per transaction. `event` is the first matching change, and
  `changes` lists every matching change of that entity in the transaction.
- **Before-values needed:** `PROPERTY p CREATE` / `UPDATE` and the `full`
  payload need the stream's before-values (the default), so they are refused
  on a stream with `beforeValues: false`.
- **Matching cost:** matching happens when a consumer reads, never on the
  commit path.
- **Not yet:** `WHERE` conditions on selectors are not supported yet.

## Payloads: full rows and keys

| `payload` | `row` |
|---|---|
| `changes` (default) | `null`; the event already holds the change and, for updates, the before-value |
| `keys` | `{labels, keys}`: the entity's labels and its unique / node-key constraint values, so a consumer can `MERGE` by key |
| `full` | `{before, after, keys}`: every label and property before and after the transaction (`before` is `null` for a create, `after` for a delete) |

Full rows are captured at commit only for entities whose label or type a
`keys` / `full` subscription names; nothing else pays. In hybrid storage this
reads the entity's unchanged properties while committing: about 3 µs per
captured entity for a 7-property node, measured on a benchmark that updates 100
captured nodes per transaction. Subscriptions without `keys` / `full` add
nothing to commits.

## Consuming

```cypher
CALL gdb.subscription.read('audit', {consumer: 'billing', max: 100, waitMs: 30000})
  YIELD deliveryId, changeId, txId, seq, deliveryCount, metadata, event;
CALL gdb.subscription.ack('audit', $deliveryIds);
CALL gdb.subscription.nack('audit', $deliveryIds, {retryIn: '10s'});
```

A consumer is created by its first `read`. `read` returns up to `max`
messages (default 100): first any that are due for redelivery, then new ones.
With `waitMs` (at most 25,000, below the drivers' default 30-second timeout) it
waits for new changes when none are ready,
holding no database lock while it waits. `event` and `metadata` are exactly
what [`gdb.stream.read`](streams.md#events) returns, so code written for the
stream works unchanged. A queue's messages also carry `changes`, and `row`
follows the payload option. Reading needs the same rights as
`gdb.stream.read`.

## A consumer program

Any Bolt driver can consume a subscription with a short loop. This one uses
the native Python driver; the pattern is the same in every language:

```python
import os
from galactus import Driver

with Driver("bolt://127.0.0.1:7687", "gdb", os.environ["GDB_PASSWORD"],
            database="galactus", timeout=30) as driver:
    while True:
        # Waits up to 25 s; returns as soon as a matching change commits.
        result = driver.execute_query(
            "CALL gdb.subscription.read('fulfilment', {consumer: 'shipper', waitMs: 25000}) "
            "YIELD deliveryId, event, row")
        done = []
        for r in result.records:
            handle(r["event"], r["row"])        # your code: make it safe to repeat
            done.append(r["deliveryId"])
        if done:
            driver.execute_query(
                "CALL gdb.subscription.ack('fulfilment', $ids)", {"ids": done})
```

Keep `waitMs` below your driver's timeout (30 seconds by default). The server
caps it at 25,000 ms.

## Delivery guarantees

- **At least once.** A message counts as handled only when acknowledged.
  Unacknowledged messages come back after `ackTimeout`, or straight away after
  a `nack` with no `retryIn`. Make consumers idempotent: `changeId` is unique
  per change event.
- **Order:** messages are offered in commit order, per consumer.
- **Dead letters:** a message delivered `maxDeliveries` times without an ack
  is set aside, counted, and recorded in `system_streams`.
- Acknowledging an id twice, or an id from before a reset, is a harmless no-op
  that reports `acked: 0`.
- After a server restart, each consumer resumes from its last acknowledged
  position. Anything delivered but not yet acknowledged is delivered again.

## Starting points and replay

| Start | Delivers |
|---|---|
| `'now'` | Changes from now on |
| `'earliest'` | Everything the stream still holds, and with no stream retention limit, everything since the stream was turned on |
| a change id | Everything after that change, if still retained |
| `'snapshot'` | Every existing node and relationship as `node.snapshot` / `rel.snapshot` events, then every change since the snapshot began. This is "from zero", including data that existed before the stream |

The snapshot is read in pages while writes continue, so an entity may appear
in it with a newer state than at the snapshot's start. The changes since the
start are then replayed in order. A consumer that applies messages
idempotently (by `elementId` or its keys) ends up exactly in step with the
database.

Replay at any time with `ALTER SUBSCRIPTION … RESET`.

## Managing subscriptions

```cypher
SHOW SUBSCRIPTIONS;
ALTER SUBSCRIPTION audit PAUSE;
ALTER SUBSCRIPTION audit RESUME;
ALTER SUBSCRIPTION audit RESET TO 'earliest';                  // every consumer
ALTER SUBSCRIPTION audit RESET CONSUMER billing TO 'snapshot'; // one consumer
ALTER SUBSCRIPTION audit DROP CONSUMER billing;
DROP SUBSCRIPTION audit [IF EXISTS];
```

`SHOW SUBSCRIPTIONS` gives one row per consumer: state, position, in-flight
count, delivered, acknowledged, redelivered and dead-lettered totals, and when
it last read. Dropping a database drops its subscriptions.

## Where state lives

- Definitions are kept in `system`.
- Consumer positions, counters and dead letters are kept in the reserved
  **`system_streams`** database, written in 100 ms batches. You can query it
  read-only, for example
  `MATCH (c:Consumer)-[:CONSUMES]->(s:Subscription) RETURN s.name, c.name, c.acked`.
  It cannot be written, dropped or re-engined, and it does not count towards
  the edition's database limit.
- Database names starting `system_` are reserved for the server. If a user
  database is already called `system_streams`, turning on
  `GDB_CHANGE_STREAMS` fails at startup with a message naming it.

## Retention and expiry

Every consumer **pins** its stream: retention (`retention`, `maxBytes`) never
removes changes a consumer has not yet acknowledged. With `maxLag`, a consumer
that falls further behind than that expires instead. It stops pinning, and its
reads fail with `Expired` until it is reset. Without `maxLag`, a consumer that
stops reading holds disk indefinitely. Use `consumerIdleTimeout` or
`DROP CONSUMER` for consumers that have gone away.

When the stream starts a new history (a restore, or the stream turned off and
on), consumers fail with `ResetRequired` rather than continue across the gap.

## Errors

| Code | When |
|---|---|
| `Galactus.ClientError.Subscriptions.NotFound` | No subscription (or consumer) by that name |
| `Galactus.ClientError.Subscriptions.AlreadyExists` | `CREATE` with a taken name |
| `Galactus.ClientError.Subscriptions.Paused` | `read` on a paused subscription |
| `Galactus.ClientError.Subscriptions.Expired` | The consumer fell behind its `maxLag`, or its start is no longer retained |
| `Galactus.ClientError.Subscriptions.ResetRequired` | The source stream started a new history |
| `Galactus.ClientError.Subscriptions.InvalidSelector` | A selector that does not parse, a `WHERE` condition, or a selector or payload that needs before-values the stream does not record |
| `Galactus.ClientError.Streams.Disabled` | The gate is off, or the source database has no stream |

## Related articles

[Change streams](streams.md) · [Server configuration](configuration.md) · [Database administration](databases.md)
