Skip to content

fix: bound each ListOffsets partition read by a deadline - #827

Open
solace-aross wants to merge 7 commits into
mainfrom
fix/sol-155091-request-deadline-bound
Open

solace-aross wants to merge 7 commits into
mainfrom
fix/sol-155091-request-deadline-bound

Conversation

@solace-aross

@solace-aross solace-aross commented Oct 3, 2026 •

Copy link
Copy Markdown
Collaborator

Summary

ListOffsets made one Storage::list_offsets call per partition in the
request, with no time limit. Engines answer partitions one after another, so
a slow storage read for one partition delayed every partition after it, and
the whole request could outlast the client's own read timeout. The client
then reconnects and re-sends the request while the abandoned one keeps
reading.

Change

  • Each partition is now read with its own Storage::list_offsets call, up
    to 4 concurrently, all bounded by one deadline 5 seconds after the request
    arrives.
  • A partition still unread at the deadline answers REQUEST_TIMED_OUT,
    which the Java consumer/admin client, librdkafka, and franz-go all retry.
    5 seconds sits under the shortest client read deadline for ListOffsets
    (franz-go, 10s); the new deadline module has the full table.
  • A storage error still fails the whole request.
  • On PostgreSQL, a read dropped at the deadline asks the server to cancel its statement and keeps its
    connection and its permit until the statement has ended, whether or not the cancel request reached the
    server, then returns the connection to the pool. The connection is closed and taken out of the pool only
    when it fails or when the statement has not ended within the connection's statement_timeout plus 5
    seconds (35 seconds by default; an operator's own statement_timeout in the URL moves it), logged at
    warn with the host the cancel went to. The engine keeps at most 8 ListOffsets partition reads in flight
    across all requests, half its pool of 16, counting an abandoned read until its statement has ended, so the
    broker never holds more than its 16 server connections. Each cancelled statement logs ERROR: canceling statement due to user request on the server. The other engines have only the per-request limit of 4.
  • nisshi_storage_read_deadline_exceeded has a stage attribute: queued is a read that was waiting for
    one of the 8 permits, for a pooled connection, or for the libSQL semaphore or request channel; reading is
    a read whose storage call was slow; not_started is a read the deadline passed before it began. A request
    with timed-out partitions logs one warning with the count per stage and a sample of the partitions.

The Storage trait is unchanged.

Test plan

  • New unit tests in nisshi-storage/src/service/deadline.rs and
    service/list_offsets.rs cover: answers returned in request order,
    a slow partition timing out alone without blocking others, the
    deadline already passed before a read starts, a storage error
    failing the whole request, and the counter under each stage.
  • New tests in nisshi-storage-sql/src/pg.rs cover the Postgres read limit without a server, the cancel and
    connection return against a real server (pg_sleep(30), dropped, then pg_stat_activity polled until the
    backend is idle and the connection is back in the pool), three tests against a real server with the cancel
    pointed at a closed port (statement outlives its cancel, statement still running at the bound, connection
    terminated by the server), the statement_timeout options parser, and the bound derived from the URL.
  • A read waiting for the libSQL semaphore or request channel counts as queued.
  • many_partitions_with_latency: 64 partitions with injected latency
    on every engine (in-memory, libSQL, SlateDB, Postgres).
  • Existing list_offsets integration tests pass across all four
    backends.
  • cargo clippy --workspace --all-features --all-targets -- -D warnings
    clean.
  • cargo fmt --all --check clean.

🤖 Generated with Claude Code

https://claude.ai/code/session_01D8pQUJmdicavN4zsumUdaD

@sgamelin sgamelin left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for this. Putting a time limit on ListOffsets is the right call, and the unit tests are tight: paused time, the not-started-after-the-deadline guard and request order are all pinned down.

I checked the client side, and REQUEST_TIMED_OUT is a safe code to send. The Java consumer retries only the failed partitions (OffsetFetcherUtils.java#L170-L173), the admin client retries it as a RetriableException (ListOffsetsHandler.java#L177), and librdkafka retries it (rdkafka_request.c#L1045). The timed-out row (offset and timestamp -1) also matches what Kafka 4.0 sends from its remote list-offsets purgatory.

I'm requesting changes for two things. Both are about what happens to a read after the deadline drops it.

  1. Merge after #836, or include its fix. In libSQL mode=mpsc, one task runs all storage. That task exits the first time it can't deliver a response (service.rs#L1011-L1013): tx.send fails because the caller dropped its receiver, and the ? ends the loop. After that, every storage call fails with UnableToSend until the broker restarts. Before this PR, that took a client disconnecting mid-request. Now any ListOffsets read slower than 5s does it. #836 fixes exactly this, so landing it first is enough.
  2. On PostgreSQL, a dropped read keeps its statement running on a connection that the pool treats as idle. Details inline.

The inline comments also cover the cost of one engine call per partition, logging, and a note on the Kafka difference.

Comment thread nisshi-storage/src/service/list_offsets.rs
Comment thread nisshi-storage/src/service/list_offsets.rs
Comment thread nisshi-storage/src/service/list_offsets.rs Outdated
Comment thread nisshi-storage/src/service/deadline.rs
@solace-aross

Copy link
Copy Markdown
Collaborator Author

On the first summary point: this PR now depends on #836 merging first, since #836 is what keeps the mpsc supervisor alive when a caller drops its receiver. Rebased onto main and updated the CHANGELOG entry in 8174f88.

@solace-aross

Copy link
Copy Markdown
Collaborator Author

Rebased on main. The Postgres list_offsets body now lives in list_offsets_on with main's no-row default for Timestamp lookups, so timestamp_no_match still holds. #836's supervisor fix is on main, so the merge-order note no longer applies.

@sgamelin sgamelin left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for turning this round so quickly. What I asked for is in: #836 is on main, the warning and the counter test look good (inverting the is_none() check now fails counts_each_abandoned_read), and the 64-partition test shows where the per-request limit sits. I'm happy to leave the trait change and a configurable deadline for later.

I'm still requesting changes, for two problems with the new limits rather than the old ones. Details inline.

  1. The shared limit of 8 applies to every backend, and the wait for it counts against the deadline. I suggested a limit shared by all ListOffsets requests but didn't say where it should live; it belongs with the Postgres pool. As a static in nisshi-storage, it lets a few partitions that always read slower than 5s hold every permit through their retries. Then every ListOffsets on the broker answers REQUEST_TIMED_OUT, on S3 and memory too.
  2. Each abandoned read now opens two new server connections, and its backend no longer counts against the pool. This load arrives when Postgres is already slow.

The third inline comment suggests a test that sends the cancel to a real server.

Comment thread nisshi-storage/src/service/list_offsets.rs Outdated
Comment thread nisshi-storage-sql/src/pg.rs Outdated
Comment thread nisshi-storage-sql/src/pg.rs Outdated
@solace-aross
solace-aross force-pushed the fix/sol-155091-request-deadline-bound branch 2 times, most recently from daac0ff to b30a3a7 Compare October 8, 2026 23:15
@solace-aross

Copy link
Copy Markdown
Collaborator Author

Rebased on main and ready for another look. The mechanism you reviewed on daac0ff is unchanged. Three things differ since then:

  • The rustdoc on cancel_and_return now says what happens when a cancel fails or outlasts the grace period: Object::take drops the client, but tokio-postgres only sends Terminate once the pending statement has answered, so that connection stays open on the server until the statement ends (up to statement_timeout), beside the pool's replacement.
  • abandoned_statement_is_cancelled_and_its_connection_returned now asserts the cancel task still holds its LIST_OFFSETS_READS permit right after the drop.
  • many_partitions_with_latency has a SlateDB leg, so it runs on all four engines.

The real-server cancel test and the Postgres leg only run in CI's Postgres matrix.

@sgamelin sgamelin left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks, this round closes most of it. Both earlier points are in. The limit lives in Postgres now, so memory, S3 and SlateDB keep only the per-request 4. The stage attribute separates queue waits from slow reads. A successful cancel now returns the connection to the pool. I checked that abandoned_statement_is_cancelled_and_its_connection_returned really runs against the server in the postgres:17 job, so it isn't silently skipped. It asserts size == 2 and available == 1, so it would catch a take on that path. Using ClosedByPeer to wait for the server's EOF, as libpq does, is a good catch.

I'm still requesting changes for one thing, and part of it is my fault: the short timeout I suggested last round leaves the gap open when a cancel misses its statement. Details inline. The other inline comments are optional.

Comment thread nisshi-storage-sql/src/pg.rs
Comment thread nisshi-storage/src/service/deadline.rs
Comment thread nisshi-storage-sql/src/pg.rs Outdated
Comment thread CHANGELOG.md Outdated
@solace-aross

Copy link
Copy Markdown
Collaborator Author

Filed the two follow-ups from the design review:

  • Give the metadata reader's DbReaderOptions an explicit retry budget that fits inside this PR's 5s deadline. Today it is unbounded in practice, so a dropped read keeps retrying against the object store after nisshi has already answered the client.
  • Once the Fetch half lands, a persistently slow partition can starve every partition after it in the same request, since Fetch reads partitions sequentially (unlike ListOffsets here). Accepted for this first cut, tracked separately.

solace-aross and others added 7 commits October 9, 2026 12:33
ListOffsetsService made one Storage::list_offsets call for every
partition in the request, with no time limit. Engines answer the
partitions one after another, so a slow storage read for one partition
delayed all the partitions after it, and the request could outlast the
client's own timeout. The client then reconnects and asks again while
the abandoned request keeps reading.

Read each partition with its own Storage::list_offsets call, up to 4 at
once, all bounded by one deadline 5 seconds after the request arrives.
A partition still unread at the deadline answers REQUEST_TIMED_OUT,
which the Java consumer and admin client, librdkafka and franz-go all
retry. 5 seconds sits under the shortest client read deadline for
ListOffsets (franz-go, 10 seconds); the table is in the new deadline
module. A storage error still fails the whole request.

The concurrency limit keeps one request from taking a large share of a
database connection pool: a read dropped at the deadline can keep its
pooled connection busy until the server finishes the statement.

The Storage trait is unchanged.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
Capitalize and punctuate a comment left as a lowercase fragment.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Signed-off-by: Andrea Ross <168456375+solace-aross@users.noreply.github.com>
…ions

On PostgreSQL, a ListOffsets read dropped at the deadline now asks the
server to cancel its statement and returns its connection to the pool
once the statement has ended, so the pool never hands out a connection
still busy with it and never holds more than its 16 server connections.
The cancel request waits for the server to close the cancel connection,
as libpq does, so a statement sent on the pooled connection afterwards
cannot be the one cancelled. The connection is closed instead when the
cancel request fails or the statement does not end within a short grace
period, and that is logged at warn.

The Postgres engine keeps at most 8 ListOffsets partition reads in
flight across all requests, half its pool, counting an abandoned read
until its statement has ended. The service keeps only its per-request
limit of 4, so the memory, object store and SlateDB engines have no
shared queue.

A request with timed-out partitions logs one warning with the deadline,
the elapsed time and a sample of the partitions. The deadline counter
gains a stage attribute (not_started, queued, reading), so a partition
that waited for a connection is told from one whose read was slow.

Tests: the deadline counter by stage, the Postgres read limit without a
server, the cancel against a real server, and a broker test that reads
64 partitions through LatencyIntroducingStorage in both isolation levels
on every engine.

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>
Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01D8pQUJmdicavN4zsumUdaD
…l task keeps its permit

`many_partitions_with_latency` now has a SlateDB leg, so the 64-partition
read runs on all four engines at the same cost.

`abandoned_statement_is_cancelled_and_its_connection_returned` asserts,
right after the connection is dropped, that the cancel task still holds
its `LIST_OFFSETS_READS` permit. Before, a task that let the permit go at
once would have passed.

The rustdoc on `cancel_and_return` said the failure path closes the
connection and keeps the pool to `POOL_MAX_SIZE` server connections.
`Object::take` drops the client, but tokio-postgres sends Terminate only
once the pending statement has answered, so that connection stays open on
the server until the statement ends, beside the pool's replacement. The
doc now says so.

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>
…until it ends

A ListOffsets read dropped at the deadline on PostgreSQL sent a cancel
request, and then gave the connection two seconds to answer an empty
query before taking it out of the pool and releasing its permit. A
cancel request that misses its statement (sent to a host other than the
one running the statement, or dropped by a proxy) therefore freed a
permit every two seconds while its backend ran on, so the broker could
have many more backends open than the pool's 16 and the 8 permits
promised to bound.

The cancel task now keeps the connection and the permit until the empty
query answers, whether or not the cancel request reached the server, so
a statement that outlives its cancel counts against the 8 and shows as
reads queued for a permit rather than as extra server connections. The
connection is closed and taken out of the pool only when that query
fails, or at a bound past the connection's statement_timeout, where the
server has ended the statement itself. The bound is read from the
connection options, so an operator's own statement_timeout in the URL
moves it; the default is 30 seconds plus a 5 second margin.

deadpool-postgres aborts a dropped connection's task at once, so closing
a taken connection does not end its statement: the backend notices the
closed socket only when it next sends a result. The rustdoc said the
client sends Terminate after the last statement answers, which does not
hold; it now says what happens.

The cancel task is spawned in the request's span, so its warnings carry
the peer, the correlation id and the partition, and each warning names
the host the cancel request went to.

Tests: the options parser, the bound derived from the URL, and three
tests against a real server with the cancel requests pointed at a closed
port. A statement that outlives its cancel keeps its connection and
permit until a monitor connection cancels it through the server; a
statement still running at the bound loses its connection and returns
its permit while its backend runs on; a connection terminated by the
server is taken out of the pool.

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>
…nnel as queued

On libSQL in semaphore mode, every storage call waits for the proxy's
single permit, and with 4 reads per request most ListOffsets reads that
miss the deadline do so in that wait. In mpsc mode the wait is for room
on the request channel. Both counted as stage=reading, which the
counter's documentation says points at a slow engine.

The proxy's acquire for list_offsets and the channel service's reserve
now run under deadline::queued, so a read abandoned in either wait is
counted as queued. The stage is a task-local of the task that called
the Storage method, so queued's rustdoc now says that a wait on a task
the engine spawns, such as the task serving the request channel in mpsc
mode, is not seen.

Tests: a read waiting behind a stalled read for the proxy's permit, and
a read waiting for a channel whose only slot is held, each miss as
queued. Both fail without the change.

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>
…e squash commit

CONTRIBUTING.md now says not to edit CHANGELOG.md in a pull request:
the release generates it from the squash commits, whose bodies are the
pull request descriptions. The entry's content moves to the pull
request description, corrected on two points: stage=queued also counts
the wait for a pooled connection, so it does not always point at the
limit of 8, and each cancelled statement logs an error on the server.

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>
@solace-aross
solace-aross force-pushed the fix/sol-155091-request-deadline-bound branch from b30a3a7 to 9ba1fc1 Compare October 9, 2026 17:23

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants