diff --git a/cmd/crossbar/main.go b/cmd/crossbar/main.go index a2b592c..edcfc17 100644 --- a/cmd/crossbar/main.go +++ b/cmd/crossbar/main.go @@ -70,7 +70,7 @@ func run() error { if h.Wake == nil { continue } - targets[name] = wake.Target{MAC: h.Wake.MAC, Broadcast: h.Wake.Broadcast, Wait: h.Wake.Wait.Duration} + targets[name] = wake.Target{MAC: h.Wake.MAC, Broadcasts: h.Wake.Addresses(), Wait: h.Wake.Wait.Duration} } waker := wake.New(targets, hosts) for name, h := range cfg.Hosts { diff --git a/docs/implementer-log.md b/docs/implementer-log.md index 7aab652..aafd711 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.2/02-broadcasts | 2026-09-25 | done | 1 | pass | The Wake struct and checkWake live in `internal/config/identity.go` (added in task 04), not `config.go`, so I edited `identity.go` rather than `config.go`; `wake.go` logs a broadcast that fails to resolve/send before continuing (task rule 2 allows "logged or ignored"). | Added `Broadcasts` to `Wake` and `Wake.Addresses()` (Broadcast then Broadcasts, never empty for a parsed config); `checkWake` errors on both-set → `.broadcasts`, neither-or-empty-list → `.broadcast`, and a non-`host:port` entry → `.broadcasts`; `Target` gains `Broadcasts` and `Wake`/`sendAll` send to Broadcast then each Broadcasts in order, logging/past a failure and returning false only when no address could be sent; `main.go` fills `Target.Broadcasts` from `Wake.Addresses()` and leaves `Target.Broadcast` empty so `sendAll` does not double-send. Given `broadcasts_test.go` and `config_v22_test.go` byte-identical, v2 `wake_test.go`/`config_v2_test.go` untouched and green; `make gate` → `gate: ok` first run. | ? | | v2.2/01-route-templates | 2026-09-25 | done | 1 | fail | `internal/config` red only on `Wake.Addresses()` (task 02), the one allowed red; `go build ./...` clean, proxy/admin/health/wake/lease/store/identity/fingerprint all pass under `-race`. New `internal/config/route.go`: `templateName` pattern `^[a-z0-9][a-z0-9-]*-\*$` and `Route()` (valid-name guard excludes `*`; exact wins; else longest `"-*"`, prefix keeps the dash, non-empty remainder required, longest-prefix wins deterministically). `config.go`: the route-name check accepts a template too (one line). `proxy.go`: `route()` and `ServeHTTP` resolve both path and `X-Crossbar-Route` header forms through `cfg.Route`, and the conflicting-route check compares concrete names via `cfg.Route` (identical to before for non-template configs). `admin.go` `routeView` lists a lease under the exact key it matches or the longest template key; `admin_ops.go` `routePin` resolves through `cfg.Route` so a concrete route under a template can be pinned before its first request and the template name 404s. `main.go` identity lookup uses `cfg.Route`. All three given tests byte-identical (`config_v22_test.go` keeps `TestWakeBroadcasts`, which is why config is red). | ? | | v2.1/02-props-loaded-only | 2026-09-25 | done | 1 | pass | `movedHeader` separator `><`→`>` and the `CtxHeader` doc comment in `proxy.go`, both forced by the given router test (`moved:small>big`) which the task text did not mention; no production code parses the separator (`forward.go` passes it straight through) so it is safe. | Implemented per-model context. `health`: added `ModelCtx` and a `Models map[string]ModelCtx` field on `Status`, plus `PerSlotCtxFor(model)` (per-model figure when present, else host-level `PerSlotCtx` for a loaded model, else 0); moved `props` into a new `props.go` and added `propsModel`/`propsModels`. Poller rules 1-4: `/v1/models` treats an entry as loaded only with no `status` or `status.value=="loaded"` (other values dropped from `Loaded`); plain `/props` with `role:router` leaves host NCtx/Slots 0; each loaded model is asked `GET /props?model=` and a failed/malformed answer leaves that id absent without failing the host; `Models` is a fresh non-nil map every successful poll, `MarkDown` leaves it. `ctxguard.go`: every `PerSlotCtx()` became `PerSlotCtxFor(model)` (leased host, candidates, wake "cannot serve" check) and `largestSlotCtx(hosts,h,model)` counts only hosts that have it loaded. `admin.go`: `HostView` gains `models` (empty object, never null). A plain single server keeps working as v2. All three given tests byte-identical; `make gate` → `gate: ok` on the first run. | ? | | v2.1/01-cancel-record | 2026-09-25 | done | 1 | pass | none | Implemented the rule: added a `cancelled` field to `forwardState`; the `ErrorHandler` sets it when it observes `context.Canceled` (client gone before any response byte) so the delivered row is no longer turned into a 499 by a pooled close after the body; removed the post-hoc `r.Context().Err()` check in the normal path, leaving the recover path's `http.ErrAbortHandler` (mid-body) check as the other 499 source. Given test failed the first run (`Errors:7`, status counts held 25×200/7×499), passes 3× under `-race`; `TestClientCancelMidStreamIsRecorded`, `TestClientCancelWhileQueuedIsRecorded` and `TestQueueFullIs503` still pass; `forward.go` 230 lines; `make gate` printed `gate: ok` on the first run. | ? | diff --git a/example.toml b/example.toml index 03f0372..4c3a1ff 100644 --- a/example.toml +++ b/example.toml @@ -19,6 +19,7 @@ models = { "ornith-1.5-35b-a3b" = { parallel = 2 } } [hosts.beta.wake] # v2: wake a sleeping host when nothing else can take a new lease mac = "aa:bb:cc:dd:ee:02" broadcast = "127.0.0.1:19082" # the LAN broadcast address, port 9, in production +# broadcasts = ["127.0.0.1:19082", "127.0.0.1:19083"] # v2.2: a roaming host on several networks; this or broadcast, not both, and at least one wait = "20s" # v1: a route is a set of candidate hosts; each conversation gets a sticky lease on the host with diff --git a/internal/config/identity.go b/internal/config/identity.go index a83ef76..c22db2e 100644 --- a/internal/config/identity.go +++ b/internal/config/identity.go @@ -7,11 +7,24 @@ import ( ) // Wake is the magic-wake pattern sent to a host to rouse it: its MAC, the -// broadcast address to aim at, and how long to wait for the answer. +// broadcast address(s) to aim at, and how long to wait for the answer. A host +// that roams between networks names several, so Broadcast (one) and Broadcasts +// (a list) are alternatives: exactly one must be set. type Wake struct { - MAC string `toml:"mac"` - Broadcast string `toml:"broadcast"` - Wait Duration `toml:"wait"` + MAC string `toml:"mac"` + Broadcast string `toml:"broadcast"` + Broadcasts []string `toml:"broadcasts"` + Wait Duration `toml:"wait"` +} + +// Addresses is Broadcast (when set) followed by Broadcasts: the ordered list to +// send wake packets to, never empty for a parsed config. +func (w *Wake) Addresses() []string { + addrs := make([]string, 0, 1+len(w.Broadcasts)) + if w.Broadcast != "" { + addrs = append(addrs, w.Broadcast) + } + return append(addrs, w.Broadcasts...) } const ( @@ -40,8 +53,19 @@ func (c *Config) checkWake() *Error { return &Error{Field: wakeField + ".mac", Msg: "must be a MAC address"} } - if _, _, err := net.SplitHostPort(w.Broadcast); err != nil || w.Broadcast == "" { + broadcastSet := w.Broadcast != "" + broadcastsSet := len(w.Broadcasts) > 0 + switch { + case broadcastSet && broadcastsSet: + return &Error{Field: wakeField + ".broadcasts", Msg: "choose broadcast or broadcasts, not both"} + case !broadcastSet && !broadcastsSet: return &Error{Field: wakeField + ".broadcast", Msg: "must be a non-empty host:port"} + default: + for _, a := range w.Broadcasts { + if _, _, err := net.SplitHostPort(a); err != nil || a == "" { + return &Error{Field: wakeField + ".broadcasts", Msg: "must be a non-empty host:port"} + } + } } if w.Wait.Duration == 0 { diff --git a/internal/wake/broadcasts_test.go b/internal/wake/broadcasts_test.go new file mode 100644 index 0000000..238ca53 --- /dev/null +++ b/internal/wake/broadcasts_test.go @@ -0,0 +1,73 @@ +package wake_test + +import ( + "net" + "testing" + "time" + + "git.wntrmute.dev/kyle/crossbar/internal/wake" +) + +// listener returns a UDP socket on 127.0.0.1 and a channel that gets one value per datagram. +func listener(t *testing.T) (string, <-chan []byte) { + t.Helper() + pc, err := net.ListenPacket("udp4", "127.0.0.1:0") + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { pc.Close() }) + got := make(chan []byte, 4) + go func() { + buf := make([]byte, 256) + for { + n, _, err := pc.ReadFrom(buf) + if err != nil { + return + } + b := make([]byte, n) + copy(b, buf[:n]) + got <- b + } + }() + return pc.LocalAddr().String(), got +} + +func expectPacket(t *testing.T, name string, got <-chan []byte) { + t.Helper() + select { + case b := <-got: + if len(b) != 102 { + t.Errorf("%s: got %d bytes, want a 102-byte magic packet", name, len(b)) + } + case <-time.After(2 * time.Second): + t.Errorf("%s: no packet within two seconds", name) + } +} + +// A target may name several broadcast addresses (a host that roams between two networks): the +// packet goes to every one of them, and one address that cannot be resolved does not stop the +// others. +func TestWakeSendsToEveryBroadcast(t *testing.T) { + a, gotA := listener(t) + b, gotB := listener(t) + h := &fakeHealth{after: 1 << 30} // never healthy + w := wake.New(map[string]wake.Target{"titan": {MAC: "aa:bb:cc:dd:ee:ff", Broadcasts: []string{a, "256.1.1.1:9", b}, Wait: 300 * time.Millisecond}}, h) + w.PollEvery(20 * time.Millisecond) + if w.Wake(t.Context(), "titan") { + t.Errorf("Wake must report false when the host never comes up") + } + expectPacket(t, "first address", gotA) + expectPacket(t, "third address, after an unresolvable second", gotB) +} + +// The single-address form keeps working, alone or together with the list. +func TestWakeBroadcastAndBroadcastsCombine(t *testing.T) { + a, gotA := listener(t) + b, gotB := listener(t) + h := &fakeHealth{after: 1 << 30} // never healthy + w := wake.New(map[string]wake.Target{"titan": {MAC: "aa:bb:cc:dd:ee:ff", Broadcast: a, Broadcasts: []string{b}, Wait: 300 * time.Millisecond}}, h) + w.PollEvery(20 * time.Millisecond) + w.Wake(t.Context(), "titan") + expectPacket(t, "Broadcast", gotA) + expectPacket(t, "Broadcasts[0]", gotB) +} diff --git a/internal/wake/wake.go b/internal/wake/wake.go index 4ae6c94..0331254 100644 --- a/internal/wake/wake.go +++ b/internal/wake/wake.go @@ -7,15 +7,25 @@ package wake import ( "context" "fmt" + "log/slog" "net" "sync" "time" ) -// Target describes how to wake one named host. +// log reports a broadcast that fails to resolve or send. Wake continues past +// such failures (rule: one dead address must not stop the others), so this is +// the only place the package logs; the address and error are not request data. +var log = slog.New(slog.Default().Handler()) + +// Target describes how to wake one named host. Broadcast is the single-address +// form (as before); Broadcasts names more than one (a host that roams between +// networks). Wake sends to Broadcast (if set) and then each of Broadcasts. type Target struct { - MAC, Broadcast string // MAC "aa:bb:cc:dd:ee:ff" (any separator, any case); Broadcast "host:port" - Wait time.Duration + MAC string + Broadcast string + Broadcasts []string + Wait time.Duration } // Health reports whether a named host is currently healthy. Implementations must @@ -67,6 +77,26 @@ func Send(mac, broadcast string) error { return nil } +// sendAll emits one magic packet for mac to broadcast (if non-empty) and then +// to each address in the rest, in order. An address that fails to resolve or +// send is logged and does not stop the others; it reports whether at least one +// packet went out. +func sendAll(mac, broadcast string, rest []string) bool { + addrs := append([]string{broadcast}, rest...) + sent := false + for _, addr := range addrs { + if addr == "" { + continue + } + if err := Send(mac, addr); err != nil { + log.Error("wake broadcast failed", "addr", addr, "err", err) + continue + } + sent = true + } + return sent +} + // 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 { @@ -95,10 +125,13 @@ func (w *Waker) PollEvery(d time.Duration) { 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. +// Wake sends a magic packet for host to every broadcast address — Broadcast +// (if set) then each of Broadcasts, in order — if none was sent in the last +// Wait, then polls health until the host is healthy, the wait elapses, or ctx is +// done. A broadcast that fails to resolve or send is logged and does not stop +// the others; Wake returns true only when the host becomes healthy, and false +// for an unknown host, on timeout, when ctx ends first, or when no address +// could be sent to. func (w *Waker) Wake(ctx context.Context, host string) bool { w.mu.Lock() target, ok := w.targets[host] @@ -110,7 +143,7 @@ func (w *Waker) Wake(ctx context.Context, host string) bool { 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 { + if !sendAll(target.MAC, target.Broadcast, target.Broadcasts) { return false } w.mu.Lock()