// 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 } // 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 }