diff --git a/README.md b/README.md index d03fefd..e1b03f0 100644 --- a/README.md +++ b/README.md @@ -33,6 +33,7 @@ serves it. Templates are embedded defaults, overridable from a directory | `GET /ipxe/{mac}` · `GET /boot/ipxe?mac=` | iPXE boot script | | `GET /ks/{ident}` | rendered kickstart (MAC or hostname) | | `POST /provisioned/{ident}` | end-of-kickstart callback (token) → clears `pxe_enabled` in NetBox | +| `POST /logs` | installer log relay (token) → VictoriaLogs `vlinsert` | | `GET /healthz` · `/readyz` · `/metrics` | health + Prometheus | The boot path is served over **plain HTTP** (PXE installers have no internal-CA @@ -64,7 +65,7 @@ rationale in [docs/endpoints.md](docs/endpoints.md). Env-based (12-factor), see [`config.example.env`](config.example.env). Key vars: `BOOTAPI_NETBOX_URL`, `BOOTAPI_NETBOX_TOKEN[_FILE]`, `BOOTAPI_BASE_URL`, `BOOTAPI_BOOT_BASE_URL`, `BOOTAPI_ROOT_PASSWORD_HASH[_FILE]`, -`BOOTAPI_UNKNOWN_MAC_FALLBACK`. +`BOOTAPI_UNKNOWN_MAC_FALLBACK`, `BOOTAPI_VLINSERT_URL`. ## Development diff --git a/cmd/bootapi/main.go b/cmd/bootapi/main.go index 021ce0e..9e644e9 100644 --- a/cmd/bootapi/main.go +++ b/cmd/bootapi/main.go @@ -47,7 +47,10 @@ func main() { slog.Warn("no NetBox token set (BOOTAPI_NETBOX_TOKEN/_FILE); NetBox reads will likely be denied") } if cfg.ProvisionToken == "" { - slog.Warn("no BOOTAPI_PROVISION_TOKEN set; the /provisioned callback is disabled (pxe_enabled will not auto-clear)") + slog.Warn("no BOOTAPI_PROVISION_TOKEN set; the /provisioned callback and /logs relay are disabled") + } + if cfg.VLInsertURL == "" { + slog.Warn("BOOTAPI_VLINSERT_URL is empty; the /logs installer log relay is disabled") } rcfg := render.RenderConfig{ @@ -92,6 +95,8 @@ func main() { Cache: cache, UnknownMACFallback: cfg.UnknownMACFallback, ProvisionToken: cfg.ProvisionToken, + VLInsertURL: cfg.VLInsertURL, + VLInsertTimeout: cfg.VLInsertTimeout, TLSAddr: cfg.TLSListenAddr, TLSCertFile: cfg.TLSCertFile, TLSKeyFile: cfg.TLSKeyFile, diff --git a/config.example.env b/config.example.env index 5d1c89b..fac1bc4 100644 --- a/config.example.env +++ b/config.example.env @@ -55,6 +55,15 @@ BOOTAPI_ARTIFACT_BASE_URL=https://artifactapi.k8s.syd1.au.unkin.net/api/v1/remot BOOTAPI_PROVISION_TOKEN= # BOOTAPI_PROVISION_TOKEN_FILE=/var/run/secrets/bootapi/provision_token +# --- installer log relay (POST /logs -> VictoriaLogs vlinsert) --- +# A host being discovered or installed is not in k8s (no vlagent) and has no +# internal-CA trust or credentials for the HTTPS log ingest, so bootapi relays +# its newline-delimited JSON logs. Guarded by BOOTAPI_PROVISION_TOKEN. +# An EMPTY value disables the endpoint (it then 503s). +BOOTAPI_VLINSERT_URL=http://vlinsert-logs.logging.svc.cluster.local:9481 +# Short on purpose: logs are best-effort, an install must not wait on the sink. +BOOTAPI_VLINSERT_TIMEOUT=5s + # --- puppet bootstrap targets (k8s puppetserver; baked into kickstart %post) --- BOOTAPI_PUPPET_SERVER=puppet.k8s.syd1.au.unkin.net BOOTAPI_PUPPET_CA_SERVER=puppetca.k8s.syd1.au.unkin.net diff --git a/docs/endpoints.md b/docs/endpoints.md index f4ae07e..fe680a0 100644 --- a/docs/endpoints.md +++ b/docs/endpoints.md @@ -35,6 +35,7 @@ of the install. | GET | `/boot/ipxe?mac=...` | Query-string alias of `/ipxe/{mac}`. | | GET | `/ks/{ident}` | Rendered kickstart. `{ident}` is a MAC (auto-detected) or a hostname; trailing `.ks`/`.cfg` is stripped. | | POST | `/provisioned/{ident}` | End-of-kickstart callback; clears `pxe_enabled` in NetBox. **Token-guarded** (`Authorization: Bearer `). | +| POST | `/logs` | Relays a host's newline-delimited JSON install logs to VictoriaLogs. **Token-guarded** (same token). | | GET | `/healthz` | Liveness: always `200 ok`. | | GET | `/readyz` | Readiness: `200` once templates parsed. Does **not** probe NetBox. | | GET | `/metrics` | Prometheus metrics (see below). | @@ -85,6 +86,33 @@ bad/missing token, `404` for an unknown host, `503` when no failure. The default kickstart templates call it from `%post` over plain HTTP (the token authenticates the call; no CA trust needed at install time). +## The installer log relay + +`POST /logs` exists because a host being PXE-discovered or installed is not in +Kubernetes, so the cluster's vlagent cannot collect its logs — if an install +fails, the only record dies with the machine. That environment also has no +internal-CA trust and no credentials for the HTTPS-only log ingest, and bootapi +is already the plain-HTTP broker it can reach, so bootapi forwards for it. + +- Body: newline-delimited JSON, one log record per line (`application/x-ndjson` + or `application/json`), capped at **1 MiB**. +- Auth: the same `BOOTAPI_PROVISION_TOKEN` as `/provisioned`, same fail-closed + behavior — `503` when no token (or no `BOOTAPI_VLINSERT_URL`) is configured, + `401` on a bad/missing token. +- Forwarded as one short-timeout POST to + `{BOOTAPI_VLINSERT_URL}/insert/jsonline?_stream_fields=serial,phase&_msg_field=msg&_time_field=time`. + `serial` and `phase` are constant for a run and low-cardinality, so they key + the stream; MAC is per-NIC and goes in `extra_fields` so it stays searchable + without multiplying streams. +- `extra_fields` carries only what bootapi *observes* rather than what the + client claims: the request's source IP (`src_ip`), plus the resolved NetBox + device name (`device`) when `?mac=` resolves. Resolution is optional — a miss + just omits the field. +- Responses: `202` once the batch is accepted, `400` on an empty or oversized + body. A vlinsert failure is logged and counted but **still returns `202`** — + no retries, no buffering: a host must never block its install because the log + sink is down. + ## Metrics All on `/metrics`, prefix `bootapi_`: @@ -95,6 +123,8 @@ All on `/metrics`, prefix `bootapi_`: - `bootapi_netbox_lookup_duration_seconds{field}` — histogram. - `bootapi_netbox_cache_hits_total` / `bootapi_netbox_cache_misses_total`. - `bootapi_provisioned_total{result}` — result = `ok|unauthorized|notfound|error|disabled`. +- `bootapi_log_relay_total{result}` — result = `ok|error|unauthorized|disabled`. +- `bootapi_log_relay_lines_total` — installer log lines successfully relayed. - `bootapi_ipxe_gated_total` — known hosts served local-boot because `pxe_enabled=false`. - `bootapi_template_sync_total` / `bootapi_template_sync_failures_total` / `bootapi_template_generation` — template git-sync (see [template-authoring.md](template-authoring.md)). - standard Go/process collectors. diff --git a/internal/config/config.go b/internal/config/config.go index 8989e8c..8a35076 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -79,6 +79,15 @@ type Config struct { // endpoint (fail closed). Prefer ProvisionTokenFile in k8s. ProvisionToken string + // VLInsertURL is the VictoriaLogs vlinsert base that POST /logs relays + // installer logs to. Empty disables the endpoint (it then 503s): a host + // being installed has no vlagent and no credentials for the HTTPS ingest, + // so bootapi relays for it. + VLInsertURL string + // VLInsertTimeout bounds the single best-effort forward attempt. Short on + // purpose — an install must not wait on the log sink. + VLInsertTimeout time.Duration + // PuppetServer / PuppetCAServer are baked into kickstart %post so the // freshly-installed host checks in to the k8s puppetserver. PuppetServer string @@ -118,6 +127,16 @@ func Load() (*Config, error) { if err != nil { return nil, fmt.Errorf("invalid BOOTAPI_TEMPLATE_GIT_INTERVAL: %w", err) } + vlTimeout, err := time.ParseDuration(getenv("BOOTAPI_VLINSERT_TIMEOUT", "5s")) + if err != nil { + return nil, fmt.Errorf("invalid BOOTAPI_VLINSERT_TIMEOUT: %w", err) + } + // Unlike the other URLs, an explicitly EMPTY value is meaningful here: it + // disables the log relay, so LookupEnv rather than getenv. + vlURL, ok := os.LookupEnv("BOOTAPI_VLINSERT_URL") + if !ok { + vlURL = "http://vlinsert-logs.logging.svc.cluster.local:9481" + } token, err := readSecret("BOOTAPI_NETBOX_TOKEN") if err != nil { @@ -169,6 +188,8 @@ func Load() (*Config, error) { ArtifactBaseURL: strings.TrimRight(getenv("BOOTAPI_ARTIFACT_BASE_URL", "https://artifactapi.k8s.syd1.au.unkin.net/api/v1/remote"), "/"), BootBaseURL: strings.TrimRight(os.Getenv("BOOTAPI_BOOT_BASE_URL"), "/"), ProvisionToken: provToken, + VLInsertURL: strings.TrimRight(vlURL, "/"), + VLInsertTimeout: vlTimeout, PuppetServer: getenv("BOOTAPI_PUPPET_SERVER", "puppet.k8s.syd1.au.unkin.net"), PuppetCAServer: getenv("BOOTAPI_PUPPET_CA_SERVER", "puppetca.k8s.syd1.au.unkin.net"), PuppetCAURL: getenv("BOOTAPI_PUPPET_CA_URL", "puppetca.k8s.syd1.au.unkin.net"), diff --git a/internal/config/config_test.go b/internal/config/config_test.go index 3faad79..ddacf75 100644 --- a/internal/config/config_test.go +++ b/internal/config/config_test.go @@ -40,6 +40,34 @@ func TestLoadDefaults(t *testing.T) { if c.ArtifactBaseURL != "https://artifactapi.k8s.syd1.au.unkin.net/api/v1/remote" { t.Errorf("ArtifactBaseURL = %q", c.ArtifactBaseURL) } + if c.VLInsertURL != "http://vlinsert-logs.logging.svc.cluster.local:9481" { + t.Errorf("VLInsertURL = %q", c.VLInsertURL) + } + if c.VLInsertTimeout != 5*time.Second { + t.Errorf("VLInsertTimeout = %v, want 5s", c.VLInsertTimeout) + } +} + +func TestVLInsertURLEmptyDisablesRelay(t *testing.T) { + clearEnv(t) + // An explicitly empty value must NOT fall back to the default: it disables + // the /logs relay. + t.Setenv("BOOTAPI_VLINSERT_URL", "") + c, err := Load() + if err != nil { + t.Fatal(err) + } + if c.VLInsertURL != "" { + t.Errorf("VLInsertURL = %q, want empty (relay disabled)", c.VLInsertURL) + } +} + +func TestVLInsertBadTimeout(t *testing.T) { + clearEnv(t) + t.Setenv("BOOTAPI_VLINSERT_TIMEOUT", "soon") + if _, err := Load(); err == nil { + t.Fatal("expected error for invalid BOOTAPI_VLINSERT_TIMEOUT") + } } func TestCallbackBaseDefaultsToBase(t *testing.T) { diff --git a/internal/server/metrics.go b/internal/server/metrics.go index e16e9d6..7f733b1 100644 --- a/internal/server/metrics.go +++ b/internal/server/metrics.go @@ -28,6 +28,8 @@ type metrics struct { netboxLookups *prometheus.CounterVec // by field,result netboxDuration *prometheus.HistogramVec provisioned *prometheus.CounterVec // by result + logRelay *prometheus.CounterVec // by result + logLines prometheus.Counter ipxeGated prometheus.Counter } @@ -56,12 +58,20 @@ func newMetrics(cache cacheStats, git gitStats) *metrics { Name: "bootapi_provisioned_total", Help: "Provisioned callbacks, by result (ok|unauthorized|notfound|error|disabled).", }, []string{"result"}), + logRelay: prometheus.NewCounterVec(prometheus.CounterOpts{ + Name: "bootapi_log_relay_total", + Help: "Installer log batches relayed to VictoriaLogs, by result (ok|error|unauthorized|disabled).", + }, []string{"result"}), + logLines: prometheus.NewCounter(prometheus.CounterOpts{ + Name: "bootapi_log_relay_lines_total", + Help: "Installer log lines successfully relayed to VictoriaLogs.", + }), ipxeGated: prometheus.NewCounter(prometheus.CounterOpts{ Name: "bootapi_ipxe_gated_total", Help: "Known hosts served the local-boot fallback because pxe_enabled=false.", }), } - reg.MustRegister(m.httpRequests, m.renders, m.netboxLookups, m.netboxDuration, m.provisioned, m.ipxeGated) + reg.MustRegister(m.httpRequests, m.renders, m.netboxLookups, m.netboxDuration, m.provisioned, m.logRelay, m.logLines, m.ipxeGated) if cache != nil { reg.MustRegister(newCacheCollector(cache)) } diff --git a/internal/server/server.go b/internal/server/server.go index fb098b8..a900a08 100644 --- a/internal/server/server.go +++ b/internal/server/server.go @@ -3,12 +3,16 @@ package server import ( + "bytes" "context" "crypto/subtle" "errors" "fmt" + "io" "log/slog" + "net" "net/http" + "net/url" "strings" "sync" "time" @@ -30,9 +34,15 @@ type Server struct { // fallback is the unknown-MAC iPXE behavior: "local" (safe default) or // "shell" (debug). fallback string - // provisionToken guards POST /provisioned; empty disables the endpoint. + // provisionToken guards POST /provisioned and POST /logs; empty disables + // both endpoints. provisionToken string + // vlinsertURL is the VictoriaLogs vlinsert base POST /logs relays to; + // empty disables the endpoint. vlClient bounds the single forward attempt. + vlinsertURL string + vlClient *http.Client + // TLS listener (optional); the plain-HTTP listener is always on. tlsAddr string tlsCert string @@ -48,6 +58,8 @@ type Options struct { GitStats gitStats UnknownMACFallback string ProvisionToken string + VLInsertURL string + VLInsertTimeout time.Duration TLSAddr string TLSCertFile string TLSKeyFile string @@ -59,12 +71,18 @@ func New(o Options) *Server { if fb == "" { fb = "local" } + timeout := o.VLInsertTimeout + if timeout <= 0 { + timeout = 5 * time.Second + } return &Server{ nb: o.NetBox, engine: o.Engine, metrics: newMetrics(o.Cache, o.GitStats), fallback: fb, provisionToken: o.ProvisionToken, + vlinsertURL: strings.TrimRight(o.VLInsertURL, "/"), + vlClient: &http.Client{Timeout: timeout}, tlsAddr: o.TLSAddr, tlsCert: o.TLSCertFile, tlsKey: o.TLSKeyFile, @@ -92,6 +110,10 @@ func (s *Server) Router() http.Handler { // End-of-kickstart callback: flips pxe_enabled off in NetBox. Token-guarded. r.Post("/provisioned/{ident}", s.handleProvisioned) + // Installer log relay: a host being installed is not in Kubernetes, so + // vlagent cannot collect its logs. Token-guarded like /provisioned. + r.Post("/logs", s.handleLogs) + return r } @@ -253,6 +275,96 @@ func (s *Server) handleProvisioned(w http.ResponseWriter, r *http.Request) { w.WriteHeader(http.StatusNoContent) } +// maxLogBody caps a relayed batch. An install's log stream is a handful of +// lines per phase; bootapi is a relay, not a log buffer. +const maxLogBody = 1 << 20 // 1 MiB + +// handleLogs relays newline-delimited JSON log records from a host that is +// PXE-discovering or installing into VictoriaLogs. Such a host is not in +// Kubernetes (no vlagent) and has no internal-CA trust or credentials for the +// HTTPS log ingest, so bootapi — the plain-HTTP broker it can already reach — +// forwards them. Best-effort: a dead log sink must never block an install, so +// every outcome after authentication is 202. +func (s *Server) handleLogs(w http.ResponseWriter, r *http.Request) { + if s.vlinsertURL == "" || s.provisionToken == "" { + http.Error(w, "log relay disabled: no vlinsert URL or no token configured", http.StatusServiceUnavailable) + s.metrics.logRelay.WithLabelValues("disabled").Inc() + return + } + if subtle.ConstantTimeCompare([]byte(bearer(r)), []byte(s.provisionToken)) != 1 { + http.Error(w, "invalid or missing provision token", http.StatusUnauthorized) + s.metrics.logRelay.WithLabelValues("unauthorized").Inc() + return + } + body, err := io.ReadAll(http.MaxBytesReader(w, r.Body, maxLogBody)) + if err != nil { + http.Error(w, "log batch unreadable or larger than 1MiB", http.StatusBadRequest) + s.metrics.logRelay.WithLabelValues("error").Inc() + return + } + if len(bytes.TrimSpace(body)) == 0 { + http.Error(w, "empty log batch", http.StatusBadRequest) + s.metrics.logRelay.WithLabelValues("error").Inc() + return + } + + lines := bytes.Count(bytes.TrimRight(body, "\n"), []byte("\n")) + 1 + if err := s.relayLogs(r.Context(), body, s.observedFields(r)); err != nil { + slog.Error("log relay to vlinsert failed; dropping batch", "lines", lines, "err", err) + s.metrics.logRelay.WithLabelValues("error").Inc() + } else { + s.metrics.logRelay.WithLabelValues("ok").Inc() + s.metrics.logLines.Add(float64(lines)) + } + s.ok(w, http.StatusAccepted, "text/plain", []byte("accepted\n"), "logs") +} + +// relayLogs makes one short-timeout POST to vlinsert's jsonline endpoint. +// serial and phase are constant for a run and low-cardinality, so they key the +// stream; MAC is per-NIC and rides in extra_fields so it stays searchable +// without multiplying streams. +func (s *Server) relayLogs(ctx context.Context, body []byte, extra string) error { + q := url.Values{ + "_stream_fields": {"serial,phase"}, + "_msg_field": {"msg"}, + "_time_field": {"time"}, + } + if extra != "" { + q.Set("extra_fields", extra) + } + req, err := http.NewRequestWithContext(ctx, http.MethodPost, s.vlinsertURL+"/insert/jsonline?"+q.Encode(), bytes.NewReader(body)) + if err != nil { + return err + } + req.Header.Set("Content-Type", "application/x-ndjson") + resp, err := s.vlClient.Do(req) + if err != nil { + return err + } + defer func() { _ = resp.Body.Close() }() + _, _ = io.Copy(io.Discard, resp.Body) + if resp.StatusCode >= http.StatusMultipleChoices { + return fmt.Errorf("vlinsert: %s", resp.Status) + } + return nil +} + +// observedFields builds the extra_fields value from what bootapi observes +// rather than what the client claims: the source IP, plus the NetBox device +// name when ?mac= resolves. Resolution is optional — a miss just omits it. +func (s *Server) observedFields(r *http.Request) string { + fields := []string{} + if ip, _, err := net.SplitHostPort(r.RemoteAddr); err == nil && ip != "" { + fields = append(fields, "src_ip="+ip) + } + if mac := r.URL.Query().Get("mac"); looksLikeMAC(mac) { + if host, err := s.lookup(r.Context(), "mac", mac); err == nil { + fields = append(fields, "device="+host.Hostname) + } + } + return strings.Join(fields, ",") +} + // bearer extracts a token from "Authorization: Bearer " or a bare "token" // header. func bearer(r *http.Request) string { diff --git a/internal/server/server_test.go b/internal/server/server_test.go index 82afaf3..3d535c6 100644 --- a/internal/server/server_test.go +++ b/internal/server/server_test.go @@ -2,10 +2,13 @@ package server import ( "context" + "io" "net/http" "net/http/httptest" + "net/url" "strings" "testing" + "time" "git.unkin.net/unkin/bootapi/internal/model" "git.unkin.net/unkin/bootapi/internal/netbox" @@ -68,6 +71,11 @@ func newTestServer(t *testing.T, nb netbox.API, fallback string) *Server { } func newTestServerToken(t *testing.T, nb netbox.API, fallback, provToken string) *Server { + t.Helper() + return newTestServerLogs(t, nb, fallback, provToken, "") +} + +func newTestServerLogs(t *testing.T, nb netbox.API, fallback, provToken, vlURL string) *Server { t.Helper() set, err := render.BuildSet(templates.FS, nil) if err != nil { @@ -80,7 +88,10 @@ func newTestServerToken(t *testing.T, nb netbox.API, fallback, provToken string) DefaultDomain: "main.unkin.net", DefaultTemplate: "almalinux9", RootPasswordHash: "$6$abc$def", }, set) - return New(Options{NetBox: nb, Engine: eng, UnknownMACFallback: fallback, ProvisionToken: provToken}) + return New(Options{ + NetBox: nb, Engine: eng, UnknownMACFallback: fallback, + ProvisionToken: provToken, VLInsertURL: vlURL, VLInsertTimeout: 2 * time.Second, + }) } func do(t *testing.T, h http.Handler, path string) *httptest.ResponseRecorder { @@ -291,3 +302,187 @@ func TestLooksLikeMAC(t *testing.T) { } } } + +// vlStub is a stand-in vlinsert that records what bootapi forwarded. +type vlStub struct { + srv *httptest.Server + hits int + path string + query url.Values + body string + ctype string + status int +} + +func newVLStub(t *testing.T, status int) *vlStub { + t.Helper() + v := &vlStub{status: status} + v.srv = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + b, err := io.ReadAll(r.Body) + if err != nil { + t.Errorf("vlinsert stub read body: %v", err) + } + v.hits++ + v.path, v.query, v.body, v.ctype = r.URL.Path, r.URL.Query(), string(b), r.Header.Get("Content-Type") + w.WriteHeader(v.status) + })) + t.Cleanup(v.srv.Close) + return v +} + +func postLogs(t *testing.T, h http.Handler, path, token, body string) *httptest.ResponseRecorder { + t.Helper() + rec := httptest.NewRecorder() + req := httptest.NewRequest(http.MethodPost, path, strings.NewReader(body)) + req.Header.Set("Content-Type", "application/x-ndjson") + if token != "" { + req.Header.Set("Authorization", "Bearer "+token) + } + h.ServeHTTP(rec, req) + return rec +} + +func TestLogRelayForwardsToVLInsert(t *testing.T) { + vl := newVLStub(t, http.StatusNoContent) + res := &fakeNB{byMAC: map[string]*model.Host{"aa:bb:cc:00:11:22": testHost()}} + h := newTestServerLogs(t, res, "local", "prov-secret", vl.srv.URL).Router() + + batch := `{"time":"2026-10-03T00:00:00Z","serial":"SN1","phase":"discovery","msg":"hello"}` + rec := postLogs(t, h, "/logs?mac=aa:bb:cc:00:11:22", "prov-secret", batch+"\n") + if rec.Code != http.StatusAccepted { + t.Fatalf("status = %d, want 202\n%s", rec.Code, rec.Body.String()) + } + if vl.hits != 1 { + t.Fatalf("vlinsert hits = %d, want 1", vl.hits) + } + if vl.path != "/insert/jsonline" { + t.Errorf("path = %q", vl.path) + } + for k, want := range map[string]string{ + "_stream_fields": "serial,phase", + "_msg_field": "msg", + "_time_field": "time", + } { + if got := vl.query.Get(k); got != want { + t.Errorf("query %s = %q, want %q", k, got, want) + } + } + // extra_fields carries only what bootapi observed, not client claims. + if got := vl.query.Get("extra_fields"); got != "src_ip=192.0.2.1,device=web01" { + t.Errorf("extra_fields = %q", got) + } + if strings.TrimSpace(vl.body) != batch { + t.Errorf("forwarded body = %q", vl.body) + } + if vl.ctype != "application/x-ndjson" { + t.Errorf("forwarded content-type = %q", vl.ctype) + } +} + +func TestLogRelayUnresolvedMACStillForwards(t *testing.T) { + vl := newVLStub(t, http.StatusNoContent) + h := newTestServerLogs(t, &fakeNB{}, "local", "prov-secret", vl.srv.URL).Router() + + rec := postLogs(t, h, "/logs?mac=de:ad:be:ef:00:00", "prov-secret", `{"msg":"x"}`+"\n") + if rec.Code != http.StatusAccepted || vl.hits != 1 { + t.Fatalf("status = %d hits = %d, want 202/1", rec.Code, vl.hits) + } + if got := vl.query.Get("extra_fields"); got != "src_ip=192.0.2.1" { + t.Errorf("unresolved MAC must omit device: extra_fields = %q", got) + } +} + +func TestLogRelayAuth(t *testing.T) { + vl := newVLStub(t, http.StatusNoContent) + h := newTestServerLogs(t, &fakeNB{}, "local", "prov-secret", vl.srv.URL).Router() + + if rec := postLogs(t, h, "/logs", "wrong", `{"msg":"x"}`); rec.Code != http.StatusUnauthorized { + t.Errorf("wrong token: status = %d, want 401", rec.Code) + } + if rec := postLogs(t, h, "/logs", "", `{"msg":"x"}`); rec.Code != http.StatusUnauthorized { + t.Errorf("no token: status = %d, want 401", rec.Code) + } + if vl.hits != 0 { + t.Errorf("unauthorized calls must not forward, hits = %d", vl.hits) + } +} + +func TestLogRelayDisabledWithoutProvisionToken(t *testing.T) { + vl := newVLStub(t, http.StatusNoContent) + h := newTestServerLogs(t, &fakeNB{}, "local", "", vl.srv.URL).Router() + if rec := postLogs(t, h, "/logs", "anything", `{"msg":"x"}`); rec.Code != http.StatusServiceUnavailable { + t.Errorf("status = %d, want 503 when no token configured", rec.Code) + } + if vl.hits != 0 { + t.Errorf("disabled relay must not forward, hits = %d", vl.hits) + } +} + +func TestLogRelayDisabledWithoutVLInsertURL(t *testing.T) { + h := newTestServerLogs(t, &fakeNB{}, "local", "prov-secret", "").Router() + if rec := postLogs(t, h, "/logs", "prov-secret", `{"msg":"x"}`); rec.Code != http.StatusServiceUnavailable { + t.Errorf("status = %d, want 503 when BOOTAPI_VLINSERT_URL is empty", rec.Code) + } +} + +func TestLogRelayEmptyBody(t *testing.T) { + vl := newVLStub(t, http.StatusNoContent) + h := newTestServerLogs(t, &fakeNB{}, "local", "prov-secret", vl.srv.URL).Router() + if rec := postLogs(t, h, "/logs", "prov-secret", "\n \n"); rec.Code != http.StatusBadRequest { + t.Errorf("status = %d, want 400 for an empty batch", rec.Code) + } + if vl.hits != 0 { + t.Errorf("empty batch must not forward, hits = %d", vl.hits) + } +} + +func TestLogRelayOversizedBody(t *testing.T) { + vl := newVLStub(t, http.StatusNoContent) + h := newTestServerLogs(t, &fakeNB{}, "local", "prov-secret", vl.srv.URL).Router() + big := strings.Repeat("x", maxLogBody+1) + if rec := postLogs(t, h, "/logs", "prov-secret", big); rec.Code != http.StatusBadRequest { + t.Errorf("status = %d, want 400 for an oversized batch", rec.Code) + } + if vl.hits != 0 { + t.Errorf("oversized batch must not forward, hits = %d", vl.hits) + } +} + +func TestLogRelayVLInsertFailureStillAccepts(t *testing.T) { + vl := newVLStub(t, http.StatusInternalServerError) + srv := newTestServerLogs(t, &fakeNB{}, "local", "prov-secret", vl.srv.URL) + h := srv.Router() + + // A dead log sink must never block an install. + if rec := postLogs(t, h, "/logs", "prov-secret", `{"msg":"x"}`+"\n"); rec.Code != http.StatusAccepted { + t.Fatalf("status = %d, want 202 even when vlinsert fails", rec.Code) + } + body := do(t, h, "/metrics").Body.String() + if !strings.Contains(body, `bootapi_log_relay_total{result="error"} 1`) { + t.Error("error result not counted") + } + if strings.Contains(body, "bootapi_log_relay_lines_total 1") { + t.Error("dropped lines must not be counted as relayed") + } +} + +func TestLogRelayLineCountMetric(t *testing.T) { + vl := newVLStub(t, http.StatusNoContent) + srv := newTestServerLogs(t, &fakeNB{}, "local", "prov-secret", vl.srv.URL) + h := srv.Router() + + batch := `{"msg":"a"}` + "\n" + `{"msg":"b"}` + "\n" + `{"msg":"c"}` + "\n" + if rec := postLogs(t, h, "/logs", "prov-secret", batch); rec.Code != http.StatusAccepted { + t.Fatalf("status = %d", rec.Code) + } + body := do(t, h, "/metrics").Body.String() + for _, want := range []string{ + `bootapi_log_relay_total{result="ok"} 1`, + "bootapi_log_relay_lines_total 3", + `bootapi_http_requests_total{endpoint="logs",status="2xx"} 1`, + } { + if !strings.Contains(body, want) { + t.Errorf("metrics missing %q", want) + } + } +}