Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
54 changes: 40 additions & 14 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
[workspace]
members = [
"crates/lance-context-ingestion",
"crates/lance-context-merge",
"crates/lance-context-core",
"crates/lance-context-metrics",
Expand Down
4 changes: 4 additions & 0 deletions crates/lance-context-core/src/registry.rs
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,10 @@ pub struct RegistryEntry {
///
/// The name says which stores exist; the dataset on object storage is the
/// data. Implementations must be safe to share across tasks (`&self`).
#[allow(
clippy::double_must_use,
reason = "async_trait adds must_use to boxed futures"
)]
#[async_trait::async_trait]
pub trait StoreRegistry: Send + Sync {
/// Whether a store named `name` exists.
Expand Down
34 changes: 34 additions & 0 deletions crates/lance-context-ingestion/Cargo.toml
Original file line number Diff line number Diff line change
@@ -0,0 +1,34 @@
[package]
name = "lance-context-ingestion"
version = "0.1.0"
edition.workspace = true
license.workspace = true
description = "Bounded session-ordered ingestion with recoverable WAL and independent consumers"

[features]
default = []
lance = ["dep:arrow-array", "dep:arrow-ipc", "dep:arrow-schema", "dep:lance", "dep:lance-file", "dep:lance-index", "dep:lance-io", "dep:lance-table", "dep:redb", "dep:tempfile"]

[dependencies]
async-trait = "0.1"
arrow-array = { version = "58", optional = true }
arrow-ipc = { version = "58", optional = true }
arrow-schema = { version = "58", optional = true }
lance = { version = "9.0.0", optional = true }
lance-file = { version = "9.0.0", optional = true }
lance-index = { version = "9.0.0", optional = true }
lance-io = { version = "9.0.0", optional = true }
lance-table = { version = "9.0.0", optional = true }
futures = "0.3"
object_store = "0.13.2"
redb = { version = "3.1.3", optional = true }
serde = { version = "1", features = ["derive"] }
serde_json = "1"
sha2 = "0.10"
tempfile = { version = "3", optional = true }
thiserror = "2"
tokio = { version = "1", features = ["macros", "rt-multi-thread", "sync", "time"] }
uuid = { version = "1", features = ["v4", "v5"] }

[dev-dependencies]
tempfile = "3"
220 changes: 220 additions & 0 deletions crates/lance-context-ingestion/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,220 @@
# lance-context-ingestion

Experimental streaming ingestion primitives. The implementation currently provides
the pipeline, journal protocol and an optional Lance table sink. Application
alignment/source adapters and deployment orchestration remain separate. This crate
does not provide a complete distributed ingestion service.

```text
replayable source / stable partition-local receipts
-> bounded concurrent history prefetch
-> ordered alignment within each virtual partition
-> bounded WAL queue / byte, count, timer flush
-> immutable payload + skip links + conditional head publication -> ACK
| |
v v
checkpoint consumer table consumer
coalesced session deltas independently grouped WAL ranges
| |
conditional session states staged fragments / fenced manifest commit
|
asynchronous indexer
```

## Implemented contracts

- `Partition::start_with_loader` overlaps history loading, alignment and WAL
publication. Different virtual partitions run independently. Within a partition,
load completion may reorder, while alignment and durable publication stay ordered.
Do not change session-to-virtual-partition routing when changing worker count.

- `enqueue` reserves bytes before admission. A reservation remains held through the
durable ACK, covering input, history, bounded output and serialization. Adapter
caches, transient adapter allocation, executor and storage-client memory need their
own budgets. Sources must also bound concurrent requests waiting for admission.

- `Aligner` owns the session cache and speculative state. Prefetched checkpoints may
be stale: reconcile their revisions with pending deltas. On failure, discard the
adapter and recover; speculative changes must never become checkpoints directly.

- `Journal` stores immutable segments binding records, state deltas, exact input
digests, receipt identities, run/schema and predecessor sequence. Conditional head
updates fence stale writers. Failed or cancelled commits poison a writer.

- ACK follows head publication. An uploaded orphan is not committed. Retry the
same partition sequence, receipt, session and input bytes after an uncertain result.
A gap is rejected; older receipts are compared with the committed journal.

- Immutable skip links support bounded chronological recovery pages and historical
receipt lookup without reading unrelated record payloads. Recovery replays all
pages after the adapter's checkpoint, not just one next WAL segment.

- `Partition::start_with_aligners` can run several session alignment lanes inside
one stable durable partition, independently of history-loader concurrency and WAL
batch size. Same-session calls stay ordered; completed lanes rejoin the original
sequence before publication. Changing lane count on restart reroutes replayed
session deltas without changing partition identity. An error or cancellation
stops all speculative lanes. Lane count must fit the queue-entry budget; adapter
caches are additional to the shared input/output byte budget. This does not
schedule workers on other machines or remove the WAL ordering barrier.

- Named `Consumer`s have independent durable cursors and coalesce producer segments
into their own batches. Sink output and input coverage must be committed together;
cursor writes can fail after output succeeds, so repeated/regrouped input must be
idempotent. The scheduler owns exclusive consumer assignment and sink-side fencing.

- `ReceiptIndex` resolves exact source receipt/session/input-digest identities after
an HTTP retry or restart. A dedicated consumer writes immutable receipt mappings;
`find_many` reads the index concurrently, then reconciles its unindexed WAL suffix
once for the whole requested batch. It handles index writes ahead of the cursor
and rejects a receipt reused at another sequence. A miss covers only the returned
committed position: the admission owner must still check its in-flight map and
serialize sequence assignment. This index does not schedule source fan-out or
replace the requirement to ACK every partition before advancing source progress.

- `SourcePartition` adds a serialized admission owner for continuous callers that
have stable source receipts but no partition sequence numbers. `enqueue_many`
checks committed and in-flight identities, assigns sequences only to new inputs,
and returns independent ACK waiters. A repeated in-flight receipt watches the
original commit without blocking later alignment dispatch. Dropped ACK waiters
do not cancel work; cancellation during admission requires recovery of the
unknown admitted prefix. Run its dedicated receipt consumer independently and
keep HTTP/source queues bounded. This API does not supply HTTP authentication,
cross-partition fan-out, or a migration from an application's previous WAL format.

- `Writer::with_backlog` optionally limits committed segments outstanding for
every required consumer. A consumer that has not started is at zero; table and
checkpoint progress are both required when both are configured. The publisher
waits before writing another segment, retaining bounded pipeline reservations;
other partitions and already durable retries remain independent. Consumers must
continue running while producers drain. Cancellation or ownership transfer fences
a paused writer. Reapply the same policy on every acquire/restart. The gate checks
durable cursor metadata and its committed ancestry; it adds storage reads and is
not itself a throughput optimization. It bounds unconsumed payload bytes by
`max_segments * max_segment_bytes`, not retained history, orphan uploads or total
storage. No WAL garbage collection or scheduling is implied.

- `SessionCheckpoints` provides an actual object-store checkpoint sink: group by
session, reduce ordered deltas, then write each session once with a conditional put.
A partially successful checkpoint batch can leave some session states ahead of the
global cursor. Restore each state with its own sequence and skip already applied
deltas when replaying the remaining global prefix.
`Reducer::apply_batch` receives only a session's ordered, unapplied deltas and
lets an adapter decode its state once and encode it once per consumer batch.
Its default preserves per-delta `apply` behavior and intermediate size checks;
overrides must preserve those semantics and bound intermediate state themselves.
The sink also checks final output size before writing. This interface alone does
not accelerate existing reducers or change the deployed ingestion adapter.

- `SessionCheckpoints::recover` reconstructs one evicted or cold session from the
checkpoint consumer cursor plus committed WAL suffix, one segment at a time.
It returns the captured WAL position separately from the session mutation
sequence and skips partially published checkpoint deltas. This is read-only;
callers must still reconcile their newer speculative state and bound the cache.
Use the checkpoint consumer name, never an unrelated table consumer cursor.

- With the `lance` feature, `lance_sink::stage` writes immutable Lance 2.2 files
using a Zstd-annotated schema. Lance's constant-valued pages use scalar
encoding before codec selection; even a single large string can take that path. `LanceTableSink::commit_staged` coalesces staged
partitions and atomically publishes rows plus each partition's covered sequence.
Fully covered retries are skipped; gaps and partial overlaps are rejected.
A caller-supplied ownership-guarded commit handler is wrapped by an exact version
pin to reject implicit rebasing. Uncertain publication poisons the sink; reopen
and reconcile the table watermarks before retrying. The sink does not acquire
table ownership or schedule index maintenance.

Use a backend supporting atomic conditional updates, such as a suitably configured
cloud object store. `object_store::local::LocalFileSystem` does not implement the
required update operation and is rejected; there is no unsafe local-lock fallback.
Tests use the real `object_store::memory::InMemory` implementation for CAS behavior,
with fault injection around writes. These tests do not establish cloud durability,
process-crash behavior on a real durable service, or production throughput.

## Bulk sources

For replayable batch sources, use `SourcePartition::start_batched` with a fresh
`BatchFlush` shared by the partition's alignment adapters. Pass the existing source
batch to `enqueue_many`; do not replace its stable record receipts when combining
batches for transport or WAL publication. The admission owner validates the entire
batch, preserves in-flight deduplication, and requests a flush through its final
assigned sequence. An all-duplicate batch does not manufacture a WAL entry.

This mode does not use `BatchPolicy::max_delay`. WAL collection continues until a
source batch boundary, the byte/count limit, output memory headroom, or shutdown.
A later already-admitted batch boundary can coalesce available batches. Limits may
split a large batch into committed prefixes: rows, state deltas and source receipts
remain together in each segment, and callers await **all** returned ACKs before
advancing source progress. Table/checkpoint consumers still group segments
independently. The bulk constructor changes scheduling, not the WAL format.

Same-session alignment sees its speculative state throughout the batch. If an
adapter evicts uncommitted state and needs to recover it, call
`BatchFlush::request_prefix(sequence)` before waiting for that prefix's durability.
This flushes available ordered outputs even when a later batch boundary is pending;
otherwise the blocked alignment could prevent the batch itself from completing.
Do not treat a stale checkpoint as the current batch's state. Cache/input/output
budgets still apply; increasing a batch target does not permit unbounded memory.

## Remaining integration

1. Adapt the existing revisioned session history/cache and cross-call alignment
implementation. Preserve its compaction/branch identity rules and source ordering;
the generic pipeline deliberately does not invent a new turn-ID algorithm.
1. Wire table ownership and independently scheduled session ZoneMap maintenance.
Staging is independent of publication; a distributed worker transport must
carry validated staged results and retain ownership through manifest commit.
1. Add source fan-out receipts and contiguous source progress. A source call spanning
partitions is complete only after every required partition ACK.
1. Add worker ownership orchestration, stage timing/queue telemetry, consumer run loops,
deployment of the backlog policy and safe WAL reclamation. No WAL files are deleted here.
1. Verify real compacted sessions, process crash/restart with durable storage, Lance
uncertain commits and source retry integration before a guarded production handoff.
Existing source-reader local audit history must survive that handoff.

Run focused checks from the workspace root:

```sh
CARGO_TARGET_DIR=/tmp/trace-streaming-target cargo test -p lance-context-ingestion --features lance --offline
CARGO_TARGET_DIR=/tmp/trace-streaming-target cargo clippy -p lance-context-ingestion --features lance --all-targets --offline -- -D warnings
cargo fmt -p lance-context-ingestion --check
```

# Local state and a Lance recovery log

The `local_lance` module (`lance` feature) stores individual binary state cells
and exact receipts in a local redb database. This is a disposable local index;
it does not require a remote KV request for each message. Session routing and
the durable partition lease remain the caller's responsibility.

An `AlignedCall` contains typed output rows and the matching state mutations.
`LocalLancePartition::commit` coalesces calls into one Lance 2.2 commit with Zstd
requested on its columns, then applies one local transaction. A cancelled or
failed commit fences the instance. Reopen and check the original receipts before
realigning or acknowledging a retry. The exact-version commit fence prevents a
stale writer from acknowledging another writer's output at the same sequence.

`checkpoint` streams the local tables in bounded batches to a separate immutable
Lance dataset. `publish_checkpoint` pins its exact URI/version in the WAL's
manifest under the same lease/version fence. Do this on a byte/time threshold,
not once per source call. Cold restart discovers that pointer, restores state,
and reads only new immutable WAL fragments. Recovery projects state and receipt
columns, excluding output bodies. No method deletes WAL or checkpoint history.

The first output schema supports flat Arrow columns. State cell encoding and
legacy migration are adapter contracts; this API does not convert JSON payloads
supplied by a caller. This is an opt-in backend: existing `Journal` and
`SessionCheckpoints` adapters keep their existing storage and recovery behavior.
Migrate their committed suffix and receipts before changing the intake log.

`local_lance_reader::OutputRange` reads only typed output from an immutable
append range after validating both version watermarks, run/schema/partition and
full fragment ancestry. It enforces physical and decoded byte limits separately.
`BatchReference` is a bounded, checksummed binary descriptor binding that range,
its generation/predecessor and application metadata. A consumer must budget the
referenced payload and its own decoded representation, never the descriptor size.
No parsing failure permits falling back to a different log format.

`batch_charge` checks the input, IPC encoding and retained Arrow allocation charge
before admission. Shared IPC buffers are charged once per allocation while their
batches remain alive. Charge computation currently encodes and decodes a call;
commit encodes again. Account for that CPU work when profiling the adapter.
Loading
Loading