Files
kyle 2fe3b9c865 Proxy: move or refuse prompts that do not fit the leased host's context
Implemented-By: OpenCode session (model recorded in docs/implementer-log.md)
2026-09-25 09:56:46 -07:00

354 lines
10 KiB
Go

// Package lease is crossbar's sticky placement table. It remembers which host each
// conversation (and each route) is on and keeps it there unless the host is
// unhealthy, the lease has been idle past lease_idle, or an operator releases or
// pins the route. Every change is written through to a Persister and replayed back
// at start so a restart does not reshuffle sessions.
package lease
import (
"errors"
"fmt"
"sort"
"sync"
"time"
"git.wntrmute.dev/kyle/crossbar/internal/store"
)
var (
ErrNoHost = errors.New("lease: no usable host")
ErrPinnedDown = errors.New("lease: pinned host is not healthy")
ErrUnknownHost = errors.New("lease: unknown host")
)
// Key identifies one conversation: a route, the fingerprint of its first request
// (empty for a request with no user message, which leases the route itself), and
// the model it wants.
type Key struct {
Route, FP, Model string
}
// Lease is one routed model on one host, in memory.
type Lease struct {
Key
Host string
State store.State
Created, LastUsed time.Time
}
// Persister is the durable half of the table: leases, the pin row, and the event
// log. *store.Store satisfies it.
type Persister interface {
SaveLease(store.Lease) error
DeleteLease(route, fp, model string) error
ListLeases() ([]store.Lease, error)
RecordEvent(store.LeaseEvent) error
}
// Hosts is the live health view the table consults before placing a lease.
type Hosts interface {
Healthy(name string) bool
Draining(name string) bool
}
// Chooser picks a host for a new lease among the eligible candidates.
type Chooser interface {
Choose(candidates []string, model string) (string, bool)
}
type Table struct {
mu sync.Mutex
leases map[Key]*Lease // active leases, keyed by (route, fp, model)
pins map[string]string // route -> pinned host
seen map[string]map[string]bool // route -> candidate hosts ever asked for or stored
p Persister
hosts Hosts
choose Chooser
idle time.Duration
}
// New builds a table and loads its state. Rows with FP=="" && Model=="" &&
// State==Pinned are pins; every other row becomes an active lease and its host
// counts as a candidate already seen for the route.
func New(p Persister, h Hosts, c Chooser, idle time.Duration) (*Table, error) {
loads, err := p.ListLeases()
if err != nil {
return nil, fmt.Errorf("lease: list leases: %w", err)
}
t := &Table{
leases: make(map[Key]*Lease),
pins: make(map[string]string),
seen: make(map[string]map[string]bool),
p: p,
hosts: h,
choose: c,
idle: idle,
}
for _, l := range loads {
if l.FP == "" && l.Model == "" && l.State == store.Pinned {
t.pins[l.Route] = l.Host
} else {
t.leases[Key{l.Route, l.FP, l.Model}] = &Lease{Key{l.Route, l.FP, l.Model}, l.Host, l.State, l.Created, l.LastUsed}
}
if t.seen[l.Route] == nil {
t.seen[l.Route] = make(map[string]bool)
}
t.seen[l.Route][l.Host] = true
}
return t, nil
}
func (l *Lease) store() store.Lease {
return store.Lease{Route: l.Route, FP: l.FP, Model: l.Model, Host: l.Host, State: l.State, Created: l.Created, LastUsed: l.LastUsed}
}
// Acquire places k, keeping it sticky. See the task's Acquire ordering: pinned
// route, existing lease, inherit the route's host, else choose. Only one lease is
// created per call, and a failed save is rolled back in memory.
func (t *Table) Acquire(k Key, candidates []string, now time.Time) (host string, reused bool, err error) {
t.mu.Lock()
defer t.mu.Unlock()
if t.seen[k.Route] == nil {
t.seen[k.Route] = make(map[string]bool)
}
for _, c := range candidates {
t.seen[k.Route][c] = true
}
// unhealthyFrom is set when an existing lease sat on a dead host; the next
// create records an unhealthy move instead of a fresh new.
var unhealthyFrom string
reason := store.ReasonNew
// 1. Pinned route.
if pin, ok := t.pins[k.Route]; ok {
if !t.hosts.Healthy(pin) {
return "", false, ErrPinnedDown
}
if _, exists := t.leases[k]; exists {
return pin, true, nil
}
l := &Lease{k, pin, store.Active, now, now}
t.leases[k] = l
if err := t.save(l); err != nil {
delete(t.leases, k)
return "", false, err
}
t.event(now, k, store.ReasonNew, "", pin)
return pin, false, nil
}
// 2. Existing lease for k.
if l, exists := t.leases[k]; exists {
if t.hosts.Healthy(l.Host) {
l.LastUsed = now
if err := t.save(l); err != nil {
return "", false, err
}
return l.Host, true, nil
}
unhealthyFrom = l.Host
}
// 3. Inherit the route's own host when a fingerprinted request can start
// where the route already lives.
if k.FP != "" {
rk := Key{k.Route, "", k.Model}
if rl, exists := t.leases[rk]; exists && t.hosts.Healthy(rl.Host) {
l := &Lease{k, rl.Host, store.Active, now, now}
t.leases[k] = l
if err := t.save(l); err != nil {
delete(t.leases, k)
return "", false, err
}
if unhealthyFrom != "" {
reason = store.ReasonUnhealthy
}
t.event(now, k, reason, unhealthyFrom, rl.Host)
return rl.Host, true, nil
}
}
// 4. Choose among healthy, non-draining candidates.
filtered := make([]string, 0, len(candidates))
for _, c := range candidates {
if t.hosts.Healthy(c) && !t.hosts.Draining(c) {
filtered = append(filtered, c)
}
}
host, ok := t.choose.Choose(filtered, k.Model)
if !ok {
return "", false, ErrNoHost
}
l := &Lease{k, host, store.Active, now, now}
t.leases[k] = l
if err := t.save(l); err != nil {
delete(t.leases, k)
return "", false, err
}
if unhealthyFrom != "" {
reason = store.ReasonUnhealthy
}
t.event(now, k, reason, unhealthyFrom, host)
return host, false, nil
}
// Move relocates k to host, dropping any existing lease for k (on any host). It re-leases k onto
// host, deletes the old row, and records a ctx event carrying the old and new hosts. It returns an
// error only from the persister, rolling back the in-memory lease on save failure.
func (t *Table) Move(k Key, host string, now time.Time) error {
t.mu.Lock()
defer t.mu.Unlock()
var from string
if l, exists := t.leases[k]; exists {
from = l.Host
delete(t.leases, k)
_ = t.p.DeleteLease(k.Route, k.FP, k.Model)
}
l := &Lease{k, host, store.Active, now, now}
t.leases[k] = l
if err := t.save(l); err != nil {
delete(t.leases, k)
return fmt.Errorf("lease: %w", err)
}
t.event(now, k, store.ReasonCtx, from, host)
return nil
}
// Candidates records hosts as seen for route (idempotent), so Pin can accept a host the route
// is configured for before any request has used it. cmd/crossbar calls it for every route at
// start; the admin handler calls it before Pin.
func (t *Table) Candidates(route string, hosts []string) {
t.mu.Lock()
defer t.mu.Unlock()
if t.seen[route] == nil {
t.seen[route] = make(map[string]bool)
}
for _, h := range hosts {
t.seen[route][h] = true
}
}
// save writes a lease through, failing the call on a persister error.
func (t *Table) save(l *Lease) error {
if err := t.p.SaveLease(l.store()); err != nil {
return fmt.Errorf("lease: %w", err)
}
return nil
}
// event appends a lease event. from is empty for a fresh placement.
func (t *Table) event(now time.Time, k Key, reason, from, to string) {
_ = t.p.RecordEvent(store.LeaseEvent{TS: now, Route: k.Route, Model: k.Model, FromHost: from, ToHost: to, Reason: reason})
}
// ExpireIdle removes leases idle longer than lease_idle (never pins, which are
// not in the lease map) and records an idle event for each. It returns how many
// it removed.
func (t *Table) ExpireIdle(now time.Time) int {
t.mu.Lock()
defer t.mu.Unlock()
var idle []Key
for k, l := range t.leases {
if now.Sub(l.LastUsed) > t.idle {
idle = append(idle, k)
}
}
for _, k := range idle {
l := t.leases[k]
delete(t.leases, k)
_ = t.p.DeleteLease(k.Route, k.FP, k.Model)
t.event(now, k, store.ReasonIdle, l.Host, "")
}
return len(idle)
}
// Pin routes route to host. host must be a candidate the table has seen for the
// route, else ErrUnknownHost. It records a pin event, stores the pin row, and
// deletes the route's existing leases on other hosts so the next turn lands on
// the pin.
func (t *Table) Pin(route, host string, now time.Time) error {
t.mu.Lock()
defer t.mu.Unlock()
if t.seen[route] == nil || !t.seen[route][host] {
return ErrUnknownHost
}
t.event(now, Key{Route: route}, store.ReasonPin, "", host)
if err := t.p.SaveLease(store.Lease{Route: route, FP: "", Model: "", Host: host, State: store.Pinned, Created: now, LastUsed: now}); err != nil {
return fmt.Errorf("lease: %w", err)
}
t.pins[route] = host
for k, l := range t.leases {
if k.Route == route && l.Host != host {
delete(t.leases, k)
_ = t.p.DeleteLease(k.Route, k.FP, k.Model)
t.event(now, k, store.ReasonPin, l.Host, host)
}
}
return nil
}
// Unpin clears the pin for route and records a release. Existing leases stay put.
func (t *Table) Unpin(route string) {
t.mu.Lock()
defer t.mu.Unlock()
delete(t.pins, route)
_ = t.p.DeleteLease(route, "", "")
t.event(time.Now(), Key{Route: route}, store.ReasonRelease, "", "")
}
// Pinned returns the pinned host for route, or "" if it is not pinned.
func (t *Table) Pinned(route string) string {
t.mu.Lock()
defer t.mu.Unlock()
return t.pins[route]
}
// Release drops every lease (not the pin) of route and records a release event
// per dropped lease. It returns how many were removed.
func (t *Table) Release(route string) int {
t.mu.Lock()
defer t.mu.Unlock()
var keys []Key
for k := range t.leases {
if k.Route == route {
keys = append(keys, k)
}
}
for _, k := range keys {
delete(t.leases, k)
_ = t.p.DeleteLease(k.Route, k.FP, k.Model)
t.event(time.Now(), k, store.ReasonRelease, "", "")
}
return len(keys)
}
// Snapshot returns copies of the active leases, sorted by route, fp, model. Pins
// are excluded.
func (t *Table) Snapshot() []Lease {
t.mu.Lock()
defer t.mu.Unlock()
out := make([]Lease, 0, len(t.leases))
for _, l := range t.leases {
out = append(out, *l)
}
sort.Slice(out, func(i, j int) bool {
if out[i].Route != out[j].Route {
return out[i].Route < out[j].Route
}
if out[i].FP != out[j].FP {
return out[i].FP < out[j].FP
}
return out[i].Model < out[j].Model
})
return out
}