Consumer¶
Consumer reads records from feed via cursor (feed-scoped) or stream sequence (stream-scoped). Feed-scoped reads walk sealed commit history. Stream-scoped reads target one stream directly, sealed or open. See Architecture.
Read Scopes¶
| Scope | API | What It Reads |
|---|---|---|
| Feed-scoped | read(from, mode) / subscription.poll(mode) |
Sealed commits after cursor |
| Stream-scoped | read(feed, streamSeq, mode) |
One stream by sequence, sealed or unsealed |
Feed-scoped reads are continuous-consumption. Cursor marks position in commit history; each page returns records from commits after that position plus end cursor for next call.
Stream-scoped reads bypass commit history. Sealed stream served from commit artifacts. Unsealed stream stitched from durable WAL tail. Inspect open ingest or replay one stream in isolation.
Cursor & Paging¶
CursorToken is feed name plus commit sequence. Sequence 0 is feed start. from cursor is exclusive: records begin in first commit after that sequence.
| API | Purpose |
|---|---|
head(feed) |
Current feed head cursor |
log(feed) |
Full sealed commit list |
read(from, mode) |
One page from from to head |
read(from, to, mode) |
One page from from to to |
Page shape is ReadPage:
| Field | Meaning |
|---|---|
records |
Page records |
end |
Cursor after last included commit |
Pass end as next from to continue. When end equals head, consumer is caught up.
Paging is bounded server-side:
| Setting | Value |
|---|---|
Default limit |
1,000 records |
Maximum limit |
10,000 records |
For RAW and DEDUP, server fills page commit by commit. For LATEST, server resolves winners across full requested range before returning.
Subscriptions¶
subscription is named, persisted cursor in metadata. Multiple consumers use different names on same feed and advance independently.
| Step | API | What Happens |
|---|---|---|
| Bind | PUT /v1/subscriptions/{name}?feed=... |
Associates name with feed; creates cursor at start if new |
| Poll | poll(mode) |
Reads from saved position to current head |
| Advance | advance(end) |
Persists new position |
Java client binds lazily on first position() or poll() call. Subscription with no saved position starts at commit 0.
Poll returns same ReadPage shape as direct feed read. After processing records, call advance(batch.end()) so next poll skips handled data. Position survives process restarts and re-binding under same name.
Long Poll¶
poll(mode, wait) blocks up to wait for new commit on feed. Server wakes waiters when seal publishes commit. If nothing arrives before timeout, poll returns empty page at current position.
Long poll avoids tight loops when waiting for next sealed commit.
Read Modes¶
Modes apply at read time over chosen scope and range. Stored data is unchanged.
| Mode | Behavior |
|---|---|
RAW |
Every record in arrival order, including duplicates |
DEDUP |
First occurrence of each fingerprint in range |
LATEST |
One winner per key by version (then commit order); keyless records pass through |
DEDUP keys on stored fingerprint bytes, not business key. Default fingerprint is SHA-256 of the payload. Identical payloads collapse unless the producer supplied different fingerprints.
LATEST requires keyed feeds for per-key collapse. Uses key index and version resolver: higher version wins, ties break toward later commit. Tombstone winners emitted as tombstone records. Point lookup of tombstoned key returns empty.
Point Lookup¶
Keyed feeds: latest(feed, key) returns current winning record for one key without scanning range. Server walks key-index metadata and fetches winning pack entry only.
Diff & Merge¶
Commit-scoped utilities compare keyed state between two cursor positions:
| API | Purpose |
|---|---|
diff(from, to) |
Added, updated, and deleted keys between two commits |
merge(left, right) |
Three-way merge resolution between two commit tips |
Operates on key metadata only; full record payloads not required.
Truncation¶
Retention and compaction remove old commits from active log. Cursor before log floor fails with CommitTruncatedException.
| Situation | Result |
|---|---|
| Cursor precedes retention floor | Read rejected; log starts at higher commit |
| Stream sequence compacted away | Stream read rejected; active log starts at later stream |
Advance subscriptions and cursors as data is processed so positions stay inside retained window.
Examples¶
CLI feed-scoped read from start:
./streamstack read orders --from 0 --mode LATEST
CLI stream-scoped read of one open or sealed stream:
./streamstack read stream orders 3 --mode RAW
CLI subscription poll with auto-advance:
./streamstack subscribe etl --feed orders --mode RAW --wait 5000
Java SDK continuous consumption:
Subscription sub = client.subscription("etl", "orders");
ReadPage page = sub.poll(ReadMode.RAW, Duration.ofSeconds(5));
for (Record record : page.records()) {
process(record);
}
sub.advance(page.end());
Java SDK manual cursor paging:
CursorToken cursor = CursorToken.start("orders");
while (true) {
ReadPage page = client.reads().read(cursor, ReadMode.DEDUP);
if (page.records().isEmpty()) {
break;
}
page.records().forEach(this::process);
cursor = page.end();
}
Java SDK stream-scoped read and point lookup:
StreamReadResponse stream = client.reads().read("orders", 3L, ReadMode.RAW);
stream.records().forEach(this::process);
client.reads().latest("orders", "policy-1").ifPresent(this::process);