> ## Documentation Index
> Fetch the complete documentation index at: https://filament.getgalaxy.io/llms.txt
> Use this file to discover all available pages before exploring further.

# NATS JetStream

> Publish acknowledged JSON messages in bounded runs or continuous epochs

`connectors/nats` registers the `nats` sink (**alpha**). It publishes one JSON
message per input row to an existing JetStream stream. Both bounded runs and
continuous execution use the same bounded asynchronous publisher. Each record
remains an individual NATS message.

## Configuration

| Field                 | Scope      | Required | Description                                                                                                           |
| --------------------- | ---------- | -------- | --------------------------------------------------------------------------------------------------------------------- |
| `url`                 | Connection | Yes      | NATS server URL; use `tls://` for TLS                                                                                 |
| `auth_mode`           | Connection | No       | `none` (default), `user_password`, `jwt`, or `token`; shows only the selected method’s fields                         |
| `username`            | Connection | No       | Username; supply with `password`                                                                                      |
| `password`            | Connection | No       | Secret password; supply with `username`                                                                               |
| `jwt`                 | Connection | No       | Secret NATS user JWT; supply with `nkey_seed`                                                                         |
| `nkey_seed`           | Connection | No       | Secret private user NKey seed matching the JWT                                                                        |
| `token`               | Connection | No       | Secret authentication token                                                                                           |
| `stream`              | Pipeline   | Yes      | Destination stream name: letters, digits, `-` and `_`, no dots. `{resource}` expands to the destination resource name |
| `subject`             | Pipeline   | Yes      | Publish subject. `{resource}` expands to the destination resource name                                                |
| `create_stream`       | Pipeline   | No       | Create a missing stream; default `true`                                                                               |
| `publish_timeout`     | Pipeline   | No       | Positive duration from each publish submission to its acknowledgment; default `5s`                                    |
| `max_in_flight`       | Pipeline   | No       | Maximum outstanding messages, from 1 through 65,536; default `1024`                                                   |
| `max_in_flight_bytes` | Pipeline   | No       | Maximum outstanding message bytes; positive integer, default `16777216` (16 MiB)                                      |

Source and sink connections share these authentication settings. Choose only one
method: username/password, JWT/NKey seed, or token. Select `none` for an unauthenticated server. Selected authentication fields are required.
An omitted `auth_mode` means no authentication; credentials require an explicit mode. JWT authentication
uses the seed to sign the server challenge; a JWT alone is insufficient.
Username/password and JWT/seed require no files mounted on workers. Passwords,
JWTs, seeds, and tokens use Filament’s secret storage.

Both `stream` and `subject` are templates. The `{resource}` placeholder, or
its Go template spelling `{{resource}}`, is replaced with the destination
resource name, so one configuration can route
every resource to its own subject, or its own stream, or both:

| `stream`     | `subject`           | Resource `orders` publishes to |
| ------------ | ------------------- | ------------------------------ |
| `EVENTS`     | `events.{resource}` | `events.orders` in `EVENTS`    |
| `EVENTS`     | `events.all`        | `events.all` in `EVENTS`       |
| `{resource}` | `{resource}.v1`     | `orders.v1` in `orders`        |

For example, provision a stream named `WAREHOUSE_EVENTS` capturing
`warehouse.>`, then configure:

```json theme={"theme":{"light":"github-light-default","dark":"github-dark-default"}}
{
  "url": "nats://localhost:4222",
  "stream": "WAREHOUSE_EVENTS",
  "subject": "warehouse.{resource}",
  "publish_timeout": "5s",
  "max_in_flight": 1024,
  "max_in_flight_bytes": 16777216
}
```

Expanded subjects must be concrete: wildcards, whitespace, and empty subject
tokens are rejected, and stream names cannot contain dots or wildcards.
Resource names are not sanitized. Use destination resource mapping when source
names are unsuitable for NATS subjects or stream names.

With `create_stream` enabled, a missing stream is created with file storage,
server defaults for retention and limits, and one subject filter derived from
the subject template: `events.{resource}` yields `events.>`, `{resource}.events`
yields `*.events`, a fixed subject yields itself, and a per-resource stream
captures that resource's exact subject. Anything beyond that, such as retention
or replicas, belongs in server-side stream management. Disable creation on
production pipelines where streams are provisioned deliberately.

The sink never modifies an existing stream. It must enable publish
acknowledgments and its subject filter must capture the expanded subject. A shared stream is checked at open using `*` for the resource, before
any write. A per-resource stream is resolved when that resource's first batch
is routed, and a missing or non-capturing stream fails the batch before any
message is sent. Note that a filter of `events.>` matches `events.orders` but
not the bare subject `events`. Each publish carries `Nats-Expected-Stream`, so a
subject routed to another stream fails rather than silently writing there.
Workers need permission to inspect the stream, publish to its subjects, and
receive publish acknowledgments. Connection testing also requires JetStream
account-info access.

## Write modes and message format

Supported bounded ingestion types are `full_append`, `incremental_append`, and
`cdc_append`. Continuous execution supports append for rows, messages, and
change events. Replace, upsert, merge, and destination deletion are unsupported.

Each message body is a JSON object containing the Arrow row's columns, without
a trailing newline. The shared JSON encoder renders bytes as base64 and temporal
values as text. Every message has:

* `Content-Type: application/json`
* `Filament-Operation: insert`, `update`, or `delete`

Updates and deletes append change-history messages; they do not modify previous
messages. Message-source envelope columns remain in the JSON object. Raw payload
and header passthrough are not implemented.

The sink validates and serializes the complete batch before publishing. Messages
exceeding the configured stream or server size limit fail; they are not split.

Within each `Apply`, messages are submitted in row order over one reused
connection. The publisher pauses whenever either outstanding limit is reached.
SDK acknowledgment callbacks can arrive in any order and release capacity for
more messages. Shared completion tracking avoids one waiting goroutine per message.
`Apply` waits for every acknowledgment before returning success. Set
`max_in_flight` to `1` for sequential publication.

The byte limit includes encoded payloads, headers, and subjects, excluding reply
subjects and protocol framing. An individual message must fit within this limit;
otherwise the entire batch is rejected before any messages are sent. The limit
bounds outstanding publications, not total process memory: the complete encoded
batch is also held until `Apply` returns. Use the run's `BatchMaxRows` and
`BatchMaxBytes` options to control input batch size; these do not combine records
into a single NATS message.

Concurrent `Apply` calls remain serialized, so the configured limits apply to
the whole sink session. No global ordering is promised for concurrently produced
batches. Raw-message aggregation and connection pooling are not required for
asynchronous publication and are not implemented.

## Durability, replay, and failures

`Apply` succeeds only after JetStream acknowledges every message. Durability
therefore follows successful `Apply`, using the destination stream's configured
storage, replication, and retention settings. `Commit` finalizes a bounded run;
`CommitEpoch` returns per-resource totals for an already acknowledged epoch.
Continuous source progress advances through Filament's normal epoch certification
flow, after destination success.

Messages become visible individually. Neither batches nor epochs are atomic.
An abort stops work and clears local accounting; it never deletes acknowledged
messages. If a batch fails partway through, any accepted messages remain visible, including
messages submitted after the one whose acknowledgment failed.

When a row includes a nonempty `_filament_event_id`, the sink derives a stable
`Nats-Msg-Id` from the event ID, tenant, pipeline, sink connection, stream,
subject, and destination resource. The ID remains stable across attempts and
rebatching. Duplicate acknowledgments count as successfully handled input rows.
Rows without event IDs are published without deduplication IDs.

Delivery is **at least once**. JetStream duplicate suppression lasts only for the
stream's configured duplicate window. Retries outside that window, or retries
of ordinary rows without event IDs, can produce duplicates. Reusing an event ID
for different content can suppress the later publication.

A lost acknowledgment leaves the write outcome uncertain. The sink stops
publishing, rejects the epoch commit, and reports unproven destination quiescence
on session closure. Any other session failure is also retained and returned by
`CloseSession`; an epoch still active at closure is dropped like an abort.
A bounded `Abort` does not repeat a failure that `Apply` already returned. On failure or cancellation, it stops submissions, closes the connection,
detaches completion tracking, and cleans up SDK publications before returning.
Late callbacks cannot complete or fail a subsequent batch.
It does not reconnect or retry publishes in the background.
It advertises no destination owner fencing, isolated epochs, or server-enforced
bound on residual effects; uncertain shutdown follows Filament's conservative
attempt recovery path.

Write receipts report payload bytes and an encoded CRC32C over concatenated JSON
message bodies. This integrity evidence excludes NATS headers and protocol framing.
