Repository navigation
Wait in basic_publish while the connection is blocked - #76
carlhoerberg wants to merge 2 commits into
Conversation
The client recorded connection.blocked and called on_blocked, but basic_publish kept publishing. Against RabbitMQ that only stalls on a full socket, but LavinMQ rejects publishes during flow control with a channel error, so blocked publishers lost their channel. basic_publish now waits until connection.unblocked, or until the connection closes. Waiters are woken last on close, after the channels are cleaned up, so they then behave like any publish on a closed connection. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_011ob4Nr8aeNr1kj3PgMosUS
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_011ob4Nr8aeNr1kj3PgMosUS
|
@codex review |
Codex Review SummaryThis comment shows the latest Codex review activity on this pull request.
ℹ️ About Codex in GitHubYour team has set up Codex to review pull requests in this repo. Reviews are triggered when you
Codex reacts with 👀 while any review is running, comments if it has suggestions, and reacts with 👍 once all reviews finish with no findings. |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: a87592cb02
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| Log.info { "Blocked by server, reason: #{f.reason}" } | ||
| @unblocked ||= ::Channel(Nil).new | ||
| @blocked = true | ||
| @on_blocked.try &.call(f.reason) |
There was a problem hiding this comment.
Invoke the blocked callback outside the read loop
When an on_blocked callback calls basic_publish (or the newly public wait_while_blocked), it waits on @unblocked from inside the connection's read_loop. That same fiber is therefore unable to read the subsequent connection.unblocked or close frame that would release the wait, permanently deadlocking the connection. Dispatch the callback outside the read-loop fiber or otherwise ensure callbacks cannot wait on an event only that loop can process.
Useful? React with 👍 / 👎.
Why
The client already handled
connection.blocked: it logged it, setblocked?and calledon_blocked. Butbasic_publishkept publishing.406 PRECONDITION_FAILED, which closes the channel. In a test under a cgroup memory limit, lavinmqperf gotconnection.blockedand then had its channel closed about 20 ms later on the next publish.What
basic_publishnow waits while the connection is blocked:connection.blocked: the connection creates a channel that publishers wait on.connection.unblocked: it closes that channel, waking every waiting publisher at once.Connection#close) also wake waiters, so a publisher never waits on a dead connection. They wake waiters last, after the channels are cleaned up. A woken publisher then behaves like any publish on a closed connection: after a serverconnection.closeit raises, and after a transport failure the publish is dropped, as before.Connection#wait_while_blockedis public, for code that wants the same wait before doing other work.The blocked/unblocked frame handling moved into
process_blocked/process_unblocked, soread_loopisn't more complex than before.Compatibility
Callers that published while blocked now wait instead. Against RabbitMQ that changes where they wait (in
basic_publishinstead of on the socket), not whether they wait.Testing
spec/blocked_spec.cruses a minimal fake AMQP server, since a real server can't be made to sendconnection.blockedon demand. It covers:Full suite against a local LavinMQ: 94 examples, 1 failure. The failure is
Error Handling Network errors wraps connection refused, which fails the same way onmainwith Crystal 1.21 ("Operation now in progress" instead of "Connection refused"), so it's unrelated.🤖 Generated with Claude Code
https://claude.ai/code/session_011ob4Nr8aeNr1kj3PgMosUS