feat: add bounded session scheduling
This commit is contained in:
@@ -0,0 +1,277 @@
|
||||
// Package session owns transport-neutral server session scheduling and fencing.
|
||||
package session
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
var (
|
||||
ErrControlLaneFull = errors.New("session control lane is full")
|
||||
ErrSessionClosed = errors.New("session is closed")
|
||||
)
|
||||
|
||||
type FrameKind uint8
|
||||
|
||||
const (
|
||||
FrameControl FrameKind = iota + 1
|
||||
FrameData
|
||||
)
|
||||
|
||||
type Frame struct {
|
||||
Kind FrameKind
|
||||
Payload []byte
|
||||
}
|
||||
|
||||
// WriterQueue is owned by one socket writer. Data saturation leaves the work
|
||||
// durable; control saturation is a session-fatal invariant.
|
||||
type WriterQueue struct {
|
||||
control chan Frame
|
||||
data chan Frame
|
||||
controlBurst uint32
|
||||
controlRun uint32
|
||||
closed chan struct{}
|
||||
closeOnce sync.Once
|
||||
}
|
||||
|
||||
func NewWriterQueue(controlCapacity, dataCapacity, controlBurst uint32) *WriterQueue {
|
||||
if controlCapacity == 0 || dataCapacity == 0 || controlBurst == 0 {
|
||||
panic("session queue capacities must be positive")
|
||||
}
|
||||
return &WriterQueue{
|
||||
control: make(chan Frame, controlCapacity), data: make(chan Frame, dataCapacity),
|
||||
controlBurst: controlBurst, closed: make(chan struct{}),
|
||||
}
|
||||
}
|
||||
|
||||
func (queue *WriterQueue) EnqueueControl(frame Frame) error {
|
||||
if frame.Kind != FrameControl {
|
||||
return ErrControlLaneFull
|
||||
}
|
||||
select {
|
||||
case <-queue.closed:
|
||||
return ErrSessionClosed
|
||||
default:
|
||||
}
|
||||
select {
|
||||
case queue.control <- frame:
|
||||
return nil
|
||||
default:
|
||||
return ErrControlLaneFull
|
||||
}
|
||||
}
|
||||
|
||||
// EnqueueData reports false when the bounded data lane is full. It never
|
||||
// creates a fallback goroutine or buffer.
|
||||
func (queue *WriterQueue) EnqueueData(frame Frame) bool {
|
||||
if frame.Kind != FrameData {
|
||||
return false
|
||||
}
|
||||
select {
|
||||
case <-queue.closed:
|
||||
return false
|
||||
default:
|
||||
}
|
||||
select {
|
||||
case queue.data <- frame:
|
||||
return true
|
||||
default:
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
func (queue *WriterQueue) Next(ctx context.Context) (Frame, error) {
|
||||
for {
|
||||
if queue.controlRun < queue.controlBurst {
|
||||
select {
|
||||
case frame := <-queue.control:
|
||||
queue.controlRun++
|
||||
return frame, nil
|
||||
default:
|
||||
}
|
||||
}
|
||||
select {
|
||||
case frame := <-queue.data:
|
||||
queue.controlRun = 0
|
||||
return frame, nil
|
||||
default:
|
||||
}
|
||||
select {
|
||||
case frame := <-queue.control:
|
||||
queue.controlRun++
|
||||
return frame, nil
|
||||
default:
|
||||
}
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return Frame{}, ctx.Err()
|
||||
case <-queue.closed:
|
||||
return Frame{}, ErrSessionClosed
|
||||
case frame := <-queue.control:
|
||||
queue.controlRun++
|
||||
return frame, nil
|
||||
case frame := <-queue.data:
|
||||
queue.controlRun = 0
|
||||
return frame, nil
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (queue *WriterQueue) Close() { queue.closeOnce.Do(func() { close(queue.closed) }) }
|
||||
|
||||
type HeartbeatAction uint8
|
||||
|
||||
const (
|
||||
HeartbeatNone HeartbeatAction = iota
|
||||
HeartbeatPing
|
||||
HeartbeatClose
|
||||
)
|
||||
|
||||
type Heartbeat struct {
|
||||
idle time.Duration
|
||||
timeout time.Duration
|
||||
lastInbound time.Duration
|
||||
pingSent bool
|
||||
}
|
||||
|
||||
func NewHeartbeat(idle, timeout, now time.Duration) *Heartbeat {
|
||||
if idle <= 0 || timeout <= idle {
|
||||
panic("invalid heartbeat timings")
|
||||
}
|
||||
return &Heartbeat{idle: idle, timeout: timeout, lastInbound: now}
|
||||
}
|
||||
|
||||
func (heartbeat *Heartbeat) ObserveInbound(now time.Duration) {
|
||||
if now >= heartbeat.lastInbound {
|
||||
heartbeat.lastInbound = now
|
||||
heartbeat.pingSent = false
|
||||
}
|
||||
}
|
||||
|
||||
func (heartbeat *Heartbeat) Check(now time.Duration) HeartbeatAction {
|
||||
if now < heartbeat.lastInbound {
|
||||
return HeartbeatNone
|
||||
}
|
||||
elapsed := now - heartbeat.lastInbound
|
||||
if elapsed >= heartbeat.timeout {
|
||||
return HeartbeatClose
|
||||
}
|
||||
if elapsed >= heartbeat.idle && !heartbeat.pingSent {
|
||||
heartbeat.pingSent = true
|
||||
return HeartbeatPing
|
||||
}
|
||||
return HeartbeatNone
|
||||
}
|
||||
|
||||
type DispatchLane uint8
|
||||
|
||||
const (
|
||||
DispatchNone DispatchLane = iota
|
||||
DispatchRunning
|
||||
DispatchQueued
|
||||
)
|
||||
|
||||
type CapacityShadow struct {
|
||||
running, queued uint32
|
||||
maxRunning, maxQueued uint32
|
||||
shadowRunning uint32
|
||||
shadowQueued uint32
|
||||
}
|
||||
|
||||
func (shadow *CapacityShadow) UpdateAdvertised(running, queued, maxRunning, maxQueued uint32) bool {
|
||||
if maxRunning == 0 || maxQueued == 0 || running > maxRunning || queued > maxQueued {
|
||||
return false
|
||||
}
|
||||
shadow.running, shadow.queued = running, queued
|
||||
shadow.maxRunning, shadow.maxQueued = maxRunning, maxQueued
|
||||
return true
|
||||
}
|
||||
|
||||
func (shadow *CapacityShadow) Reserve() DispatchLane {
|
||||
if shadow.running+shadow.shadowRunning < shadow.maxRunning {
|
||||
shadow.shadowRunning++
|
||||
return DispatchRunning
|
||||
}
|
||||
if shadow.queued+shadow.shadowQueued < shadow.maxQueued {
|
||||
shadow.shadowQueued++
|
||||
return DispatchQueued
|
||||
}
|
||||
return DispatchNone
|
||||
}
|
||||
|
||||
func (shadow *CapacityShadow) Release(lane DispatchLane) bool {
|
||||
switch lane {
|
||||
case DispatchRunning:
|
||||
if shadow.shadowRunning == 0 {
|
||||
return false
|
||||
}
|
||||
shadow.shadowRunning--
|
||||
case DispatchQueued:
|
||||
if shadow.shadowQueued == 0 {
|
||||
return false
|
||||
}
|
||||
shadow.shadowQueued--
|
||||
default:
|
||||
return false
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
func (shadow CapacityShadow) Pending() (running, queued uint32) {
|
||||
return shadow.shadowRunning, shadow.shadowQueued
|
||||
}
|
||||
|
||||
type Handle struct {
|
||||
ClientID string
|
||||
SessionID [16]byte
|
||||
Generation uint64
|
||||
Context context.Context
|
||||
cancel context.CancelFunc
|
||||
}
|
||||
|
||||
type Registry struct {
|
||||
mu sync.Mutex
|
||||
clients map[string]*Handle
|
||||
}
|
||||
|
||||
func NewRegistry() *Registry { return &Registry{clients: make(map[string]*Handle)} }
|
||||
|
||||
// Install replaces only an older in-memory handle for the same client. The
|
||||
// caller must have already durably fenced it in the store.
|
||||
func (registry *Registry) Install(clientID string, sessionID [16]byte, generation uint64) (*Handle, error) {
|
||||
if clientID == "" || sessionID == [16]byte{} || generation == 0 {
|
||||
return nil, ErrSessionClosed
|
||||
}
|
||||
registry.mu.Lock()
|
||||
defer registry.mu.Unlock()
|
||||
if old := registry.clients[clientID]; old != nil {
|
||||
old.cancel()
|
||||
}
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
handle := &Handle{ClientID: clientID, SessionID: sessionID, Generation: generation, Context: ctx, cancel: cancel}
|
||||
registry.clients[clientID] = handle
|
||||
return handle, nil
|
||||
}
|
||||
|
||||
func (registry *Registry) Remove(handle *Handle) bool {
|
||||
if handle == nil {
|
||||
return false
|
||||
}
|
||||
registry.mu.Lock()
|
||||
defer registry.mu.Unlock()
|
||||
current := registry.clients[handle.ClientID]
|
||||
if current != handle || current.Generation != handle.Generation || current.SessionID != handle.SessionID {
|
||||
return false
|
||||
}
|
||||
delete(registry.clients, handle.ClientID)
|
||||
handle.cancel()
|
||||
return true
|
||||
}
|
||||
|
||||
func (registry *Registry) Get(clientID string) *Handle {
|
||||
registry.mu.Lock()
|
||||
defer registry.mu.Unlock()
|
||||
return registry.clients[clientID]
|
||||
}
|
||||
Reference in New Issue
Block a user