Repository navigation
fix(dynostore): ListOffsets by timestamp now answers from record time, not object write time - #840
solace-aross wants to merge 14 commits into
Conversation
sgamelin
left a comment
There was a problem hiding this comment.
Thanks for this. Answering from record time instead of object write time is the right fix, and the ceiling-plus-scan logic is sound for a complete index: every batch before the ceiling entry has a max_timestamp below the target, so the first match can't come earlier. The offset-ordered replay in merge_time_index and the serde compatibility with existing watermark documents also hold up, and rolling back to an older binary is safe.
I'm requesting changes for one correctness problem: the backfill turns a transient object-store error into a permanently wrong index (inline). The other main point is where the index lives. It grows with nearly every produce from the first one, inside the document that every produce and every fetch reads (inline).
Also:
- The client-visible part of the no-match change depends on #838. Until that lands, the service layer still maps
Noneto offset 0 (list_offsets.rs#L168). CHANGELOG.mdon main now has an[Unreleased]section and an advisory check. Consumers that seek by time on S3 or memory storage land on different offsets after upgrading, and the first lookup on an existing partition runs the backfill, so this deserves a### Fixedentry.- On #835: agreed that its
retainmust be deleted, not renamed. With a ceiling lookup, dropping the entries below the cutoff skips the surviving records of a batch that straddles it. If the lookup becomes a floor (see the growth comment), pruning belowlowbecomes safe again. - Merging with #838 needs one decision in
nisshi-broker/tests/it/list_offsets.rs. #838 makes every backend answerNone, so once both land the newtimestamp_no_match_offsetparameter has noSome(0)callers left and can go (#838's newtursomodule calls the three-argument form). #838'sTimestamp(_) => Nonearm in the listing path also becomes unreachable after this PR's earlycontinue.
56cc234 to
1e5767c
Compare
675e8b5 to
803df26
Compare
sgamelin
left a comment
There was a problem hiding this comment.
Thanks for the rework. Every point from the last round is addressed, and each has a test:
- Read errors: a failed GET or decode now fails the request instead of being skipped (1e5767c).
- Backfill: concurrent lookups in one process share a single backfill, it logs at start and end, and a missing watermark document no longer gets created.
- Lookup: it lists only from the scan's start (
list_with_offset), applies theREAD_COMMITTEDbound, and the running-maximum guard has a test.
The sparse index with a floor lookup and a cap settles the growth question. The backfill token closes the mixed-version race. The in-flight check, where the listing has to tile the log up to high, catches a race I missed. #838 has landed, so the None answer now reaches clients as -1.
I'm requesting changes for one point: the index trusts the batch header's max_timestamp, and a complete index is never rebuilt, so one inconsistent batch affects the partition for good (inline). Two cost points are non-blocking (inline): an incomplete index re-runs the whole backfill on every lookup, and the backfill reads whole batch objects when it needs only the header.
Also:
- A failed batch read now fails the whole
ListOffsetsrequest through?. That takes down every partition in the request (list_offsets.rs#L185), where Kafka answers the error on the failing partition only. The LIST could already do this before, but this path now reads every batch. A stored batch that can't be inflated now fails every lookup whose scan crosses it. Could the error go in that partition'serror_codeinstead? - Every concurrency test uses one
DynoStore, so the commit never hits a CAS conflict, and the path that merges again against a refetched document never runs. TwoDynoStores over oneInMemorywould cover it, with the second producing between the first's collect and commit. - On #835: its
timestamps.retain(...)no longer compiles against this branch. Pruningtime_indexbelow the log start is safe now but optional, so it can just be dropped.
| // A produce that waited for a concurrent backfill would | ||
| // leave a gap: the backfill reads the index before this | ||
| // entry exists, and merges nothing in for it. | ||
| watermark.time_index.append( |
There was a problem hiding this comment.
The index takes max_timestamp from the batch header, and so do the backfill (L814) and the no-match short-circuit (L1854). The client writes that header, and nothing checks it against the records. Kafka doesn't trust it: LogValidator recomputes the maximum from the records and sets it on the batch (LogValidator.java#L282, #L415).
The index is persisted, and a complete index is never rebuilt. So a batch whose header disagrees with its records stays in effect for the life of the partition:
- Header below a record's timestamp: a lookup for a target between the two takes the short-circuit and answers no match, though the record exists. Nisshi writes such a batch itself today. For
LogAppendTime, the produce service setsbase_timestampandmax_timestampto now and keeps the client's deltas (produce.rs#L326-L337). Kafka also reports every record of aLogAppendTimebatch at the batch's max timestamp, where the scan answers base plus delta (L972). - Header above every record: the maximum never rises again, so no later batch gets an entry. Each lookup then scans from the last entry before that batch, a range that grows with the partition.
Could the maximum come from the records? One way is to recompute it in the produce service and set it on the batch, as Kafka does, which also covers the other backends. The other is for this path and the backfill to derive it from the records. A test with a header below its records' timestamps would pin it.
There was a problem hiding this comment.
Agreed that the header is the weak point, and the right place to fix it is the produce path, for every backend. #847 does that: for each CreateTime batch, the shared produce service decodes the records and rewrites max_timestamp to their actual maximum (recomputing the CRC) before anything is written, and rejects a record more than an hour ahead of the broker's clock, so a header can neither understate the records nor run far ahead of them. Once it merges, the index, the backfill and the short-circuit here all read a header that equals the records for every new batch. I'm leaving this path as is rather than adding a second derivation here. What #847 does not cover, and I'd track as a follow-up: batches stored before it merges keep the header the client wrote, so a legacy partition whose headers understate the records answers no match for those targets, and a complete index is never rebuilt; and LogAppendTime batches, where the header is the broker's clock and a positive delta puts a record above it, which also means the scan at L972 answers base plus delta where Kafka answers the batch's max timestamp. A follow-up ticket has been raised for a dynostore-side defence: derive the maximum from the records in the backfill (it reads the whole object anyway) and in the scan, and answer LogAppendTime records at the batch's max.
| // so it is in flight until a later batch is written and ages. A | ||
| // partition whose newest produce failed in between therefore lists | ||
| // again on each lookup until its next produce. | ||
| if batches > 0 && next_offset.is_some_and(|expected| expected < backfill.high) { |
There was a problem hiding this comment.
While the index is incomplete, every lookup runs the whole backfill again: a LIST from offset 0, a GET of every batch, and two CAS writes to watermark.json. Three states keep it incomplete:
- this trailing gap on a partition that gets no more produces, for example after a broker stopped between the offset CAS and the batch PUT;
- an interior gap, for up to
BATCH_SETTLED_AFTER; - produces from an older broker during a rolling upgrade.
On a large partition from before the upgrade, each offsetsForTimes costs a full read until the state clears. That also makes the CHANGELOG's "reads every stored batch of that partition once" inaccurate. Two changes would bound it:
- Read only the header. The backfill uses
batch_length,last_offset_deltaandmax_timestamp, all in the first 61 bytes, so a ranged GET saves nearly all the bytes. This needs the header's maximum to be trustworthy (see L1514). - Keep progress. The batches before the first unsettled gap already give valid entries. The commit could save them with the offset they reach, and the next backfill could list from there instead of from 0.
No test leaves a gap after the last listed batch, so deleting this check keeps the suite green.
There was a problem hiding this comment.
Agreed on the cost, and the "reads every stored batch once" claim was wrong for an index that stays incomplete. The CHANGELOG entry is gone (CONTRIBUTING now generates the changelog from the squash commit), and the PR description says what actually happens: the first lookup reads every batch, and the backfill runs again on each lookup while the index stays incomplete. The missing test is backfill_is_left_incomplete_by_a_gap_after_the_last_listed_batch: an aged listing whose last assigned offset is unwritten answers from the listed batches, leaves the index incomplete, and the lookup after the batch lands completes it; deleting the check fails it. The two bounds you propose (a ranged GET of the 61-byte header, and saving the entries before the first unsettled gap with the offset they reach so the next backfill lists from there) are deferred to a follow-up ticket that has been raised: the header GET depends on the header being trustworthy (the thread above), and saved progress is a format change I'd rather make on its own.
Replace the per-partition `timestamps` field's dead last_modified-based ListOffsets(Timestamp) matching with a real time index, keyed by each batch's max_timestamp and maintained by Kafka's own monotonic TimeIndex.maybeAppend rule, floored at NO_TIMESTAMP (-1). produce always inserts into the index under the monotonic rule, with no guard for a legacy or still-backfilling partition: deferring the insert until a concurrent backfill catches up would create a permanent index gap, since the backfill can read the index while it's still absent, merge nothing in, and never pick up that entry on either side. A new time_index_complete flag (serde-defaulted, so a pre-existing document decodes cleanly) tracks whether the index covers a partition from offset 0 onward; it's set on a brand-new partition's first produce (nothing predates offset 0) and on a successful backfill commit. A read that finds the index incomplete runs a backfill first: list every real batch object, replay the monotonic rule over their headers to get a candidate, then commit it via a CAS write that merges the candidate with whatever is live at write time. The merge is not a plain BTreeMap union - that can let a lower, more-recently inserted entry outrank an already-established higher one - but a reduction to (offset, timestamp) pairs, sorted by offset, replayed through the monotonic rule from scratch. ListOffsets(Timestamp) does a ceiling lookup in the index to find a starting offset, then a sequential scan of real batches forward from there, decoding individual records rather than trusting a batch header's max_timestamp alone (it's a ceiling hint that can be stale after compaction elsewhere in this codebase). The scan skips any record whose logical offset is below the partition's low watermark, so a batch that straddles a DeleteRecords cutoff never answers with already-deleted data. An empty index falls back to scanning from low forward rather than reporting no match, a path this PR's own delete_records (still todo!()) can't yet reach but #835's will, once it starts pruning time_index entries - that PR's current `timestamps.retain(...)` call is currently a no-op (the field was always null) but will become a real per-entry prune against this field's new semantics once both land; this PR owns the "never prune, filter on read" rule, so whichever of the two merges second should drop that retain call. The lake-sink produce branch never writes a `.batch` object, so it never touches the index either; the sequential scan already handles an empty listing by returning None rather than erroring, which this relies on. Leaves the unbounded growth of this index (GET+PUT on every produce, write amplification scaling with index size rather than batch count) as a separate follow-up ticket: there is no retention path for DynoStore at all yet, so it isn't reachable today regardless. Also updates two pre-existing shared ListOffsets integration tests (new_topic, single_record) whose Timestamp-no-match assertions depended on a shared `Some(0)` default this PR's DynoStore fix no longer produces for a genuine no-match; pg/lite/slatedb keep the old expectation via a new parameter, since their own Timestamp handling is unchanged. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01D2qdgPhMyLGZN5dLR8CVHs Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
Round 2 review follow-up: - Rewrite the empty-index fallback comment to state the actual end-state rule plainly: time_index entries are never pruned, and the per-record `low` filter in the sequential scan is what handles logical deletion - not index pruning. The prior "defensive path a future delete_records would need" framing was wrong: independent mutation testing against PR #835's `timestamps.retain(...)` logic showed pruning (even partial) produces wrong answers, both when the pruned index is still non-empty (a ceiling lookup can skip a surviving matching record whose own entry was pruned) and when it's pruned to fully empty (the `low`-based fallback excludes a straddling batch below `low` that still holds a surviving match). Also change the empty-index fallback's start_offset from `low` to `0`, which fixes the pruned-to-empty case specifically; under "never prune" the two coincide for any partition whose index can genuinely be empty, so this costs nothing in the common case. Drop "dead today" from the comment: `empty_partition_no_match` genuinely exercises this branch. - Replace the `Watermark::time_index` doc comment's reference to `dynostore::tests::schema_change` (which never actually touches `Watermark`) with a new, real regression test, `watermark_decodes_pre_time_index_documents`, that decodes three pre-time_index watermark document shapes (`timestamps: null`, key omitted, and real old data under the key) into `Watermark` and confirms `time_index_complete` defaults to `false` on all three. - Add `equal_and_mid_batch_timestamps_match_inside_a_batch`, covering the ticket's two scenarios that weren't pinned by a real test yet: an inclusive timestamp-equality boundary and a genuine mid-batch match, both inside one multi-record batch. Also coordinates with PR #835 (DeleteRecords): left a review comment on its `timestamps.retain(...)` line in nisshi-storage-dynostore/src/dynostore.rs explaining why that line must be deleted (not merely renamed) once this change lands - #835 (comment) And splits the backfill-performance deferral out of the growth follow-up ticket (which only covers steady-state unbounded index growth) into a new, separate ticket for the one-time full-GET-per-batch cost of backfilling a legacy partition's index. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01D2qdgPhMyLGZN5dLR8CVHs Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
A time index backfill or a timestamp scan that cannot read a batch now fails the request instead of skipping the batch. Only a batch deleted after the LIST is skipped. A failed backfill leaves the index incomplete, so the next lookup retries it. Also: concurrent lookups in one process share one backfill, which logs its start and end; a lookup does not create a watermark document; the scan LISTs from its start offset; READ_COMMITTED lookups stop at the last stable offset; and a test covers the produce-path monotonic guard. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_011P97dHLhTRdJ35fMJFpYPg Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
…lookup The time index lived in every partition's watermark document and gained an entry on nearly every produce, so the document that every produce rewrites and every fetch reads grew without bound. The index is now a `TimeIndex` value in the document, appended under Kafka's rule: a batch is indexed only when it raises the partition's maximum timestamp, and only once `index.interval.bytes` (4096 bytes) of batches have been appended since the last entry. The index holds at most 1024 entries; past that it drops every other entry and doubles the interval, so the document stays about 25 KiB however long the partition lives. A lookup now takes the greatest entry at or below the target (the floor) and scans forward from its offset, as `LogSegment.findOffsetByTimestamp` does, instead of the smallest entry at or above it. Every entry keeps the guarantee that each batch before it is below its timestamp, and removing an entry keeps that guarantee, so thinning the index, or pruning it below a new log start, cannot make a lookup answer wrong. A target above the greatest timestamp of every batch answers no match without a scan. The watermark document no longer carries a `timestamps` key; a document with that key, which was always `null`, decodes unchanged. A document written before the index decodes to an incomplete index that the first lookup backfills, as before. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01D8pQUJmdicavN4zsumUdaD Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
… an older broker's write A produce assigns its offset in the watermark document before it writes its batch object, so a backfill's listing can see a later batch and not an earlier one still in flight. A candidate entry for the later batch was then the floor of a lookup inside the earlier one, and the scan started past it for good. The backfill now reads the high watermark when it begins and checks that the listed batches tile the log up to it. The batches past a gap that is not settled (one below a batch written less than ten minutes ago) raise the candidate's maximum timestamp but get no entry, and the index is left incomplete, so the next lookup lists again. A binary without the index rewrites the document without it, and the batch it produced is then in neither the listing nor the live index. The backfill now sets a token in the index before its listing, and its commit marks the index complete only while that token is still there. The scan lists lazily, so a lookup stops listing at its match, and a batch key shorter than its offset no longer panics the listing. The CHANGELOG entry states that an older broker's produce discards the index, and asks for every broker sharing a bucket to be upgraded before relying on these lookups, as the group-id key change in the same section already does. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01D8pQUJmdicavN4zsumUdaD Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
The summary of `TimeIndex::merge` said that the merge returns a complete index. The code, and the last paragraph of the same comment, leave the result incomplete when the candidate is incomplete or the backfill's token is gone. The summary now states the guarantee the result keeps and defers to that paragraph for completeness. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01D8pQUJmdicavN4zsumUdaD Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
The changelog is generated from each squash commit's body when a release is cut, so a pull request no longer edits `CHANGELOG.md`. The entry's content, with its upgrade note, moves to the pull request description. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01D8pQUJmdicavN4zsumUdaD Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
The backfill leaves the index incomplete when the listed batches stop short of the high watermark, however old they are, because no later batch settles that gap. No test exercised the check, so deleting it kept the suite green. The new test lists an aged partition with the last assigned offset still unwritten: the lookup answers from the listed batches, the index stays incomplete, and the lookup after the batch lands completes it. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01D8pQUJmdicavN4zsumUdaD Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
…when one is settled A backfill takes a gap in its listing as a failed produce when the batch after it is older than `BATCH_SETTLED_AFTER`. That age was the broker's `SystemTime::now()` minus the write time the store reported, so it depended on the two clocks agreeing. The backfill now reads the store's write time of the watermark document its begin CAS just wrote, and ages each batch against that, so both times come from the store. A gap taken as settled is indexed past for good, so the backfill now logs a warning naming the gap and the age. The constant's comment states what bounds the window between a produce's offset CAS and its batch write: nothing in the broker; the object store client's retry timeout bounds the batch PUT, and a lake commit or a transaction CAS in between has no limit of its own. The test store can report every write time from a clock behind the test's own. A gap between two batches written just now stays in flight under that clock, where an age taken from the broker's clock would settle it. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01D8pQUJmdicavN4zsumUdaD Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
A `ListOffsets(Timestamp)` lookup that could not read a batch failed the whole request, so every other partition in the request got no answer. Kafka answers the error in the failing partition's `error_code` only. The lookup of one partition now lives in `timestamp_lookup`, and `list_offsets` turns its error into that partition's response: the error code of an API error, else `UNKNOWN_SERVER_ERROR`, with no offset. The other partitions answer as before. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01D8pQUJmdicavN4zsumUdaD Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
…r process Every concurrency test ran one `DynoStore`, so a backfill commit never met a conditional PUT that failed, and the path that reads the document again and merges against it never ran. The new test runs two stores over one object store: the second produces between the first's listing and its commit, the commit's PUT fails once, and the merged index holds the produced batch's entry. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01D8pQUJmdicavN4zsumUdaD Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
A timestamp lookup that could not list, read or inflate a batch answered `UNKNOWN_SERVER_ERROR`, which a Kafka client does not retry, so a transient object store failure made `offsetsForTimes` throw. Kafka answers a log it cannot read with `KAFKA_STORAGE_ERROR`, which a client retries until its own timeout. The lookup now answers that code. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01D8pQUJmdicavN4zsumUdaD Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
The scan's doc said it decodes records because a batch header's `max_timestamp` can be stale after pruning or compaction; nothing prunes or compacts a partition. It decodes records because the answer is the first record at or after the target, which a header cannot place inside a batch. The re-read of the log start after a backfill was explained by a concurrent DeleteRecords, which does not exist yet; the comment now states the requirement it keeps. The gap check names the batch it skips after a deleted batch, and the empty-listing rule names the one case it misses. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01D8pQUJmdicavN4zsumUdaD Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
803df26 to
45da6a8
Compare
|
Summary
On S3 and in-memory storage,
ListOffsetsby timestamp (a consumer'soffsetsForTimes) answered from the wrong data. It matched an object'slast_modified(write time) instead of any record's own timestamp, used>where Kafka's contract is>=, returned a batch's base offset instead of the first matching record inside it, and listed every batch of the partition on every call. It now answers the first record whose own timestamp is at or after the target, with a per-partition time index that follows Kafka's time-index semantics.What an operator sees differently
KAFKA_STORAGE_ERRORin that partition'serror_code, which a client retries; the other partitions of the request answer normally.Design
TimeIndexis kept per partition as atime_indexfield in the existingwatermark.jsondocument, so there is no new storage shape (nisshi-storage-dynostore/src/dynostore/time_index.rs). The oldtimestampskey is no longer written and is ignored on read. A document withouttime_indexdecodes to an incomplete index, and an older binary can read the new document.interval_bytesof batches have been appended since the last entry. The interval starts at 4096 bytes (Kafka's defaultindex.interval.bytes). Timestamps of -1 or below are never indexed.max_timestamp. Removing an entry only lengthens a scan, so thinning the index and pruning entries below the log start are both safe.produceassigns a batch's offset before it writes the batch object, so a listing can see a later batch and not an earlier one in flight. The backfill checks that the listed batches tile the log up to the high watermark. A gap below a batch written more than ten minutes earlier is a produce that failed in between and holds no records; the backfill logs a warning for it. Past any other gap, batches raise the candidate's maximum timestamp but get no entry, and the index is left incomplete, so the next lookup lists again. Ages are measured between two write times the store reports (the batch object's and the watermark document's), so the broker's clock does not enter into it.READ_COMMITTED, a match at or past the last stable offset is no match.Known limits, tracked separately
max_timestampfrom its header. fix(storage): rewrite a produced batch's max_timestamp and reject far-future CreateTime records #847 rewrites that header from the records for everyCreateTimebatch in the shared produce path, which covers every backend for new batches. Batches stored before fix(storage): rewrite a produced batch's max_timestamp and reject far-future CreateTime records #847, andLogAppendTimebatches, keep a header that can disagree with the records; a dynostore-side defence is a follow-up.Cross-PR coordination
delete_records) callstimestamps.retain(...)on the old field. That no longer compiles against this branch and can be dropped: pruning below the log start is safe but optional.Noneis the no-match answer on every backend, and the broker sends -1.max_timestamprewrite) is the fix for the header trust point above; this PR does not depend on it to merge.Testing
cargo nextest run -p nisshi-broker -p nisshi-storage -p nisshi-storage-dynostore --all-features --no-fail-fast -E '!test(/::pg::/)': 535 passed; the 5 failures are Postgres tests that need a server the local run did not have, which CI runs.cargo fmt --all --check,cargo clippy --workspace --all-features --all-targets -- -D warningsandcargo doc --workspace --all-features --no-deps --document-private-itemswithRUSTDOCFLAGS="-D warnings"are clean.🤖 Generated with Claude Code
https://claude.ai/code/session_01D8pQUJmdicavN4zsumUdaD