Repository navigation
stream, conn: split channels by direction and simplify stream lifecycle - #339
Conversation
|
|
Overall Grade |
Security Reliability Complexity Hygiene |
Code Review Summary
| Analyzer | Status | Updated (UTC) | Details |
|---|---|---|---|
| Go | Oct 3, 2026 4:24p.m. | Review ↗ | |
| Shell | Oct 3, 2026 4:24p.m. | Review ↗ |
Important
AI Review is run only on demand for your team. We're only showing results of static analysis review right now. To trigger AI Review, comment @deepsourcebot review on this thread.
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
Session-bound replies can migrate to replacement streams, and several lifecycle and error-contract issues remain unresolved.
Review effort: Balanced
Findings: 1
Open (5)
Keep handler replies bound to their originating session · New Clean up expired calls without requiring subsequent traffic · New Preserve all pending calls for retry after session loss · New Propagate invalid peer configuration errors from inbound manager · New Remove unrelated agent attribution configuration · New
What changed in this PR
Refactors Gorums’ stream and connection layers around directional channels and per-stream sessions, with lifecycle, dispatch, testing, and documentation updates.
Changes:
- Introduces outbound, inbound, and local channels with session-scoped pending calls and dispatch queues.
- Improves stream draining, cancellation, stalled-send reporting, graceful shutdown, and local test servers.
- Expands generator validation, tests, and API documentation.
| File | Description |
|---|---|
stream_dedup_inflight_test.go |
Tests dropped shared-stream calls. |
stream_dedup_dispatch_test.go |
Tests deduplicated dispatch and nested calls. |
server.go |
Updates server construction and shutdown lifecycle. |
server_testing.go |
Clarifies test server interface documentation. |
server_test.go |
Tests usable peer-change configurations. |
server_handler.go |
Documents handler ordering and release semantics. |
server_graceful_stop_test.go |
Tests graceful shutdown behavior. |
server_e2e_test.go |
Updates nested-call documentation. |
README.md |
Removes obsolete Ansible requirement. |
Makefile |
Corrects release version path. |
local_servers.go |
Adds listener injection and error cleanup. |
local_servers_test.go |
Tests local-server errors and listener closure. |
internal/testutils/servers/servers.go |
Clarifies server stop signaling. |
internal/testutils/servers/integration.go |
Adds integration listener factory. |
internal/testutils/servers/bufconn.go |
Adds concurrent-safe bufconn listeners. |
internal/testutils/servers/bufconn_race_test.go |
Tests bufconn registry concurrency. |
internal/tests/oneway/oneway_test.go |
Updates test lifecycle commentary. |
internal/testprotos/failing/reservedrpc/reserved.proto |
Adds reserved RPC fixture. |
internal/testprotos/failing/reservedenum/reserved.proto |
Adds reserved enum fixture. |
internal/testprotos/failing_test.go |
Expands reserved-name generator tests. |
internal/stream/transport.go |
Refactors transport around channel interfaces. |
internal/stream/testhelpers.go |
Updates stream test channel helpers. |
internal/stream/session.go |
Implements per-stream session lifecycle. |
internal/stream/session_test.go |
Tests session routing, draining, and health. |
internal/stream/server.go |
Delegates accepted streams to inbound channels. |
internal/stream/send_queue.go |
Adds closable bounded send queues. |
internal/stream/send_queue_test.go |
Tests send queue closure and backpressure. |
internal/stream/router.go |
Removes the former message router. |
internal/stream/request.go |
Defines requests and stream errors. |
internal/stream/request_test.go |
Tests cancellation-safe error delivery. |
internal/stream/pending.go |
Adds session-local pending-call tracking. |
internal/stream/pending_test.go |
Tests pending-call expiry and draining. |
internal/stream/outbound.go |
Adds reconnecting outbound channels. |
internal/stream/local.go |
Adds in-process local channels. |
internal/stream/latency.go |
Extracts latency estimation. |
internal/stream/inbound.go |
Adds accepted-stream inbound channels. |
internal/stream/gorums_message.go |
Exports direct payload message construction. |
internal/stream/doc.go |
Documents the directional channel architecture. |
internal/stream/dispatch.go |
Adds ordered handler dispatching. |
internal/stream/dispatch_test.go |
Tests dispatch ordering and back-channel routing. |
internal/stream/channel_goaway_test.go |
Tests GOAWAY draining and replacement. |
internal/stream/channel_cancel_test.go |
Tests cancellation during stalled sends. |
internal/impl/quorumcall.go |
Updates lazy dispatch documentation. |
internal/impl/call.go |
Preserves buffered one-way confirmations. |
internal/impl/call_context.go |
Updates dispatch documentation. |
internal/impl/call_context_test.go |
Adapts tests to channel-based transports. |
internal/conn/outbound_manager.go |
Configures new outbound channel options. |
internal/conn/node.go |
Refactors nodes for directional channels. |
internal/conn/node_trysend_test.go |
Removes obsolete router-based tests. |
internal/conn/node_test.go |
Updates node lifecycle and transport tests. |
internal/conn/node_source.go |
Updates implementation documentation. |
internal/conn/inbound_manager.go |
Creates and manages inbound channels. |
internal/conn/inbound_manager_test.go |
Tests inbound registration and cleanup. |
internal/conn/dial_opts.go |
Updates dial option semantics. |
inbound_manager_test.go |
Adapts public inbound-manager tests. |
gorumstest/servers.go |
Clarifies default test server behavior. |
gorumstest/servers_test.go |
Tests nil server factory handling. |
gorumstest/options.go |
Renames and documents test dial options. |
gorumstest/gorumstest.go |
Uses bufconn-backed local server groups. |
errors.go |
Exposes stalled-send errors. |
doc/release-guide.md |
Corrects release version path. |
dial_options.go |
Updates send-buffer documentation. |
config_test.go |
Renames configuration tests and helpers. |
cmd/protoc-gen-gorums/gengorums/gorums.go |
Expands reserved identifier validation. |
cmd/protoc-gen-gorums/gengorums/gorums_dev.go |
Passes plugin context to validation. |
call_test.go |
Updates interceptor test commentary. |
call_async_test.go |
Tests one-way confirmations after cancellation. |
benchkit/harness_test.go |
Uses the unreachable configuration helper. |
backchannel_dispatch_test.go |
Tests nested back-channel calls. |
AGENTS.md |
Documents permitted internal forwarding wrappers. |
.claude/settings.json |
Adds agent attribution settings. |
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
8d3884a to
85e29c9
Compare
85e29c9 to
600fb9b
Compare
600fb9b to
e536c66
Compare
f747906 to
75ab74f
Compare
ensureConnectedNodeStream replaced the stream whenever the connection left Ready, and the pending calls on the cancelled stream were never requeued. Leave a stream in place until the receive path clears it.
pauseReconnect drained streamReady before waiting, so a signal for a stream the sender had just created was discarded and the receiver slept the full backoff while that stream went unread.
GracefulStop only stopped the gRPC server, so WaitForPeers kept blocking and the peer configuration's sender and receiver kept redialing. Stop the inbound manager and close the peer configuration after in-flight RPCs finish.
A full send queue does not block replies. Two-way requests fail fast, one-way requests wait, and a reply with no response channel is dropped. TestTrySendDoesNotBlockOnFullQueue and TestChannelTrySendDoesNotBlockOnFullQueue already cover that.
The receive loop waited until a handler released before reading the next frame, so a nested call to the peer that sent the request never saw its reply. Requests now enter a per-stream queue drained by one dispatcher; replies are delivered on the reader immediately.
A stream kept open across a server GOAWAY held the drain open, so GracefulStop never returned and MaxConnectionAge never moved the channel to a new connection. Mark the stream draining when its connection leaves Ready and clear it once no call is pending, so the next send opens a stream on a new connection.
ensureConnectedNodeStream created a stream context and, when NodeStream failed, left it live until the next attempt overwrote it. With eager reconnect against a peer that stays down, one child context of the connection leaked per backoff tick. Cancel it when creation fails.
The cancel watcher cleared the stream to unblock a Send stalled by flow control. An inbound channel cannot open a replacement stream, so the clear severed it for good while ConnectedPeers still listed the peer: under dedup a borrower call with a short deadline stopped all traffic on the shared stream. Arm the watcher on outbound channels only.
A one-way call collected its send confirmations with a select over the reply channel and the context. When both were ready Go chose at random, so Wait on an Async handle often returned the context error for sends that had completed. Take a buffered confirmation first and fall back to the context only when none is available.
…ends Closing an inbound channel cancelled its pending calls with ErrNodeClosed, so a borrower call in flight when a shared stream dropped did not match the documented ErrStreamDown, and a caller that retries on it gave up. An inbound channel closes only with its stream, so report the stream drop.
NewLocalServers panicked on an invalid peer configuration, such as missing dial options, and leaked the listeners it had allocated. Its stop function also left a preallocated listener open once a server had been started on another listener with Serve. Return the error, stop the servers built so far, and close every preallocated listener.
SetPeerConfig swaps the inbound view's own node for the peer configuration's node, but change detection compared node IDs only, so the WithPeerChange snapshot delivered during NewServer held a node with no outbound manager and Add on it panicked. Compare the nodes themselves so the callback receives the usable configuration.
Servers documents that a nil srvFn selects a default server, but it passed the nil function on and panicked. Substitute DefaultServer.
serverOptions held the WithPeers node source twice, as peerNodes and outboundNodes, and both were always set together; it is now one field. The outbound stream-transition hook had three names (onStreamChange, StreamState, peerStreamChanged) and passed a node ID and state that its only consumer ignored, since the connected-peer rebuild reads the state from the nodes. The hook is now a plain func() named onStreamChange at every layer, and Channel.StreamUp reports the new state.
The Configuration type is now Config, so tests named TestConfiguration* no longer follow the TestFileNameFeatureName convention. Rename them to TestConfig* in the root package and internal/conn; no test changes behavior.
The field holds the gorums.DialOption values passed to a helper, and gorums has no manager type, so name it after the type it holds.
The forwarding-function rule reads as forbidding the gorumsimpl call constructors and the root MapRequest, MapResponse, NewConfig, and WithNodeList wrappers. They exist because generic functions cannot be aliased and a function variable would be reassignable, so record them as an allowed exception.
Several snippets no longer compiled or failed at run time against the current API: - quorum-call functions return *gorums.Call, so helpers that take *gorums.Responses are now passed the call's Responses field, and the guide says so; - ResponseSeq methods such as IgnoreErrors are reached through Results(), not directly on *gorums.Responses; - MapRequest and MapResponse cannot infer all type arguments, so the snippets spell them out, and one-way calls use *emptypb.Empty; - MetadataEntry_builder is not exported, so the metadata interceptor uses the MetadataEntry setters; - WithPeers and NewLocalServers need dial options with transport credentials, or NewServer panics and NewLocalServers fails.
The dev guide's alias block named Configuration, but generated code declares Config. The release guide pointed at the root version.go, which moved to runtime/gorumsimpl/version.go. The README listed Ansible as a requirement for a benchmark script that no longer uses it; the sweep command deploys over ssh.
prepare-release told the releaser to edit the root version.go, which moved to runtime/gorumsimpl/version.go.
The sender armed a context.AfterFunc per request that cleared the stream when the request's context ended during Send. The stream is shared by every caller of the node, so one caller's short deadline tore it down for all of them, and a late watcher could still clear a healthy stream after Send returned. A cancelled caller already returns on its own context; the request is sent if the peer resumes reading, and its late reply is discarded. Detecting a peer that stops reading is left to gRPC keepalive, which the user guide now recommends. This removes the per-request AfterFunc registration, sendGuard, and cancelInflightSend, and the tests that pinned the watcher's decision logic. The cancelled-send test now covers outbound channels too.
A two-way request could pass the closed check in trySend, then land in the send channel after the sender had already drained it and exited, leaving the request queued forever. A caller with a background context never returned. sendQueue guards pushes with a read lock and close with the write lock, so every push either sees the queue closed or completes before close drains it. The sender also closes the queue when the channel's context ends without Close. This replaces the two-stage selects in Enqueue and trySend and drainSendQ with one push method.
Handler ordering had three implementations: a channel-and-runner requestDispatch per outbound channel, a second one per server NodeStream, and the router's dispatchMu for local channels. Inbound channels also started a runner they never used. dispatcher is a mutex-guarded queue that starts the next handler from the running handler's release, or when it returns, so it needs no goroutine of its own and no selects on its hot path. Local channels now queue requests on it instead of blocking the caller on dispatchMu: a one-way call is confirmed once queued, a two-way call fails fast with ErrSendQueueFull when the queue is full, and a handler that calls its own node before releasing no longer deadlocks. Because the dispatcher already runs each handler in its own goroutine, the inbound routers now call the handler on that goroutine. The unreachable back-channel branch of RouteMessage is removed and its remaining delivery inlined into dispatchInbound.
One Channel type played three roles, told apart by a nil send queue or a nil connection, and several goroutines could open, clear, or replace the outbound stream. Guards against stale streams (clearStream's comparison, owner tokens on a shared router, streamReady and the backoff wait that drained it) existed only because no goroutine owned a stream's lifetime. The server side ran its own copy of the receive loop through PeerNode, peerNode, nilPeerNode, and a reply forwarder. Channel is now an interface with three implementations: - OutboundChannel: one goroutine is the only sender and the only opener of streams. Each stream is a session with its own pending calls, receive loop, and GOAWAY watcher; a lost session's calls are sent again on the next one. - InboundChannel: one session over an accepted stream. Serve runs the receive loop, and Server.NodeStream reduces to AcceptPeer and Serve. - LocalChannel: the in-process handler path. Both stream kinds route received frames through one session.handle, with the request ID space chosen by direction. MessageRouter is gone: pending calls live per session, and the latency estimate lives in the node's Transport, shared with borrower transports. Behavior changes: - A stream that receives GOAWAY stops taking new requests at once, so it retires even under sustained load; before, new calls kept it busy and the server's drain waited. - Untracked inbound streams (plain clients) queue replies through a send queue that waits for space, which keeps their old guarantee that replies are not dropped; this server sends them no requests, so waiting cannot deadlock. Tracked peers still drop replies on a full queue. - WithBufferSizes' receiveSize now sets the per-stream request dispatch queue capacity; the reply forwarder queue it used to size no longer exists. - A channel with no server records its failed first stream in LastErr without waiting for a request. Tests that pinned removed mechanisms (ensureStream, clearStream, pauseReconnect, stale receivers, the router's owner filtering, peerNode.TrySend) are deleted; behavior tests are ported, and new tests cover sessions, pending calls, retry outcomes, waiting replies, and GOAWAY under load.
Transport.Enqueue silently dropped a request when an owned transport had no channel, as a known peer's node has while the peer is disconnected. A caller holding an older connected-peer configuration then got no node error, and its call waited until its context ended. Every transport without a channel now fails the request with ErrStreamDown, as shared transports already did.
InboundChannel.Close waits for its send loop, which can be blocked in a server-side Send until the peer reads. The stream cleanup closed the channel inside detach, under the node's inbound lock and the InboundManager's lock, so one stalled peer could block AcceptPeer, onStreamChange, and WaitForPeers for every other peer. detach now only removes the channel and re-points the node's active channel; the cleanup closes the channel after releasing the locks.
- A session that ended itself no longer records the resulting cancellation as the channel's LastErr, which overwrote the real error or made a node look failing after a GOAWAY drain. - A request taken from the queue when its session ended concurrently with registration is carried to the next session, keeping its queue position. - A request held while Close interrupts stream creation fails with ErrNodeClosed, not the cancellation error. - The inbound send loop cancels the channel before closing its queue, so waiting pushes are released first. - A waiting reply dropped because its context ended is counted in DroppedReplies, as one dropped because the channel closed is. The user guide no longer claims that gRPC keepalive detects a peer that stops reading a stream; it detects lost connections only. Doc comments now state the exceptions for waiting replies and local channels, and WithBufferSizes states that receiveSize applies to inbound streams.
TestOnewayAsyncWaitAfterContextEnds slept 10 ms before cancelling the call's context, assuming the sends were confirmed by then. Under -race with packages running in parallel the sleep was sometimes too short, and Wait correctly reported the context error for an unconfirmed send. Each channel confirms a one-way send before it sends the next request, so the test now completes a two-way call to the same nodes after Async and only then cancels.
A streaming (correctable) call stays pending until its session ends, since only the caller knows when it is done, and so does a two-way call whose reply was dropped. Such calls grew the pending table on a long-lived stream and kept a stream that received GOAWAY from ever retiring, so the server's graceful stop waited until the client closed the stream. The pending table now removes calls whose context has ended once it has doubled in size since the last sweep, which bounds it to about twice the live calls at amortized constant cost. A draining session also watches its pending calls' contexts, so it ends as soon as each call has completed or its caller is done. The watch is registered only when draining starts, keeping per-request context callbacks off the send path.
NewLocalServers always listened on random localhost TCP ports. The new option supplies the listener for each server instead, so local server groups can run over other transports, such as in-memory connections in tests.
LocalServers used real TCP even in the default build, unlike the other gorumstest helpers. Each run of TestStreamDedupWaitForAllConcurrent with 50 servers opens about 1,225 localhost connections; repeated runs left enough of them in TIME_WAIT to exhaust macOS's ephemeral ports, so dials failed with "can't assign requested address" and gRPC's dial backoff pushed some peers past the test's deadline. The test had failed this way since it was added. LocalServers now uses bufconn listeners and the matching dialer, and real TCP only under the integration build tag, as the other helpers do. The test passes 100 of 100 -race runs.
A live peer that stops reading a stream, for example because a handler never returns, blocks the channel's send loop on flow control. gRPC keepalive does not notice, since the transport still answers pings, so the node looked healthy. Each channel records when its current send began, and LastErr reports the new ErrSendStalled while that send has been blocked for a second or longer, so ByLastError orders the node last. The check runs when LastErr is read; a send costs one clock read and two atomic stores.
Replace manual WaitGroup accounting around the outbound supervisor and session workers with WaitGroup.Go while preserving the existing lifecycle.
Inbound cleanup now cancels the session and closes its queue without waiting for the sender. This breaks the half-close deadlock where a flow-controlled transport send can finish only after the RPC returns.
The server reads a stream only while its request dispatch queue has room. The docs promised that a handler may call the sending peer before Release and that the call completes. That holds only until the peer fills the queue: then the reply waits behind the blocked read, and the call waits until its context ends. State the condition in the Release and WithBufferSizes comments and in the user guide, and recommend calling Release before a nested call to a peer that may send many requests meanwhile.
75ab74f to
1f1767f
Compare
The dispatcher started every handler on a new goroutine. A new goroutine gets the runtime's starting stack, which is 2 KB when the process's average stack use is low. It is low on a node that dials many peers, because each dialed ClientConn parks several goroutines with shallow stacks. The request path, HandleRequest and then unmarshaling, needs more than 2 KB, so every request grew and copied its stack, and the grown stack was freed when the goroutine exited. In cluster benchmarks with 29 nodes, this stack copying took 24 to 29 percent of each node's CPU. When a handler returns without calling release, its goroutine now takes the next queued handler, so a busy stream keeps one goroutine and the stack it has grown. A release before return still starts the next handler on a new goroutine, so a handler that releases early still runs concurrently with the next one. Handler order and the queue bound are unchanged. HandleRequest no longer defers ServerContext.Release. That release fired inside the handler run, before it returned, which made the dispatcher start every next request on a new goroutine. Both callers run HandleRequest through the dispatcher, which releases when the handler returns, as RequestHandler specifies. TestDispatcherGoroutineReuse checks both reuse cases, and TestServerHandleRequestRelease checks that HandleRequest releases only when the handler does. On the cluster, this raised throughput by 8 to 28 percent and lowered p50 latency by 8 to 32 percent, in both stream modes.
AppendToIncomingContext runs for every received request and copied the incoming metadata twice: metadata.FromIncomingContext returns a copy, and the code then copied that copy before appending the message's entries. On a server-side stream the context carries the stream's gRPC headers, so every request paid for two map copies, even when the message had no entries to add, which is the common case. In cluster benchmarks, this took up to 15 percent of the CPU on nodes that receive most of their requests on server-side streams. A message without entries now returns the context unchanged. A message with entries appends to the single copy that FromIncomingContext returns. The context's own metadata is never modified. One corner changes: a request without entries, on a context that has no incoming metadata, no longer gets an empty incoming metadata map. The callers in this repository treat a missing map and an empty one alike. TestMessageAppendToIncomingContext checks the four combinations of entries and incoming metadata, that the original metadata is left unchanged, and that a message without entries allocates nothing. On the cluster, on top of the dispatcher's goroutine reuse, this raised throughput by 6 to 22 percent and lowered p99 latency by 5 to 30 percent.
sendShared enqueued a call's request to the nodes in configuration order, which is sorted by node ID. Every caller therefore reached the highest-ID node last in every call. In dual stream mode that node waits longest for its share of each call, and under load it fell behind: in cluster benchmarks with 29 nodes and 32 workers per node, the highest-ID node ran at 0.79 of the median node's throughput for quorum calls, and the spread across nodes was twice as wide as needed. The fan-out now starts at a random node and wraps around, so each node is reached last in about one call in N. Each node still receives a caller's requests in call order, since only the order across nodes changes. The cost is one rand.IntN per call. TestCallContextSendSharedFanOutStart checks that every node receives every call and that each node is visited first in some call. On the cluster, the highest-ID node's relative throughput rose from 0.79-0.95 to 0.97-0.99 in dual mode, and the spread for quorum calls halved. Total throughput and latency did not change measurably.



Restructures the stream layer so that each stream has exactly one owner, and fixes the connection-lifecycle bugs found while reviewing the stack below.
Channels by direction.
stream.Channelis now an interface with three implementations:OutboundChannelruns over streams this node dials. One goroutine is the only sender and the only opener of streams. Each stream is asessionwith its own pending calls, receive loop, and GOAWAY watcher, and a lost session's calls are sent again on the next one.InboundChannelruns over one stream the server accepted.Server.NodeStreamasks aPeerAcceptor(theInboundManager) for the channel and runs its session; the session ends with the stream.LocalChannelserves requests to the local node in-process through the handler.Both stream kinds route received frames through one
session.handle, with the message-ID space chosen by direction.MessageRouteris gone: pending calls live per session, and the latency estimate lives in the node'sTransport, which deduplicated borrower nodes share. The send queue is a closablesendQueue, so a request is either accepted before close or rejected, never stranded. Request handlers on every channel kind run through one serialdispatcherthat starts the next handler when the running one releases or returns.Server handlers.
Handlers for requests from the same stream start one at a time, in arrival order. The server keeps reading the stream while a handler runs, so replies arrive whether or not the handler has called
ServerContext.Release, and a handler may call the peer that sent the request beforeReleasewithout deadlocking. Deduplicated replies bypass the dispatch queue.Stream lifecycle.
ErrStreamDown; a request to a node with no channel fails instead of waiting.Server.GracefulStopnow also unblocksWaitForPeers,WaitForClients, andWaitForAlland closes the peer configuration built byWithPeers.InboundManagerlock, and inbound cleanup returns without waiting for a send blocked on flow control.WithPeerChangecallback with a usable configuration.API changes.
ErrSendStalled(new):Node.LastErrreports it while a send to the node has been blocked for a second or longer, soByLastErrororders a node that stopped reading its stream last.WithBufferSizes:receiveSizenow sets the capacity of each inbound stream's dispatch queue, with a default of 4096 when zero. The server stops reading a stream while its queue is full. A reply never waits for send-queue space; it fails fast, or is dropped and counted byNode.DroppedReplies.WithLocalListeners(new) letsNewLocalServerscreate its listeners with a caller-supplied function.NewLocalServersreturns an error for an invalid peer configuration instead of panicking, and its stop function closes every preallocated listener.Config,Node,NodeContext,ConfigContext) used as the Go name of an enum or RPC method, and of a message in any file of the service's Go package, not only the service's own file.Test framework.
gorumstest.LocalServersconnects its servers in memory in the default build and over localhost TCP under theintegrationtag.gorumstest.Serversstarts the default server when given a nil server function. The bufconn dialer's address map is guarded against concurrent starts.Documentation.
The developer guide gains a runtime-architecture section with package-layering and per-layer class diagrams for
internal/streamandinternal/conn. The user guide covers the handler ordering andReleasesemantics, cancellation and stalled sends, keepalive, the dispatch queue, and the reserved names; its call, interceptor, and server snippets now compile against the current API. Doc comments across the touched packages describe current behavior, and stale references to the version file, Ansible, and old aliases are fixed.Verification:
go test ./... -count=1,go test -race ./internal/stream/... ./internal/conn/... -count=1,gofmt -l.Top of the stack, on top of #337.
Closes #205