Repository navigation
fix: ListOffsets by timestamp should answer offset -1, not 0, when no record matches - #838
Conversation
|
Fixed in c5933d8. You were right: Pushed Thanks for reproducing it directly, that made this quick to nail down. |
sgamelin
left a comment
There was a problem hiding this comment.
Thanks, good catch, and the fix matches Kafka. I checked it against 3.9.1: when fetchOffsetForTimestamp finds nothing, handleListOffsetRequest answers buildErrorResponse(Errors.NONE, partition), which sends offset -1 and timestamp -1 (KafkaApis.scala#L1171-L1177, #L1214-L1226). The Java client drops an entry whose offset is -1 and builds an OffsetAndTimestamp from any other, and that constructor rejects a negative timestamp (OffsetFetcherUtils.java#L135, OffsetAndTimestamp.java#L34-L39). So the old offset 0 really did make offsetsForTimes() throw. Every backend with a no-match fallback is covered, null included, and ListOffsetsService is the only production code that reads ListOffsetResponse.offset, so nothing else in the tree relied on the old 0. Reverting the service default, or any of the four tested backends' None, fails at least one test that runs in CI.
Nothing here blocks. The inline comments cover the details. Three points that don't fit on a line:
Merge order with #839 and #840. This PR sets the convention: a no-match is offset: None, sent as -1. #839 still returns Some(0) on both of its no-match paths, and #840 parametrizes new_topic/single_record over Some(0) or None per backend, so whichever lands second has to reconcile them. There is also a behaviour change to weigh. On slatedb and dynostore, today's timestamp lookup can report no match when one exists: slatedb keys its index by each batch's base_timestamp (storage.rs#L900-L903), and dynostore compares the object's last_modified with a strict > (dynostore.rs#L1307-L1311). Today a false miss answers 0 and the consumer replays the partition. After this PR it answers -1 and the consumer starts at the end, so it skips records. Both answers are wrong, but skipping is worse for an at-least-once consumer. Landing this in the same release as #839 and #840 avoids shipping that window. I'm not asking for a change here, only flagging the ordering.
CHANGELOG. Clients see this change on the wire, so it deserves an entry under [Unreleased] (the advisory changelog job flags it too). For example:
Fixed
- ListOffsets by timestamp now answers offset -1 and timestamp -1 with error NONE when no record has a timestamp at or after the target, as Apache Kafka does. The broker answered offset 0 before. That made the Java consumer's
offsetsForTimes()throwIllegalArgumentException: Invalid negative timestamp, and it sent a client that seeks to the returned offset back to the start of the partition. A partition answered with an error code now also carries offset -1 instead of 0.
Tests (optional). The case in the description, a non-empty partition with a target after its last record, is checked only at the storage layer in single_record. At the wire, every -1 assertion is on an empty partition. A fourth lookup in multiple_record at third + 1s, asserting that partition 0 gets offset -1 and timestamp -1, would pin it on all four backends. single_record's after lookup could also assert that timestamp is None. The null backend change has no test, and the new turso module is #[ignore]d like the existing ones, so limbo.rs compiles but nothing runs it. The description doesn't mention the turso module either.
5dd4e84 to
984735d
Compare
|
@sgamelin following up on the points in your review body:
Rebased on main and pushed as 984735d. |
984735d to
2cbd40a
Compare
|
Small follow-up push (2cbd40a): the null backend test's tokio dev-dependency needed its Cargo.lock line, which cargo-deny caught. No code change beyond that. |
2cbd40a to
4b19ef3
Compare
Kafka's real broker answers a timestamp-based ListOffsets lookup with offset=-1/timestamp=-1 (UNKNOWN_OFFSET/UNKNOWN_TIMESTAMP) when no record matches, via KafkaApis.scala's buildErrorResponse. Nisshi was answering offset=0/timestamp=None instead, which the Java client's offsetsForTimes() treats as a valid hit (it only filters out entries where offset == UNKNOWN_OFFSET), so callers would then crash constructing OffsetAndTimestamp from a negative timestamp. Earliest/Latest on an empty partition still correctly return offset 0; the bug was specific to Timestamp lookups with no match. Fixed across all six storage backends (pg, lite, limbo, dynostore, slatedb, null) and corrected the service-layer default from 0 to -1. Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com> 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>
`after` was captured with SystemTime::now() immediately after produce returned. The record's stored timestamp is captured just before produce as integer milliseconds (to_timestamp(&SystemTime::now())), so when both captures land in the same millisecond, lite's `r.timestamp >= $4` lookup legitimately matches at offset 0 instead of returning no match. The new L752 assertion (Some(0) -> None) made this previously-invisible race observable, failing ~5-10% of the time under concurrent test load (0/30 isolated). Push `after` a full second past SystemTime::now() so it's deterministically later than any record's millisecond timestamp, regardless of scheduling delay between capturing the record's timestamp and this one. 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>
The offset=-1 fallback comment in list_offsets.rs named a specific backend file and the current paths that leave offset: None on it. That's a stale-the-day-it's-wrong remote comment (.claude/rules/comments.md): nothing in the named file would notify an editor that this comment exists, and the enumeration goes wrong the moment another backend's behavior changes. State the durable requirement instead. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
Document on ListOffsetResponse::offset when an answer has no offset and that ListOffsetsService sends it as -1, and shrink each backend's copy of the rule to one line that says what the code does. Pin the Kafka citation to the 3.9.1 tag, note the leader-epoch difference (Kafka sends -1, nisshi 0), trim the test comment, and add a CHANGELOG entry. 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>
…backend multiple_record now looks up one second past its last record and asserts every partition, including the one holding records, answers offset -1 and timestamp -1 with no error. single_record also asserts the no-match answer carries no timestamp, and the null backend gets a unit test for its Earliest/Latest/Timestamp answers. Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com> Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
4b19ef3 to
aba4286
Compare
sgamelin
left a comment
There was a problem hiding this comment.
Thanks for working through all of it. I re-reviewed at aba4286 by range-diff against 4ea4dcb, my last review. The first three commits are unchanged, and the last push is a clean rebase onto main.
Earlier points:
- The no-offset rule: addressed.
ListOffsetResponse::offsetnow states it,Storage::list_offsetspoints to it, and each backend's comment says what its own code does. I checked every backend against that doc: Earliest and Latest with no error always answerSome, so the rule holds. - The "empty partition" comment: addressed. It now describes what the code does, and the retention case is tracked in #864.
- Leader epoch: addressed. The difference from Kafka is now noted next to
.leader_epoch(Some(0))and in the doc comment ontimestamp_no_match. - The
single_recordcomment: addressed. - CHANGELOG: addressed.
- Tests: addressed.
multiple_recordnow checks for offset -1 and timestamp -1 past the last record of a non-empty partition, on pg, lite, in_memory and slatedb.single_recordchecks that the no-match answer has no timestamp, and the null backend has its own unit test. The description now covers the ignored turso module. - Merge order with #839 and #840: agreed in the thread, so nothing more is needed here.
I found nothing new in the latest commits. Approving.
…uce, consume, offsets and restarts (nisshi-io#902) Stacked on nisshi-io#862: review that one first. This PR adds the next set of smoke tests on top of its harness. ## What this adds Tests that run the Kafka CLI tools the way a user does, on every storage engine: - `broker`: ApiVersions names node 111, and the cluster id is the one the broker started with. - `topics`: describe, auto-create on produce, delete and re-create. - `configs`: describe, add and delete topic configs, and broker defaults. - `produce`: every `acks` setting, compression codecs, tombstones, `max.message.bytes`, idempotent and concurrent producers, `LogAppendTime`. - `consume`: reading from an offset, a group resuming where it stopped, a pattern subscription that picks up a topic created later. - `offsets`: lookups by time. - `restart`: topics, records, committed offsets, topic configs and SCRAM users survive a SIGTERM restart on PostgreSQL and SQLite, and an in-memory broker starts again empty. - `storage_url`: an unparsable or invalid storage URL stops the broker with an error that names it. The harness now has one module per Kafka tool. Each command whose output a test reads has its own `Output<Marker>` type, so a reader can only be called on its own command's output. Brokers can restart on the same storage, and tests can start isolated brokers. The CI change gives a smoke leg that fails before `just smoke` runs (setup, toolchain) a FAIL row, so it still shows in the report. `run.sh` also fails the leg if no nextest status line parses, instead of showing a green leg with no tests. `.claude/rules/smoke-tests.md` sets the rules for writing a smoke test. ## Ignored tests A test that fails because of a broker bug is ignored with the bug in its reason, so CI stays green and the fix enables it. Every reason names the issue or PR that fixes it, e.g. nisshi-io#798 for the five produce tests that read the latest offset on in-memory storage, and nisshi-io#849 for the `max-timestamp` lookup. ## Testing - `just clippy`-equivalent for `nisshi-smoke-test` (all features, and each engine feature alone), `cargo fmt`, the crate's unit tests, rustdoc with warnings denied, ShellCheck on `run.sh` and actionlint on `ci.yml` all pass. - `just smoke postgres`, `sqlite` and `memory` on this tree, rebased on main: every test that isn't ignored passes, and the shared broker passes its checks. Local runs also run the ignored tests: postgres 55 of 64 pass, sqlite 58 of 67, memory 45 of 58, and every failure is an ignored test. - The `maintenance_interval`, timestamp-lookup and `acks=0` tests were ignored for bugs that nisshi-io#850, nisshi-io#838 and nisshi-io#845 fixed, and are now enabled. `acks=0` stays ignored on in-memory storage for nisshi-io#798, like the other produce tests. - `s3` is unchanged: it still runs only `topic_lifecycle`. 🤖 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>
Summary
When a timestamp-based ListOffsets lookup finds no matching record, nisshi was answering
offset=0/timestamp=None. Kafka's real broker answersoffset=-1/timestamp=-1(UNKNOWN_OFFSET/UNKNOWN_TIMESTAMP) on this path, viaKafkaApis.scala'sbuildErrorResponse.This matters because the Java client's
offsetsForTimes()keeps any entry whereoffset != UNKNOWN_OFFSET. Nisshi'soffset=0answer slips past that filter as if it were a real hit, and the client then throws constructingOffsetAndTimestampfrom the negative timestamp.Earliest/Latest lookups on an empty partition are unaffected and still correctly return offset 0; the bug was specific to Timestamp lookups with no match.
Changes
Fixed across all six storage backends, plus the service-layer default:
nisshi-storage-sql/src/pg.rsnisshi-storage-sql/src/lite.rsnisshi-storage-sql/src/limbo.rsnisshi-storage-dynostore/src/dynostore.rsnisshi-storage-slatedb/src/storage.rsnisshi-storage-null/src/lib.rsnisshi-storage/src/service/list_offsets.rs(default corrected from0to-1)Tests:
nisshi-broker/tests/it/list_offsets.rs: wire-level checks on every backend, including a lookup past the last record of a non-empty partition. Adds atursomodule mirroring the existing ones; it is#[ignore]d like them, so it compiles but does not run in CI.nisshi-broker/tests/it/storage_list_offsets.rs: storage-layer checks for the no-match answer.nisshi-storage-null: unit test for the Earliest/Latest/Timestamp answers.Each backend now distinguishes "empty partition" (Earliest/Latest, offset 0) from "no record matched" (Timestamp with no hit, offset/timestamp
Noneso the caller applies Kafka's not-found default of-1).Test plan
list_offsetsandstorage_list_offsetssuites)just build-allcleanjust fmtcleanjust clippycleannisshi-schema::lake::berg) gap from no docker-compose services in this environment, unrelated to this change🤖 Generated with Claude Code
https://claude.ai/code/session_01D2qdgPhMyLGZN5dLR8CVHs