Keyboard shortcuts

Press ← or → to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

Execution Semantics

How Rivet behaves under normal execution, retries, crashes, resume, repair, and reconcile. This is the contract a downstream pipeline can build on.

This page is a user-facing summary. The binding contracts live in the ADRs, which this document links into.


Scope

This document describes guarantees and known non-guarantees for:

  • single-table exports (full, incremental, chunked, time-window, cdc),
  • destination writes (local, S3, GCS, Azure Blob Storage, stdout),
  • state and journal updates,
  • automatic retries and --resume after crashes,
  • rivet reconcile and rivet repair,
  • in-process and subprocess parallelism.

It does not cover database / network / disk failures whose mode is external to Rivet (e.g. a corrupted Parquet file caused by a bad disk).


Core concepts

TermMeaning
RunOne invocation of rivet run. Has a unique run_id. Produces zero or more output files.
ExportA named entry in rivet.yaml — a query + cursor + destination triple. One run executes one or many exports.
Modefull, incremental, chunked, time_window, cdc. See docs/modes/.
BatchOne FETCH worth of rows materialized as an Arrow RecordBatch. Streamed; never accumulated in memory.
ChunkA row range (chunked mode) processed as one unit. Produces one output file. Has a checkpoint row in the state DB.
CursorThe last extracted value for incremental exports. Stored in export_state.last_cursor_value.
File logThe per-export ledger of files written to the destination (file_log table, renamed from file_manifest in schema v8).
JournalA per-run event log, persisted as a single JSON document in the state DB (run_journal table); inspect with rivet journal. Authoritative for “what happened when” during one run.
ProgressionThe committed / verified boundary per export (export_progression table). Advisory only.

Normal execution

The pipeline is one straight line per export (see architecture.md § Data flow):

begin_query → FETCH batch → write batch to temp file
                          → next FETCH
                          → ...
            → finalize writer
            → destination.write(temp_file)
            → record manifest entry
            → advance cursor (incremental) / record chunk completion (chunked)
            → record metric (at end of run)

The exact state-write ordering is defined by ADR-0001 — State Update Invariants (I1–I7). The pipeline source code references these IDs at the call sites.


Retry semantics

Retries are classified by error type in src/pipeline/retry.rs:

The classifier (RetryClass) has two outcomes — Transient (retry) and Permanent (propagate):

ClassExamplesRetried?
Transientconnection reset, “server has gone away”, lock-wait timeout, deadlock / serialization failure, “too many connections”, cloud (temporary) writes (S3 / GCS)Yes — exponential backoff up to tuning.max_retries. The variant carries needs_reconnect (reopen the source connection first — e.g. a network reset or 08xxx SQLSTATE) and extra_delay_ms (an added settling delay for capacity errors like “too many connections” / “database system is starting up”)
Permanentsyntax error, auth failure, missing table / column, a statement-duration timeout (statement_timeout / max_execution_time)No — propagates immediately. Uncategorized errors default to Permanent (they are not retried)

A retried batch starts from the same cursor position as the failed attempt — see ADR-0001 I3 (Write Before Cursor). At-least-once delivery to the destination is therefore possible: on retry after a destination write succeeded but the cursor failed to advance, the same rows are written again, producing a duplicate file.

Retry-safety per destination is declared by capabilities().retry_safe — see ADR-0004 — Destination Write Contracts:

Destinationretry_safe
S3, GCS, Azuretrue (no partial visible objects; safe to retry)
Local filesystemtrue (staged temp file + atomic rename; a failure leaves nothing at the final path)
stdoutfalse (no commit boundary; retry produces duplicate/corrupt output)

When retries occur against a non-retry-safe destination, the pipeline logs a WARN so operators see the mismatch.


Crash semantics

If the process is killed (SIGKILL, OOM, host reboot), the next state depends on where the crash landed. The full failure-point map is in ADR-0001 § Failure Point Map. Summarised:

Crash pointFiles at destinationManifestCursorNext run does
Mid-extraction (before dest.write)noneno entrynot advancedre-extract from last cursor
After dest.write, before manifestfile presentno entrynot advancedre-extract → duplicate file at destination
After manifest, before cursorfile presententrynot advancedre-extract → second duplicate + manifest entry
After cursor updatefile presententryadvancednext run starts from new cursor; metric may be missing
Clean error (Err return)noneno entrynot advancednormal retry

Two strong invariants hold across every crash point:

  • No row is silently skipped. The cursor only advances after the corresponding write succeeds (I3).
  • No file at the destination is incomplete. Writers are finalized before destination upload (I1).

The trade-off is at-least-once at the destination: a crash between write and cursor advancement produces a duplicate file. Downstream consumers must tolerate this — see Known non-guarantees below.


Resume semantics

rivet run --resume consults the state DB to decide what work is outstanding. A plain rivet run (no --resume) never skips the chunks of a run that FINISHED — it does a fresh full pass. A checkpointed run holds its export’s run lease (an flock beside a SQLite state DB, a session advisory lock on a Postgres one) for as long as the process lives, and the OS releases it when the process dies, kill -9 included. So when a plain run finds a chunk-checkpoint run still in progress:

  • its process is gone (the lease is free): the run resumes it, exactly as --resume would. A chunk_dense plan cannot be resumed, so it starts over.
  • its process is alive (the lease is held): the run is refused, and so is --resume — wait for the live run to finish.

What --resume does:

  • Incremental exports resume from export_state.last_cursor_value.
  • Chunked exports consult the chunk_task table: tasks in completed are skipped; tasks in pending or running (the latter reset to pending on resume) are re-issued; tasks in failed are retried while attempts < max_chunk_attempts.
  • Full and time-window modes do not resume — they restart from the beginning. The previous run’s output files remain at the destination unless cleaned manually.

Chunk task transitions are strictly forward (pending → running → {completed | failed}). A completed chunk is never re-claimed, even after a crash — see ADR-0001 I5 (Chunk Task Acyclicity).


Repair semantics

rivet repair re-exports specific chunks identified as mismatch or unknown by a prior rivet reconcile. The full contract is ADR-0009 — Reconcile and Targeted Repair.

Key properties:

  • Repair derives only from a reconcile report (RR1). There is no operator-typed chunk range.
  • Dry-run by default (RR2). --execute is required to perform writes.
  • Repair writes new files alongside originals (RR5). Rivet does not delete or overwrite prior destination files. Downstream dedup is the operator’s responsibility.
  • Repair does not advance the committed boundary (RR4). Run rivet reconcile again afterwards to advance last_verified_*.

Reconcile semantics

rivet reconcile compares per-chunk row counts between source and destination for the latest chunked run:

Per-partition outcomeMeaning
matchCounts equal
mismatchBoth counts known and different
unknownEither count missing (e.g. chunk never completed)

Scope and limits:

  • Chunked mode only in v1. time_window and incremental exports surface a clear “not supported” error.
  • COUNT(*) based. No hash-based partition verification yet.
  • Verified boundary advances only on a fully clean report — zero mismatches and zero unknowns (ADR-0008 § PG5).
  • Exit code gates on mismatch. rivet reconcile exits non-zero when any partition is a mismatch, so rivet reconcile && <next step> does not proceed on disagreeing data (mirrors rivet validate). unknown partitions (an incomplete chunk, or a non-integer keyset key with no source re-count) are surfaced as a warning but do not fail the command — “could not verify” is not “verified wrong”, and a keyset export is structurally all-unknown. The mismatch detail is always in the printed report regardless of exit code.

Reconcile reads from the same source and never writes files itself.


Parallel execution semantics

Rivet runs two distinct parallel engines — see ADR-0010 — Two Parallel Engines:

EngineUse caseCrash isolation
In-process scoped threadsChunked export of a single table, split into N concurrent chunksNone — a panicking worker can fail the run
Subprocess fan-out (--parallel-export-processes)Many independent exports concurrently, one child per exportOS-level — a failing child exits non-zero; the parent aggregates and returns non-zero

Both engines honour the same state invariants. Chunk checkpoints serialise the parallel threads’ state writes; subprocess children share one state DB — the parent migrates it once before spawning, and SQLite WAL (or the PostgreSQL state backend) handles the children’s concurrent writes.


Quality gates

When quality: is configured, the pipeline evaluates row-count, null-ratio, and uniqueness checks before advancing committed progression. A failing gate aborts the run with a non-zero exit code; the destination files remain (manual cleanup or replacement is the operator’s call). See docs/best-practices/quality-checks.md.


Destination commit boundaries

DestinationCommit protocolWhat “Ok” means
S3 / GCS / AzureFinalizeOnCloseObject is committed only after writer close; a mid-upload failure leaves nothing visible
Local filesystemAtomicOk means the full file is present; staged temp file + atomic rename, so a failure leaves nothing at the final path (retry_safe: true, partial_write_risk: false)
stdoutStreamingNo atomic commit boundary; partial output may be observable before write() returns

Full per-backend table and rationale: ADR-0004 — Destination Write Contracts.


Known non-guarantees

Rivet does not currently guarantee:

  • Exactly-once delivery to the destination. Crashes between destination write and cursor advancement can produce duplicate files. Plan downstream dedup or idempotent ingestion — the manifest’s per-part content_fingerprint is the supported dedup key: identical rows produce byte-identical parts (and the same fingerprint) across rivet releases, so a duplicate is safely droppable by fingerprint. See recipes/idempotent-warehouse-load.md.
  • Continuous / near-real-time replication. Rivet does capture CDC to files (mode: cdc — inserts/updates/deletes via a Postgres logical replication slot / MySQL binlog / SQL Server CDC change tables / MongoDB change streams, into typed Parquet/CSV — or the JSON-blob document image for MongoDB — resuming from the last committed log position each run), but it is not a continuously-running stream — changes are captured per invocation, not delivered live. For always-on near-real-time replication use Debezium or Estuary.
  • Completeness of incremental cursors that can tie. Incremental resume uses a strict WHERE cursor > last_value. If two rows share the high-watermark value and the second becomes visible only after the run that advanced the watermark past it — e.g. a low-resolution updated_at (second granularity) or rows committed at the same timestamp after the read snapshot — the next run skips them and they are never exported. (Keyset pagination is unaffected: its key is planner-enforced unique + NOT NULL.) Use a strictly per-row-distinct, monotonic cursor (a sequence/identity id, or a timestamp with sub-value uniqueness); when the cursor can tie, re-snapshot the affected window with full/chunked mode.
  • Automatic cleanup of an interrupted write’s temp file. A crash mid-write on the local destination may leave a dot-prefixed temp file in the target directory — never a partial final file (the commit is an atomic rename, so the final path is the complete file or absent). The stray temp file is harmless and can be removed manually.
  • Schema migration handling. If the source schema changes between runs, Rivet does not migrate the destination; it surfaces a schema-drift error (see tests/live/live_schema_drift.rs).
  • Correctness of user-authored SQL. Rivet executes query: verbatim. A query that omits a WHERE clause or selects from the wrong table will export the wrong data — there is no semantic validation.
  • Protection from poorly indexed source queries. Preflight (rivet doctor, rivet check) warns about missing cursor indexes and unbounded ORDER BY, but it does not refuse to run. The operator decides.
  • Stdout state safety. Using stdout with cursor or manifest state is technically allowed but not meaningful; plan validation rejects stdout + chunked and stdout + max_file_size before execution (ADR-0004 Known Gap).
  • Atomicity across exports in one run. If a run has three exports and the second fails, the first export’s writes are already committed — the run does not roll back.
  • Cross-run ordering when running in parallel from multiple machines against the same state DB. The default state DB is a local SQLite file; concurrent processes on one machine are handled via WAL, but multi-machine deployments should point RIVET_STATE_URL at a shared PostgreSQL state backend — cross-run ordering across machines is still not guaranteed.

Test coverage

The invariants on this page are exercised by:

Invariants and recovery suites run as named semantic gates in PR CI (.github/workflows/ci.yml). Branch protection blocks merges on regression. See reliability-matrix.md for the full coverage matrix.