Repository navigation
fix: implement DeleteRecords on lite, limbo and dynostore; fix Postgres column names - #835
solace-aross wants to merge 6 commits into
Conversation
sgamelin
left a comment
There was a problem hiding this comment.
Thanks for taking this on. Removing the todo!() panics and the Postgres query that deleted the wrong range are real fixes, and delete_records_cutoff matches Kafka 3.9.1: -1 resolves to the high watermark, any other negative offset or one above the high watermark gets OFFSET_OUT_OF_RANGE, and an offset at or below the log start is a no-op that reports the current log start. I checked it against Partition.deleteRecordsOnLeader and UnifiedLog.maybeIncrementLogStartOffset at the 3.9.1 tag.
I'm requesting changes for one thing (1). The rest are strong suggestions or optional.
1. ListOffsets(Earliest) stays below the new log start on four of the five engines (blocking)
Kafka answers EARLIEST with logStartOffset and NO_TIMESTAMP (UnifiedLog.scala#L1326-L1336). Here:
- pg, lite and limbo answer from
list_earliest_offset.sql, which takes the lowestoffset_idstill stored. After a-1delete that is the retained row athigh_watermark - 1, as the doc comment ondelete_to_high_watermark_keeps_latestalready notes. - DynoStore answers with the base offset of the oldest object (dynostore.rs#L1454-L1460). So it is below the log start after
-1(the last batch's base offset, which can be far belowhigh_watermark - 1) and also after any cutoff that falls inside a batch. Indelete_offset_inside_a_batch_keeps_the_whole_batch, Earliest would answer 0 for a log start of 2. - SlateDB reads
watermark.lowand is correct.
A consumer that resets to earliest, or a new group, therefore reads records that the DeleteRecords response reported as deleted. Once #826 lands (Fetch answers OFFSET_OUT_OF_RANGE below the log start), this becomes a loop: the fetch at high_watermark - 1 is out of range, the consumer resets to Earliest, gets high_watermark - 1 again, and repeats.
The fix is small enough for this PR: floor Earliest at the watermark's low. For SQL, join watermark in list_earliest_offset.sql, take the first record with offset_id >= w.low, and answer w.low with no timestamp when none remains. For DynoStore, floor the answer at w.low the same way. Please also assert Earliest in delete_to_high_watermark_keeps_latest and in the straddling-batch test.
Related, optional: once Earliest reads the watermark, the record kept at high_watermark - 1 exists only because the SQL Latest queries join the record at w.high - 1 (list_latest_offset_uncommitted.sql), and DynoStore's Latest reads the newest object (#798). SlateDB's Latest reads watermark.high, so protect_active_batch protects nothing there, and SlateDB's own policy_delete already removes the last batch. Answering Latest from the watermark would let the five hand-copied clamps go; if they stay, the comments should name the query that needs the retained row.
2. pg, lite and limbo delete the rows inside one request-wide transaction that holds the watermark locks (strongly suggested)
On pg, delete_records takes for no key update on each partition's watermark row in request order and holds every lock until the single commit, while an unbounded record_delete_by_offset.sql (plus the header cascade) runs for each partition. Three effects:
- Produce to every partition already processed waits on the row lock.
- The delete runs under
DEFAULT_STATEMENT_TIMEOUT_MS(30 s).policy_deleteopts out withSET LOCAL statement_timeout = 0(pg.rs#L1695-L1698) for this reason. On a large partition the delete times out, the whole request rolls back, and each retry does the same, so-1cannot succeed. - EndTxn locks watermarks in
order by t.name, tp.partition(txn_select_produced_topitions.sql), and this locks them in request order, so overlapping requests can deadlock and fail the whole request.
On lite and limbo the IMMEDIATE transaction holds the database's writer lock for the duration.
Kafka only advances logStartOffset in the request path; segment removal happens later in the retention task. Could the request commit only the watermark advance, one short transaction per partition, and leave the row removal to a bounded or background step (for example policy_delete also removing offset_id < w.low)? If it stays in the request for now, sorting the topitions before locking at least removes the deadlock.
3. Smaller points (inline)
- DynoStore: the reclaim is described as best-effort, but a listing error returns
Errafter the watermark is committed. - DynoStore and #840: the timestamp pruning is a no-op on main but conflicts with #840's design. Details inline; this decides whether merge order matters.
- SlateDB: every candidate batch is decoded inside the transaction.
- Tests: the SlateDB physical test and the limbo
-1test pass with the guard broken. - On error, Kafka answers
low_watermark = -1(ReplicaManager.scala#L1127-L1142); this PR answers 0 for an unknown topic or partition and the current log start forOFFSET_OUT_OF_RANGE. Clients act on the error code, so this is minor, but please either return -1 or add a comment at the code site explaining the difference.
4. Follow-ups
The description says the out-of-scope items are filed as follow-ups. I couldn't find GitHub issues for them (Earliest from the watermark, POLICY_VIOLATION for compact topics, broker tests for turso). Could you open them and link them from the description? The Fetch item is #826.
Merge note for #839: it replaces SlateDB's timestamps.retain with prune_time_index_below(new_low_watermark) and adds clear_latest_if_partition_empty. Whichever PR merges second should call both inside this PR's Ok(cutoff) arm with the validated cutoff (not the raw request offset), and re-run #839's delete_records tests against the end-offset deletion rule.
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 SOL-155296/#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 follow-up, filed as SOL-155343: 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. SOL-155076 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 on SOL-155076: - 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 (SOL-155296, 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 SOL-155076 lands - #835 (comment) And splits the backfill-performance deferral out of SOL-155343 (which only covers steady-state unbounded index growth) into a new, separate ticket, SOL-155344, for the one-time full-GET-per-batch cost of backfilling a legacy partition's index. SOL-155076 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>
c7cfc6f to
bd52b8d
Compare
|
Summary points, in bd52b8d after a rebase onto main:
|
…es column names
SOL-155296: DeleteRecords panicked with todo!() on libSQL (lite) and
DynoStore, and failed at SQL-prepare time on Postgres because the query
referenced record.id/record.partition, columns that don't exist (the real
columns are record.offset_id/topition, with partition living on topition).
The Postgres query also deleted the wrong range (offset_id >= cutoff
instead of < cutoff) and never persisted the new low watermark to the
watermark table, so a "successful" delete had no lasting effect on Fetch,
ListOffsets or offset_stage.
A fifth backend, limbo (turso://), had the identical todo!() and wasn't
mentioned by the ticket; it gets the same treatment as lite since it
shares the same schema and watermark helpers.
Adds nisshi_storage::delete_records_cutoff, a single pure function
encoding Kafka's DeleteRecords validation (offset -1 means "delete up to
the high watermark", any other negative or above-high-watermark offset is
OffsetOutOfRange, at-or-below the current log start is a no-op that
reports the existing log start unchanged) so all five backends -
Postgres, lite, limbo, DynoStore and SlateDB - validate identically and
can't drift from one another.
Every backend now physically deletes only up to high_watermark - 1: the
record/batch holding the last committed offset is Kafka's active segment
and must never be removed, even for the common `offset == -1` ("delete
everything") request, or ListOffsets(Latest) would regress to 0 once
nothing remains to derive it from. SlateDB already had this deletion path
but was missing the validation and the active-segment guard; DynoStore's
physical reclaim derives the same guarantee for free by never deleting the
last object in a partition's listing (its successor's base offset is what
bounds a batch's span, and the last object has no successor).
No backend returns a hard Err for a per-partition problem (unknown topic,
out-of-range partition or offset): each is reported as a per-partition
ErrorCode so a sibling partition in the same request still succeeds.
Replaces the single ignored, assertion-light test in
nisshi-broker/tests/it/delete_records.rs with the repo's generic-fn +
per-backend-module pattern, covering pg/dynostore/lite/slatedb end to end:
cutoff resolution, the -1/keep-latest-watermark case, out-of-range
rejection, the log-start no-op case, and that an unknown topic or
partition never fails a sibling partition in the same request. limbo gets
equivalent coverage as #[ignore]'d unit tests in its own test module,
consistent with every other engine-level limbo test (turso can't yet run
the full produce flow in this environment/CI).
Also fixes two pre-existing nisshi-storage-slatedb tests that produced
batches via a hand-built Batch literal with an empty record_data: harmless
before, but DeleteRecords now decodes every candidate batch to find its
end offset, and that literal doesn't decode. Replaced with a batch built
through the real encode path.
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>
…teRecords SOL-155296 follow-up from review: - Add a CHANGELOG [Unreleased] entry for the DeleteRecords fix (the repo's advisory CI check flags a commit with none). - Add delete_offset_inside_a_batch_keeps_the_whole_batch to the committed nisshi-broker/tests/it/delete_records.rs suite, run on all four backend modules (pg/in_memory/lite/slatedb). Every other scenario in this file produces single-record batches, so a cutoff always lands on a batch boundary; this one produces a single four-record batch and deletes at an offset strictly inside its range, directly exercising the end-vs- start deletion rule the previous commit's fix depends on. Confirmed by mutation: reverting DynoStore's successor-offset check, or SlateDB's end-offset check back to a start-offset check, each makes this new test fail while leaving every pre-existing test green; both pass again once reverted. - Add test_delete_records_to_high_watermark_keeps_one_physical_batch to nisshi-storage-slatedb/src/tests.rs. SlateDB's public ListOffsets(Latest) reads watermark.high directly, never by scanning batches, so no public- API test can observe whether protect_active_batch actually keeps the active batch's object in the underlying store. This test builds the Engine from a Db handle it keeps for itself so it can scan the batch-key prefix directly after a -1 DeleteRecords and assert exactly one physical batch object remains. Confirmed by mutation: hard-coding protect_active_batch's condition to false makes this test fail (1 expected, 0 found) while all 34 other nisshi-storage-slatedb tests and all 24 delete_records::slatedb::* broker tests stay green; reverted after confirming. 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>
A test-helper doc comment framed a detail about DeleteRecords' batch-span
check as a change in time ("now decodes") instead of a present fact.
Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
ListOffsets(Earliest) on pg, lite, limbo and DynoStore now floors its answer at the watermark's low, so a consumer that resets to earliest no longer reads records a DeleteRecords response reported as deleted. Also from review: - pg locks the watermark rows in topic name, partition order, matching EndTxn, so overlapping requests do not deadlock. - DynoStore logs and skips a failed reclaim listing instead of failing a request whose watermark advance has already committed, batches the object deletes, and no longer prunes the watermark timestamps. - SlateDB finds deletable batches from the next batch's base offset instead of decoding every candidate inside the transaction. - A failed partition reports low_watermark -1, as Kafka does. - Tests assert Earliest after a delete, which batch SlateDB keeps, and ListOffsets on limbo after a -1 delete. 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>
DeleteRecords on Postgres ran the whole request in one transaction that held every partition's watermark row lock while it deleted the records under the 30s statement timeout. A large `-1` delete timed out and rolled back on every retry, produce to the partitions already processed waited on the lock, and the lock order could deadlock with EndTxn. Each partition now commits its log start in a short transaction of its own, so the request holds at most one watermark lock at a time. Removing the records below the new log start follows in a separate transaction and is best effort. `maintain` removes any rows a failed removal left, keeping the record at high - 1 that Latest reads, with the statement timeout lifted as the other sweeps do. A storage failure on a partition answers KAFKA_STORAGE_ERROR for that partition instead of failing the whole request. Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com> Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
bd52b8d to
ec7a750
Compare
|
@sgamelin closing out the rest of your review body, in ec7a750 (rebased on main):
Could you take another look when you have a moment? |
The DeleteRecords change made list_earliest_offset a union of the first record at or above watermark.low and watermark.low itself, with a single outer order by and limit. Plain union dedupes, and the constant branch column hides the index order, so both Postgres and SQLite read every record at or above the log start before choosing one. Plain union all does not fix it on SQLite: EXPLAIN QUERY PLAN with scan stats on a 300,000-record partition shows 299,000 rows read and sorted either way. The first branch now takes its own order by and limit 1, the same shape as the query before this PR, and the branches are joined with union all. SQLite reads 2 record rows for the same partition. The result is unchanged: rows tied on (o, offset_id) are identical tuples, so deduplication never changed the first row. The outer order by keeps offset_id so a partition with more than one watermark row still answers with the smallest low. Co-Authored-By: Claude Sonnet 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>
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 SOL-155296/#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 follow-up, filed as SOL-155343: 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. SOL-155076 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 on SOL-155076: - 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 (SOL-155296, 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 SOL-155076 lands - #835 (comment) And splits the backfill-performance deferral out of SOL-155343 (which only covers steady-state unbounded index growth) into a new, separate ticket, SOL-155344, for the one-time full-GET-per-batch cost of backfilling a legacy partition's index. SOL-155076 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>
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 SOL-155296/#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 follow-up, filed as SOL-155343: 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. SOL-155076 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 on SOL-155076: - 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 (SOL-155296, 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 SOL-155076 lands - #835 (comment) And splits the backfill-performance deferral out of SOL-155343 (which only covers steady-state unbounded index growth) into a new, separate ticket, SOL-155344, for the one-time full-GET-per-batch cost of backfilling a legacy partition's index. SOL-155076 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>
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>
…te-records and SCRAM (nisshi-io#923) Stacked on nisshi-io#902: review that one first. ## What changes, and why? It adds smoke tests that run the Kafka CLI tools against Nisshi for the broker's maintenance and authentication: - **Retention** (PostgreSQL and SQLite): expired records deleted and newer ones kept, the next offset after retention, a topic without `cleanup.policy`, `retention.ms=-1`, a 30-day `retention.ms`, `kafka-configs --add-config` and `--delete-config retention.ms`, `retention.bytes`, offsets after retention empties a partition. - **Compaction** (PostgreSQL and SQLite): the latest value per key, a tombstone as the key's only record, and `compact,delete`. - **`kafka-delete-records`** (every engine): the new low watermark, the earliest offset, the latest offset. - **SCRAM:** users created with `kafka-configs` and described with it; login with SCRAM-SHA-256 and SCRAM-SHA-512, a wrong password, and a client without credentials, all on PostgreSQL and SQLite, because a broker with `--authentication` can't create its first user, so each test restarts the broker. That restart duplicated three SCRAM tests in `restart.rs`, which this PR moves into `scram.rs`. - **SQLite `vacuum_into`:** a broker started on the snapshot, after its own database has been deleted, has the topics and records. Harness changes: `kafka-delete-records`, `kafka-configs --describe --entity-type users`, a producer that sets old timestamps, a broker that runs maintenance every 2 seconds, and a restart onto other storage that can delete files first. `Broker::file_exists` now copies the file out with `docker cp`. It used `docker exec … test`, and the broker image has no `test` command, so it returned `false` for every file in a broker container. ### Ignored tests Each is waiting for: - nisshi-io#918: a topic without `cleanup.policy` is never cleaned - nisshi-io#919: `retention.ms=-1` deletes every record - nisshi-io#920: a `retention.ms` above 2147483647 stops retention on PostgreSQL (2 tests) - nisshi-io#921: `retention.bytes` is ignored - nisshi-io#922: `kafka-configs` can't describe users without `DescribeClientQuotas` - nisshi-io#864: offsets go back to 0 when retention empties a partition - nisshi-io#904: `kafka-configs --delete-config` refuses a topic's own config - PR nisshi-io#835, which fixes DeleteRecords: the three `delete_records` tests ## Upgrade impact None. This changes only the smoke tests. ## How was this tested? Local runs of `just smoke <engine>` with every test, the ignored ones included, and the Kafka 3.9 tools: | Engine | Tests run | Passed | Failed, all ignored | |---|---|---|---| | SQLite | 86 | 69 | 17 | | Memory | 62 | 46 | 16 | | PostgreSQL | 82 | 63 | 19 | On each engine, the failed tests are exactly the tests ignored on that engine, so CI, which skips them, passes. The failures include the ignored tests from nisshi-io#902. - [x] `cargo clippy -p nisshi-smoke-test --all-features --all-targets -- -D warnings`, and with `--features memory` - [x] `cargo fmt --check`, and rustdoc with warnings denied - [x] each commit builds and passes the crate's unit tests on its own - [ ] S3 and the Kafka 3.7 and 4.3 tools: CI's smoke jobs 🤖 Generated with [Claude Code](https://claude.com/claude-code) --------- Signed-off-by: William Kourlas <156007774+solace-wkourlas@users.noreply.github.com> Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
What does this change do, and why?
DeleteRecords panicked with
todo!()on libSQL (lite) and DynoStore, and failed at SQL-prepare time on Postgres because the query referencedrecord.id/record.partition, columns that don't exist (the real columns arerecord.offset_id/topition, with partition living ontopition). The Postgres query also deleted the wrong range (offset_id >= cutoffinstead of< cutoff) and never persisted the new low watermark, so a "successful" delete had no lasting effect on Fetch, ListOffsets oroffset_stage.A fifth backend, limbo (
turso://, a default feature), had the identicaltodo!()and wasn't mentioned by the ticket; it gets the same treatment as lite since it shares the same schema and watermark helpers.Shared validation
Adds
nisshi_storage::delete_records_cutoff, a single pure function encoding Kafka's DeleteRecords validation (confirmed against realapache/kafkatrunk source): offset-1means "delete up to the high watermark" (the common case, not an edge case); any other negative offset, or one above the high watermark, isOffsetOutOfRange; at or below the current log start is a no-op that reports the existing log start unchanged. All five backends — Postgres, lite, limbo, DynoStore and SlateDB — call this before any mutation, so none can drift from the others.Corrected deletion semantics
Every backend now physically deletes only up to
high_watermark - 1: the record/batch holding the last committed offset is Kafka's active segment and must never be removed, even for the commonoffset == -1("delete everything") request — otherwiseListOffsets(Latest)would regress to 0 until the next produce. SlateDB already had a deletion path but was missing the validation entirely and deleted by batch start offset rather than end offset (over-deleting a batch that straddles the cutoff, and deleting the active segment on the common-1call). DynoStore's physical reclaim gets the same "never delete the active segment" guarantee for free: an object is only deleted when its successor's base offset is also below the cutoff, and the last object in a partition's listing has no successor.No backend returns a hard
Errfor a per-partition problem (unknown topic, out-of-range partition or offset) — each is reported as a per-partition error code, so a sibling partition in the same request still succeeds.Earliest after a delete
ListOffsets(Earliest)answers the new log start on every engine.list_earliest_offset.sql(pg, lite, limbo) returns the first record at or abovewatermark.low, orwatermark.lowwith no timestamp when none remains, and DynoStore floors its answer at the low watermark. SlateDB already readwatermark.low.Postgres transactions
Each partition commits its new log start in a short transaction of its own, so a request holds at most one watermark row lock at a time: produce waits only while its own partition advances, and the request cannot deadlock with EndTxn. Removing the records below the log start follows in a separate transaction and is best effort. If it fails or times out,
maintainremoves the rows on its next run (record_delete_below_log_start.sql, statement timeout lifted like the other sweeps), always keeping the record athigh_watermark - 1. A storage failure on one partition answersKAFKA_STORAGE_ERRORfor that partition instead of failing the request.Until #826 lands, Fetch does not reject an offset below the log start on any engine. On Postgres that means a fetch below the new log start can return deleted records between a timed-out removal and the next
maintainrun (10 minutes by default). The window opens only on a delete large enough to time out.Follow-ups
ListOffsets(Latest)from the high watermark on every engine, which lets thehigh_watermark - 1clamps go.POLICY_VIOLATIONfor compact-only topics.OFFSET_OUT_OF_RANGE.How was this tested?
nisshi-broker/tests/it/delete_records.rs: cutoff resolution end-to-end (response,offset_stage, ListOffsets, Fetch all agree), the-1/keep-latest-watermark case, out-of-range rejection (both directions), the log-start no-op case, an unknown topic/partition never failing a sibling partition in the same request, and a batch that straddles the cutoff keeping its earlier-offset records reachable (directly exercises the end-vs-start deletion fix).#[ignore]d unit tests in its own test module, matching the pre-existing pattern for every other limbo engine-level test in this repo.maintain_removes_records_left_below_log_start(pg unit test, runs in CI) covers the Postgres sweep: rows below the log start stay untilmaintain, then go, and the record athigh_watermark - 1stays.just clippy,cargo fmt --all --check,just build-allclean.Note: the limbo (turso) DeleteRecords tests are
#[ignore]d like the other limbo engine tests, so they don't run in CI.🤖 Generated with Claude Code
https://claude.ai/code/session_011P97dHLhTRdJ35fMJFpYPg