Skip to content

Commit fcfa2e4

Browse files
drpcstream: gate data-frame sends on the per-stream send window
Wire the sendWindow credit gate into rawWriteLocked: a KindMessage frame now acquires len(frame) bytes of per-stream send credit before it is handed to the writer. Control frames (invoke/metadata) bypass the gate, and write stats are recorded only after credit is acquired, next to the WriteFrame they describe. The window is opt-in: a stream has no send window by default, so data writes stay ungated (unlimited) and behavior is unchanged until one is installed. terminate closes the window with the send-side error (sigs.send is first-wins, holding io.EOF when a cancel/error path pre-set it), so a send parked on credit returns the same error as one parked in WriteFrame or a later send. Close and SendError now close the send window before taking the write lock, mirroring SendCancel's signal-first ordering. Without this, both would block on the write lock behind a send parked on credit and never reach the terminate that frees it. SendError pre-closes with io.EOF to match the send signal it sets moments later; Close pre-closes with termClosed for the same reason. CloseSend is a graceful half-close, so it waits for a parked send to complete rather than aborting it -- but it must not hold s.mu while waiting for the write lock: the parked writer is only freed by terminate, and every terminate path (Cancel/Close/SendError/an inbound terminal frame) needs s.mu first. CloseSend now releases s.mu before waiting, then re-locks and re-checks state before sending KindCloseSend. Per-stream only; the connection-level window and the enablement path that installs the window come in later commits. Co-Authored-By: roachdev-claude <roachdev-claude-bot@cockroachlabs.com>
1 parent 4d08d0b commit fcfa2e4

2 files changed

Lines changed: 313 additions & 0 deletions

File tree

Lines changed: 269 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,269 @@
1+
// Copyright (C) 2026 Cockroach Labs.
2+
// See LICENSE for copying information.
3+
4+
package drpcstream
5+
6+
import (
7+
"context"
8+
"errors"
9+
"io"
10+
"testing"
11+
"time"
12+
13+
"github.com/zeebo/assert"
14+
"github.com/zeebo/errs"
15+
16+
"storj.io/drpc/drpcwire"
17+
)
18+
19+
// newGateStream builds a stream writing to io.Discard with an explicit
20+
// SplitSize so small payloads are a single frame.
21+
func newGateStream(t *testing.T) *Stream {
22+
mw := testMuxWriter(t)
23+
return NewWithOptions(context.Background(), 1, mw, NewBufferPool(), Options{SplitSize: 64 << 10})
24+
}
25+
26+
// By default no send window is installed, so data writes are ungated
27+
// (unlimited) and behavior is unchanged.
28+
func TestStream_SendWindowDefaultUngated(t *testing.T) {
29+
st := newGateStream(t)
30+
assert.That(t, st.sendw == nil)
31+
assert.NoError(t, st.RawWrite(drpcwire.KindMessage, []byte("hello")))
32+
}
33+
34+
// With a finite send window, a data write blocks until enough credit is
35+
// granted.
36+
func TestStream_SendWindowGatesDataWrite(t *testing.T) {
37+
st := newGateStream(t)
38+
st.sendw = newSendWindow(4) // 4 bytes of credit
39+
40+
done := make(chan error, 1)
41+
go func() { done <- st.RawWrite(drpcwire.KindMessage, []byte("hello")) }() // 5 bytes > 4
42+
43+
select {
44+
case <-done:
45+
t.Fatal("data write returned before sufficient credit")
46+
case <-time.After(blockShort):
47+
}
48+
49+
st.sendw.grant(1) // 4 + 1 = 5 >= 5
50+
51+
select {
52+
case err := <-done:
53+
assert.NoError(t, err)
54+
case <-time.After(time.Second):
55+
t.Fatal("data write did not complete after grant")
56+
}
57+
}
58+
59+
// Control kinds (here, invoke) are not flow-controlled: they proceed even with
60+
// zero send credit.
61+
func TestStream_SendWindowControlKindsBypassGate(t *testing.T) {
62+
st := newGateStream(t)
63+
st.sendw = newSendWindow(0) // no credit at all
64+
65+
assert.NoError(t, st.WriteInvoke("service.Method", nil))
66+
}
67+
68+
// SendCancel preempts a send parked on credit: terminate (which closes the
69+
// window) runs before SendCancel takes the write lock, so the parked write
70+
// wakes, releases the lock, and the cancel frame goes out.
71+
func TestStream_SendWindowSendCancelPreemptsParkedWrite(t *testing.T) {
72+
st := newGateStream(t)
73+
st.sendw = newSendWindow(0) // send will park immediately
74+
75+
done := make(chan error, 1)
76+
go func() { done <- st.RawWrite(drpcwire.KindMessage, []byte("data")) }()
77+
78+
select {
79+
case <-done:
80+
t.Fatal("data write returned before cancel")
81+
case <-time.After(blockShort):
82+
}
83+
84+
assert.NoError(t, st.SendCancel(context.Canceled))
85+
86+
select {
87+
case err := <-done:
88+
// Same error as a send parked in WriteFrame or a later send would see.
89+
assert.That(t, errors.Is(err, io.EOF))
90+
case <-time.After(time.Second):
91+
t.Fatal("parked data write was not preempted by SendCancel")
92+
}
93+
94+
// A subsequent send observes the same error as the parked one.
95+
assert.That(t, errors.Is(st.RawWrite(drpcwire.KindMessage, []byte("more")), io.EOF))
96+
}
97+
98+
// Close preempts a send parked on credit the same way: it closes the send
99+
// window before taking the write lock rather than waiting on a grant that
100+
// may never come.
101+
func TestStream_SendWindowClosePreemptsParkedWrite(t *testing.T) {
102+
st := newGateStream(t)
103+
st.sendw = newSendWindow(0) // send will park immediately
104+
105+
done := make(chan error, 1)
106+
go func() { done <- st.RawWrite(drpcwire.KindMessage, []byte("data")) }()
107+
108+
select {
109+
case <-done:
110+
t.Fatal("data write returned before close")
111+
case <-time.After(blockShort):
112+
}
113+
114+
assert.NoError(t, st.Close())
115+
116+
select {
117+
case err := <-done:
118+
// Same error later sends see: terminate sets sigs.send to termClosed.
119+
assert.That(t, errors.Is(err, termClosed))
120+
case <-time.After(time.Second):
121+
t.Fatal("parked data write was not preempted by Close")
122+
}
123+
}
124+
125+
// SendError preempts a send parked on credit, like Close: it closes the send
126+
// window before taking the write lock so reporting an error is never stuck
127+
// behind a slow consumer.
128+
func TestStream_SendWindowSendErrorPreemptsParkedWrite(t *testing.T) {
129+
st := newGateStream(t)
130+
st.sendw = newSendWindow(0) // send will park immediately
131+
132+
done := make(chan error, 1)
133+
go func() { done <- st.RawWrite(drpcwire.KindMessage, []byte("data")) }()
134+
135+
select {
136+
case <-done:
137+
t.Fatal("data write returned before error")
138+
case <-time.After(blockShort):
139+
}
140+
141+
assert.NoError(t, st.SendError(errs.New("boom")))
142+
143+
select {
144+
case err := <-done:
145+
// io.EOF, matching sigs.send: parked and later sends agree.
146+
assert.That(t, errors.Is(err, io.EOF))
147+
case <-time.After(time.Second):
148+
t.Fatal("parked data write was not preempted by SendError")
149+
}
150+
}
151+
152+
// CloseSend is a graceful half-close, not a termination: it must NOT preempt
153+
// a send parked on credit. It waits for the write lock; once credit arrives
154+
// the parked write completes successfully and CloseSend follows it out.
155+
func TestStream_SendWindowCloseSendWaitsForParkedWrite(t *testing.T) {
156+
st := newGateStream(t)
157+
st.sendw = newSendWindow(0) // send will park immediately
158+
159+
write := make(chan error, 1)
160+
go func() { write <- st.RawWrite(drpcwire.KindMessage, []byte("data")) }()
161+
162+
select {
163+
case <-write:
164+
t.Fatal("data write returned before credit")
165+
case <-time.After(blockShort):
166+
}
167+
168+
closeSend := make(chan error, 1)
169+
go func() { closeSend <- st.CloseSend() }()
170+
171+
// Neither may make progress yet: the write is parked on credit and
172+
// CloseSend is parked behind it on the write lock.
173+
select {
174+
case <-write:
175+
t.Fatal("data write returned without credit")
176+
case <-closeSend:
177+
t.Fatal("CloseSend preempted a parked data write")
178+
case <-time.After(blockShort):
179+
}
180+
181+
st.sendw.grant(uint64(len("data")))
182+
183+
select {
184+
case err := <-write:
185+
assert.NoError(t, err) // the parked write completed, not aborted
186+
case <-time.After(time.Second):
187+
t.Fatal("parked data write did not complete after grant")
188+
}
189+
select {
190+
case err := <-closeSend:
191+
assert.NoError(t, err)
192+
case <-time.After(time.Second):
193+
t.Fatal("CloseSend did not complete after the parked write finished")
194+
}
195+
}
196+
197+
// While CloseSend waits behind a credit-parked send, termination must still
198+
// be able to proceed: Cancel needs s.mu to terminate, so CloseSend must not
199+
// hold s.mu while waiting for the write lock.
200+
func TestStream_SendWindowCancelUnwedgesCloseSendBehindParkedWrite(t *testing.T) {
201+
st := newGateStream(t)
202+
st.sendw = newSendWindow(0) // send will park immediately
203+
204+
write := make(chan error, 1)
205+
go func() { write <- st.RawWrite(drpcwire.KindMessage, []byte("data")) }()
206+
207+
select {
208+
case <-write:
209+
t.Fatal("data write returned before credit")
210+
case <-time.After(blockShort):
211+
}
212+
213+
closeSend := make(chan error, 1)
214+
go func() { closeSend <- st.CloseSend() }()
215+
216+
select {
217+
case <-closeSend:
218+
t.Fatal("CloseSend preempted a parked data write")
219+
case <-time.After(blockShort):
220+
}
221+
222+
// Cancel terminates the stream, closing the send window: the parked write
223+
// wakes with an error and CloseSend unblocks as a no-op.
224+
canceled := make(chan struct{})
225+
go func() { st.Cancel(errs.New("boom")); close(canceled) }()
226+
227+
select {
228+
case <-canceled:
229+
case <-time.After(time.Second):
230+
t.Fatal("Cancel blocked behind CloseSend waiting for the write lock")
231+
}
232+
select {
233+
case err := <-write:
234+
assert.That(t, errors.Is(err, io.EOF))
235+
case <-time.After(time.Second):
236+
t.Fatal("parked data write did not wake on termination")
237+
}
238+
select {
239+
case err := <-closeSend:
240+
assert.NoError(t, err)
241+
case <-time.After(time.Second):
242+
t.Fatal("CloseSend did not unblock after termination")
243+
}
244+
}
245+
246+
// Terminating the stream wakes a send parked on credit.
247+
func TestStream_SendWindowTerminateWakesParkedWrite(t *testing.T) {
248+
st := newGateStream(t)
249+
st.sendw = newSendWindow(0) // send will park immediately
250+
251+
done := make(chan error, 1)
252+
go func() { done <- st.RawWrite(drpcwire.KindMessage, []byte("data")) }()
253+
254+
select {
255+
case <-done:
256+
t.Fatal("data write returned before termination")
257+
case <-time.After(blockShort):
258+
}
259+
260+
st.Cancel(errs.New("boom"))
261+
262+
select {
263+
case err := <-done:
264+
// Cancel pre-sets sigs.send to io.EOF; the window closes with it.
265+
assert.That(t, errors.Is(err, io.EOF))
266+
case <-time.After(time.Second):
267+
t.Fatal("parked data write did not wake on termination")
268+
}
269+
}

drpcstream/stream.go

Lines changed: 44 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -55,6 +55,10 @@ type Stream struct {
5555
recvQueue ringBuffer
5656
wbuf []byte
5757

58+
// sendw is the per-stream send-side flow-control window. It is nil when
59+
// flow control is not enabled, in which case data writes are ungated.
60+
sendw *sendWindow
61+
5862
mu sync.Mutex // protects state transitions
5963
sigs struct {
6064
send drpcsignal.Signal // set when done sending messages
@@ -352,6 +356,12 @@ func (s *Stream) terminate(err error) {
352356
s.sigs.recv.Set(err)
353357
s.sigs.term.Set(err)
354358
s.recvQueue.Close(err)
359+
if s.sendw != nil {
360+
// Close with the send-side error: sigs.send is first-wins, so when a
361+
// caller pre-set it (io.EOF for cancel/error), a send parked on credit
362+
// returns the same error as one parked in WriteFrame or a later send.
363+
s.sendw.close(s.sigs.send.Err())
364+
}
355365
s.checkFinished()
356366
}
357367

@@ -402,6 +412,15 @@ func (s *Stream) rawWriteLocked(kind drpcwire.Kind, data []byte) (err error) {
402412
fr.Data, data = drpcwire.SplitData(data, n)
403413
fr.Done = len(data) == 0
404414

415+
// Only data frames consume send credit; a nil window (flow control
416+
// disabled) leaves sends ungated. acquire parks until credit arrives,
417+
// the ctx is canceled, or the window closes (stream termination).
418+
if kind == drpcwire.KindMessage && s.sendw != nil {
419+
if err := s.sendw.acquire(s.Context(), int64(len(fr.Data))); err != nil {
420+
return err
421+
}
422+
}
423+
405424
drpcopts.GetStreamStats(&s.opts.Internal).AddWritten(uint64(len(fr.Data)))
406425
s.log("SEND", fr.String)
407426

@@ -502,6 +521,13 @@ func (s *Stream) SendError(serr error) (err error) {
502521
}
503522

504523
defer s.checkFinished()
524+
525+
// Close the send window before taking the write lock so a send parked on
526+
// credit wakes and releases the lock (same ordering as Close). io.EOF to
527+
// match the sigs.send error set below, so parked and later sends agree.
528+
if s.sendw != nil {
529+
s.sendw.close(io.EOF)
530+
}
505531
s.write.Lock()
506532
defer s.write.Unlock()
507533

@@ -558,6 +584,13 @@ func (s *Stream) Close() (err error) {
558584
}
559585

560586
defer s.checkFinished()
587+
588+
// Close the send window before taking the write lock so a send parked on
589+
// credit wakes and releases the lock, instead of stalling the close on a
590+
// grant that may never come (mirrors SendCancel's signal-first ordering).
591+
if s.sendw != nil {
592+
s.sendw.close(termClosed)
593+
}
561594
s.write.Lock()
562595
defer s.write.Unlock()
563596

@@ -581,11 +614,22 @@ func (s *Stream) CloseSend() (err error) {
581614
s.mu.Unlock()
582615
return nil
583616
}
617+
s.mu.Unlock()
584618

585619
defer s.checkFinished()
620+
621+
// Wait for the write lock without holding s.mu: the writer holding it may
622+
// be parked on send credit, and every path that frees it (terminate, via
623+
// Cancel/Close/SendError or an inbound terminal frame) needs s.mu first.
586624
s.write.Lock()
587625
defer s.write.Unlock()
588626

627+
s.mu.Lock()
628+
// Re-check: sending may have ended while we waited for the write lock.
629+
if s.sigs.send.IsSet() || s.sigs.term.IsSet() {
630+
s.mu.Unlock()
631+
return nil
632+
}
589633
s.sigs.send.Set(sendClosed)
590634
s.terminateIfBothClosed()
591635
s.mu.Unlock()

0 commit comments

Comments
 (0)