diff --git a/docs/implementer-log.md b/docs/implementer-log.md index 5e4c46a..4aae4b3 100644 --- a/docs/implementer-log.md +++ b/docs/implementer-log.md @@ -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:>` 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 /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 /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. | ? | diff --git a/internal/wake/wake.go b/internal/wake/wake.go new file mode 100644 index 0000000..4ae6c94 --- /dev/null +++ b/internal/wake/wake.go @@ -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: + } + } +} diff --git a/internal/wake/wake_test.go b/internal/wake/wake_test.go new file mode 100644 index 0000000..53e6179 --- /dev/null +++ b/internal/wake/wake_test.go @@ -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)) + } +}