> ## 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.

# Streaming connectors

> Implement reusable source sessions and sink epochs with explicit recovery boundaries

A streaming connector keeps a session open across many commits. Its most
important job is to explain when a portion of the input is complete and when
the corresponding output is durable. The runner uses those two promises to
decide where a later attempt may safely resume.

Start with the [streaming lifecycle](/pages/guides/concepts/streaming) for the
end-to-end picture. This page focuses on the connector's part in that lifecycle.

## Declare what the pair can support

Sources describe whether they emit rows, changes, or messages, which operations
they produce, and what ordering and replay guarantees they offer. Continuous
planning currently requires replayable, at-least-once input. A best-effort feed
cannot promise to recover the messages lost during a worker interruption.

Sinks declare continuous write policies separately from bounded policies.
Supporting bounded upserts does not automatically provide streaming upserts.
The planner checks operations, ordering, durability, and required primary keys
before writing. Replace is unavailable for continuous execution because it
needs a finite snapshot and a completion boundary.

The current planner accepts an order-dependent sink policy only when the
source promises strict global order. Partition or transaction ordering alone
does not satisfy that requirement. Declare the guarantee the connector can
actually maintain through extraction and writing.

## Keep the source session alive

Implement `StreamSource` to open a reusable `StreamSession`. Opening receives
the selected resources, certified resume positions, resource membership, and
the admitted worker identity. Use those positions to resume the existing feed;
do not silently replace an unavailable position with the latest available one.

The runner alternates reads and acknowledgements on the session. A read writes
typed rows through `StreamRecordSink` and returns the source positions covered
by that read. Those positions are candidates until the runner has finished the
pipeline and committed the epoch.

Backpressure can block a row flush for much longer than a network read. Keep
protocol heartbeats and any required delivery renewal alive while waiting.
Honor cancellation, and join an outstanding read or acknowledgement before
closing the session. Closing concurrently with a read is not part of the
contract.

## Show that the source boundary is complete

Consider a source transaction that updates an order and inserts its line
items. Reaching a row limit after the order update is not a safe place to save
progress if resuming there would skip the line items. Finish the transaction
and report its end before returning coverage.

Sources express these boundaries with ordered controls alongside the rows.
Transaction begin and end markers describe transaction boundaries; a progress
boundary can describe a safe position outside a transaction. Sending a control
flushes the affected builders and waits for preceding sink writes and integrity
checks. It does not commit the epoch or acknowledge the provider.

Before an epoch can be sealed, every covered domain needs an accepted safe
boundary after its last row, and no transaction may remain open. Returning a
position from `Read` without the corresponding pipeline evidence is not enough.
The source is also responsible for proving it did not skip relevant input on
the way to that position. The pipeline cannot discover a missing broker message
from the rows it received.

Treat record, byte, idle, and age limits as targets for finding that boundary.
Cancellation is an interruption, never evidence that a transaction completed.
The current runner accepts domain-position coverage; the inbox-claim types in
the contracts do not provide an implemented inbox execution path.

## Give bookmarks a stable meaning

An ordering domain identifies the part of the source whose positions can be
compared, together with its incarnation. A partition offset is useful only if
it still refers to the same partition history. Detecting a recreated source
prevents an old offset from being mistaken for progress in a new feed.

Implement a position codec that validates, compares, and canonicalizes those
bookmarks. Codec lookup must work without connector configuration or network
access. Register the codecs for runtime persistence as well as exposing them
through the source. The datastore needs the same interpretation when it checks
for regression or compares a retried certificate.

Where the source has a stable event identity, preserve it across retries. It
combines the domain, position, and an ordinal for events sharing a position.
Do not build that identity from the run or worker attempt: doing so would make
a replay look like a new event to a deduplicating sink.

## Acknowledge only certified work

`Acknowledge` receives coverage after destination commit and durable epoch
certification. A source must not acknowledge speculative progress during a
read. Check current ownership through the supplied authority callback before
provider acknowledgements, including redeliveries that you suppress because
they are already certified.

This check matters when an old worker wakes up after losing its lease. Knowing
the correct bookmark does not give it permission to change the provider's
consumer state.

## Implement the sink's epoch lifecycle

A streaming sink opens once, then repeats three steps: begin an epoch, apply
its batches, and commit it. `StreamingSink` supplies that lifecycle through
`BeginEpoch`, `CommitEpoch`, and `AbortEpoch`; `CloseSession` ends the reusable
session. The runner does not adapt bounded `Commit` and `Abort` methods into
streaming behavior.

Bind each batch to the active epoch and reject mismatched identities. Keep the
same Arrow and encoded integrity obligations as a bounded sink. When commit
returns, the writes must have reached the durability promised by the selected
policy, and the per-resource receipts must describe the completed work.

For PostgreSQL append, this means starting a transaction, copying the batches
into it, and returning receipts after transaction commit. For NATS append,
messages are already acknowledged by the destination during batch application;
the epoch commit summarizes that accepted work. These implementations have
different visibility and rollback behavior even though they use the same
runner lifecycle.

Aborting discards only tentative effects. It cannot undo messages already
published or earlier committed epochs. Document that boundary so users can
reason about retries and duplicates.

## Make shutdown an honest promise

A successful `CloseSession` means no later destination effects remain possible
from that session. Releasing a local client or reaching a timeout does not
establish that promise if remote requests could still finish.

Return an error when shutdown is uncertain. The runtime records that
uncertainty and blocks automatic takeover. Similarly, matching an epoch ID in
memory does not prove destination-side owner fencing: that stronger capability
requires the destination to check ownership atomically with its effects.

## Test interruptions at the boundaries

Successful delivery is only part of streaming support. Exercise a slow sink,
cancellation during a read, a partial source transaction, and loss of ownership
before acknowledgement. Verify that a source cannot certify a position whose
rows have not finished writing.

For a sink, test failed writes, uncertain commit responses, repeated lifecycle
calls, and shutdown with outstanding work. An interruption after destination
commit but before certificate persistence is especially useful: recovery may
replay that epoch, and the result must match the connector's documented
duplicate-handling behavior.

The public contracts live in `stream.go` and `epoch.go`. The coordinator is in
`runner/continuous.go`, and pipeline completion is enforced in
`pipeline/epoch.go`. The Kafka and NATS sources and PostgreSQL and NATS sinks
provide concrete implementations to read alongside those contracts.


This documentation is built and hosted on [Mintlify](https://mintlify.com), a developer documentation platform.