feat: add client reconnect scheduling primitives
This commit is contained in:
@@ -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
|
||||
}
|
||||
}
|
||||
@@ -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())
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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"
|
||||
|
||||
Reference in New Issue
Block a user