Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
37 commits
Select commit Hold shift + click to select a range
37de5a6
perf(wire): frame without per-packet allocations
tphakala Sep 25, 2026
fc0c160
test(etservertest): add fake etserver and in-memory network
tphakala Sep 25, 2026
5804e5b
feat(etcp): add reconnect backoff schedule
tphakala Sep 25, 2026
29b6d7f
feat(etcp): add idle-timeout conn with chunked writes
tphakala Sep 25, 2026
1e28074
feat(etcp): add replay ring of sealed packets
tphakala Sep 25, 2026
1206ca9
feat(etcp): add reconnecting encrypted packet connection
tphakala Sep 25, 2026
3629a95
feat(etcp): detect dead links with keepalive probes
tphakala Sep 25, 2026
4817a15
feat(etcp): block writes while the unsent backlog is full
tphakala Sep 25, 2026
ce3567c
test(etcp): pin long outages, slow catchups and random cuts
tphakala Sep 25, 2026
6dfc3b3
docs: add etcp and etservertest to the AGENTS.md layout table
tphakala Sep 25, 2026
d97876d
fix(etcp): send no keepalive probe before the caller's first packet
tphakala Sep 25, 2026
ab837e0
fix(etcp): keep the recover exchange alive while the peer's catchup a…
tphakala Sep 25, 2026
3713ee1
fix(etcp): never drop the backlog wakeup a blocked writer was handed
tphakala Sep 25, 2026
3793e8c
fix(etcp): refuse unusable Dialer options and oversized catchups
tphakala Sep 25, 2026
74af615
perf(etcp): bound each writer pass instead of copying the whole backlog
tphakala Sep 25, 2026
59b114a
docs(etcp): state that a slow upload can cost a reconnect, and pin de…
tphakala Sep 25, 2026
b647075
test(etcp): check the property test's tail for late duplicates
tphakala Sep 25, 2026
cae7c8b
fix(etcp): report the context cause when Dial's context ends mid-hand…
tphakala Sep 25, 2026
2aeb676
test(etcp): cover the payload bound, first-dial recover and redial st…
tphakala Sep 25, 2026
6744cab
test(etcp): make the suite fail fast and cover untested guards
tphakala Sep 25, 2026
580550f
fix: tighten three guards and correct docs and upstream citations
tphakala Sep 25, 2026
c37f439
fix(etcp): end the Conn on an oversized or undecodable handshake message
tphakala Sep 25, 2026
404dacf
test(etcp): fix a flaky fake-server test and make the late-replay che…
tphakala Sep 25, 2026
b9dac3f
fix(etcp): report the dial-step context cause and correct overstated …
tphakala Sep 25, 2026
f2e9f91
fix(etcp): report a bad recover message instead of the stuck write's …
tphakala Sep 25, 2026
5e6df35
test(etcp): pin the definitive-answer rule and tidy the last comments
tphakala Sep 25, 2026
0a0eef5
fix(etcp): let a bad recover message win over a racing connection reset
tphakala Sep 25, 2026
1bc05d0
test(etcp): pin a bad catchup after our half of recover completes
tphakala Sep 25, 2026
3211aac
fix(etcp): trim the replay ring when recover advances flushed
tphakala Sep 25, 2026
3bb1262
test(etcp): pin recover's trim boundary and harden the flapping test
tphakala Sep 25, 2026
2bb9391
test(etcp): pin recover's backlog accounting and writer wakeup
tphakala Sep 25, 2026
aba76d0
fix(etcp): apply the replay limit to written packets alone
tphakala Sep 25, 2026
8ccd298
fix(etcp): let a dead link go while the caller is not reading
tphakala Sep 25, 2026
373e1a9
fix(etcp): return a pending packet when the Conn ends before a new link
tphakala Sep 25, 2026
86d2e82
fix(etcp): refuse a write that reaches the queue after the Conn ended
tphakala Sep 25, 2026
6977802
fix(etcp): never quote the server's rejection text in errors
tphakala Sep 25, 2026
0c380b8
fix(etcp): bound probes behind a stuck writer and oversized catchup e…
tphakala Sep 25, 2026
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
16 changes: 9 additions & 7 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -13,13 +13,15 @@ Module path: `github.com/tphakala/et-go`. Go version: see `go.mod`.

## Layout

| Path | Role |
| ------------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| `cmd/et` | Entry point and flag parsing. |
| `internal/protocol` | Wire messages generated from upstream's `.proto` files (opaque API), plus `Header`, `Version` and `Packet`. Regenerate with `go generate ./internal/protocol` (needs protoc 3.21.12). |
| `internal/seal` | One direction of the libsodium-compatible encrypted stream: secretbox with a counter nonce. |
| `internal/wire` | Handshake message framing, stream frame framing, packet layout and size limits. |
| `rules/` | ruleguard matchers used by golangci-lint (build tag `ruleguard`). |
| Path | Role |
| ----------------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| `cmd/et` | Entry point and flag parsing. |
| `internal/etcp` | Reliable, ordered, encrypted packet connection over replaceable TCP links: replay ring, recover exchange, liveness probes, reconnect backoff, write backpressure. |
| `internal/etservertest` | Test-only fake etserver (written independently from upstream semantics) and an in-memory `net.Pipe` network with cut and refuse controls, for synctest-driven etcp tests. |
| `internal/protocol` | Wire messages generated from upstream's `.proto` files (opaque API), plus `Header`, `Version` and `Packet`. Regenerate with `go generate ./internal/protocol` (needs protoc 3.21.12). |
| `internal/seal` | One direction of the libsodium-compatible encrypted stream: secretbox with a counter nonce. |
| `internal/wire` | Handshake message framing, stream frame framing, packet layout and size limits. |
| `rules/` | ruleguard matchers used by golangci-lint (build tag `ruleguard`). |

Update this table when a package is added.

Expand Down
39 changes: 39 additions & 0 deletions internal/etcp/backoff.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,39 @@
package etcp

import (
"math/rand/v2"
"time"
)

const (
backoffBase = 250 * time.Millisecond
backoffMax = 5 * time.Second
backoffReset = 30 * time.Second // a link that lived this long resets the schedule
)

// backoff yields reconnect delays: none before the first attempt, then 250 ms
// doubling to a 5 s cap, each with +/-20% jitter, never above the cap.
type backoff struct {
attempt int
jitter func() float64 // returns [0, 1); nil means math/rand/v2
}

func (b *backoff) next() time.Duration {
n := b.attempt
b.attempt++
if n == 0 {
return 0
}
// The shift is clamped because backoffBase<<36 overflows int64 and would
// yield a negative delay, a hot redial loop; 5 is the first shift at
// which backoffBase exceeds backoffMax, so the clamp never lowers a delay.
d := min(backoffBase<<min(n-1, 5), backoffMax)
j := rand.Float64
if b.jitter != nil {
j = b.jitter
}
d = time.Duration(float64(d) * (0.8 + 0.4*j()))
return min(d, backoffMax)
}

func (b *backoff) reset() { b.attempt = 0 }
54 changes: 54 additions & 0 deletions internal/etcp/backoff_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,54 @@
package etcp

import (
"testing"
"time"
)

func TestBackoffSchedule(t *testing.T) {
b := backoff{jitter: func() float64 { return 0.5 }} // factor 1.0: no jitter
want := []time.Duration{
0,
250 * time.Millisecond,
500 * time.Millisecond,
time.Second,
2 * time.Second,
4 * time.Second,
5 * time.Second,
5 * time.Second,
}
for i, w := range want {
if got := b.next(); got != w {
t.Fatalf("attempt %d: next() = %v, want %v", i, got, w)
}
}
b.reset()
if got := b.next(); got != 0 {
t.Fatalf("after reset: next() = %v, want 0", got)
}
}

func TestBackoffJitterBounds(t *testing.T) {
tests := []struct {
name string
jitter float64
second time.Duration // delay before attempt 2 (base 250 ms)
}{
{name: "low", jitter: 0, second: 200 * time.Millisecond},
{name: "high", jitter: 0.999, second: 299_900_000 * time.Nanosecond},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
b := backoff{jitter: func() float64 { return tt.jitter }}
b.next()
if got := b.next(); got != tt.second {
t.Fatalf("second delay = %v, want %v", got, tt.second)
}
for range 100 { // far past the attempt where an unclamped shift overflows
if got := b.next(); got <= 0 || got > backoffMax {
t.Fatalf("delay %v outside (0, %v]", got, backoffMax)
}
}
})
}
}
123 changes: 123 additions & 0 deletions internal/etcp/backpressure_internal_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,123 @@
package etcp

import (
"context"
"strings"
"testing"
"testing/synctest"
"time"

"github.com/tphakala/et-go/internal/protocol"
)

// blockedConn returns a Conn with no link whose backlog is over its limit, so
// every WritePacket parks until the test makes room.
func blockedConn(t *testing.T) *Conn {
t.Helper()
var d Dialer
c := d.newConn("et.example:2022", "XXXtestclient001", strings.Repeat("k", 32))
c.mu.Lock()
c.unsent = c.limit + 1
c.mu.Unlock()
return c
}

// makeRoom empties the backlog and hands out one wakeup, as a link writer
// does after a drain.
func makeRoom(c *Conn) {
c.mu.Lock()
c.unsent = 0
c.mu.Unlock()
signal(c.space)
}

func waitWrite(t *testing.T, name string, done <-chan error) {
t.Helper()
select {
case err := <-done:
if err != nil {
t.Fatalf("writer %s: %v", name, err)
}
case <-time.After(time.Minute):
t.Fatalf("writer %s still blocked a minute after room was made", name)
}
}

// One drain wakes one blocked writer, which passes the turn on, so every
// writer that fits proceeds.
func TestBlockedWritersAllProceed(t *testing.T) {
synctest.Test(t, func(t *testing.T) {
c := blockedConn(t)
defer c.cancel(nil)

doneA := make(chan error, 1)
doneB := make(chan error, 1)
go func() { doneA <- c.WritePacket(t.Context(), protocol.Packet{}) }()
synctest.Wait()
go func() { doneB <- c.WritePacket(t.Context(), protocol.Packet{}) }()
synctest.Wait()

makeRoom(c)
waitWrite(t, "A", doneA)
waitWrite(t, "B", doneB)
})
}

// A writer that is handed the wakeup and whose context is cancelled before it
// runs must not drop the wakeup: the other blocked writer would then stay
// parked with room available. Channel receivers are served in order, so the
// wakeup goes to A, the first writer to park.
func TestWakeupNotLostOnCancel(t *testing.T) {
for range 20 {
synctest.Test(t, func(t *testing.T) {
c := blockedConn(t)
defer c.cancel(nil)

ctxA, cancelA := context.WithCancel(t.Context())
doneA := make(chan error, 1)
doneB := make(chan error, 1)
go func() { doneA <- c.WritePacket(ctxA, protocol.Packet{}) }()
synctest.Wait()
go func() { doneB <- c.WritePacket(t.Context(), protocol.Packet{}) }()
synctest.Wait()

makeRoom(c) // A is handed the wakeup
cancelA() // and is cancelled before it gets to run
select { // A may enqueue or return its cause; either is fine
case <-doneA:
case <-time.After(time.Minute):
t.Fatal("writer A still blocked a minute after it was cancelled")
}
waitWrite(t, "B", doneB)
})
}
}

// A write that reaches the queue lock after the Conn ended is refused, not
// queued on the dead Conn. The test holds the lock so the writer parks on
// it, then ends the Conn. It uses real time because synctest cannot wait
// for a goroutine blocked on a mutex; a writer that has not reached the lock
// in time sees the ended Conn anyway, so the test cannot fail spuriously.
func TestWritePacketRefusedOnceConnEnded(t *testing.T) {
var d Dialer
c := d.newConn("et.example:2022", "XXXtestclient001", strings.Repeat("k", 32))
c.mu.Lock()
done := make(chan error, 1)
go func() { done <- c.WritePacket(t.Context(), protocol.Packet{}) }()
time.Sleep(50 * time.Millisecond) // let the writer park on c.mu
c.cancel(errClosed)
c.mu.Unlock()
select {
case err := <-done:
if err == nil {
t.Fatal("WritePacket = nil after the Conn ended, want its cause")
}
case <-time.After(time.Minute):
t.Fatal("WritePacket did not return")
}
c.mu.Lock()
defer c.mu.Unlock()
if n := c.ring.next(); n != 0 {
t.Fatalf("%d packets queued on the ended Conn, want 0", n)
}
}
91 changes: 91 additions & 0 deletions internal/etcp/backpressure_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,91 @@
package etcp_test

import (
"context"
"errors"
"net"
"testing"
"testing/synctest"
"time"

"github.com/tphakala/et-go/internal/etcp"
)

// A writer blocked on a full backlog is released by its own context, with
// that context's cause, and by Close, with net.ErrClosed. An already
// cancelled context is refused without queueing.
func TestBackpressureWaitEnds(t *testing.T) {
synctest.Test(t, func(t *testing.T) {
h := newHarness(t, etcp.Dialer{ReplayLimit: 1 << 10})
defer h.close()

synctest.Wait()
h.net.SetRefuse(true)
h.net.CutAll()
synctest.Wait()
errCause := errors.New("caller gave up")
ctx, cancel := context.WithCancelCause(t.Context())
cancel(errCause)
if err := h.conn.WritePacket(ctx, numbered(0, 10)); !errors.Is(err, errCause) {
t.Fatalf("write with a cancelled context = %v, want %v", err, errCause)
}
if err := h.conn.WritePacket(t.Context(), numbered(0, 2048)); err != nil {
t.Fatalf("first write: %v", err)
}

ctx, cancel = context.WithCancelCause(t.Context())
done := make(chan error, 1)
go func() { done <- h.conn.WritePacket(ctx, numbered(1, 10)) }()
synctest.Wait()
cancel(errCause)
if err := within(t, done); !errors.Is(err, errCause) {
t.Fatalf("blocked write after cancel = %v, want %v", err, errCause)
}

go func() { done <- h.conn.WritePacket(t.Context(), numbered(1, 10)) }()
synctest.Wait()
if err := h.conn.Close(); err != nil {
t.Fatalf("Close: %v", err)
}
if err := within(t, done); !errors.Is(err, net.ErrClosed) {
t.Fatalf("blocked write after Close = %v, want net.ErrClosed", err)
}
})
}

func TestBackpressureWhileDisconnected(t *testing.T) {
synctest.Test(t, func(t *testing.T) {
h := newHarness(t, etcp.Dialer{ReplayLimit: 4 << 10})
defer h.close()

synctest.Wait()
h.net.SetRefuse(true)
h.net.CutAll()
synctest.Wait() // the dead link is noticed and the supervisor is backing off

// Each packet seals to 1024+2+16 = 1042 bytes. Writes are admitted
// while the unsent backlog is at most 4 KiB, so the fourth takes it
// past the limit and the fifth must wait.
for i := range 4 {
if err := h.conn.WritePacket(t.Context(), numbered(i, 1024)); err != nil {
t.Fatalf("WritePacket %d: %v", i, err)
}
}
done := make(chan error, 1)
go func() { done <- h.conn.WritePacket(t.Context(), numbered(4, 1024)) }()
synctest.Sleep(time.Minute)
select {
case err := <-done:
t.Fatalf("fifth write returned %v while the backlog was full", err)
default:
}

h.net.SetRefuse(false)
if err := within(t, done); err != nil {
t.Fatalf("fifth write after reconnect: %v", err)
}
if err := expectNumbered(t.Context(), 5, h.srv.Recv); err != nil {
t.Fatal(err)
}
})
}
Loading
Loading