Repository navigation
v2.9.0
Minor release. Operations now wait for a reconnect that is already in flight instead of failing at once, which also lets a synchronous worker drive a reconnect the heartbeat started. Plus the outcome of several review rounds over drain(), drainSubscription(), SlowConsumerPolicy::Error and the read path, most of them cases where a message could be lost, a reconnect could stall, or a process could freeze.
Upgrading is recommended for anyone running with reconnects on (the default), anyone using drain() or drainSubscription(), SlowConsumerPolicy::Error, or a processIncoming() loop next to other fibers. Read the upgrade notes first: an operation that waits for a reconnect fails with a different exception when the wait runs out.
Upgrade notes
- An operation whose wait for a reconnect runs out fails with a
TimeoutException(for exampleSubscribe to "orders" timed out waiting for the connection to be re-established), where it used to fail at once withConnection is not open.TimeoutExceptionis not aConnectionException, so acatch (ConnectionException)aroundsubscribe(),flush(),rtt(),request(),requestWithHeaders()orrequestMany()no longer catches that case. processIncoming()andreadIncoming()without aCancellationnow wait for a reconnect in flight to end, and for another fiber's socket read to end, instead of returning or failing at once. Pass aCancellationto bound them.SubscriptionQueue::fetch(), andnext()without a timeout, returnnullduring a reconnect instead of throwing;next()andfetchAll()with a timeout wait for the reconnect within it.drain()now emitsConnectionEvent::Closed, so a listener that reconnects on everyClosednow also reconnects after a gracefuldrain().- Under
SlowConsumerPolicy::Error, an overflow that an operation's read runs into is reported through the error listener, logged at error level, instead of failing that operation. Register anerrorListener, or setslowConsumerErrorsFailOperations: trueto have operations fail as before.
Added
[feature]NatsOptions::$waitForReconnect(defaulttrue). While a reconnect is in flight, requests,subscribe(),flush(),rtt(),drain(),drainSubscription()and the reads (processIncoming(),readIncoming(),SubscriptionQueuepolling) wait for it within their own timeout, then run on the new connection. Setfalseto fail fast again.[feature]SlowConsumerException, aConnectionExceptionthat names the subscription whose queue overflowed, so a subscriber that fell behind can be told apart from a failed connection.[feature]NatsOptions::$slowConsumerErrorsFailOperations(defaultfalse) restores the previous behavior, where an overflow of any subscription failed whichever operation's read ran into it.
Fixed, the highlights
Reconnects:
- A synchronous application could leave a reconnect stalled until the process restarted: a reconnect only advances while something waits on the event loop, and every operation failed at once. Operations now wait for it, and a publish buffered during a reconnect yields one event-loop tick.
- Under
SlowConsumerPolicy::Error, a subscriber that could not keep up made the connection close for good on its next reconnect: the first message after the re-subscription overflowed its queue and failed every attempt. The overflow is now reported and the reconnect carries on. - A request, fetch or
processIncoming()whose read failed while another fiber was already reconnecting waited for the whole reconnect, ignoring its own timeout. It now gives up at its own deadline. - A request made by a handler of a message delivered right after a reconnect timed out, and
drain()called from aReconnectedlistener waited out its whole budget: the read that started the reconnect held the socket read until it was over. disconnect()during a reconnect now cuts its backoff short and wins over aconnect()that is still dialling, and a failure is tied to the connection it happened on, so it cannot tear down a connection the application has since reopened.
Freezes and lost messages:
- A
processIncoming()orreadIncoming()loop could freeze the process at 100% CPU while another fiber held the socket read (a request waiting for its reply, the heartbeat): the call returned 0 at once without ever letting the event loop run. It now waits for that read. drain()during a reconnect threw without closing anything, so the reconnect could reopen the connection being shut down; it now waits for the reconnect within its budget, bounds its backlog pass by the same budget, and reports what it discards.drainSubscription()could leave the subscription alive on the server, dropped messages already received while reconnecting or duringdrain(), lost messages still in flight when a handler threw inside its flush, and could hang on a stalled socket. All fixed.- Under
SlowConsumerPolicy::Error, an operation no longer fails with another subscription's overflow: arequest()failed although its reply had arrived, and a fetch discarded what it had collected. ASubscriptionQueue's polling buffer follows the same rules, and a read delivers everything else before it throws. - Key/Value
keys()andhistory()could return a cut-short list as if complete (under the defaultDropOldest) when one read brought more records than the pending limit. SubscriptionQueue::fetchAll()failing on its own overflow lost the messages it had taken;subscribeQueue()could throw while replaying early messages and leave the subscription feeding a queue nobody held.
Close and drain lifecycle, and reporting:
drain()emits oneClosedevent once it is over, and every close is announced once.connect()is refused while a drain runs, because its teardown would close a connection opened meanwhile.- What a read reports about its frames (a drop policy's drop, a recoverable
-ERR, a malformed asyncINFO, an overflow) is reported once the whole chunk is queued, so an error listener that reads on the connection no longer reorders messages. A fatal-ERRoutranks an overflow met earlier in the same read. - A logger that throws no longer breaks a cleanup, a reconnect or a read, and no longer keeps a report from the error listener.
- A handler that threw during the heartbeat's own read was swallowed without a trace; it is now reported, and delivery goes on.
Quality gates
PHPStan level 8, 2304 unit tests (with data sets), 149 live integration tests, 48 Behat scenarios, 45 runnable examples executed against a live server, 99.29% unit statement coverage (95% unit and 97% combined floors enforced in CI) and ~94.5% Infection covered MSI (90% floor enforced in CI). Every fix has a test that fails when the fix alone is reverted. The full per-change record is in the CHANGELOG.