// 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 // Written is closed by the sole socket writer after a successful write. // It is used only for protocol barriers such as reconciliation-before-work. Written chan<- struct{} } // 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] }