§09 Links & Workers

Declarative materialization (zero-ETL), exactly-once apply, background task scheduling

Status: Approved · 2026-09-22 · amended 2026-09-26 (M1.2 as built) · amended 2026-09-26 (collection write backpressure, D86)

Links are Operon's zero-ETL mechanism: declared, continuously maintained materializations from streams into tables, collections and graphs. Workers are the stateless pool that executes links and all other background work.


LinkExampleTarget commit
stream → tableStream llm_calls → Iceberg tableIceberg snapshot (via Lakekeeper)
stream → collectionStream support_tickets → searchable, embedded collectionCollection manifest
table → collectionSearch projection of products(title, description)Collection manifest
table/collection → graphEdge table knows → CSR/CSC sidecars of a mapped graph (§07 §2.1)Graph manifest
table → tableMaterialized view / rollup (§08 §4)Iceberg snapshot
table/collection → changelog streamRow-level changes of tickets as a stream (§02 §8.1)Fenced append to the changelog stream, before the target commit
durable_events → table/graphDurable-execution search index and execution graph (§14 Phase B)Iceberg snapshot / graph manifest
implicitEvery table/collection/graph's own implicit stream → itselfper target
CREATE LINK tickets_search
  FROM STREAM support_tickets
  INTO COLLECTION tickets
  TRANSFORM (
    SELECT key AS _pk,
           json_get_str(value, '$.subject') AS subject,
           json_get_str(value, '$.body')    AS body,
           embed('endpoint:text-embedding-3', json_get_str(value, '$.body')) AS body_vec,
           timestamp AS ts
  )
  WITH (batch_bytes = '64MiB', batch_interval = '2s', on_error = 'dead_letter');

2. Transforms

  • SQL (DataFusion) over the decoded record batch: projections, filters, JSON extraction, casts, scalar UDFs.
  • Decoders: JSON, Avro/Protobuf (with a registered schema), CSV, raw bytes.
  • embed() UDF (optional): calls an external embedding endpoint (OpenAI-compatible HTTP, or a self-hosted model server) with batching, retries and rate limiting; results cached by content hash. Off by default; configured per namespace. (Operon does not host models.)
  • Mergeable aggregate states for MV-style links (§08 §4).
  • Not supported: stateful joins across streams, windowed aggregations with watermarks (§00 §7).

3. Exactly-once semantics

  1. A link task leases (link_id, source_partition_range) in meta with an epoch.
  2. It reads from the last applied offset recorded in the target's latest commit (Iceberg snapshot summary property / collection or graph manifest field).
  3. It builds target files, then commits the target including the new applied offsets.
  4. A zombie task (stale epoch) fails its commit: Iceberg optimistic concurrency / manifest-pointer CAS rejects it; leases are also fenced by epoch.
  5. Restart = resume from the committed applied offset. Duplicate work after crashes is discarded, never double-applied.
  6. A link with a changelog appends change records before the target commit, fenced by epoch and by the changelog's recorded source_upto; the retried task skips what is already appended (§02 §8.1).

4. Errors, dead letters, schema evolution

  • on_error: fail (pause link, alert), skip, or dead_letter (write offending records + error to a DLQ stream).
  • Source schema changes: additive fields auto-propagate if the target allows (Iceberg schema evolution, collection dynamic mapping); incompatible changes pause the link with a clear error.
  • Link lag (offsets and seconds) exported as metrics and exposed in system.links.

5. Worker task catalog

TaskTriggerOutput
SegmenterWAL objects pendingStream segments; WAL deletion
Stream compactionCompacted stream dirty ratioNew segments
RetentionPolicy scheduleSegment/file deletions
Link applySource lag > 0Target commits
Split mergeSplit count/size tiersMerged Tantivy splits
Lance compaction / vector index optimizeFragment count, unindexed rows, centroid driftNew Lance version / index
Hot artifact buildPromotion or pin; new manifest versionHNSW / projection artifacts
Iceberg compactionSmall files, delete ratioRewritten data files (nimtable/iceberg-compaction)
Iceberg maintenanceScheduleSnapshot expiry, orphan cleanup, manifest rewrite
Graph sidecar buildNew/compacted edge segmentsCSR/CSC sidecars
Graph algorithmsUser jobResult tables/columns
GCScheduleUnreferenced object deletion (§03 §7)
Meta snapshotRaft log size / scheduleSnapshot to S3
Durable timer sweep (§14 Phase B)Timer shard leaseFires due promise/task/schedule deadlines
Durable retention (§14)Policy scheduleDeletes settled origin documents past retention

As built in M1.3 (task keys under the lease task/<key>, with their priorities):

TaskKeyPriorityTriggerOutput
Split mergecollection-merge/<cid> (per namespace)CompactionStableLogMergePolicy plans a merge of splits with the same schema version, or a split with ≥ 30 % (and at least 1 000) deleted docs is rewritten alone; bounded by max_merge_docs (10 M) and max_merge_bytes (1 GiB)One split re-indexed from _source, committed as a Maintenance manifest with rebase
Lance compactioncollection-compact/<cid> (per namespace)CompactionEnough small fragments, or a fragment's deleted share over the thresholdLance's rewrite committed as a detached version, then a Maintenance manifest
Hot artifact buildhot-build/<cid>/<i> (per namespace; i the vector index)HotBuildA column pinned or promoted with no artifact, or a stale artifact past the rebuild threshold (M1.3 Ruling 16)An HNSW artifact under hot/hnsw/…, referenced by a Maintenance manifest
Learner evictionmeta-membership (cluster-wide)MaintenanceEvery 60 s in cluster modeRemoves metastore learners whose node/<id> lease has been expired for 10 min

Known limits carried to M2's resource budgets (§6): a merge streams its documents into the split writer, but the built index and the bundle are still in memory, so its peak memory is about twice the output split (bounded by max_merge_bytes, 1 GiB); a hot build's work directory, which can hold gigabytes, is removed synchronously on the task's thread (M1.3 rows R32.3, 5.6).

6. Scheduling

  • Tasks are leases in meta (task_key, epoch, owner, deadline); workers pull tasks, renew leases, and are fenced by epoch.
  • Priorities: link apply (freshness SLO) > segmenting > merges/compaction > hot builds > maintenance > GC.
  • Fair share per namespace with weights; per-namespace caps on concurrent tasks and bytes/s to prevent noisy neighbors.
  • Autoscaling signals: total link lag (seconds), compaction debt (bytes), queue age per priority.
  • Resource budgets: CPU/memory per task class; object-store request budgets (PUT/GET rate) per node to stay under prefix limits.

As built in M0 (operon-worker, M0.4):

  • A task source proposes task keys with work (the segmenter one per partition with a due WAL run, retention and GC one singleton each, link apply one per link, D30). Every node runs one Worker, which polls its sources (default every 1 s), orders the candidates by priority, then round-robin across namespaces within a priority, and starts what fits under max_concurrent (16) and max_per_namespace (4).
  • A task runs under the metastore lease task/<key> (TTL 30 s) with a fence at the lease's epoch and a cancellation token. The worker renews every third of the TTL; a renewal that finds the lease expired re-takes it at the same epoch if nobody else took it (ReacquireLease), so a slow run keeps its fence; a lease someone else took cancels the run, and every metastore commit it attempts with its fence is rejected. Leases are released when a run ends; a crashed worker's leases expire.
  • Weights, byte-rate caps, autoscaling signals and per-class resource budgets are not built yet.

7. Backpressure

  • Collection links (M1.3, D86): each collection has a budget on its unapplied backlog, the records past applied (max_unapplied_records, default 1 000 000) and their log bytes (max_unapplied_bytes, default 128 MiB, at most half of tail.max_bytes). While the backlog is at either budget, collection writes are refused with HTTP 429 or gRPC RESOURCE_EXHAUSTED and Retry-After (estimated from the link's recent apply rate, 1–30 s); every write response reports the backlog (Operon-Unapplied-Records, Operon-Unapplied-Bytes). A bulk load may send Operon-Backpressure: off to write up to 4× the budget. M2 makes the budget a per-namespace and per-collection quota (D65) with metrics.
  • Other links (M4 tables, M3 graphs) may configure a max_lag that throttles their source stream's producers the same way, or keep accepting (the log absorbs the backlog).
  • Tail memory is bounded per object on query nodes; when exceeded, strong reads fall back to a range tail read directly from the log (bounded; Unavailable beyond the bound) — slower but correct (M1.2).

On this page