2.8 KiB
ADR-GRPC-005: One writer per stream, a bounded queue, and resume that refuses to guess
- Status: accepted
- Date: 2026-08-30
- Scope:
:grpc:grpc-policy
Context
StreamObserver is not thread-safe, and the failure when two producers call onNext concurrently is
not an exception — it is interleaved bytes, which a client decodes as a corrupt message or, worse, as
a valid one it should never have received.
Two further properties of server streams are easy to get wrong in ways that look healthy. A consumer that falls behind either terminates the stream or silently loses messages, and the second leaves a client with a stream that appears fine and is missing changes. And a reconnect either continues from a position the server can still replay, or skips whatever is no longer there.
Decision
A bounded queue drained by one writer. GrpcSerializedStreamWriter accepts messages from any
thread and hands them to the transport only from flush, which is synchronized. write returning
ACCEPTED means queued, and the name is deliberately not sent: the transport call returns as soon
as bytes are handed over, so no method here can honestly report delivery.
Both a message bound and a byte bound. Either alone is unbounded in the other dimension.
GrpcFlowControlPolicy also takes the transport's own readiness signal, because a writer that relies
only on its queue bound produces as fast as it can allocate.
Termination is the default for a slow consumer. GrpcSlowConsumerPolicy.DROP_OLDEST exists for
feeds whose business meaning tolerates loss, and is not the default, because a client cannot detect
dropped messages: the sequence numbers it sees are the ones it was sent.
Resume is refused rather than faked. GrpcStreamGapDetector requires a signed, unexpired token
whose caller and filter fingerprints match the current request, refuses one whose snapshot version
moved, and returns FULL_RESYNC_REQUIRED when the cursor predates retained history. GrpcResumeToken
carries a key id so the signing key can rotate without invalidating every outstanding token.
Consequences
A stream carries an envelope, not a bare payload. GrpcStreamEnvelope holds the stream id,
generation, sequence, snapshot version and resume token, because resume, gap detection and drain all
need a position and a generation.
Four clocks, not one. GrpcStreamLifetimePolicy separates setup deadline, idle timeout, max
duration and heartbeat interval, and refuses combinations where one can never fire. Merging any pair
produces a familiar bug: an idle timeout used as a max duration kills healthy busy streams.
A heartbeat is a liveness signal and nothing else. It is not an application acknowledgement and
not an ordering guarantee; GrpcStreamHeartbeat says so in the place somebody would otherwise reuse
it.