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
ImplementStreamSource 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 fromRead 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 successfulCloseSession 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 instream.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.