add POST /logs installer log relay to VictoriaLogs
ci/woodpecker/pr/pre-commit Pipeline was successful
ci/woodpecker/pr/build Pipeline was successful
ci/woodpecker/pr/test Pipeline was successful

A host being PXE-discovered or installed is not in Kubernetes, so vlagent
cannot collect its logs and a failed install leaves no record; the installer
environment has no internal-CA trust or credentials for the HTTPS log ingest,
and bootapi is already the plain-HTTP broker it can reach.

- add POST /logs, token-guarded like POST /provisioned, relaying ndjson to
  vlinsert's jsonline endpoint keyed on serial+phase
- stamp observed source IP and resolved NetBox device name into extra_fields
- return 202 on a sink failure so logs never block an install
- add BOOTAPI_VLINSERT_URL/_TIMEOUT and bootapi_log_relay metrics
This commit is contained in:
2026-10-03 19:44:27 +10:00
parent 0f0fb7fa8e
commit 3c77895788
9 changed files with 416 additions and 5 deletions
+11 -1
View File
@@ -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))
}
+113 -1
View File
@@ -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 <t>" or a bare "token"
// header.
func bearer(r *http.Request) string {
+196 -1
View File
@@ -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)
}
}
}