From e3d22d256e15e797887f57dba891c1232fb02598 Mon Sep 17 00:00:00 2001 From: cabbage Date: Sun, 6 Sep 2026 06:57:47 +0000 Subject: [PATCH] feat: add client reconnect scheduling primitives --- internal/client/runtime/queue.go | 124 ++++++++++++++++++ internal/client/runtime/runtime_test.go | 94 ++++++++++++++ internal/client/runtime/session.go | 160 ++++++++++++++++++++++++ test/coverage.toml | 12 ++ 4 files changed, 390 insertions(+) create mode 100644 internal/client/runtime/queue.go create mode 100644 internal/client/runtime/runtime_test.go create mode 100644 internal/client/runtime/session.go diff --git a/internal/client/runtime/queue.go b/internal/client/runtime/queue.go new file mode 100644 index 0000000..d99d181 --- /dev/null +++ b/internal/client/runtime/queue.go @@ -0,0 +1,124 @@ +package runtime + +import ( + "errors" + "sync" +) + +var ErrSendQueueFull = errors.New("client session send queue is full") + +type FramePriority uint8 + +const ( + PriorityControl FramePriority = iota + 1 // Ping, Pong, Close, and session fencing. + PriorityEssential // Lifecycle, acknowledgement, signal, and loss markers. + PriorityData // Replayed output and ordinary durable data. +) + +type Frame struct { + Priority FramePriority + Payload []byte +} + +type QueueOptions struct { + MaxBytes uint64 + ControlReserveBytes uint64 + EssentialReserveBytes uint64 + ControlBurst uint32 +} + +func (options QueueOptions) Validate() error { + if options.MaxBytes == 0 || options.ControlReserveBytes == 0 || options.EssentialReserveBytes == 0 || options.ControlReserveBytes+options.EssentialReserveBytes >= options.MaxBytes || options.ControlBurst == 0 { + return errors.New("invalid client send queue options") + } + return nil +} + +// SendQueue has no fallback goroutine or unbounded channel. Output cannot use +// bytes reserved for essential state, and essential state cannot use the final +// control reserve needed for heartbeat/close traffic. +type SendQueue struct { + mu sync.Mutex + options QueueOptions + bytes uint64 + controlRun uint32 + control []Frame + essential []Frame + data []Frame +} + +func NewSendQueue(options QueueOptions) (*SendQueue, error) { + if err := options.Validate(); err != nil { + return nil, err + } + return &SendQueue{options: options}, nil +} + +func (queue *SendQueue) Enqueue(frame Frame) error { + if frame.Priority < PriorityControl || frame.Priority > PriorityData { + return ErrSendQueueFull + } + copyFrame := Frame{Priority: frame.Priority, Payload: append([]byte(nil), frame.Payload...)} + needed := uint64(len(copyFrame.Payload)) + queue.mu.Lock() + defer queue.mu.Unlock() + limit := queue.limitForLocked(copyFrame.Priority) + if queue.bytes > limit || needed > limit-queue.bytes { + return ErrSendQueueFull + } + queue.bytes += needed + switch copyFrame.Priority { + case PriorityControl: + queue.control = append(queue.control, copyFrame) + case PriorityEssential: + queue.essential = append(queue.essential, copyFrame) + case PriorityData: + queue.data = append(queue.data, copyFrame) + } + return nil +} + +func (queue *SendQueue) Next() (Frame, bool) { + queue.mu.Lock() + defer queue.mu.Unlock() + var selected *Frame + if len(queue.control) > 0 && (queue.controlRun < queue.options.ControlBurst || (len(queue.essential) == 0 && len(queue.data) == 0)) { + selected = &queue.control[0] + queue.control = queue.control[1:] + queue.controlRun++ + } else if len(queue.essential) > 0 { + selected = &queue.essential[0] + queue.essential = queue.essential[1:] + queue.controlRun = 0 + } else if len(queue.data) > 0 { + selected = &queue.data[0] + queue.data = queue.data[1:] + queue.controlRun = 0 + } else if len(queue.control) > 0 { + selected = &queue.control[0] + queue.control = queue.control[1:] + queue.controlRun++ + } + if selected == nil { + return Frame{}, false + } + queue.bytes -= uint64(len(selected.Payload)) + return *selected, true +} + +func (queue *SendQueue) Bytes() uint64 { + queue.mu.Lock() + defer queue.mu.Unlock() + return queue.bytes +} + +func (queue *SendQueue) limitForLocked(priority FramePriority) uint64 { + switch priority { + case PriorityControl: + return queue.options.MaxBytes + case PriorityEssential: + return queue.options.MaxBytes - queue.options.ControlReserveBytes + default: + return queue.options.MaxBytes - queue.options.ControlReserveBytes - queue.options.EssentialReserveBytes + } +} diff --git a/internal/client/runtime/runtime_test.go b/internal/client/runtime/runtime_test.go new file mode 100644 index 0000000..e45e9d5 --- /dev/null +++ b/internal/client/runtime/runtime_test.go @@ -0,0 +1,94 @@ +package runtime + +import ( + "errors" + "testing" + "time" +) + +func TestSessionMachineBackoffAndStableReset_HP_FLOW_01(t *testing.T) { + t.Parallel() + machine, err := NewSessionMachine(BackoffOptions{Initial: time.Second, Maximum: 8 * time.Second, StableReset: 60 * time.Second}, func(cap time.Duration) time.Duration { return cap / 2 }) + if err != nil { + t.Fatal(err) + } + if err := machine.StartConnecting(0); err != nil { + t.Fatal(err) + } + if delay := machine.Failed(0); delay != 500*time.Millisecond || machine.NextConnectAt() != delay { + t.Fatalf("first failure delay/next = %s/%s", delay, machine.NextConnectAt()) + } + if err := machine.StartConnecting(499 * time.Millisecond); !errors.Is(err, ErrInvalidTransition) { + t.Fatalf("early reconnect error = %v, want ErrInvalidTransition", err) + } + if err := machine.StartConnecting(500 * time.Millisecond); err != nil { + t.Fatal(err) + } + if delay := machine.Failed(500 * time.Millisecond); delay != time.Second { + t.Fatalf("second failure delay = %s, want 1s", delay) + } + if err := machine.StartConnecting(1500 * time.Millisecond); err != nil { + t.Fatal(err) + } + if err := machine.TransportConnected(); err != nil { + t.Fatal(err) + } + if err := machine.Welcomed(); err != nil { + t.Fatal(err) + } + if err := machine.Reconciled(2 * time.Second); err != nil { + t.Fatal(err) + } + if delay := machine.Failed(62 * time.Second); delay != 500*time.Millisecond { + t.Fatalf("stable session did not reset backoff: %s", delay) + } +} + +func TestSendQueuePriorityReserveAndFairness_HP_FLOW_02(t *testing.T) { + t.Parallel() + queue, err := NewSendQueue(QueueOptions{MaxBytes: 100, ControlReserveBytes: 20, EssentialReserveBytes: 30, ControlBurst: 2}) + if err != nil { + t.Fatal(err) + } + if err := queue.Enqueue(Frame{Priority: PriorityData, Payload: make([]byte, 50)}); err != nil { + t.Fatal(err) + } + if err := queue.Enqueue(Frame{Priority: PriorityData, Payload: []byte("x")}); !errors.Is(err, ErrSendQueueFull) { + t.Fatalf("data consumed reserved bytes: %v", err) + } + if err := queue.Enqueue(Frame{Priority: PriorityEssential, Payload: make([]byte, 30)}); err != nil { + t.Fatal(err) + } + if err := queue.Enqueue(Frame{Priority: PriorityControl, Payload: make([]byte, 10)}); err != nil { + t.Fatal(err) + } + if err := queue.Enqueue(Frame{Priority: PriorityControl, Payload: make([]byte, 10)}); err != nil { + t.Fatal(err) + } + if got := queue.Bytes(); got != 100 { + t.Fatalf("queue bytes = %d, want 100", got) + } + for want := 0; want < 2; want++ { + frame, ok := queue.Next() + if !ok || frame.Priority != PriorityControl { + t.Fatalf("control frame %d = %#v, %t", want, frame, ok) + } + if want == 0 { + if err := queue.Enqueue(Frame{Priority: PriorityControl, Payload: []byte("c")}); err != nil { + t.Fatal(err) + } + } + } + frame, ok := queue.Next() + if !ok || frame.Priority != PriorityEssential { + t.Fatalf("essential frame starved by controls: %#v, %t", frame, ok) + } + frame, ok = queue.Next() + if !ok || frame.Priority != PriorityControl { + t.Fatalf("remaining control frame = %#v, %t", frame, ok) + } + frame, ok = queue.Next() + if !ok || frame.Priority != PriorityData || queue.Bytes() != 0 { + t.Fatalf("data frame/final bytes = %#v, %t/%d", frame, ok, queue.Bytes()) + } +} diff --git a/internal/client/runtime/session.go b/internal/client/runtime/session.go new file mode 100644 index 0000000..e3570aa --- /dev/null +++ b/internal/client/runtime/session.go @@ -0,0 +1,160 @@ +// Package runtime owns shared client session state. It contains no socket or +// OS calls, so supervisor workers can outlive a replaceable network session. +package runtime + +import ( + "errors" + "math" + "sync" + "time" +) + +var ErrInvalidTransition = errors.New("invalid client session transition") + +type State uint8 + +const ( + StateBackoff State = iota + 1 + StateConnecting + StateHello + StateReconciling + StateActive + StateClosing +) + +type BackoffOptions struct { + Initial time.Duration + Maximum time.Duration + StableReset time.Duration +} + +func (options BackoffOptions) Validate() error { + if options.Initial <= 0 || options.Maximum < options.Initial || options.StableReset <= 0 { + return errors.New("invalid client reconnect backoff options") + } + return nil +} + +// Jitter returns a delay in [0, cap]. It is injected so all backoff boundaries +// are testable without sleeps or probabilistic assertions. +type Jitter func(cap time.Duration) time.Duration + +type SessionMachine struct { + mu sync.Mutex + options BackoffOptions + jitter Jitter + state State + failures uint32 + nextConnect time.Duration + activeSince time.Duration +} + +func NewSessionMachine(options BackoffOptions, jitter Jitter) (*SessionMachine, error) { + if err := options.Validate(); err != nil { + return nil, err + } + if jitter == nil { + return nil, errors.New("client reconnect jitter is required") + } + return &SessionMachine{options: options, jitter: jitter, state: StateBackoff}, nil +} + +func (machine *SessionMachine) State() State { + machine.mu.Lock() + defer machine.mu.Unlock() + return machine.state +} + +func (machine *SessionMachine) NextConnectAt() time.Duration { + machine.mu.Lock() + defer machine.mu.Unlock() + return machine.nextConnect +} + +func (machine *SessionMachine) StartConnecting(now time.Duration) error { + machine.mu.Lock() + defer machine.mu.Unlock() + if machine.state != StateBackoff || now < machine.nextConnect { + return ErrInvalidTransition + } + machine.state = StateConnecting + return nil +} + +func (machine *SessionMachine) TransportConnected() error { + return machine.transition(StateConnecting, StateHello) +} + +func (machine *SessionMachine) Welcomed() error { + return machine.transition(StateHello, StateReconciling) +} + +func (machine *SessionMachine) Reconciled(now time.Duration) error { + machine.mu.Lock() + defer machine.mu.Unlock() + if machine.state != StateReconciling { + return ErrInvalidTransition + } + machine.state = StateActive + machine.activeSince = now + return nil +} + +func (machine *SessionMachine) BeginClosing() error { + machine.mu.Lock() + defer machine.mu.Unlock() + if machine.state == StateBackoff || machine.state == StateClosing { + return ErrInvalidTransition + } + machine.state = StateClosing + return nil +} + +// Failed moves every in-flight session state to backoff. Only a continuous +// active period reaching StableReset clears exponential history. +func (machine *SessionMachine) Failed(now time.Duration) time.Duration { + machine.mu.Lock() + defer machine.mu.Unlock() + if machine.state == StateActive && now >= machine.activeSince && now-machine.activeSince >= machine.options.StableReset { + machine.failures = 0 + } + if machine.failures < math.MaxUint32 { + machine.failures++ + } + cap := machine.capLocked() + delay := machine.jitter(cap) + if delay < 0 { + delay = 0 + } + if delay > cap { + delay = cap + } + machine.state = StateBackoff + machine.nextConnect = now + delay + machine.activeSince = 0 + return delay +} + +func (machine *SessionMachine) transition(from, to State) error { + machine.mu.Lock() + defer machine.mu.Unlock() + if machine.state != from { + return ErrInvalidTransition + } + machine.state = to + return nil +} + +func (machine *SessionMachine) capLocked() time.Duration { + cap := machine.options.Initial + for index := uint32(1); index < machine.failures && cap < machine.options.Maximum; index++ { + if cap > machine.options.Maximum/2 { + return machine.options.Maximum + } + cap *= 2 + } + if cap > machine.options.Maximum { + return machine.options.Maximum + } + return cap +} diff --git a/test/coverage.toml b/test/coverage.toml index 8ad2f89..f799e51 100644 --- a/test/coverage.toml +++ b/test/coverage.toml @@ -227,6 +227,18 @@ layer = "unit" status = "implemented" tests = ["internal/client/spool/output_test.go:TestAppendOutputNeverEvictsAssignedEvent_BH_OUTFLOW_02"] +[[requirements]] +id = "HP-FLOW-01" +layer = "unit" +status = "implemented" +tests = ["internal/client/runtime/runtime_test.go:TestSessionMachineBackoffAndStableReset_HP_FLOW_01"] + +[[requirements]] +id = "HP-FLOW-02" +layer = "unit" +status = "implemented" +tests = ["internal/client/runtime/runtime_test.go:TestSendQueuePriorityReserveAndFairness_HP_FLOW_02"] + [[requirements]] id = "BH-RET-01" layer = "unit"