Compare commits

...

8 Commits

Author SHA1 Message Date
benvin 3d124eac46 Merge pull request 'Relay installer logs to VictoriaLogs via POST /logs' (#9) from benvin/log-relay into main
Reviewed-on: #9
2026-10-03 20:46:07 +10:00
unkin-agent 3c77895788 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
2026-10-03 19:44:27 +10:00
benvin 0f0fb7fa8e Merge pull request 'ci: switch Go steps to gobuilder, enable S3 build cache' (#8) from benvin/gocache into main
Reviewed-on: #8
2026-10-02 23:49:34 +10:00
unkin-agent f7119e2361 ci: switch Go steps to gobuilder, enable S3 build cache
ci/woodpecker/pr/pre-commit Pipeline was successful
ci/woodpecker/pr/build Pipeline was successful
ci/woodpecker/pr/test Pipeline was successful
golang:1.25 and the stale almalinux9-gobuilder pin are replaced by
gobuilder:0.1.2-alma9 on every step that invokes the Go toolchain, with
GOCACHEPROG wired to the baked-in go-cache-plugin against the shared
S3 cache bucket.

- pre-commit, test, release build/test: image -> gobuilder:0.1.2-alma9
- add GOCACHE_* env + AWS creds from org secrets, absolute cache-dir
- lint step (test.yaml), buildx steps, release upload step untouched
2026-10-02 23:45:30 +10:00
benvin edc90f01d9 Merge pull request 'Fetch templates over HTTP instead of shelling out to git' (#7) from benvin/templates-http into main
ci/woodpecker/tag/docker Pipeline was successful
ci/woodpecker/tag/release Pipeline was successful
Reviewed-on: #7
2026-09-26 23:50:25 +10:00
unkin-agent a414918350 Fetch templates over HTTP instead of shelling out to git
ci/woodpecker/pr/build Pipeline was successful
ci/woodpecker/pr/pre-commit Pipeline was successful
ci/woodpecker/pr/test Pipeline was successful
The runtime image is distroless and has no git binary, so every sync
failed and bootapi silently served the stale embedded templates.

- fetch the branch tarball (<repo>/archive/<branch>.tar.gz) and extract
  it into an in-memory FS; no checkout, no writable volume
- digest the extracted tree, not the archive bytes, so a recompressed
  identical archive is not a change
- skip entries that would escape the tree
- log the source commit from Gitea's immutable Link header
2026-09-26 18:59:07 +10:00
benvin ca86d9cedc Merge pull request 'ci: add buildkit_config CA trust for artifactapi push' (#6) from benvin/buildx-ca-config into main
Reviewed-on: #6
2026-08-15 18:46:50 +10:00
unkin-agent fae0cde692 ci: add buildkit_config CA trust for artifactapi push
ci/woodpecker/pr/build Pipeline was successful
ci/woodpecker/pr/pre-commit Pipeline was successful
ci/woodpecker/pr/test Pipeline was successful
2026-08-15 18:30:34 +10:00
17 changed files with 818 additions and 175 deletions
+3
View File
@@ -10,6 +10,9 @@ steps:
repo: artifactapi.k8s.syd1.au.unkin.net/docker-internal/bootapi repo: artifactapi.k8s.syd1.au.unkin.net/docker-internal/bootapi
build_args: build_args:
VERSION: ${CI_COMMIT_TAG} VERSION: ${CI_COMMIT_TAG}
buildkit_config: |
[registry."artifactapi.k8s.syd1.au.unkin.net"]
ca = ["/etc/docker/certs.d/artifactapi.k8s.syd1.au.unkin.net/ca.crt"]
tags: tags:
- ${CI_COMMIT_TAG} - ${CI_COMMIT_TAG}
- latest - latest
+12 -1
View File
@@ -3,7 +3,18 @@ when:
steps: steps:
- name: pre-commit - name: pre-commit
image: git.unkin.net/unkin/almalinux9-gobuilder:20260606 image: "artifactapi.k8s.syd1.au.unkin.net/docker-internal/gobuilder:0.1.2-alma9"
environment:
GOCACHE_S3_BUCKET: gocache
GOCACHE_S3_REGION: us-east-1
GOCACHE_S3_ENDPOINT_URL: "https://s3.ceph.unkin.net"
GOCACHE_S3_PATH_STYLE: "true"
GOCACHE_KEY_PREFIX: ci-bootapi
GOCACHEPROG: "go-cache-plugin --cache-dir=/tmp/gocache"
AWS_ACCESS_KEY_ID:
from_secret: GOCACHE_AWS_ACCESS_KEY_ID
AWS_SECRET_ACCESS_KEY:
from_secret: GOCACHE_AWS_SECRET_ACCESS_KEY
commands: commands:
- uvx pre-commit run --all-files - uvx pre-commit run --all-files
backend_options: backend_options:
+24 -2
View File
@@ -6,7 +6,18 @@ when:
# container image is built+pushed separately by docker.yaml. # container image is built+pushed separately by docker.yaml.
steps: steps:
- name: test - name: test
image: golang:1.25 image: "artifactapi.k8s.syd1.au.unkin.net/docker-internal/gobuilder:0.1.2-alma9"
environment:
GOCACHE_S3_BUCKET: gocache
GOCACHE_S3_REGION: us-east-1
GOCACHE_S3_ENDPOINT_URL: "https://s3.ceph.unkin.net"
GOCACHE_S3_PATH_STYLE: "true"
GOCACHE_KEY_PREFIX: ci-bootapi
GOCACHEPROG: "go-cache-plugin --cache-dir=/tmp/gocache"
AWS_ACCESS_KEY_ID:
from_secret: GOCACHE_AWS_ACCESS_KEY_ID
AWS_SECRET_ACCESS_KEY:
from_secret: GOCACHE_AWS_SECRET_ACCESS_KEY
commands: commands:
- go test -race ./... - go test -race ./...
backend_options: backend_options:
@@ -21,7 +32,18 @@ steps:
cpu: 2 cpu: 2
- name: build - name: build
image: git.unkin.net/unkin/almalinux9-gobuilder:20260606 image: "artifactapi.k8s.syd1.au.unkin.net/docker-internal/gobuilder:0.1.2-alma9"
environment:
GOCACHE_S3_BUCKET: gocache
GOCACHE_S3_REGION: us-east-1
GOCACHE_S3_ENDPOINT_URL: "https://s3.ceph.unkin.net"
GOCACHE_S3_PATH_STYLE: "true"
GOCACHE_KEY_PREFIX: ci-bootapi
GOCACHEPROG: "go-cache-plugin --cache-dir=/tmp/gocache"
AWS_ACCESS_KEY_ID:
from_secret: GOCACHE_AWS_ACCESS_KEY_ID
AWS_SECRET_ACCESS_KEY:
from_secret: GOCACHE_AWS_SECRET_ACCESS_KEY
commands: commands:
- make release-binaries VERSION=${CI_COMMIT_TAG} - make release-binaries VERSION=${CI_COMMIT_TAG}
depends_on: [test] depends_on: [test]
+12 -1
View File
@@ -18,7 +18,18 @@ steps:
cpu: 2 cpu: 2
- name: test - name: test
image: golang:1.25 image: "artifactapi.k8s.syd1.au.unkin.net/docker-internal/gobuilder:0.1.2-alma9"
environment:
GOCACHE_S3_BUCKET: gocache
GOCACHE_S3_REGION: us-east-1
GOCACHE_S3_ENDPOINT_URL: "https://s3.ceph.unkin.net"
GOCACHE_S3_PATH_STYLE: "true"
GOCACHE_KEY_PREFIX: ci-bootapi
GOCACHEPROG: "go-cache-plugin --cache-dir=/tmp/gocache"
AWS_ACCESS_KEY_ID:
from_secret: GOCACHE_AWS_ACCESS_KEY_ID
AWS_SECRET_ACCESS_KEY:
from_secret: GOCACHE_AWS_SECRET_ACCESS_KEY
commands: commands:
- go test -race ./... - go test -race ./...
backend_options: backend_options:
+3 -2
View File
@@ -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 /ipxe/{mac}` · `GET /boot/ipxe?mac=` | iPXE boot script |
| `GET /ks/{ident}` | rendered kickstart (MAC or hostname) | | `GET /ks/{ident}` | rendered kickstart (MAC or hostname) |
| `POST /provisioned/{ident}` | end-of-kickstart callback (token) → clears `pxe_enabled` in NetBox | | `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 | | `GET /healthz` · `/readyz` · `/metrics` | health + Prometheus |
The boot path is served over **plain HTTP** (PXE installers have no internal-CA 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: Env-based (12-factor), see [`config.example.env`](config.example.env). Key vars:
`BOOTAPI_NETBOX_URL`, `BOOTAPI_NETBOX_TOKEN[_FILE]`, `BOOTAPI_BASE_URL`, `BOOTAPI_NETBOX_URL`, `BOOTAPI_NETBOX_TOKEN[_FILE]`, `BOOTAPI_BASE_URL`,
`BOOTAPI_BOOT_BASE_URL`, `BOOTAPI_ROOT_PASSWORD_HASH[_FILE]`, `BOOTAPI_BOOT_BASE_URL`, `BOOTAPI_ROOT_PASSWORD_HASH[_FILE]`,
`BOOTAPI_UNKNOWN_MAC_FALLBACK`. `BOOTAPI_UNKNOWN_MAC_FALLBACK`, `BOOTAPI_VLINSERT_URL`.
## Development ## Development
@@ -88,7 +89,7 @@ internal/model/ Host/Interface data model (incl. pxe_enabled gate)
internal/netbox/ NetBox client (reads + pxe_enabled write) + TTL cache, behind an interface internal/netbox/ NetBox client (reads + pxe_enabled write) + TTL cache, behind an interface
internal/catalog/ distro catalog: NetBox host -> boot images/kickstart internal/catalog/ distro catalog: NetBox host -> boot images/kickstart
internal/render/ text/template engine (swappable Set), selection, loader internal/render/ text/template engine (swappable Set), selection, loader
internal/gitsync/ periodic git pull + atomic template reload (last-good) internal/gitsync/ periodic templates-tarball fetch + atomic reload (last-good)
internal/server/ chi HTTP handlers + Prometheus metrics internal/server/ chi HTTP handlers + Prometheus metrics
templates/ embedded defaults: kickstart, iPXE, catalog/*.yaml templates/ embedded defaults: kickstart, iPXE, catalog/*.yaml
docs/ see above docs/ see above
+11 -8
View File
@@ -9,7 +9,6 @@ import (
"log/slog" "log/slog"
"os" "os"
"os/signal" "os/signal"
"path/filepath"
"syscall" "syscall"
"git.unkin.net/unkin/bootapi/internal/config" "git.unkin.net/unkin/bootapi/internal/config"
@@ -48,7 +47,10 @@ func main() {
slog.Warn("no NetBox token set (BOOTAPI_NETBOX_TOKEN/_FILE); NetBox reads will likely be denied") slog.Warn("no NetBox token set (BOOTAPI_NETBOX_TOKEN/_FILE); NetBox reads will likely be denied")
} }
if cfg.ProvisionToken == "" { 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{ rcfg := render.RenderConfig{
@@ -93,6 +95,8 @@ func main() {
Cache: cache, Cache: cache,
UnknownMACFallback: cfg.UnknownMACFallback, UnknownMACFallback: cfg.UnknownMACFallback,
ProvisionToken: cfg.ProvisionToken, ProvisionToken: cfg.ProvisionToken,
VLInsertURL: cfg.VLInsertURL,
VLInsertTimeout: cfg.VLInsertTimeout,
TLSAddr: cfg.TLSListenAddr, TLSAddr: cfg.TLSListenAddr,
TLSCertFile: cfg.TLSCertFile, TLSCertFile: cfg.TLSCertFile,
TLSKeyFile: cfg.TLSKeyFile, TLSKeyFile: cfg.TLSKeyFile,
@@ -133,10 +137,10 @@ func runValidate(dir string) int {
return 0 return 0
} }
// buildEngine constructs the render Engine and, when a templates git repo is // buildEngine constructs the render Engine and, when a templates repo is
// configured, a Syncer that reloads it periodically. Precedence: git repo → // configured, a Syncer that reloads it periodically. Precedence: templates repo
// local override dir → embedded defaults only. Git/dir failures degrade to the // → local override dir → embedded defaults only. Fetch/dir failures degrade to
// embedded defaults rather than failing startup. // the embedded defaults rather than failing startup.
func buildEngine(ctx context.Context, cfg *config.Config, rcfg render.RenderConfig) (*render.Engine, *gitsync.Syncer, error) { func buildEngine(ctx context.Context, cfg *config.Config, rcfg render.RenderConfig) (*render.Engine, *gitsync.Syncer, error) {
switch { switch {
case cfg.TemplateGitURL != "": case cfg.TemplateGitURL != "":
@@ -145,11 +149,10 @@ func buildEngine(ctx context.Context, cfg *config.Config, rcfg render.RenderConf
Branch: cfg.TemplateGitBranch, Branch: cfg.TemplateGitBranch,
Token: cfg.TemplateGitToken, Token: cfg.TemplateGitToken,
Interval: cfg.TemplateGitInterval, Interval: cfg.TemplateGitInterval,
WorkDir: filepath.Join(os.TempDir(), "bootapi-templates"),
}, templates.FS) }, templates.FS)
set, gerr := syncer.Bootstrap(ctx) set, gerr := syncer.Bootstrap(ctx)
if gerr != nil { if gerr != nil {
slog.Warn("template git bootstrap degraded to embedded defaults", "err", gerr) slog.Warn("template bootstrap degraded to embedded defaults", "err", gerr)
} }
engine := render.NewEngine(rcfg, set) engine := render.NewEngine(rcfg, set)
syncer.SetEngine(engine) syncer.SetEngine(engine)
+9
View File
@@ -55,6 +55,15 @@ BOOTAPI_ARTIFACT_BASE_URL=https://artifactapi.k8s.syd1.au.unkin.net/api/v1/remot
BOOTAPI_PROVISION_TOKEN= BOOTAPI_PROVISION_TOKEN=
# BOOTAPI_PROVISION_TOKEN_FILE=/var/run/secrets/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) --- # --- puppet bootstrap targets (k8s puppetserver; baked into kickstart %post) ---
BOOTAPI_PUPPET_SERVER=puppet.k8s.syd1.au.unkin.net BOOTAPI_PUPPET_SERVER=puppet.k8s.syd1.au.unkin.net
BOOTAPI_PUPPET_CA_SERVER=puppetca.k8s.syd1.au.unkin.net BOOTAPI_PUPPET_CA_SERVER=puppetca.k8s.syd1.au.unkin.net
+3 -4
View File
@@ -44,10 +44,9 @@ Create `apps/base/bootapi/` following the argocd-apps `AGENTS.md` pattern:
as `BOOTAPI_NETBOX_TOKEN_FILE` / `BOOTAPI_PROVISION_TOKEN_FILE` / as `BOOTAPI_NETBOX_TOKEN_FILE` / `BOOTAPI_PROVISION_TOKEN_FILE` /
`BOOTAPI_ROOT_PASSWORD_HASH_FILE` (mount the Secret). Least-privilege `BOOTAPI_ROOT_PASSWORD_HASH_FILE` (mount the Secret). Least-privilege
securityContext (`runAsNonRoot`, `drop: [all]`). Baseline resources: requests securityContext (`runAsNonRoot`, `drop: [all]`). Baseline resources: requests
`512Mi`/`1`, limits `2Gi`/`2` cpu. The pod needs `git` on PATH for template `512Mi`/`1`, limits `2Gi`/`2` cpu. No `git` binary and no writable volume
sync (the distroless image includes only the static binary — either add a git are needed: template sync is an HTTP fetch of the repo's branch tarball, held
layer, use an initContainer that seeds the checkout, or fall back to a in memory.
ConfigMap; simplest is a small alpine+git base for this service).
6. **Service + exposure**: see the Gateway section below. 6. **Service + exposure**: see the Gateway section below.
7. Register in `argocd/applicationsets/platform.yaml` (`apps/overlays/*/bootapi`) 7. Register in `argocd/applicationsets/platform.yaml` (`apps/overlays/*/bootapi`)
and the platform AppProject destinations. and the platform AppProject destinations.
+30
View File
@@ -35,6 +35,7 @@ of the install.
| GET | `/boot/ipxe?mac=...` | Query-string alias of `/ipxe/{mac}`. | | 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. | | 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 <BOOTAPI_PROVISION_TOKEN>`). | | POST | `/provisioned/{ident}` | End-of-kickstart callback; clears `pxe_enabled` in NetBox. **Token-guarded** (`Authorization: Bearer <BOOTAPI_PROVISION_TOKEN>`). |
| POST | `/logs` | Relays a host's newline-delimited JSON install logs to VictoriaLogs. **Token-guarded** (same token). |
| GET | `/healthz` | Liveness: always `200 ok`. | | GET | `/healthz` | Liveness: always `200 ok`. |
| GET | `/readyz` | Readiness: `200` once templates parsed. Does **not** probe NetBox. | | GET | `/readyz` | Readiness: `200` once templates parsed. Does **not** probe NetBox. |
| GET | `/metrics` | Prometheus metrics (see below). | | 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 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 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 ## Metrics
All on `/metrics`, prefix `bootapi_`: All on `/metrics`, prefix `bootapi_`:
@@ -95,6 +123,8 @@ All on `/metrics`, prefix `bootapi_`:
- `bootapi_netbox_lookup_duration_seconds{field}` — histogram. - `bootapi_netbox_lookup_duration_seconds{field}` — histogram.
- `bootapi_netbox_cache_hits_total` / `bootapi_netbox_cache_misses_total`. - `bootapi_netbox_cache_hits_total` / `bootapi_netbox_cache_misses_total`.
- `bootapi_provisioned_total{result}` — result = `ok|unauthorized|notfound|error|disabled`. - `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_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)). - `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. - standard Go/process collectors.
+2 -1
View File
@@ -8,7 +8,8 @@ bootapi ships an embedded default set and lets you override or extend it.
`templates/ipxe/*.ipxe.tmpl` and `templates/catalog/*.yaml`, compiled into the `templates/ipxe/*.ipxe.tmpl` and `templates/catalog/*.yaml`, compiled into the
binary (`templates/embed.go`). These are the always-available startup fallback. binary (`templates/embed.go`). These are the always-available startup fallback.
- **Template git repo** (preferred in prod): `BOOTAPI_TEMPLATE_GIT_URL`. bootapi - **Template git repo** (preferred in prod): `BOOTAPI_TEMPLATE_GIT_URL`. bootapi
clones it at startup and re-pulls every `BOOTAPI_TEMPLATE_GIT_INTERVAL` fetches the branch tarball (`<repo>/archive/<branch>.tar.gz`) over HTTP at
startup and re-fetches every `BOOTAPI_TEMPLATE_GIT_INTERVAL`
(default 3m, like argocd), atomically swapping the loaded set on change. A (default 3m, like argocd), atomically swapping the loaded set on change. A
parse failure keeps the **last-good** set and is only logged + counted parse failure keeps the **last-good** set and is only logged + counted
(`bootapi_template_sync_failures_total`), so a bad push can't take bootapi (`bootapi_template_sync_failures_total`), so a bad push can't take bootapi
+24 -2
View File
@@ -46,14 +46,15 @@ type Config struct {
DefaultTemplate string DefaultTemplate string
// --- template git-sync (preferred over TemplateDir) --- // --- template git-sync (preferred over TemplateDir) ---
// TemplateGitURL, when set, makes bootapi clone a templates repo and re-pull // TemplateGitURL, when set, makes bootapi fetch a templates repo's branch
// tarball over HTTP (<repo>/archive/<branch>.tar.gz) and re-fetch
// it every TemplateGitInterval, atomically swapping the loaded set on change // it every TemplateGitInterval, atomically swapping the loaded set on change
// and keeping the last-good set on a parse failure. // and keeping the last-good set on a parse failure.
TemplateGitURL string TemplateGitURL string
TemplateGitBranch string TemplateGitBranch string
TemplateGitInterval time.Duration TemplateGitInterval time.Duration
// TemplateGitToken is an optional token for a private templates repo, // TemplateGitToken is an optional token for a private templates repo,
// injected into the HTTPS clone URL. Empty for a public repo. // sent as a Gitea token header. Empty for a public repo.
TemplateGitToken string TemplateGitToken string
// BaseURL is the http:// base PXE clients use to reach bootapi. It is baked // BaseURL is the http:// base PXE clients use to reach bootapi. It is baked
@@ -78,6 +79,15 @@ type Config struct {
// endpoint (fail closed). Prefer ProvisionTokenFile in k8s. // endpoint (fail closed). Prefer ProvisionTokenFile in k8s.
ProvisionToken string 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 // PuppetServer / PuppetCAServer are baked into kickstart %post so the
// freshly-installed host checks in to the k8s puppetserver. // freshly-installed host checks in to the k8s puppetserver.
PuppetServer string PuppetServer string
@@ -117,6 +127,16 @@ func Load() (*Config, error) {
if err != nil { if err != nil {
return nil, fmt.Errorf("invalid BOOTAPI_TEMPLATE_GIT_INTERVAL: %w", err) 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") token, err := readSecret("BOOTAPI_NETBOX_TOKEN")
if err != nil { if err != nil {
@@ -168,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"), "/"), 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"), "/"), BootBaseURL: strings.TrimRight(os.Getenv("BOOTAPI_BOOT_BASE_URL"), "/"),
ProvisionToken: provToken, ProvisionToken: provToken,
VLInsertURL: strings.TrimRight(vlURL, "/"),
VLInsertTimeout: vlTimeout,
PuppetServer: getenv("BOOTAPI_PUPPET_SERVER", "puppet.k8s.syd1.au.unkin.net"), PuppetServer: getenv("BOOTAPI_PUPPET_SERVER", "puppet.k8s.syd1.au.unkin.net"),
PuppetCAServer: getenv("BOOTAPI_PUPPET_CA_SERVER", "puppetca.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"), PuppetCAURL: getenv("BOOTAPI_PUPPET_CA_URL", "puppetca.k8s.syd1.au.unkin.net"),
+28
View File
@@ -40,6 +40,34 @@ func TestLoadDefaults(t *testing.T) {
if c.ArtifactBaseURL != "https://artifactapi.k8s.syd1.au.unkin.net/api/v1/remote" { if c.ArtifactBaseURL != "https://artifactapi.k8s.syd1.au.unkin.net/api/v1/remote" {
t.Errorf("ArtifactBaseURL = %q", c.ArtifactBaseURL) 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) { func TestCallbackBaseDefaultsToBase(t *testing.T) {
+169 -79
View File
@@ -1,42 +1,63 @@
// Package gitsync keeps bootapi's template Set in step with a git repo. It // Package gitsync keeps bootapi's template Set in step with the templates repo.
// clones the templates repo at startup and re-pulls it every interval (default // It fetches the repo's branch tarball over plain HTTP (Gitea's
// 3m, like argocd), atomically swapping the Engine's active Set when the repo // /archive/<branch>.tar.gz) every interval (default 3m, like argocd), holds the
// changes. A parse failure keeps the last-good Set and is only logged/counted, // template files in memory, and atomically swaps the Engine's active Set when
// so a bad template push can never take bootapi down. The embedded defaults // the content changes. A parse failure keeps the last-good Set and is only
// remain the fallback when git is unreachable at startup. // logged/counted, so a bad template push can never take bootapi down. The
// embedded defaults remain the fallback when the repo is unreachable at startup.
//
// The repo is only ever read as a file tree plus a change signal, so no git
// binary is involved: the runtime image stays distroless and the pod needs no
// writable volume.
package gitsync package gitsync
import ( import (
"archive/tar"
"compress/gzip"
"context" "context"
"crypto/sha256"
"encoding/hex"
"fmt" "fmt"
"io"
"io/fs" "io/fs"
"log/slog" "log/slog"
"os" "net/http"
"os/exec" "path"
"sort"
"strings" "strings"
"sync/atomic" "sync/atomic"
"testing/fstest"
"time" "time"
"git.unkin.net/unkin/bootapi/internal/render" "git.unkin.net/unkin/bootapi/internal/render"
) )
// maxArchiveBytes caps the downloaded tarball. The templates repo is a handful
// of text files (~12KB compressed); this only exists to bound a hostile or
// broken response.
const maxArchiveBytes = 32 << 20
// maxFileBytes caps a single extracted template.
const maxFileBytes = 1 << 20
// Options configures the syncer. // Options configures the syncer.
type Options struct { type Options struct {
URL string URL string // repo URL, e.g. https://git.unkin.net/unkin/bootapi-templates.git
Branch string Branch string
Token string // optional; injected into the HTTPS URL for a private repo Token string // optional; sent as a Gitea token header for a private repo
Interval time.Duration Interval time.Duration
WorkDir string // local checkout path Client *http.Client // optional; defaults to a 30s-timeout client
} }
// Syncer pulls a templates repo and reloads an Engine on change. // Syncer fetches a templates repo and reloads an Engine on change.
type Syncer struct { type Syncer struct {
opt Options opt Options
embedded fs.FS embedded fs.FS
engine *render.Engine engine *render.Engine
digest string // content digest of the last fetched tree
syncs atomic.Int64 // successful reloads (Set swapped) syncs atomic.Int64 // successful reloads (Set swapped)
failures atomic.Int64 // pull or parse failures (last-good kept) failures atomic.Int64 // fetch or parse failures (last-good kept)
generation atomic.Int64 // increments on every successful swap generation atomic.Int64 // increments on every successful swap
} }
@@ -50,6 +71,9 @@ func New(opt Options, embedded fs.FS) *Syncer {
if opt.Interval <= 0 { if opt.Interval <= 0 {
opt.Interval = 3 * time.Minute opt.Interval = 3 * time.Minute
} }
if opt.Client == nil {
opt.Client = &http.Client{Timeout: 30 * time.Second}
}
return &Syncer{opt: opt, embedded: embedded} return &Syncer{opt: opt, embedded: embedded}
} }
@@ -61,35 +85,46 @@ func (s *Syncer) Syncs() int64 { return s.syncs.Load() }
func (s *Syncer) Failures() int64 { return s.failures.Load() } func (s *Syncer) Failures() int64 { return s.failures.Load() }
func (s *Syncer) Generation() int64 { return s.generation.Load() } func (s *Syncer) Generation() int64 { return s.generation.Load() }
// Bootstrap clones the repo and builds the initial Set from embedded + the // ArchiveURL is the branch tarball URL derived from the repo URL.
// checkout. On any git/parse failure it returns an embedded-only Set plus a func (o Options) ArchiveURL() string {
// non-nil error (which the caller logs but treats as non-fatal, so bootapi base := strings.TrimSuffix(strings.TrimRight(o.URL, "/"), ".git")
return base + "/archive/" + o.Branch + ".tar.gz"
}
// Bootstrap fetches the repo and builds the initial Set from embedded + the
// fetched tree. On any fetch/parse failure it returns an embedded-only Set plus
// a non-nil error (which the caller logs but treats as non-fatal, so bootapi
// always starts with at least the embedded defaults). // always starts with at least the embedded defaults).
func (s *Syncer) Bootstrap(ctx context.Context) (*render.Set, error) { func (s *Syncer) Bootstrap(ctx context.Context) (*render.Set, error) {
if err := s.clone(ctx); err != nil { tree, digest, commit, err := s.fetch(ctx)
set, berr := render.BuildSet(s.embedded, nil)
if berr != nil {
return nil, berr // embedded defaults broken: genuinely fatal
}
return set, fmt.Errorf("git clone failed, using embedded defaults: %w", err)
}
set, err := render.BuildSet(s.embedded, os.DirFS(s.opt.WorkDir))
if err != nil { if err != nil {
emb, berr := render.BuildSet(s.embedded, nil) return s.embeddedSet(fmt.Errorf("template fetch failed, using embedded defaults: %w", err))
if berr != nil {
return nil, berr
} }
return emb, fmt.Errorf("git templates failed to parse, using embedded defaults: %w", err) s.digest = digest
set, err := render.BuildSet(s.embedded, tree)
if err != nil {
return s.embeddedSet(fmt.Errorf("fetched templates failed to parse, using embedded defaults: %w", err))
} }
s.generation.Add(1) s.generation.Add(1)
slog.Info("templates loaded", "commit", commit, "digest", digest)
return set, nil return set, nil
} }
// embeddedSet returns the embedded-only Set alongside the degrade reason. A
// broken embedded set is genuinely fatal.
func (s *Syncer) embeddedSet(reason error) (*render.Set, error) {
set, err := render.BuildSet(s.embedded, nil)
if err != nil {
return nil, err
}
return set, reason
}
// Run polls the repo every interval until ctx is cancelled. // Run polls the repo every interval until ctx is cancelled.
func (s *Syncer) Run(ctx context.Context) { func (s *Syncer) Run(ctx context.Context) {
t := time.NewTicker(s.opt.Interval) t := time.NewTicker(s.opt.Interval)
defer t.Stop() defer t.Stop()
slog.Info("template git-sync started", "url", s.opt.URL, "branch", s.opt.Branch, "interval", s.opt.Interval) slog.Info("template sync started", "url", s.opt.ArchiveURL(), "interval", s.opt.Interval)
for { for {
select { select {
case <-ctx.Done(): case <-ctx.Done():
@@ -101,78 +136,133 @@ func (s *Syncer) Run(ctx context.Context) {
} }
func (s *Syncer) pollOnce(ctx context.Context) { func (s *Syncer) pollOnce(ctx context.Context) {
changed, head, err := s.pull(ctx) tree, digest, commit, err := s.fetch(ctx)
if err != nil { if err != nil {
s.failures.Add(1) s.failures.Add(1)
slog.Error("template git pull failed; keeping last-good set", "err", err) slog.Error("template fetch failed; keeping last-good set", "err", err)
return return
} }
if !changed { if digest == s.digest {
return return
} }
set, err := render.BuildSet(s.embedded, os.DirFS(s.opt.WorkDir)) // Record the new digest before parsing so an unchanged bad push is counted
// once, not on every poll.
s.digest = digest
set, err := render.BuildSet(s.embedded, tree)
if err != nil { if err != nil {
s.failures.Add(1) s.failures.Add(1)
slog.Error("template reload failed to parse; keeping last-good set", "commit", head, "err", err) slog.Error("template reload failed to parse; keeping last-good set", "commit", commit, "err", err)
return return
} }
s.engine.Swap(set) s.engine.Swap(set)
s.syncs.Add(1) s.syncs.Add(1)
s.generation.Add(1) s.generation.Add(1)
slog.Info("templates reloaded from git", "commit", head, "generation", s.generation.Load()) slog.Info("templates reloaded", "commit", commit, "digest", digest, "generation", s.generation.Load())
} }
// authURL injects a token into the HTTPS clone URL when configured. // fetch downloads the branch tarball and extracts it into an in-memory FS. The
func (s *Syncer) authURL() string { // digest is taken over the extracted tree (not the gzip bytes) so a
if s.opt.Token == "" { // re-compressed but identical archive is not treated as a change. commit is the
return s.opt.URL // source commit Gitea advertises in its immutable Link header, for logging only.
} func (s *Syncer) fetch(ctx context.Context) (fs.FS, string, string, error) {
if rest, ok := strings.CutPrefix(s.opt.URL, "https://"); ok { req, err := http.NewRequestWithContext(ctx, http.MethodGet, s.opt.ArchiveURL(), nil)
return "https://" + s.opt.Token + "@" + rest
}
return s.opt.URL
}
func (s *Syncer) clone(ctx context.Context) error {
if err := os.RemoveAll(s.opt.WorkDir); err != nil {
return err
}
return run(ctx, "", "git", "clone", "--depth", "1", "--branch", s.opt.Branch, s.authURL(), s.opt.WorkDir)
}
// pull fetches origin/branch and hard-resets to it, reporting whether HEAD moved.
func (s *Syncer) pull(ctx context.Context) (changed bool, head string, err error) {
old, _ := s.head(ctx)
if err := run(ctx, s.opt.WorkDir, "git", "fetch", "--depth", "1", "origin", s.opt.Branch); err != nil {
return false, "", err
}
if err := run(ctx, s.opt.WorkDir, "git", "reset", "--hard", "origin/"+s.opt.Branch); err != nil {
return false, "", err
}
newHead, err := s.head(ctx)
if err != nil { if err != nil {
return false, "", err return nil, "", "", err
} }
return old != newHead, newHead, nil if s.opt.Token != "" {
req.Header.Set("Authorization", "token "+s.opt.Token)
}
resp, err := s.opt.Client.Do(req)
if err != nil {
return nil, "", "", err
}
defer func() { _ = resp.Body.Close() }()
if resp.StatusCode != http.StatusOK {
return nil, "", "", fmt.Errorf("GET %s: %s", s.opt.ArchiveURL(), resp.Status)
} }
func (s *Syncer) head(ctx context.Context) (string, error) { tree, err := extract(io.LimitReader(resp.Body, maxArchiveBytes))
out, err := output(ctx, s.opt.WorkDir, "git", "rev-parse", "HEAD") if err != nil {
return strings.TrimSpace(out), err return nil, "", "", err
}
if len(tree) == 0 {
return nil, "", "", fmt.Errorf("archive contained no files")
}
return tree, digest(tree), commitFromLink(resp.Header.Get("Link")), nil
} }
func run(ctx context.Context, dir, name string, args ...string) error { // extract reads a gzipped tar and returns its regular files keyed by path with
cmd := exec.CommandContext(ctx, name, args...) // the archive's single top-level directory stripped (Gitea prefixes every entry
cmd.Dir = dir // with "<repo>/"). Entries that would escape the tree are skipped rather than
if out, err := cmd.CombinedOutput(); err != nil { // trusted: the archive is a network input.
return fmt.Errorf("%s %s: %w: %s", name, strings.Join(args, " "), err, strings.TrimSpace(string(out))) func extract(r io.Reader) (fstest.MapFS, error) {
gz, err := gzip.NewReader(r)
if err != nil {
return nil, fmt.Errorf("gzip: %w", err)
}
defer func() { _ = gz.Close() }()
out := fstest.MapFS{}
tr := tar.NewReader(gz)
for {
h, err := tr.Next()
if err == io.EOF {
return out, nil
}
if err != nil {
return nil, fmt.Errorf("tar: %w", err)
}
if h.Typeflag != tar.TypeReg {
continue
}
name := stripRoot(h.Name)
if name == "" || !fs.ValidPath(name) {
continue
}
b, err := io.ReadAll(io.LimitReader(tr, maxFileBytes))
if err != nil {
return nil, fmt.Errorf("tar %s: %w", h.Name, err)
}
out[name] = &fstest.MapFile{Data: b, Mode: 0o444}
} }
return nil
} }
func output(ctx context.Context, dir, name string, args ...string) (string, error) { // stripRoot removes the archive's leading directory component.
cmd := exec.CommandContext(ctx, name, args...) func stripRoot(name string) string {
cmd.Dir = dir clean := path.Clean(strings.TrimPrefix(name, "./"))
out, err := cmd.Output() if strings.HasPrefix(clean, "/") || strings.HasPrefix(clean, "..") {
return string(out), err return ""
}
_, rest, ok := strings.Cut(clean, "/")
if !ok {
return ""
}
return rest
}
// digest hashes the extracted tree: every path and its contents, in path order.
func digest(tree fstest.MapFS) string {
names := make([]string, 0, len(tree))
for n := range tree {
names = append(names, n)
}
sort.Strings(names)
h := sha256.New()
for _, n := range names {
_, _ = fmt.Fprintf(h, "%s\x00%d\x00", n, len(tree[n].Data))
_, _ = h.Write(tree[n].Data)
}
return hex.EncodeToString(h.Sum(nil))[:16]
}
// commitFromLink pulls the commit SHA out of Gitea's immutable-archive Link
// header: <.../archive/<sha>.tar.gz?rev=<sha>>; rel="immutable".
func commitFromLink(link string) string {
_, rev, ok := strings.Cut(link, "rev=")
if !ok {
return ""
}
sha, _, _ := strings.Cut(rev, ">")
return strings.TrimSpace(sha)
} }
+164 -68
View File
@@ -1,11 +1,14 @@
package gitsync package gitsync
import ( import (
"archive/tar"
"bytes"
"compress/gzip"
"context" "context"
"os" "net/http"
"os/exec" "net/http/httptest"
"path/filepath"
"strings" "strings"
"sync/atomic"
"testing" "testing"
"time" "time"
@@ -14,31 +17,81 @@ import (
"git.unkin.net/unkin/bootapi/templates" "git.unkin.net/unkin/bootapi/templates"
) )
// gitRepo creates a real git repo at dir with an initial almalinux9 override. // archive builds a gzipped tar shaped like Gitea's: every entry under a single
func gitRepo(t *testing.T, dir string) { // "<repo>/" root, plus the directory entries Gitea includes.
func archive(t *testing.T, files map[string]string) []byte {
t.Helper() t.Helper()
gitCmd(t, "", "git", "init", "-b", "main", dir) var buf bytes.Buffer
gitCmd(t, dir, "git", "config", "user.email", "t@example.net") gz := gzip.NewWriter(&buf)
gitCmd(t, dir, "git", "config", "user.name", "test") tw := tar.NewWriter(gz)
writeKS(t, dir, "GITSYNC-V1 {{ .Hostname }}\n") if err := tw.WriteHeader(&tar.Header{Name: "bootapi-templates/", Typeflag: tar.TypeDir, Mode: 0o755}); err != nil {
gitCmd(t, dir, "git", "add", "-A") t.Fatal(err)
gitCmd(t, dir, "git", "commit", "-m", "v1")
} }
for name, body := range files {
func writeKS(t *testing.T, dir, body string) { h := &tar.Header{Name: "bootapi-templates/" + name, Typeflag: tar.TypeReg, Mode: 0o644, Size: int64(len(body))}
t.Helper() if err := tw.WriteHeader(h); err != nil {
if err := os.WriteFile(filepath.Join(dir, "almalinux9.ks.tmpl"), []byte(body), 0o600); err != nil { t.Fatal(err)
}
if _, err := tw.Write([]byte(body)); err != nil {
t.Fatal(err) t.Fatal(err)
} }
} }
if err := tw.Close(); err != nil {
func gitCmd(t *testing.T, dir, name string, args ...string) { t.Fatal(err)
t.Helper()
cmd := exec.Command(name, args...)
cmd.Dir = dir
if out, err := cmd.CombinedOutput(); err != nil {
t.Fatalf("%s %v: %v: %s", name, args, err, out)
} }
if err := gz.Close(); err != nil {
t.Fatal(err)
}
return buf.Bytes()
}
// repo serves a mutable archive at Gitea's /archive path and counts requests.
type repo struct {
t *testing.T
srv *httptest.Server
body atomic.Value // []byte
status atomic.Int64
fetches atomic.Int64
}
func newRepo(t *testing.T, files map[string]string) *repo {
t.Helper()
r := &repo{t: t}
r.body.Store(archive(t, files))
r.status.Store(http.StatusOK)
mux := http.NewServeMux()
mux.HandleFunc("/unkin/bootapi-templates/archive/main.tar.gz", func(w http.ResponseWriter, _ *http.Request) {
r.fetches.Add(1)
if code := int(r.status.Load()); code != http.StatusOK {
w.WriteHeader(code)
return
}
w.Header().Set("Link", `<https://git.example.net/api/v1/repos/unkin/bootapi-templates/archive/abc123.tar.gz?rev=abc123>; rel="immutable"`)
_, _ = w.Write(r.body.Load().([]byte))
})
r.srv = httptest.NewServer(mux)
t.Cleanup(r.srv.Close)
return r
}
func (r *repo) url() string { return r.srv.URL + "/unkin/bootapi-templates.git" }
func (r *repo) push(files map[string]string) { r.body.Store(archive(r.t, files)) }
func ks(body string) map[string]string { return map[string]string{"almalinux9.ks.tmpl": body} }
func syncer(t *testing.T, r *repo) *Syncer {
t.Helper()
return New(Options{URL: r.url(), Branch: "main", Interval: time.Hour}, templates.FS)
}
func engineFor(t *testing.T, s *Syncer) *render.Engine {
t.Helper()
set, err := s.Bootstrap(context.Background())
if err != nil {
t.Fatalf("Bootstrap: %v", err)
}
eng := render.NewEngine(render.RenderConfig{DefaultTemplate: "almalinux9", ArtifactBase: "https://af"}, set)
s.SetEngine(eng)
return eng
} }
func renderKS(t *testing.T, e *render.Engine) string { func renderKS(t *testing.T, e *render.Engine) string {
@@ -51,33 +104,32 @@ func renderKS(t *testing.T, e *render.Engine) string {
return string(out) return string(out)
} }
func TestArchiveURL(t *testing.T) {
for _, tc := range []struct{ in, want string }{
{"https://git.unkin.net/unkin/bootapi-templates.git", "https://git.unkin.net/unkin/bootapi-templates/archive/main.tar.gz"},
{"https://git.unkin.net/unkin/bootapi-templates", "https://git.unkin.net/unkin/bootapi-templates/archive/main.tar.gz"},
{"https://git.unkin.net/unkin/bootapi-templates/", "https://git.unkin.net/unkin/bootapi-templates/archive/main.tar.gz"},
} {
if got := (Options{URL: tc.in, Branch: "main"}).ArchiveURL(); got != tc.want {
t.Errorf("ArchiveURL(%q) = %q, want %q", tc.in, got, tc.want)
}
}
}
func TestBootstrapAndReload(t *testing.T) { func TestBootstrapAndReload(t *testing.T) {
if _, err := exec.LookPath("git"); err != nil { r := newRepo(t, ks("SYNC-V1 {{ .Hostname }}\n"))
t.Skip("git not available") s := syncer(t, r)
} eng := engineFor(t, s)
src := t.TempDir()
gitRepo(t, src)
s := New(Options{URL: src, Branch: "main", Interval: time.Hour, WorkDir: filepath.Join(t.TempDir(), "co")}, templates.FS) if got := renderKS(t, eng); !strings.Contains(got, "SYNC-V1 web01") {
set, err := s.Bootstrap(context.Background())
if err != nil {
t.Fatalf("Bootstrap: %v", err)
}
eng := render.NewEngine(render.RenderConfig{DefaultTemplate: "almalinux9", ArtifactBase: "https://af"}, set)
s.SetEngine(eng)
if got := renderKS(t, eng); !contains(got, "GITSYNC-V1 web01") {
t.Fatalf("initial render missing v1 override:\n%s", got) t.Fatalf("initial render missing v1 override:\n%s", got)
} }
gen1 := s.Generation() gen1 := s.Generation()
// Commit v2 upstream, then poll: the engine must swap to the new content. r.push(ks("SYNC-V2 {{ .Hostname }}\n"))
writeKS(t, src, "GITSYNC-V2 {{ .Hostname }}\n")
gitCmd(t, src, "git", "add", "-A")
gitCmd(t, src, "git", "commit", "-m", "v2")
s.pollOnce(context.Background()) s.pollOnce(context.Background())
if got := renderKS(t, eng); !contains(got, "GITSYNC-V2 web01") {
if got := renderKS(t, eng); !strings.Contains(got, "SYNC-V2 web01") {
t.Fatalf("after reload, render missing v2:\n%s", got) t.Fatalf("after reload, render missing v2:\n%s", got)
} }
if s.Generation() <= gen1 { if s.Generation() <= gen1 {
@@ -88,30 +140,33 @@ func TestBootstrapAndReload(t *testing.T) {
} }
} }
func TestReloadKeepsLastGoodOnParseError(t *testing.T) { func TestUnchangedContentDoesNotReload(t *testing.T) {
if _, err := exec.LookPath("git"); err != nil { r := newRepo(t, ks("SYNC-V1 {{ .Hostname }}\n"))
t.Skip("git not available") s := syncer(t, r)
} engineFor(t, s)
src := t.TempDir()
gitRepo(t, src)
s := New(Options{URL: src, Branch: "main", Interval: time.Hour, WorkDir: filepath.Join(t.TempDir(), "co")}, templates.FS)
set, err := s.Bootstrap(context.Background())
if err != nil {
t.Fatal(err)
}
eng := render.NewEngine(render.RenderConfig{DefaultTemplate: "almalinux9", ArtifactBase: "https://af"}, set)
s.SetEngine(eng)
// Push a template that fails to parse.
writeKS(t, src, "BROKEN {{ .Hostname \n")
gitCmd(t, src, "git", "add", "-A")
gitCmd(t, src, "git", "commit", "-m", "broken")
// Re-archiving the same files yields fresh gzip bytes; the digest is taken
// over the extracted tree, so this must NOT count as a change.
r.push(ks("SYNC-V1 {{ .Hostname }}\n"))
s.pollOnce(context.Background()) s.pollOnce(context.Background())
// The last-good v1 set must still be served, and a failure recorded. if s.Syncs() != 0 {
if got := renderKS(t, eng); !contains(got, "GITSYNC-V1 web01") { t.Errorf("syncs = %d, want 0 (identical content is not a change)", s.Syncs())
}
if s.Failures() != 0 {
t.Errorf("failures = %d, want 0", s.Failures())
}
}
func TestReloadKeepsLastGoodOnParseError(t *testing.T) {
r := newRepo(t, ks("SYNC-V1 {{ .Hostname }}\n"))
s := syncer(t, r)
eng := engineFor(t, s)
r.push(ks("BROKEN {{ .Hostname \n"))
s.pollOnce(context.Background())
if got := renderKS(t, eng); !strings.Contains(got, "SYNC-V1 web01") {
t.Fatalf("last-good not kept after parse failure:\n%s", got) t.Fatalf("last-good not kept after parse failure:\n%s", got)
} }
if s.Failures() != 1 { if s.Failures() != 1 {
@@ -120,11 +175,19 @@ func TestReloadKeepsLastGoodOnParseError(t *testing.T) {
if s.Syncs() != 0 { if s.Syncs() != 0 {
t.Errorf("syncs = %d, want 0 (bad push must not count as a sync)", s.Syncs()) t.Errorf("syncs = %d, want 0 (bad push must not count as a sync)", s.Syncs())
} }
// The same bad content on the next poll must not be counted again.
s.pollOnce(context.Background())
if s.Failures() != 1 {
t.Errorf("failures = %d, want 1 (an unchanged bad push is counted once)", s.Failures())
}
} }
func TestBootstrapDegradesToEmbedded(t *testing.T) { func TestBootstrapDegradesToEmbedded(t *testing.T) {
// A bogus URL must not fail startup: Bootstrap returns the embedded set. r := newRepo(t, ks("SYNC-V1 {{ .Hostname }}\n"))
s := New(Options{URL: "/nonexistent/repo", Branch: "main", Interval: time.Hour, WorkDir: filepath.Join(t.TempDir(), "co")}, templates.FS) r.status.Store(http.StatusNotFound)
s := syncer(t, r)
set, err := s.Bootstrap(context.Background()) set, err := s.Bootstrap(context.Background())
if err == nil { if err == nil {
t.Error("expected a non-nil (non-fatal) error describing the degrade") t.Error("expected a non-nil (non-fatal) error describing the degrade")
@@ -133,10 +196,43 @@ func TestBootstrapDegradesToEmbedded(t *testing.T) {
t.Fatal("expected the embedded fallback Set, got nil") t.Fatal("expected the embedded fallback Set, got nil")
} }
eng := render.NewEngine(render.RenderConfig{DefaultTemplate: "almalinux9", ArtifactBase: "https://af"}, set) eng := render.NewEngine(render.RenderConfig{DefaultTemplate: "almalinux9", ArtifactBase: "https://af"}, set)
// Embedded almalinux9 template still renders. if got := renderKS(t, eng); !strings.Contains(got, "rootpw") {
if got := renderKS(t, eng); !contains(got, "rootpw") {
t.Errorf("embedded fallback did not render a real kickstart:\n%s", got) t.Errorf("embedded fallback did not render a real kickstart:\n%s", got)
} }
} }
func contains(s, sub string) bool { return strings.Contains(s, sub) } func TestExtractStripsRootAndSkipsEscapes(t *testing.T) {
tree, err := extract(bytes.NewReader(archive(t, map[string]string{
"catalog/almalinux9.yaml": "name: almalinux9\n",
"../escape.ks.tmpl": "nope\n",
})))
if err != nil {
t.Fatalf("extract: %v", err)
}
if _, ok := tree["catalog/almalinux9.yaml"]; !ok {
t.Errorf("root not stripped; got keys %v", keys(tree))
}
for k := range tree {
if strings.Contains(k, "escape") {
t.Errorf("traversal entry was kept: %q", k)
}
}
}
func TestCommitFromLink(t *testing.T) {
link := `<https://git.unkin.net/api/v1/repos/unkin/bootapi-templates/archive/5df1894.tar.gz?rev=5df1894>; rel="immutable"`
if got := commitFromLink(link); got != "5df1894" {
t.Errorf("commitFromLink = %q, want 5df1894", got)
}
if got := commitFromLink(""); got != "" {
t.Errorf("commitFromLink(\"\") = %q, want empty", got)
}
}
func keys[V any](m map[string]V) []string {
out := make([]string, 0, len(m))
for k := range m {
out = append(out, k)
}
return out
}
+11 -1
View File
@@ -28,6 +28,8 @@ type metrics struct {
netboxLookups *prometheus.CounterVec // by field,result netboxLookups *prometheus.CounterVec // by field,result
netboxDuration *prometheus.HistogramVec netboxDuration *prometheus.HistogramVec
provisioned *prometheus.CounterVec // by result provisioned *prometheus.CounterVec // by result
logRelay *prometheus.CounterVec // by result
logLines prometheus.Counter
ipxeGated prometheus.Counter ipxeGated prometheus.Counter
} }
@@ -56,12 +58,20 @@ func newMetrics(cache cacheStats, git gitStats) *metrics {
Name: "bootapi_provisioned_total", Name: "bootapi_provisioned_total",
Help: "Provisioned callbacks, by result (ok|unauthorized|notfound|error|disabled).", Help: "Provisioned callbacks, by result (ok|unauthorized|notfound|error|disabled).",
}, []string{"result"}), }, []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{ ipxeGated: prometheus.NewCounter(prometheus.CounterOpts{
Name: "bootapi_ipxe_gated_total", Name: "bootapi_ipxe_gated_total",
Help: "Known hosts served the local-boot fallback because pxe_enabled=false.", 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 { if cache != nil {
reg.MustRegister(newCacheCollector(cache)) reg.MustRegister(newCacheCollector(cache))
} }
+113 -1
View File
@@ -3,12 +3,16 @@
package server package server
import ( import (
"bytes"
"context" "context"
"crypto/subtle" "crypto/subtle"
"errors" "errors"
"fmt" "fmt"
"io"
"log/slog" "log/slog"
"net"
"net/http" "net/http"
"net/url"
"strings" "strings"
"sync" "sync"
"time" "time"
@@ -30,9 +34,15 @@ type Server struct {
// fallback is the unknown-MAC iPXE behavior: "local" (safe default) or // fallback is the unknown-MAC iPXE behavior: "local" (safe default) or
// "shell" (debug). // "shell" (debug).
fallback string fallback string
// provisionToken guards POST /provisioned; empty disables the endpoint. // provisionToken guards POST /provisioned and POST /logs; empty disables
// both endpoints.
provisionToken string 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. // TLS listener (optional); the plain-HTTP listener is always on.
tlsAddr string tlsAddr string
tlsCert string tlsCert string
@@ -48,6 +58,8 @@ type Options struct {
GitStats gitStats GitStats gitStats
UnknownMACFallback string UnknownMACFallback string
ProvisionToken string ProvisionToken string
VLInsertURL string
VLInsertTimeout time.Duration
TLSAddr string TLSAddr string
TLSCertFile string TLSCertFile string
TLSKeyFile string TLSKeyFile string
@@ -59,12 +71,18 @@ func New(o Options) *Server {
if fb == "" { if fb == "" {
fb = "local" fb = "local"
} }
timeout := o.VLInsertTimeout
if timeout <= 0 {
timeout = 5 * time.Second
}
return &Server{ return &Server{
nb: o.NetBox, nb: o.NetBox,
engine: o.Engine, engine: o.Engine,
metrics: newMetrics(o.Cache, o.GitStats), metrics: newMetrics(o.Cache, o.GitStats),
fallback: fb, fallback: fb,
provisionToken: o.ProvisionToken, provisionToken: o.ProvisionToken,
vlinsertURL: strings.TrimRight(o.VLInsertURL, "/"),
vlClient: &http.Client{Timeout: timeout},
tlsAddr: o.TLSAddr, tlsAddr: o.TLSAddr,
tlsCert: o.TLSCertFile, tlsCert: o.TLSCertFile,
tlsKey: o.TLSKeyFile, 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. // End-of-kickstart callback: flips pxe_enabled off in NetBox. Token-guarded.
r.Post("/provisioned/{ident}", s.handleProvisioned) 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 return r
} }
@@ -253,6 +275,96 @@ func (s *Server) handleProvisioned(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusNoContent) 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" // bearer extracts a token from "Authorization: Bearer <t>" or a bare "token"
// header. // header.
func bearer(r *http.Request) string { func bearer(r *http.Request) string {
+196 -1
View File
@@ -2,10 +2,13 @@ package server
import ( import (
"context" "context"
"io"
"net/http" "net/http"
"net/http/httptest" "net/http/httptest"
"net/url"
"strings" "strings"
"testing" "testing"
"time"
"git.unkin.net/unkin/bootapi/internal/model" "git.unkin.net/unkin/bootapi/internal/model"
"git.unkin.net/unkin/bootapi/internal/netbox" "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 { 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() t.Helper()
set, err := render.BuildSet(templates.FS, nil) set, err := render.BuildSet(templates.FS, nil)
if err != 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", DefaultDomain: "main.unkin.net", DefaultTemplate: "almalinux9",
RootPasswordHash: "$6$abc$def", RootPasswordHash: "$6$abc$def",
}, set) }, 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 { 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)
}
}
}