Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
31 changes: 20 additions & 11 deletions internal/etcp/dialer.go
Original file line number Diff line number Diff line change
Expand Up @@ -69,10 +69,15 @@ type Dialer struct {
DialContext(ctx context.Context, network, address string) (net.Conn, error)
}
// KeepAlive is the quiet period after which a probe is sent; after two
// quiet periods the link is declared dead. Only inbound frames count as
// proof of life, and the probe queues behind any unsent backlog, so an
// upload that takes longer than two periods to drain while the server
// sends nothing costs a reconnect (the data survives it). Zero means 5 s
// quiet periods the link is declared dead. Inbound frames count as proof
// of life, and so do our writes while a probe waits behind them: the
// probe queues behind the unsent backlog, so a slow upload while the
// server sends nothing would otherwise be taken for a dead link. Writes
// are observed in 4 KiB pieces, so an uplink too slow to send one piece
// in two periods (under about 400 B/s at the default) still costs a
// reconnect, and so does data the kernel has already accepted, which
// drains out of sight: a send buffer holding more than two periods'
// worth ahead of the probe delays its echo past the deadline. Zero means 5 s
// (upstream's maximum, src/base/Headers.hpp:180 at et-v7.0.0); Dial
// refuses a negative value or one below 100 ms. Probing starts only
// after the first WritePacket, because etserver aborts when a session's
Expand All @@ -90,13 +95,17 @@ type Dialer struct {
// bound applied separately to the two kinds of retained packets: packets
// already written to a socket are trimmed down to it, and WritePacket
// blocks while the not-yet-written backlog exceeds it (a single packet
// may take the backlog over, and probes are never blocked). While
// disconnected the ring can therefore hold about twice ReplayLimit plus
// one packet. Packets count as sent once written to the socket, and
// written packets are trimmed first, so a window smaller than what the
// kernel and the network can hold in flight turns a reconnect into
// ErrReplayExceeded; so does a catchup too large for one handshake
// message. Small values are for tests only.
// may take the backlog over, and probes are never blocked). After a
// reconnect, written packets the server has not acknowledged are kept
// until it first sends on the new link, up to twice ReplayLimit: the
// server counts our catchup only once it has decoded all of it, so a
// link lost before then must be able to replay it again. The ring can
// therefore hold about three times ReplayLimit plus one packet. Packets
// count as sent once written to the socket, and written packets are
// trimmed first, so a window smaller than what the kernel and the network
// can hold in flight turns a reconnect into ErrReplayExceeded; so does a
// catchup too large for one handshake message. Small values are for
// tests only.
ReplayLimit int
// Logger receives connection events. Nil discards them.
Logger *slog.Logger
Expand Down
2 changes: 1 addition & 1 deletion internal/etcp/drain_bench_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@ func BenchmarkDrainBacklog(b *testing.B) {
}
ctx, cancel := context.WithCancel(b.Context())
w := &countingWriter{want: entries * (4 + 18), done: cancel}
if err := c.writeLoop(ctx, w); err == nil || ctx.Err() == nil {
if err := c.writeLoop(ctx, newLink(), w); err == nil || ctx.Err() == nil {
b.Fatalf("writeLoop = %v before draining %d bytes (took %d)", err, w.want, w.n)
}
c.cancel(nil)
Expand Down
53 changes: 53 additions & 0 deletions internal/etcp/helpers_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -297,3 +297,56 @@ func (s *rawServer) open(b []byte) (protocol.Packet, error) {
}
return protocol.Packet{Header: h, Payload: plain}, nil
}

// lostCatchup completes the exchange from the client's point of view, then
// drops the link as if it died before the server counted the client's
// catchup.
func (s *rawServer) lostCatchup() {
if !s.startRecover() {
return
}
var theirs protocol.CatchupBuffer
_ = wire.ReadMessage(s.br, &theirs)
}

// recoverAll completes the exchange and returns the numbers of the packets
// in the client's catchup, or nil when the exchange fails.
func (s *rawServer) recoverAll() []int {
if !s.startRecover() {
return nil
}
var theirs protocol.CatchupBuffer
if wire.ReadMessage(s.br, &theirs) != nil {
return nil
}
nums := []int{}
for _, b := range theirs.GetBuffer() {
p, err := s.open(b)
if err != nil {
return nil
}
n, err := number(p)
if err != nil {
return nil
}
nums = append(nums, n)
}
return nums
}

// startRecover answers RETURNING_CLIENT and runs the server's half of the
// recover exchange up to reading the client's catchup, in upstream's order,
// claiming to have received nothing and sending an empty catchup.
func (s *rawServer) startRecover() bool {
if s.respond(protocol.ConnectStatus_RETURNING_CLIENT) != nil {
return false
}
if wire.WriteMessage(s.conn, &protocol.SequenceHeader{}) != nil {
return false
}
var mine protocol.SequenceHeader
if wire.ReadMessage(s.br, &mine) != nil {
return false
}
return wire.WriteMessage(s.conn, &protocol.CatchupBuffer{}) == nil
}
99 changes: 82 additions & 17 deletions internal/etcp/link.go
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,10 @@ const (
// frame more than this many into writeBufSize. Taking the whole backlog
// instead would copy it under c.mu on every pass, quadratic in its length.
maxBatchEntries = writeBufSize/(4+2+secretbox.Overhead) + 1

// progressChunk is the piece size the writer sends a batch in, so that
// a slow uplink shows progress well within one keepAlive period.
progressChunk = 4 << 10
)

var (
Expand All @@ -34,10 +38,13 @@ var (

// link is the per-TCP-connection state shared by its goroutines.
type link struct {
alive chan struct{} // cap 1: a frame arrived
alive chan struct{} // cap 1: a frame arrived, or a write progressed ahead of a waiting probe
delivering atomic.Bool // the reader is blocked handing a packet to ReadPacket
}

// newLink returns the shared state for one TCP connection's goroutines.
func newLink() *link { return &link{alive: make(chan struct{}, 1)} }

// runLink runs a reader, a writer and a liveness watcher on nc until one of
// them fails, then joins all three and returns the cause.
func (c *Conn) runLink(nc net.Conn, catchup [][]byte) error {
Expand All @@ -46,18 +53,19 @@ func (c *Conn) runLink(nc net.Conn, catchup [][]byte) error {
stop := context.AfterFunc(ctx, func() { _ = nc.Close() })
defer stop()

l := &link{alive: make(chan struct{}, 1)}
l := newLink()
var wg sync.WaitGroup
wg.Go(func() { cancel(c.readLoop(ctx, l, nc, catchup)) })
wg.Go(func() { cancel(c.writeLoop(ctx, nc)) })
wg.Go(func() { cancel(c.writeLoop(ctx, l, nc)) })
wg.Go(func() { cancel(c.watch(ctx, l)) })
wg.Wait()
_ = nc.Close()
return context.Cause(ctx)
}

// readLoop hands over packets an earlier link left pending, then delivers the
// peer's catchup, then frames from the link.
// peer's catchup, then frames from the link. The first frame ends the hold
// recover placed on our catchup (see releaseHold).
func (c *Conn) readLoop(ctx context.Context, l *link, r io.Reader, catchup [][]byte) error {
for len(c.pending) > 0 {
if err := c.handOver(ctx, l, c.pending[0]); err != nil {
Expand All @@ -73,7 +81,10 @@ func (c *Conn) readLoop(ctx context.Context, l *link, r io.Reader, catchup [][]b
}
}
br := bufio.NewReaderSize(r, readBufSize)
var buf []byte
var (
buf []byte
released bool
)
for {
frame, err := wire.ReadFrame(br, buf)
if err != nil {
Expand All @@ -84,6 +95,10 @@ func (c *Conn) readLoop(ctx context.Context, l *link, r io.Reader, catchup [][]b
}
buf = frame
signal(l.alive)
if !released {
c.releaseHold()
released = true
}
if err := c.deliver(ctx, l, frame); err != nil {
return err
}
Expand Down Expand Up @@ -138,20 +153,44 @@ func (c *Conn) handOver(ctx context.Context, l *link, p protocol.Packet) error {
}
}

// releaseHold ends the hold recover placed on our catchup and trims the ring
// back to ReplayLimit. It is called on a link's first frame: upstream writes
// on a recovered socket only after it has decoded our whole catchup, since
// BackedWriter::write waits on the recover mutex that Connection::recover
// holds until then (src/base/BackedWriter.cpp:17-18 and
// src/base/Connection.cpp:109,134-142 at et-v7.0.0).
func (c *Conn) releaseHold() {
c.mu.Lock()
defer c.mu.Unlock()
if c.ring.held {
c.ring.held = false
c.ring.trim(c.flushed, c.unsent)
}
}

// writeLoop sends ring entries from flushed onwards, trimming the ring and
// releasing blocked writers as the backlog drains. It frames a batch of
// entries (at least one, then up to writeBufSize bytes) into one reused
// buffer and sends it with a single Write, so the steady state allocates
// buffer and sends it in progressChunk pieces, so the steady state allocates
// nothing per packet. Entries are immutable once sealed, so they are framed
// outside the lock.
func (c *Conn) writeLoop(ctx context.Context, w io.Writer) error {
//
// A piece written while a probe still waits behind it signals l.alive. The
// probe's echo cannot come before the probe is sent, and once the socket's
// send buffer is full a write completes only as the peer acknowledges data,
// so a backlog draining ahead of the probe proves the path is alive while
// the server sends nothing. Other writes prove nothing (a dead link's send
// buffer still takes them), so data written with no probe behind it, the
// probe itself included, never counts.
func (c *Conn) writeLoop(ctx context.Context, l *link, w io.Writer) error {
var (
batch [][]byte
buf []byte
)
for {
c.mu.Lock()
batch = c.ring.appendRange(batch[:0], c.flushed, min(c.ring.next(), c.flushed+maxBatchEntries))
base, probe := c.flushed, c.lastProbe
batch = c.ring.appendRange(batch[:0], base, min(c.ring.next(), base+maxBatchEntries))
c.mu.Unlock()
if len(batch) == 0 {
select {
Expand All @@ -163,10 +202,14 @@ func (c *Conn) writeLoop(ctx context.Context, w io.Writer) error {
}
buf = buf[:0]
sent, n := 0, 0
probeAt := -1 // offset in buf of the probe's frame, if it is in this batch
for _, f := range batch {
if sent > 0 && len(buf)+4+len(f) > writeBufSize {
break
}
if base+int64(sent) == probe {
probeAt = len(buf)
}
var err error
// WritePacket refuses packets above wire.MaxFrameSize, so an
// error here means the ring is corrupt and no link can help.
Expand All @@ -177,13 +220,20 @@ func (c *Conn) writeLoop(ctx context.Context, w io.Writer) error {
n += len(f)
}
clear(batch) // drop references so trimmed entries can be collected
// A short count without an error breaks the io.Writer contract;
// counting the whole batch as sent would drop its tail from replay.
if m, err := w.Write(buf); err != nil || m != len(buf) {
if err == nil {
err = io.ErrShortWrite
for written := 0; written < len(buf); {
piece := buf[written:min(len(buf), written+progressChunk)]
// A short count without an error breaks the io.Writer contract;
// counting the whole batch as sent would drop its tail from replay.
if m, err := w.Write(piece); err != nil || m != len(piece) {
if err == nil {
err = io.ErrShortWrite
}
return fmt.Errorf("etcp: write: %w", err)
}
written += len(piece)
if c.probeBehind(base+int64(sent), probeAt, written) {
signal(l.alive)
}
return fmt.Errorf("etcp: write: %w", err)
}
c.mu.Lock()
c.flushed += int64(sent)
Expand All @@ -197,10 +247,25 @@ func (c *Conn) writeLoop(ctx context.Context, w io.Writer) error {
}
}

// probeBehind reports whether a probe still waits behind the first written
// bytes of the current batch: in the batch past them (probeAt is its frame's
// offset, or -1), or queued after the batch, whose entries end before
// sequence end. Only the watcher queues probes, so one queued after the batch
// was framed is always behind it.
func (c *Conn) probeBehind(end int64, probeAt, written int) bool {
if probeAt >= 0 {
return written <= probeAt
}
c.mu.Lock()
defer c.mu.Unlock()
return c.lastProbe >= end
}

// watch declares the link dead after two quiet keepAlive periods, sending the
// probe after the first. Time the reader spends blocked on a slow ReadPacket
// caller does not count as silence, and neither does the time before the
// caller's first packet, when no probe may be sent.
// probe after the first. A frame from the server, or writer progress ahead
// of a waiting probe (see writeLoop), breaks the quiet. Time the reader spends blocked on
// a slow ReadPacket caller does not count as silence, and neither does the
// time before the caller's first packet, when no probe may be sent.
func (c *Conn) watch(ctx context.Context, l *link) error {
t := time.NewTimer(c.keepAlive)
defer t.Stop()
Expand Down
4 changes: 2 additions & 2 deletions internal/etcp/link_internal_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,7 @@ func TestWriteLoopShortWrite(t *testing.T) {
// more data forever.
ctx, cancel := context.WithTimeout(t.Context(), time.Minute)
defer cancel()
err := c.writeLoop(ctx, shortWriter{})
err := c.writeLoop(ctx, newLink(), shortWriter{})
if !errors.Is(err, io.ErrShortWrite) {
t.Fatalf("writeLoop = %v, want io.ErrShortWrite", err)
}
Expand Down Expand Up @@ -104,7 +104,7 @@ func TestWriteLoopKeepsWrittenWhileBacklogFull(t *testing.T) {
ctx, cancel := context.WithTimeout(t.Context(), time.Minute)
defer cancel()
done := make(chan error, 1)
go func() { done <- c.writeLoop(ctx, &backlogWriter{ctx: ctx, c: c}) }()
go func() { done <- c.writeLoop(ctx, newLink(), &backlogWriter{ctx: ctx, c: c}) }()
synctest.Wait() // the first batch is written, the second is stuck

c.mu.Lock()
Expand Down
Loading
Loading