Add the wake package: magic packets and a waiter
Implemented-By: OpenCode session (model recorded in docs/implementer-log.md)
This commit is contained in:
@@ -5,6 +5,7 @@ owner fills in the Model column. The reviewer adds findings under "Reviews" once
|
||||
|
||||
| Task | Date | Status | Gate runs | First gate | Deviations | Notes | Model |
|
||||
|---|---|---|---|---|---|---|---|
|
||||
| v2/03-wake | 2026-09-25 | done | 1 | pass | none | Implemented wake-on-LAN in new `internal/wake/wake.go`: `MagicPacket` builds the 102-byte frame via `net.ParseMAC` (six `0xff` bytes plus the MAC repeated sixteen times) and rejects bad MACs; `Send` emits one UDP4 datagram to the resolved broadcast address, returning parse/resolve/write errors; `Waker` tracks last-sent per host under a mutex and sends at most once per `Wait` window, polling health every second (`PollEvery` is a test hook) until healthy, on `Wait` timeout, or on ctx cancellation, returning false for an unknown host without sending. Copied `internal/wake/wake_test.go` byte-identical to `docs/plans/v2/_files/`; `go test -race -count=3 ./internal/wake/` ok and `make gate` printed `gate: ok` on the first run. | llama.cpp/ornith-1.5-35b-a3b |
|
||||
| v2/02-ctxguard | 2026-09-25 | done | 1 | pass | none | Implemented the context-size guard in new `internal/proxy/ctxguard.go` (estimate `int(float64(len(body))/4*1.2)`; rule 2 skip on unknown/fit; rule 3 move via `leases.Move` with a `moved:<old>><new>` header; rule 4 400 with `{"error":"prompt too large","estimate":E,"max":M}` and a status-400 accounting row, no forward, no mark-down) and wired it into `ServeHTTP` between the lease and the slot; added `Move` to `internal/lease/lease.go` (re-leases, deletes the old row, records a `ctx` event) and `ReasonCtx = "ctx"` to `internal/store`. Copied `internal/proxy/ctxguard_test.go` byte-identical to `docs/plans/v2/_files/`; `go test -race -count=2 ./internal/proxy/ ./internal/lease/` ok and `make gate` printed `gate: ok` on the first run. | llama.cpp/ornith-1.5-35b-a3b |
|
||||
| v2/01-props | 2026-09-25 | done | 1 | pass | none | Implemented /props learning in `internal/health/health.go`: added `Status.NCtx`/`Status.Slots`, `PerSlotCtx()`, and a best-effort `GET <base>/props` appended to the poll after `/v1/models`, setting NCtx/Slots to 0 (negative → 0) on any failure without counting the poll as failed; exposed them in `internal/admin/admin.go` `HostView`. Copied `internal/health/props_test.go` and the replacement `internal/proxy/helpers_test.go` byte-identical to `docs/plans/v2/_files/`. `go test -race ./...` and `make gate` pass on the first run. | ? |
|
||||
| v2/01-props | 2026-09-25 | stopped | 1 | fail | none | Implemented /props learning in `internal/health/health.go` (added `Status.NCtx`/`Status.Slots`, `PerSlotCtx`, and a best-effort `GET <base>/props` appended to the poll; 0/unknown on any failure without failing the poll) and exposed them in `internal/admin/admin.go` `HostView`; copied `internal/health/props_test.go` byte-identical to `docs/plans/v2/_files/`. `go test -race ./internal/health/ ./internal/admin/` ok. `make gate` fails on two GIVEN v1 proxy tests — `TestConversationIsStickyAndLeaseHeaderTellsWhy` (alpha 1/beta 7, want 0/6) and `TestDifferentConversationsSpreadByFreeSlots` (beta 3/alpha 2, want 2/1) — which assert exact upstream hit counts; the task-required `/props` poll now lands on that scaffold's `/` catch-all and bumps the counter by exactly 1 per host (deterministic, confirmed over 3 repeated runs, not a flake). `internal/proxy/helpers_test.go` is byte-identical to `docs/plans/v1/_files/` (protected) and cannot be updated here; the `/props` request is unavoidable per the task, so the owner must hand over a scaffold that registers `/props` without counting it as a hit. Code left uncommitted for review. | ? |
|
||||
|
||||
@@ -0,0 +1,142 @@
|
||||
// Package wake sends wake-on-LAN magic packets and waits for a sleeping host to
|
||||
// appear healthy in the health table. A sleeping host takes tens of seconds to
|
||||
// come up, so the Waker remembers when it last sent and wakes a host at most once
|
||||
// per wait window.
|
||||
package wake
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"net"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Target describes how to wake one named host.
|
||||
type Target struct {
|
||||
MAC, Broadcast string // MAC "aa:bb:cc:dd:ee:ff" (any separator, any case); Broadcast "host:port"
|
||||
Wait time.Duration
|
||||
}
|
||||
|
||||
// Health reports whether a named host is currently healthy. Implementations must
|
||||
// be safe for concurrent use.
|
||||
type Health interface{ Healthy(name string) bool }
|
||||
|
||||
// magicPacketLen is six sync bytes plus the MAC repeated sixteen times.
|
||||
const magicPacketLen = 6 + 6*16
|
||||
|
||||
// MagicPacket builds a wake-on-LAN magic packet: six 0xff bytes followed by the
|
||||
// target MAC sixteen times, a 102-byte frame.
|
||||
func MagicPacket(mac string) ([]byte, error) {
|
||||
m, err := net.ParseMAC(mac)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("wake: parse MAC %q: %w", mac, err)
|
||||
}
|
||||
if len(m) != 6 {
|
||||
return nil, fmt.Errorf("wake: MAC %q is not six bytes", mac)
|
||||
}
|
||||
pkt := make([]byte, magicPacketLen)
|
||||
for i := range pkt[:6] {
|
||||
pkt[i] = 0xff
|
||||
}
|
||||
for i := 0; i < 16; i++ {
|
||||
copy(pkt[6+i*6:], m)
|
||||
}
|
||||
return pkt, nil
|
||||
}
|
||||
|
||||
// Send emits one magic packet for mac to the broadcast address as a single UDP4
|
||||
// datagram, reporting parse, resolve and write errors.
|
||||
func Send(mac, broadcast string) error {
|
||||
pkt, err := MagicPacket(mac)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
remote, err := net.ResolveUDPAddr("udp4", broadcast)
|
||||
if err != nil {
|
||||
return fmt.Errorf("wake: resolve broadcast %q: %w", broadcast, err)
|
||||
}
|
||||
conn, err := net.DialUDP("udp4", nil, remote)
|
||||
if err != nil {
|
||||
return fmt.Errorf("wake: dial broadcast %q: %w", broadcast, err)
|
||||
}
|
||||
defer conn.Close()
|
||||
if _, err := conn.Write(pkt); err != nil {
|
||||
return fmt.Errorf("wake: write packet to %q: %w", broadcast, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// Waker wakes named hosts at most once per wait window and waits for the health
|
||||
// table to report them healthy. It is safe for concurrent Wake calls.
|
||||
type Waker struct {
|
||||
mu sync.Mutex
|
||||
targets map[string]Target
|
||||
health Health
|
||||
lastSent map[string]time.Time
|
||||
poll time.Duration
|
||||
}
|
||||
|
||||
// New returns a Waker for the given targets, polling health every second.
|
||||
func New(targets map[string]Target, h Health) *Waker {
|
||||
return &Waker{
|
||||
targets: targets,
|
||||
health: h,
|
||||
lastSent: make(map[string]time.Time),
|
||||
poll: time.Second,
|
||||
}
|
||||
}
|
||||
|
||||
// PollEvery sets how often Wake re-checks health; it is a test hook. Production
|
||||
// keeps the 1 s default from New.
|
||||
func (w *Waker) PollEvery(d time.Duration) {
|
||||
w.mu.Lock()
|
||||
w.poll = d
|
||||
w.mu.Unlock()
|
||||
}
|
||||
|
||||
// Wake sends a magic packet for host if none was sent in the last Wait, then
|
||||
// polls health until the host is healthy, the wait elapses, or ctx is done. It
|
||||
// returns true only when the host becomes healthy, and false for an unknown
|
||||
// host, on timeout, or when ctx ends first.
|
||||
func (w *Waker) Wake(ctx context.Context, host string) bool {
|
||||
w.mu.Lock()
|
||||
target, ok := w.targets[host]
|
||||
if !ok {
|
||||
w.mu.Unlock()
|
||||
return false
|
||||
}
|
||||
now := time.Now()
|
||||
if last, sent := w.lastSent[host]; !sent || now.Sub(last) >= target.Wait {
|
||||
w.lastSent[host] = now
|
||||
w.mu.Unlock()
|
||||
if err := Send(target.MAC, target.Broadcast); err != nil {
|
||||
return false
|
||||
}
|
||||
w.mu.Lock()
|
||||
}
|
||||
poll := w.poll
|
||||
deadline := now.Add(target.Wait)
|
||||
w.mu.Unlock()
|
||||
|
||||
timer := time.NewTimer(poll)
|
||||
defer timer.Stop()
|
||||
for {
|
||||
if w.health.Healthy(host) {
|
||||
return true
|
||||
}
|
||||
if time.Now().After(deadline) {
|
||||
return false
|
||||
}
|
||||
d := poll
|
||||
if rem := time.Until(deadline); rem < d {
|
||||
d = rem
|
||||
}
|
||||
timer.Reset(d)
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return false
|
||||
case <-timer.C:
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,114 @@
|
||||
package wake_test
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"net"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"git.wntrmute.dev/kyle/crossbar/internal/wake"
|
||||
)
|
||||
|
||||
func listen(t *testing.T) (*net.UDPConn, string) {
|
||||
conn, err := net.ListenUDP("udp4", &net.UDPAddr{IP: net.IPv4(127, 0, 0, 1)})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(func() { conn.Close() })
|
||||
return conn, conn.LocalAddr().String()
|
||||
}
|
||||
|
||||
func TestMagicPacket(t *testing.T) {
|
||||
pkt, err := wake.MagicPacket("aa:bb:cc:dd:ee:ff")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(pkt) != 102 || !bytes.Equal(pkt[:6], bytes.Repeat([]byte{0xff}, 6)) {
|
||||
t.Fatalf("packet = % x", pkt)
|
||||
}
|
||||
mac := []byte{0xaa, 0xbb, 0xcc, 0xdd, 0xee, 0xff}
|
||||
for i := 0; i < 16; i++ {
|
||||
if !bytes.Equal(pkt[6+6*i:12+6*i], mac) {
|
||||
t.Fatalf("repetition %d wrong: % x", i, pkt[6+6*i:12+6*i])
|
||||
}
|
||||
}
|
||||
for _, bad := range []string{"", "aa:bb", "zz:bb:cc:dd:ee:ff", "aabbccddeeff00"} {
|
||||
if _, err := wake.MagicPacket(bad); err == nil {
|
||||
t.Errorf("MagicPacket(%q) must fail", bad)
|
||||
}
|
||||
}
|
||||
if p2, _ := wake.MagicPacket("AA-BB-CC-DD-EE-FF"); !bytes.Equal(p2, pkt) {
|
||||
t.Errorf("dash-separated upper-case MAC must give the same packet")
|
||||
}
|
||||
}
|
||||
|
||||
func TestSendReachesTheBroadcastAddress(t *testing.T) {
|
||||
conn, addr := listen(t)
|
||||
if err := wake.Send("aa:bb:cc:dd:ee:ff", addr); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
buf := make([]byte, 200)
|
||||
_ = conn.SetReadDeadline(time.Now().Add(time.Second))
|
||||
n, _, err := conn.ReadFromUDP(buf)
|
||||
if err != nil || n != 102 {
|
||||
t.Fatalf("received %d bytes, err %v", n, err)
|
||||
}
|
||||
if err := wake.Send("aa:bb:cc:dd:ee:ff", "256.1.1.1:9"); err == nil {
|
||||
t.Error("an unresolvable broadcast address must be an error")
|
||||
}
|
||||
}
|
||||
|
||||
// fakeHealth flips to healthy after `after` calls to Healthy.
|
||||
type fakeHealth struct{ calls, after int }
|
||||
|
||||
func (f *fakeHealth) Healthy(name string) bool { f.calls++; return f.calls > f.after }
|
||||
|
||||
func TestWakerSendsOncePerWindowAndWaitsForHealth(t *testing.T) {
|
||||
conn, addr := listen(t)
|
||||
h := &fakeHealth{after: 3}
|
||||
w := wake.New(map[string]wake.Target{"titan": {MAC: "aa:bb:cc:dd:ee:ff", Broadcast: addr, Wait: 2 * time.Second}}, h)
|
||||
w.PollEvery(20 * time.Millisecond) // test hook: how often Wake re-checks health
|
||||
start := time.Now()
|
||||
ok := w.Wake(context.Background(), "titan")
|
||||
if !ok {
|
||||
t.Fatal("Wake must return true once the host reports healthy")
|
||||
}
|
||||
if time.Since(start) > time.Second {
|
||||
t.Errorf("Wake waited %v for a host that came up after 3 checks", time.Since(start))
|
||||
}
|
||||
_ = conn.SetReadDeadline(time.Now().Add(200 * time.Millisecond))
|
||||
buf := make([]byte, 200)
|
||||
if n, _, err := conn.ReadFromUDP(buf); err != nil || n != 102 {
|
||||
t.Fatalf("no magic packet received: %d %v", n, err)
|
||||
}
|
||||
// A second Wake inside the same window does not send again (the host is booting).
|
||||
_ = w.Wake(context.Background(), "titan")
|
||||
_ = conn.SetReadDeadline(time.Now().Add(150 * time.Millisecond))
|
||||
if n, _, err := conn.ReadFromUDP(buf); err == nil {
|
||||
t.Errorf("a second packet (%d bytes) was sent inside the wait window", n)
|
||||
}
|
||||
if w.Wake(context.Background(), "nobody") {
|
||||
t.Errorf("unknown host: Wake must return false")
|
||||
}
|
||||
}
|
||||
|
||||
func TestWakeGivesUpAfterWait(t *testing.T) {
|
||||
_, addr := listen(t)
|
||||
h := &fakeHealth{after: 1 << 30}
|
||||
w := wake.New(map[string]wake.Target{"titan": {MAC: "aa:bb:cc:dd:ee:ff", Broadcast: addr, Wait: 300 * time.Millisecond}}, h)
|
||||
w.PollEvery(20 * time.Millisecond)
|
||||
start := time.Now()
|
||||
if w.Wake(context.Background(), "titan") {
|
||||
t.Fatal("Wake must return false when the host never comes up")
|
||||
}
|
||||
if d := time.Since(start); d < 250*time.Millisecond || d > 900*time.Millisecond {
|
||||
t.Errorf("Wake returned after %v, want about the 300ms wait", d)
|
||||
}
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond)
|
||||
defer cancel()
|
||||
start = time.Now()
|
||||
if w.Wake(ctx, "titan") || time.Since(start) > 200*time.Millisecond {
|
||||
t.Errorf("a cancelled context must end the wait early (took %v)", time.Since(start))
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user