Compare commits

..
Author SHA1 Message Date
mrhid6 e6fe463216 feat: Updated for api ingress routes
Chart Release / chart (push) Successful in 25s
Server Deploy / deploy (push) Successful in 4m26s
2026-07-31 11:21:15 +01:00
mrhid6 8528f14ed7 feat: Added ingress to chart
Chart Release / chart (push) Successful in 10s
Server Deploy / deploy (push) Successful in 1m24s
2026-07-31 10:49:20 +01:00
mrhid6 df1d9658f5 fix: Chart build
Chart Release / chart (push) Successful in 10s
Server Deploy / deploy (push) Successful in 1m20s
2026-07-31 10:41:41 +01:00
mrhid6 9f9b384481 fix: Fixed chart version
Chart Release / chart (push) Failing after 13s
2026-07-31 10:37:04 +01:00
mrhid6 165114471f fix: Fixes to running on kubernetes
Chart Release / chart (push) Failing after 13s
Server Deploy / deploy (push) Successful in 6m35s
2026-07-31 10:34:10 +01:00
mrhid6 de78688093 feat: authenticate the server's Redis connection
InitRedis now takes a username and password, read from REDIS_USERNAME and
REDIS_PASSWORD, matching what admin has always done. Both empty keeps an
unauthenticated Redis working; a password with an empty username is what a
legacy requirepass instance needs, since go-redis then sends AUTH with one
argument instead of two.

This is what lets a Kubernetes install point at a managed Redis instead of
the bundled one.
2026-07-31 09:36:31 +01:00
33 changed files with 2174 additions and 287 deletions
+204
View File
@@ -0,0 +1,204 @@
name: Chart Release
on:
# Every push that touches the chart is validated. Publishing is separate and
# deliberate: a chart version is immutable in the registry once pushed, so
# it must come from a tag someone chose, not from whatever landed on main.
# No `paths` filter on push, deliberately. A paths filter applies to tag
# pushes too, so tagging a commit that happened not to touch the chart
# would skip the publish entirely — a release that silently does nothing.
# Validation is seconds of helm rendering; running it on every push to main
# is cheaper than that failure mode.
push:
branches:
- main
tags:
- "chart/v*"
pull_request:
paths:
- "deploy/chart/**"
workflow_dispatch:
env:
CHART_DIR: deploy/chart/vantage
HELM_VERSION: v3.16.3
jobs:
chart:
runs-on: ubuntu-docker
container: alpine:3.21
steps:
# git for actions/checkout, curl for both the Helm download and the
# registry upload, tar because the Helm tarball is not self-extracting.
- name: Setup
run: apk add --no-cache bash curl git tar nodejs npm
- name: Install Helm
run: |
set -eu
curl -fsSL "https://get.helm.sh/helm-${HELM_VERSION}-linux-amd64.tar.gz" \
| tar -xz -C /tmp linux-amd64/helm
mv /tmp/linux-amd64/helm /usr/local/bin/helm
helm version --short
- name: Checkout
uses: actions/checkout@v4
- name: Lint
run: helm lint "$CHART_DIR"
# Rendering is the real test. `helm lint` accepts a chart whose
# templates fail to execute, and every guard in this chart is a
# template `fail` that only fires during rendering.
- name: Render default values
run: helm template test "$CHART_DIR" > /dev/null
- name: Render a multi-replica install
run: |
helm template test "$CHART_DIR" \
--set server.replicaCount=3 \
--set web.replicaCount=3 > /dev/null
- name: Render against external Redis and MongoDB
run: |
helm template test "$CHART_DIR" \
--set redis.enabled=false \
--set redis.addr=redis.example.com:6379 \
--set mongo.enabled=false \
--set server.env.mongoUri=mongodb://mongo.example.com:27017/vantage > /dev/null
- name: Render with the Traefik ingress
run: |
helm template test "$CHART_DIR" \
--set ingress.enabled=true \
--set ingress.web.host=vantage.example.com \
--set ingress.grpc.host=agents.example.com \
--set ingress.tls.certResolver=letsencrypt \
--set server.env.grpcHost=agents.example.com:443 > /dev/null
# The shape the cloud deployment actually uses: a wildcard tenant
# namespace, /api and /auth routed at the edge, and no apex — that
# belongs to the marketing site, which this chart does not deploy.
- name: Render a wildcard host with edge-routed API paths
run: |
helm template test "$CHART_DIR" \
--set ingress.enabled=true \
--set 'ingress.web.host=*.vantage.example.com' \
--set ingress.api.enabled=true \
--set ingress.grpc.host=agents.example.com \
--set server.env.grpcHost=agents.example.com:443 \
--set ingress.tls.secretName=vantage-tls \
--set ingress.tls.grpcSecretName=agents-tls > /dev/null
# The guards are load-bearing, so their absence is a regression the
# same way a broken render is. Each of these must fail.
- name: Check the guards still refuse bad values
run: |
set -eu
refuses() {
desc="$1"; shift
if helm template test "$CHART_DIR" "$@" > /dev/null 2>&1; then
echo "GUARD MISSING: $desc was accepted"
exit 1
fi
echo "ok: refused $desc"
}
refuses "mongo disabled with an in-chart URI" \
--set mongo.enabled=false
refuses "redis disabled with no external address" \
--set redis.enabled=false
refuses "multiple replicas on a ReadWriteOnce volume" \
--set server.replicaCount=2 --set server.persistence.enabled=true
refuses "ingress with no web host" \
--set ingress.enabled=true
refuses "edge-routed API with an empty path list" \
--set ingress.enabled=true \
--set ingress.web.host=vantage.example.com \
--set ingress.grpc.enabled=false \
--set ingress.api.enabled=true \
--set 'ingress.api.paths=null'
refuses "gRPC ingress with no host" \
--set ingress.enabled=true \
--set ingress.web.host=vantage.example.com \
--set server.env.grpcHost=agents.example.com:443
refuses "gRPC ingress while grpcHost is still in-cluster" \
--set ingress.enabled=true \
--set ingress.web.host=vantage.example.com \
--set ingress.grpc.host=agents.example.com
- name: Read the chart version
id: chart
run: |
set -eu
VERSION="$(grep '^version:' "$CHART_DIR/Chart.yaml" | awk '{print $2}')"
echo "version=$VERSION" >> "$GITHUB_OUTPUT"
echo "chart version is $VERSION"
# Chart.yaml is the source of truth for the version; the tag only
# says "publish this one". A mismatch is a mistake worth stopping
# for — the alternative is stamping the tag over Chart.yaml, which
# leaves the repository disagreeing with what was published.
- name: Check the tag matches Chart.yaml
if: startsWith(github.ref, 'refs/tags/chart/v')
run: |
set -eu
TAG_VERSION="${GITHUB_REF_NAME#chart/v}"
CHART_VERSION="${{ steps.chart.outputs.version }}"
if [ "$TAG_VERSION" != "$CHART_VERSION" ]; then
echo "tag chart/v$TAG_VERSION does not match Chart.yaml version $CHART_VERSION"
echo "bump version: in $CHART_DIR/Chart.yaml, or retag."
exit 1
fi
- name: Package
run: |
set -eu
mkdir -p dist
helm package "$CHART_DIR" --destination dist
ls -l dist
- name: Publish to the Gitea chart registry
if: startsWith(github.ref, 'refs/tags/chart/v')
env:
# github.server_url is this Gitea instance, so the registry
# host needs no variable of its own and cannot drift from it.
REGISTRY: ${{ github.server_url }}/api/packages/${{ github.repository_owner }}/helm/api/charts
# The same pair server-deploy.yml uses for `docker login`.
# RELEASE_TOKEN, not REGISTRY_PASSWORD: the latter is named in
# the docs but set by no workflow, and an unset secret becomes
# an empty password, which Gitea reports as "Failed to
# authenticate user" rather than as a missing credential.
REGISTRY_USER: ${{ secrets.REGISTRY_USER }}
REGISTRY_TOKEN: ${{ secrets.RELEASE_TOKEN }}
CHART_VERSION: ${{ steps.chart.outputs.version }}
run: |
set -eu
PKG="dist/vantage-${CHART_VERSION}.tgz"
test -f "$PKG"
# Checked explicitly, because the failure it prevents is a
# 401 that looks like a permissions problem on the token that
# was never sent.
if [ -z "${REGISTRY_USER}" ] || [ -z "${REGISTRY_TOKEN}" ]; then
echo "REGISTRY_USER or RELEASE_TOKEN is not set on this repository."
echo "RELEASE_TOKEN needs the write:package scope to publish a chart."
exit 1
fi
echo "publishing to ${REGISTRY} as ${REGISTRY_USER}"
# --fail-with-body so an HTTP error is a failed step with the
# server's explanation, rather than a green run that published
# nothing. A repeated version is rejected by the registry;
# that is the intended behaviour, not something to retry past.
curl --fail-with-body -sS \
--user "${REGISTRY_USER}:${REGISTRY_TOKEN}" \
-X POST \
--upload-file "$PKG" \
"$REGISTRY"
echo "published vantage ${CHART_VERSION}"
echo " helm repo add vantage ${{ github.server_url }}/api/packages/${{ github.repository_owner }}/helm"
echo " helm install vantage vantage/vantage --version ${CHART_VERSION}"
+100 -7
View File
@@ -120,7 +120,7 @@ Upload a public key, assign it per server, revoke softly. The agent diffs desire
### Workflows
A library of reusable **steps** (bash or PowerShell scripts with declared inputs, outputs, and secret refs) composed into **workflows** targeting a set of servers. Running one snapshots the resolved steps into a `WorkflowRun`, then dispatches `RunStepCmd` over the agent command stream. Step stdout/stderr streams back as `StepOutputChunk` and is written to a log file on disk; the UI streams it live. Steps support `on_failure: stop|continue|retry`, per-run env passed between steps via `output_env`, and a per-run workspace directory the agent cleans up at the end.
A library of reusable **steps** (bash or PowerShell scripts with declared inputs, outputs, and secret refs) composed into **workflows** targeting a set of servers. Running one snapshots the resolved steps into a `WorkflowRun`, then dispatches `RunStepCmd` over the agent command stream. Step stdout/stderr streams back as `StepOutputChunk` and is written to MongoDB (`workflow_log_lines`, one document per line); the UI streams it live. Steps support `on_failure: stop|continue|retry`, per-run env passed between steps via `output_env`, and a per-run workspace directory the agent cleans up at the end.
Default steps are seeded per org at boot (`SeedDefaultSteps`) from `VANTAGE_DEFAULT_STEPS_DIR`, which `server/Dockerfile` bakes to `/opt/default-steps` from the repo's `default_steps/`. Deliberately **not** under `/data` — that is a bind mount, so the library would be editable from the host. Adding a step there means committing a file and rebuilding, which is why `default_steps/` is in the `server` rebuild trigger. **Steps with `source: "default"` are read-only**: `UpdateStep`/`DeleteStep` refuse with `ErrDefaultStep` (409), because seeding rewrites them on every boot, so an edit would silently revert and a delete would come back. `web/` mirrors this — the step modal opens read-only, Delete is hidden, and the designer's per-step script override is `readOnly` for a default library step — but as elsewhere, the API is the boundary and the UI is the courtesy. Seeding writes straight to the collection rather than through `UpdateStep`, so the guard does not lock out the seeder. Logs are swept by retention (`workflow_log_retention_days`; nil = 30 days, 0 = forever).
@@ -155,6 +155,56 @@ encrypted, single-use, and consumed when the tunnel opens. None of them reach
the agent — the session is negotiated end-to-end between guacd and the target
daemon, so the agent relays bytes it cannot read.
### Running more than one server replica
An agent's `CommandStream` terminates on exactly **one** server process. Every
piece of coordination below exists because of that single fact: with several
replicas, the process asked to do something to an agent is almost never the
process holding that agent's stream.
`server/internal/bus` is the Redis message bus that closes the gap. It adds no
infrastructure — Redis was already required for sessions — and it is **not
optional on a single-replica deployment**: dispatch takes the bus path always,
so the code running in production is the code running everywhere, rather than a
rare cross-pod branch that only fails under load.
| Concern | How it crosses replicas |
| ---------------------- | --------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| Which pod owns an agent | `vantage:agent:<server_id>` holds the owner's node ID with a 30s TTL, renewed every 10s. `Dispatcher.IsConnected` is an `EXISTS` on it |
| Sending a command | published to `vantage:cmd:<server_id>`; the owner pod acks on `vantage:ack:<command_id>`. **Request/ack, not a queue** — a command whose owner died must fail loudly (503) rather than queue |
| Step results | the owner pod publishes to `vantage:res:<command_id>`; the pod driving the run subscribes **before** dispatching, or a fast agent answers into a channel nobody has joined |
| Step output | never crosses. The dispatch envelope carries the secret mask list, so the owner pod masks and writes lines itself — unmasked bytes stay off the bus |
| Console relay | the envelope asks the owner pod to bind the listener, and the ack returns **that pod's** address for guacd. The relay's failure reason comes back on `vantage:proxyend:<proxy_id>` |
| Background jobs | `bus.RunAsLeader` — one Redis lock named `housekeeping` |
**Workflow logs are in MongoDB** (`workflow_log_lines`, one document per line,
with a `workflow_log_seq` counter document per run/server). Two pods write the
same log concurrently — the run's pod emits markers, the agent's pod emits
output — so ordering only means anything if both draw sequence numbers from the
same counter. `StepRun.log_offset` is that sequence number now, not a byte
offset. Writes are batched (128 lines or 250ms) and capped: 8 KB per line,
200k lines per server-run, after which one final `[vantage] log truncated`
marker is written and the rest is dropped. Without that cap a `yes` in a step
is a database incident. **Nothing writes to `/data` any more**, which is why
`server.persistence` now defaults to off and `VANTAGE_WORKFLOW_LOG_DIR` is gone.
**The leader lock is not an optimisation.** N replicas each running the monitor
scheduler means each check fires N times, each incident notification reaches the
customer N times, and each hourly rollup is written N times; N reapers race to
purge the same Free instance. `monitorsched`, `StartReaper`, `StartLogSweeper`,
`StartAuditSweeper` and the offline sweep therefore all run inside one
`RunAsLeader("housekeeping", …)` — one role, one lock. Each takes a context
cancelled the instant leadership is lost, and must return when it is.
Redis rather than a Kubernetes `Lease` so Compose takes the identical path: one
implementation to reason about, not two.
Two deployment requirements come with `replicaCount > 1`: every replica must
share **one** Redis (a per-pod Redis partitions the bus and every agent looks
offline to two thirds of the fleet), and `POD_IP` must be set — the chart does
it from the downward API — because `PROXY_ADVERTISE_HOST` names the Service, and
a Service cannot address the one pod holding a console listener.
### Inventory and OS updates
Agents report CPU/memory/swap/partitions/kernel — metrics every 30s, full static snapshot every 15 min. They also check for pending OS package updates hourly and can apply them on command (`ApplyUpdatesCmd`).
@@ -330,6 +380,7 @@ Key-state polling stays on the 30s `SyncKeys` interval. Full message definitions
Unauthenticated:
```
GET /healthz /readyz # liveness / readiness probes
GET /install /install.ps1 # dynamic agent install scripts
GET /update /update.ps1
GET /auth/bootstrap-status
@@ -436,7 +487,7 @@ Paddle is merchant of record; `admin/internal/paddle` is a thin REST client (no
## MongoDB Collections
`servers` · `keys` · `assignments` · `orgs` · `users` · `org_oidc` · `settings` · `secrets` · `workflows` · `workflow_steps` · `workflow_runs` · `monitors` · `incidents` · `monitor_rollups` · `notification_channels` · `console_sessions` · `audit_logs` · `migrations`
`servers` · `keys` · `assignments` · `orgs` · `users` · `org_oidc` · `settings` · `secrets` · `workflows` · `workflow_steps` · `workflow_runs` · `workflow_log_lines` · `workflow_log_seq` · `monitors` · `incidents` · `monitor_rollups` · `notification_channels` · `console_sessions` · `audit_logs` · `migrations`
Every document except `migrations` carries `org_id`. Struct definitions are the source of truth — see `server/internal/models/`.
@@ -448,6 +499,7 @@ Notes that are not obvious from the structs:
- `assignments.revoked_at: null` means active. Revocation is soft, preserving audit history.
- `workflow_runs.steps_snapshot` freezes the resolved steps so editing the library never rewrites history.
- `console_sessions.token_consumed_at` is set atomically to enforce one-time use.
- `workflow_log_lines` is keyed `(run_id, server_id, seq)` — the index is not an optimisation, every read is a range scan over it. `workflow_log_seq` holds one counter document per `run_id/server_id`, which is what lets two pods interleave into one ordered log. Neither carries `instance_id`: they are reached only through a run, and a run is already scoped.
- `users.auth_source` is `local`, `oidc` or `hq`. An `hq` user was projected from a Vantage HQ account and carries `hq_user_id`; HQ owns its role, password and existence.
Admin's own database is separate and holds `accounts` · `admin_instances` · `licenses` · `subscriptions` · `plans` · `catalogue` · `entitlements` · `paddle_events` · `staff_users` · `customer_users` · `instance_members` · `admin_audit`. `paddle_events` is the webhook idempotency log, unique on `event_id`: an event is claimed there before processing, and a duplicate of a handled event is a 200 no-op. `instance_members` is unique on `(instance_id, customer_user_id)` — one person holds at most one user in one instance, which makes a grant idempotent-by-refusal rather than silently doubling a projection. It is an _index_ of the control-plane rows, not the authority (see "Grants project, they do not federate"). Admin has no migrations collection; `models.Backfill` runs on every boot and is idempotent by filtering on the absence of what it writes.
@@ -529,13 +581,17 @@ Windows: MSI built by CI (WiX), or `installer/setup.ps1` registering the agent a
| `GRPC_HOST` | **yes** | `host:port` agents dial. Boot fails without it — there is no safe default; falling back to the web host would hand agents a port that does not speak gRPC. |
| `MONGO_URI` | no | default `mongodb://localhost:27017` |
| `MONGO_DB` | no | default `vantage` |
| `REDIS_USERNAME` | no | Redis 6+ ACL user. Leave empty for a legacy `requirepass` instance — go-redis then sends AUTH with one argument instead of two |
| `REDIS_PASSWORD` | no | empty for an unauthenticated Redis |
| `REDIS_ADDR` | no | default `localhost:6379` |
| `KEY_ENCRYPTION_KEY` | yes in practice | 64-char hex (32 bytes) for AES-256-GCM. Required for private keys, secrets, OIDC secrets, RDP credentials. |
| `GUACD_ADDR` | no | default `guacd:4822` |
| `PROXY_ADVERTISE_HOST` | no | default `server`; the hostname guacd resolves the control plane by, handed to guacd as the relay's address. Wrong here and every console session fails at connect |
| `PROXY_LISTEN_HOST` | no | default `0.0.0.0`; the interface the ephemeral relay listener binds |
| `APP_ROOT_LABEL` | no | default `vantage`; wrong value disables the host/session org guard |
| `VANTAGE_WORKFLOW_LOG_DIR` | no | where run logs are written |
| `POD_IP` | no | this pod's own address, set by the Helm chart from the downward API. **Takes precedence over `PROXY_ADVERTISE_HOST`** — a console relay listener belongs to one replica, and a Service address names all of them |
| `VANTAGE_MIGRATE_ONLY` | no | run schema setup (migrations, index builders, default-step seeding) and exit without serving. `GRPC_HOST` is not required in this mode. Set by the Helm chart's pre-upgrade Job |
| `VANTAGE_SKIP_MIGRATIONS` | no | serve without running schema setup, on the assumption a Job already did. Set by the chart's Deployment whenever `server.migrationJob.enabled`. Unset under Compose, where one process still migrates and then serves |
| `FREE_INSTANCE_REAP_AFTER` | no | duration past a Free licence's expiry before the instance and all its data are deleted. **Empty disables the reaper, and empty is the default.** Set to `336h` in `docker-compose.site.yml` only — a self-hosted deployment must never reap. Must match admin's value, which only names the date in warning emails |
**sitesvc** (`deploy/docker-compose.site.yml` only):
@@ -550,11 +606,33 @@ Windows: MSI built by CI (WiX), or `installer/setup.ps1` registering the agent a
| `SITE_ORIGIN` | yes in practice | comma-separated allowed origins; unset refuses every cross-origin browser request |
| `TRUST_PROXY` | no | only `true` behind a proxy that overwrites `X-Forwarded-For`, or clients spoof past the rate limiter |
### Ingress (Helm, Traefik)
`ingress.enabled` publishes **two** hostnames, because the two audiences arrive over different protocols:
| Values | Route |
| -------------------------------------------- | --------------------------------------------------------------------------------------------- |
| `ingress.web.host` (+ `web.extraHosts`) | browsers → `web:3000` |
| `ingress.api.paths` (when `api.enabled`) | `/api`, `/auth` → `<release>-server:8080`, bypassing the Next proxy |
| `ingress.grpc.host` | agents → a dedicated `<release>-server-grpc` Service on 9090, annotated `serversscheme: h2c` |
**`ingress.web.host` is normally a wildcard.** `*.vantage.example.com` is the per-tenant instance namespace — `APP_ROOT_LABEL` resolves the instance from the label. A Kubernetes wildcard host matches **exactly one** label, so it does not match the apex, and here that is correct rather than a gap: `vantage.hostxtra.co.uk` is the marketing site (`site/`, in `docker-compose.site.yml`), which this chart does not deploy. `extraHosts` is for a genuine second name; adding the apex to it would put the control plane on the marketing host. Every host in the list gets identical paths.
**`ingress.api.enabled` routes `/api` and `/auth` straight to the server.** Both arrangements work — without it `web` proxies those prefixes onward itself (`web/next.config.ts`) — but edge routing is one hop shorter and matches what the Nginx Proxy Manager in front of the Docker deployment already does, so leaving it off makes the request path a different shape on Kubernetes than in production. It stays **off by default** because it only helps where the server is reachable on the same host and certificate as `web`; turning it on blindly moves the whole API onto a route that may not be provisioned. Traefik derives router priority from rule length, so `PathPrefix(/api)` outranks the catch-all `/` with no priority annotation needed.
**The gRPC route needs its own Service.** The server terminates no TLS; it speaks plain h2c and always has, with TLS terminated by whatever sits in front. Traefik will not use h2c to a backend unless the *Service* says so, and that annotation applies to every port on the Service — so annotating the shared two-port `<release>-server` would force h2c on its HTTP port too.
**`server.env.grpcHost` is not derived from `ingress.grpc.host`, and the chart refuses to render if they disagree.** Agents dial whatever `grpcHost` says, and it is baked into every install one-liner; left pointing at the in-cluster Service while agents arrive through the ingress, every install succeeds and every agent then fails to connect, with nothing in the control plane explaining why. Guessing at the port (443? 9090?) would be worse than stopping.
TLS is `ingress.tls.secretName` / `grpcSecretName` (pre-existing certificates) **or** `certResolver` (Traefik ACME). Setting neither while `tls.enabled` produces a TLS router with no certificate, so Traefik serves its self-signed default — which looks valid and is trusted by nothing. NOTES.txt warns on install rather than the chart failing, since it is a real if unusual choice behind another terminator.
---
`deploy/docker-compose.yml` runs four services: `redis`, `guacd`, `server` (8080 + 9090), `web` (3000). MongoDB is external. `deploy/docker-compose.site.yml` adds five more — `site` (3003), `sitesvc` (8082), `admin` (8083), `adminsite` (3004) and `docsite` (3005) — and is only used on vantage.hostxtra.co.uk.
`docsite` is the odd one: a **static** build served by `nginx:alpine-slim`, not a Node runtime, and it listens on `80` rather than `3000`. It is reached at **`vantage.hostxtra.co.uk/docs`** — a path on the marketing host, routed by its own Nginx Proxy Manager location, which must sort **above** the catch-all forwarding to `site:3003` or Next answers the 404. A path and not a subdomain because `*.vantage.hostxtra.co.uk` is the per-tenant instance namespace and `APP_ROOT_LABEL` would read a `docs.` label as a tenant slug. NPM forwards the **full** path upstream — it does not strip `/docs` — so `DOCS_BASE_URL`, the proxy location and the directory the image copies the build into (`/usr/share/nginx/html/docs`) must all agree. When they do not, the HTML loads and every asset 404s.
`LICENSE_SIGNING_KEY` appears in **exactly one service in exactly one compose file**: `admin` in `docker-compose.site.yml`. It must never be added to `server`, and the self-hosted `docker-compose.yml` must never mention `admin` or `adminsite` at all. Admin uses an external Redis via `REDIS_ADDR`/`REDIS_USERNAME`/`REDIS_PASSWORD`; the base compose hardcodes `redis:6379` for `server`, so those variables reach admin only.
`LICENSE_SIGNING_KEY` appears in **exactly one service in exactly one compose file**: `admin` in `docker-compose.site.yml`. It must never be added to `server`, and the self-hosted `docker-compose.yml` must never mention `admin` or `adminsite` at all. Admin uses an external Redis via `REDIS_ADDR`/`REDIS_USERNAME`/`REDIS_PASSWORD`. `server` now reads the same three, so a Kubernetes install can point at a managed Redis; the base compose still hardcodes an unauthenticated `redis:6379` for it, so in Docker those credentials remain admin's alone.
---
@@ -668,10 +746,24 @@ cd /opt/vantage && docker compose -f docker-compose.yml -f docker-compose.site.y
The gap this leaves: **changing a repo variable pushes no commit, so nothing rebuilds.** After editing `ADMIN_API_URL`, `HQ_URL` or `ADMIN_ENV`, run the workflow manually — that is what `workflow_dispatch` is there for. Base images also stop being refreshed on a service nobody touches; a periodic manual run covers that.
### `chart-release.yml` — validates on every chart change, publishes on `chart/v*` tags
Two jobs' worth of work in one, split by trigger. Any push or PR touching `deploy/chart/` lints the chart and renders it four ways: defaults, a multi-replica install, external Redis and MongoDB, and a set of values that **must be refused**. That last one is the point — every safety rail in this chart is a template `fail`, and `helm lint` happily accepts a chart whose templates never execute, so only rendering proves they still fire.
Publishing runs only on a `chart/v*` tag, to the Gitea Helm registry at `/api/packages/<owner>/helm/api/charts`. **`Chart.yaml` is the source of truth for the version**; the tag only selects which one to publish, and a tag that disagrees with `Chart.yaml` fails rather than stamping over it — the alternative leaves the repository disagreeing with what shipped. A version already in the registry is rejected by Gitea, which is intended: published chart versions are immutable.
The registry host comes from `github.server_url`, so it cannot drift from the instance the workflow is running on. It authenticates with `REGISTRY_USER` + **`RELEASE_TOKEN`** — the pair `server-deploy.yml` actually uses for `docker login`. `REGISTRY_PASSWORD` is listed in the secrets table below but set by no workflow; passing an unset secret yields an empty password and Gitea answers `401 Failed to authenticate user`, which reads like a scope problem on a token that was never sent. The publish step therefore checks both are non-empty before it calls curl. `RELEASE_TOKEN` needs `write:package` in addition to `write:release`.
```bash
helm repo add vantage https://gitea.hostxtra.co.uk/api/packages/mrhid6/helm
helm install vantage vantage/vantage --version 0.1.0
```
### Tagging
```bash
git tag agent/v1.0.0 && git push origin agent/v1.0.0 # agent release
git tag chart/v0.1.0 && git push origin chart/v0.1.0 # helm chart package
git push origin main # server + web deploy
```
@@ -679,9 +771,9 @@ git push origin main # server + web deploy
| Name | Type | Value |
| ----------------------- | -------- | ---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| `RELEASE_TOKEN` | Secret | Gitea API token, `write:release` |
| `REGISTRY_USER` | Secret | Gitea username |
| `REGISTRY_PASSWORD` | Secret | Gitea token, `write:packages` |
| `RELEASE_TOKEN` | Secret | Gitea API token. Needs `write:release` (agent releases), `write:package` (container images and the Helm chart). **This is the only token any workflow authenticates with** — `docker login` and the chart publish both pair it with `REGISTRY_USER` |
| `REGISTRY_USER` | Secret | Gitea username. Must own `RELEASE_TOKEN`, or basic auth is rejected |
| ~~`REGISTRY_PASSWORD`~~ | — | **Not used.** Named here historically; no workflow reads it. Referencing an unset secret yields an empty password and a `401 Failed to authenticate user` that looks like a token scope problem. Use `RELEASE_TOKEN` |
| `DOCKER_HOST` | Variable | registry host used for image tags |
| `API_URL` | **not** a CI variable | `web` reads it at **runtime**, from the container environment — `next.config.ts` is evaluated when `server.js` boots in standalone mode, and the rewrites it feeds are server-side, never browser-side. Default `http://localhost:8080`; compose sets `http://server:8080`. `NEXT_PUBLIC_API_URL` is still honoured as a fallback for existing deployments. |
| `SITE_API_URL` | Variable | **browser-reachable** sitesvc URL, baked into the `site` image. Required — if empty, both forms report "not connected" and submit nowhere. Must also be in sitesvc's `SITE_ORIGIN`. |
@@ -715,6 +807,7 @@ git push origin main # server + web deploy
- **`org_id` on every document** — isolation enforced at the query layer, not by separate databases.
- **root only** — manages `/root/.ssh/authorized_keys`; no per-user key management.
- **Windows agents are second-class by design** — register, heartbeat, run steps, report inventory; no `authorized_keys` management.
- **Both `server` and `web` scale horizontally** — see "Running more than one server replica" below. `web` holds nothing; `server` holds per-agent state that is routed between replicas over Redis rather than duplicated.
- **Deletion lives in the control plane** — admin sends the warnings because it knows the billing address; the control plane performs the delete because it is the only service that knows which collections carry `instance_id`. Mirroring that list into admin would drift, and a drift there deletes the wrong rows.
## graphify
+2 -2
View File
@@ -2,5 +2,5 @@ apiVersion: v2
name: vantage
description: Helm chart for the Vantage stack (Redis, MongoDB, guacd, server, web)
type: application
version: 0.1.0
appVersion: "1.0.0"
version: 1.0.3
appVersion: "1.0.3"
+55
View File
@@ -1,12 +1,67 @@
Vantage has been deployed as release "{{ .Release.Name }}" in namespace "{{ .Release.Namespace }}".
Services created:
{{- if .Values.redis.enabled }}
- {{ .Release.Name }}-redis (ClusterIP {{ .Values.redis.port }})
{{- else }}
- Redis: not deployed, using external {{ .Values.redis.addr }}
{{- end }}
{{- if .Values.mongo.enabled }}
- {{ .Release.Name }}-mongo (ClusterIP {{ .Values.mongo.port }})
{{- else }}
- MongoDB: not deployed, using the external server.env.mongoUri
{{- end }}
- {{ .Release.Name }}-guacd ({{ .Values.guacd.service.type }} {{ .Values.guacd.service.port }})
- {{ .Release.Name }}-server ({{ .Values.server.service.type }} http:{{ .Values.server.service.httpPort }} grpc:{{ .Values.server.service.grpcPort }})
- {{ .Release.Name }}-web ({{ .Values.web.service.type }} {{ .Values.web.service.port }})
Scaling (server.replicaCount / web.replicaCount):
- Both scale. Pin the image tags first — replicas on different builds serve
mismatched web asset hashes, and mixed server versions share one bus.
- server replicas route agent commands, step results and console relays to
each other over Redis, so every replica must use the SAME Redis. Workflow
logs are in MongoDB, not on a volume.
- Background work (monitor scheduler, Free reaper, log and audit retention,
the offline sweep) runs on one replica at a time under a Redis leader lock.
- server.persistence must be off to scale past one replica on a ReadWriteOnce
volume. Nothing writes to it any more.
{{- if gt (int .Values.server.replicaCount) 1 }}
- Console relays are reached by pod IP; guacd must be able to dial pod IPs
directly (it can, inside the cluster network).
{{- end }}
{{- if .Values.server.migrationJob.enabled }}
- Migrations run in the {{ .Release.Name }}-migrate Job before each upgrade;
the pods skip them. Its logs are kept: kubectl logs job/{{ .Release.Name }}-migrate
{{- end }}
{{- if .Values.ingress.enabled }}
Ingress (Traefik):
- Browsers: https://{{ .Values.ingress.web.host }}
{{- range .Values.ingress.web.extraHosts }}
https://{{ . }}
{{- end }}
{{- if .Values.ingress.api.enabled }}
{{ join ", " .Values.ingress.api.paths }} go straight to the server; everything else to web.
{{- else }}
Everything goes to web, which proxies /api and /auth onward. Set
ingress.api.enabled=true to route them at the edge instead.
{{- end }}
{{- if .Values.ingress.grpc.enabled }}
- Agents: {{ .Values.ingress.grpc.host }} (gRPC, h2c behind TLS)
Agents dial server.env.grpcHost, currently {{ tpl .Values.server.env.grpcHost . }}.
Point DNS for both hostnames at the Traefik load balancer.
{{- if not .Values.ingress.tls.enabled }}
- WARNING: ingress.tls.enabled is false. Agent tokens and session cookies
would cross the network in clear.
{{- else if and (not .Values.ingress.tls.certResolver) (not .Values.ingress.tls.secretName) }}
- WARNING: TLS is on but neither ingress.tls.secretName nor
ingress.tls.certResolver is set, so Traefik will serve its self-signed
default certificate, which no agent and no browser will trust.
{{- end }}
{{- end }}
{{- end }}
By default the server/web/guacd services are ClusterIP only (no host port publishing,
unlike the original docker-compose file). To expose them externally, set
server.service.type / web.service.type / guacd.service.type to NodePort or LoadBalancer,
@@ -9,3 +9,76 @@ Common name helpers
app.kubernetes.io/instance: {{ .Release.Name }}
app.kubernetes.io/managed-by: {{ .Release.Service }}
{{- end -}}
{{/*
vantage.server.env renders the server container's environment.
It lives here because two workloads need it identically: the Deployment and the
pre-upgrade migration Job. A Job that connected to a different database than the
pods it migrates for would be worse than no Job at all, so there is one copy and
both read it.
*/}}
{{- define "vantage.server.env" -}}
- name: MONGO_URI
{{- $mongoUri := tpl .Values.server.env.mongoUri . }}
{{- if and (not .Values.mongo.enabled) (contains (printf "%s-mongo" .Release.Name) $mongoUri) }}
{{- fail "mongo.enabled is false, so server.env.mongoUri must point at an external MongoDB rather than the in-chart one" }}
{{- end }}
value: {{ $mongoUri | quote }}
- name: REDIS_ADDR
{{- if .Values.redis.enabled }}
value: "{{ .Release.Name }}-redis:{{ .Values.redis.port }}"
{{- else }}
{{- if not .Values.redis.addr }}
{{- fail "redis.enabled is false, so redis.addr must be set to an external Redis host:port" }}
{{- end }}
value: {{ .Values.redis.addr | quote }}
{{- end }}
{{- if .Values.redis.auth.existingSecret }}
- name: REDIS_USERNAME
valueFrom:
secretKeyRef:
name: {{ .Values.redis.auth.existingSecret }}
key: {{ .Values.redis.auth.usernameKey }}
optional: true
- name: REDIS_PASSWORD
valueFrom:
secretKeyRef:
name: {{ .Values.redis.auth.existingSecret }}
key: {{ .Values.redis.auth.passwordKey }}
{{- else }}
{{- if .Values.redis.auth.username }}
- name: REDIS_USERNAME
value: {{ .Values.redis.auth.username | quote }}
{{- end }}
{{- if .Values.redis.auth.password }}
- name: REDIS_PASSWORD
value: {{ .Values.redis.auth.password | quote }}
{{- end }}
{{- end }}
- name: GRPC_HOST
value: {{ .Values.server.env.grpcHost | quote }}
- name: GRPC_PORT
value: {{ .Values.server.service.grpcPort | quote }}
- name: HTTP_PORT
value: {{ .Values.server.service.httpPort | quote }}
- name: KEY_ENCRYPTION_KEY
value: {{ .Values.server.env.keyEncryptionKey | quote }}
- name: GUACD_ADDR
value: "{{ .Release.Name }}-guacd:{{ .Values.guacd.service.port }}"
- name: APP_ROOT_LABEL
value: {{ .Values.server.env.appRootLabel | quote }}
- name: PROXY_ADVERTISE_HOST
value: {{ .Values.server.env.proxyAdvertiseHost | quote }}
- name: PROXY_LISTEN_HOST
value: {{ .Values.server.env.proxyListenHost | quote }}
# The address guacd dials to reach a console relay. It must name one pod, not
# the Service: the relay listener is bound by whichever pod holds that agent's
# command stream, and a Service would send guacd to a different one. POD_IP
# takes precedence over PROXY_ADVERTISE_HOST in the server for exactly this
# reason, so the setting above stays meaningful only outside Kubernetes.
- name: POD_IP
valueFrom:
fieldRef:
fieldPath: status.podIP
{{- end -}}
+173
View File
@@ -0,0 +1,173 @@
{{- if .Values.ingress.enabled }}
{{/*
Two hostnames, because the two audiences arrive over different protocols.
Browsers reach the web host. What answers there depends on the path: with
ingress.api.enabled, /api and /auth go straight to the server and everything
else goes to `web`. Without it, everything goes to `web`, which proxies those
prefixes onward itself (web/next.config.ts).
Both work. Routing at the edge is one hop shorter and is what the Nginx Proxy
Manager deployment in front of the Docker install already does, so leaving it
off changes the shape of the request path between the two deployments. It is
still off by default, because turning it on where `web` is the only thing with
a public certificate would strand /api behind a route nobody can reach.
The web host is normally a wildcard — `*.vantage.example.com` — because that is
the per-tenant instance namespace; APP_ROOT_LABEL resolves the instance from the
label. Kubernetes wildcard hosts match exactly one label, so this does not match
the apex, and on the deployment this chart was written for it must not: the apex
is the marketing site, a separate application (see site/ and
docker-compose.site.yml). extraHosts exists for a genuine second name, not for
reclaiming the apex.
Agents reach the server's gRPC port, which is plain h2c — the server holds no
certificates of its own, TLS has always been terminated by whatever sits in
front. Traefik will not speak h2c to a backend unless told to, and it is told
per Service, which is why the gRPC route gets a Service of its own below rather
than reusing the two-port one. Annotating the shared Service would force h2c on
its HTTP port too.
*/}}
{{- $tls := .Values.ingress.tls }}
{{- $webHost := required "ingress.enabled requires ingress.web.host" .Values.ingress.web.host }}
{{- $hosts := prepend .Values.ingress.web.extraHosts $webHost }}
{{- $apiPaths := .Values.ingress.api.paths }}
{{- if and .Values.ingress.api.enabled (not $apiPaths) }}
{{- fail "ingress.api.enabled requires at least one path in ingress.api.paths" }}
{{- end }}
apiVersion: networking.k8s.io/v1
kind: Ingress
metadata:
name: {{ .Release.Name }}-web
labels:
{{- include "vantage.labels" . | nindent 4 }}
app.kubernetes.io/component: web
annotations:
traefik.ingress.kubernetes.io/router.entrypoints: {{ .Values.ingress.entrypoint | quote }}
{{- if and $tls.enabled $tls.certResolver }}
traefik.ingress.kubernetes.io/router.tls: "true"
traefik.ingress.kubernetes.io/router.tls.certresolver: {{ $tls.certResolver | quote }}
{{- else if $tls.enabled }}
traefik.ingress.kubernetes.io/router.tls: "true"
{{- end }}
{{- with .Values.ingress.annotations }}
{{- toYaml . | nindent 4 }}
{{- end }}
spec:
{{- if .Values.ingress.className }}
ingressClassName: {{ .Values.ingress.className }}
{{- end }}
{{- if and $tls.enabled $tls.secretName }}
tls:
- hosts:
{{- range $hosts }}
- {{ . | quote }}
{{- end }}
secretName: {{ $tls.secretName }}
{{- end }}
rules:
{{- range $host := $hosts }}
- host: {{ $host | quote }}
http:
paths:
{{- /*
The API paths come first and, more importantly, are longer. Traefik
derives router priority from rule length, so Host(x) &&
PathPrefix(/api) outranks Host(x) && PathPrefix(/) without anyone
having to assign priorities by hand. Order within the list is for
the reader; the length is what decides.
*/}}
{{- if $.Values.ingress.api.enabled }}
{{- range $apiPaths }}
- path: {{ . | quote }}
pathType: Prefix
backend:
service:
name: {{ $.Release.Name }}-server
port:
number: {{ $.Values.server.service.httpPort }}
{{- end }}
{{- end }}
- path: /
pathType: Prefix
backend:
service:
name: {{ $.Release.Name }}-web
port:
number: {{ $.Values.web.service.port }}
{{- end }}
{{- if .Values.ingress.grpc.enabled }}
{{- $grpcHost := required "ingress.grpc.enabled requires ingress.grpc.host" .Values.ingress.grpc.host }}
{{/*
GRPC_HOST is what an agent is told to dial, and it is baked into every install
one-liner. Left pointing at the in-cluster Service while agents are expected to
arrive through the ingress, every install would succeed and every agent would
fail to connect — with nothing in the control plane saying why.
*/}}
{{- $grpcEnv := tpl .Values.server.env.grpcHost . }}
{{- if contains (printf "%s-server" .Release.Name) $grpcEnv }}
{{- fail (printf "ingress.grpc.enabled routes agents through %s, but server.env.grpcHost is still the in-cluster address %q. Agents dial the value of grpcHost, so set it to the public gRPC address (for example %q)." $grpcHost $grpcEnv (printf "%s:443" $grpcHost)) }}
{{- end }}
---
apiVersion: v1
kind: Service
metadata:
name: {{ .Release.Name }}-server-grpc
labels:
{{- include "vantage.labels" . | nindent 4 }}
app.kubernetes.io/component: server
annotations:
# The server speaks h2c: it terminates no TLS itself. Without this Traefik
# dials the backend as HTTP/1.1 and every agent handshake fails.
traefik.ingress.kubernetes.io/service.serversscheme: h2c
spec:
type: ClusterIP
selector:
app.kubernetes.io/instance: {{ .Release.Name }}
app.kubernetes.io/component: server
ports:
- name: grpc
port: {{ .Values.server.service.grpcPort }}
targetPort: {{ .Values.server.service.grpcPort }}
---
apiVersion: networking.k8s.io/v1
kind: Ingress
metadata:
name: {{ .Release.Name }}-grpc
labels:
{{- include "vantage.labels" . | nindent 4 }}
app.kubernetes.io/component: server
annotations:
traefik.ingress.kubernetes.io/router.entrypoints: {{ .Values.ingress.entrypoint | quote }}
{{- if and $tls.enabled $tls.certResolver }}
traefik.ingress.kubernetes.io/router.tls: "true"
traefik.ingress.kubernetes.io/router.tls.certresolver: {{ $tls.certResolver | quote }}
{{- else if $tls.enabled }}
traefik.ingress.kubernetes.io/router.tls: "true"
{{- end }}
{{- with .Values.ingress.grpc.annotations }}
{{- toYaml . | nindent 4 }}
{{- end }}
spec:
{{- if .Values.ingress.className }}
ingressClassName: {{ .Values.ingress.className }}
{{- end }}
{{- if and $tls.enabled $tls.grpcSecretName }}
tls:
- hosts:
- {{ $grpcHost | quote }}
secretName: {{ $tls.grpcSecretName }}
{{- end }}
rules:
- host: {{ $grpcHost | quote }}
http:
paths:
- path: /
pathType: Prefix
backend:
service:
name: {{ .Release.Name }}-server-grpc
port:
number: {{ .Values.server.service.grpcPort }}
{{- end }}
{{- end }}
@@ -0,0 +1,68 @@
{{- if .Values.server.migrationJob.enabled }}
{{/*
Schema setup, lifted out of the serving pods.
Every server process used to run migrations, index builders and default-step
seeding at boot. With one replica that is fine. With two it is not: 0004 renames
the orgs collection to instances, and a sibling reading it mid-rename is a
corruption, not a retry.
A Helm hook Job runs it once, before any pod of the new version starts. The
Deployment then sets VANTAGE_SKIP_MIGRATIONS, which is what makes the Job's
existence load-bearing rather than decorative — if you disable the Job, the
pods go back to migrating themselves and you must go back to one replica.
hook-weight orders this after the dependency waits; before-hook-creation deletes
the previous Job so a repeat upgrade is not blocked by an immutable object. The
Job is deliberately NOT deleted on success: its logs are the record of what the
upgrade did to the database.
*/}}
apiVersion: batch/v1
kind: Job
metadata:
name: {{ .Release.Name }}-migrate
labels:
{{- include "vantage.labels" . | nindent 4 }}
app.kubernetes.io/component: migrate
annotations:
"helm.sh/hook": pre-install,pre-upgrade
"helm.sh/hook-weight": "0"
"helm.sh/hook-delete-policy": before-hook-creation
spec:
backoffLimit: {{ .Values.server.migrationJob.backoffLimit }}
# A migration that has not finished in this long is stuck, and a stuck
# migration should fail the upgrade rather than hold it open forever.
activeDeadlineSeconds: {{ .Values.server.migrationJob.activeDeadlineSeconds }}
template:
metadata:
labels:
app.kubernetes.io/instance: {{ .Release.Name }}
app.kubernetes.io/component: migrate
spec:
restartPolicy: Never
{{- if .Values.imagePullSecrets }}
imagePullSecrets:
{{- toYaml .Values.imagePullSecrets | nindent 8 }}
{{- end }}
{{- if .Values.mongo.enabled }}
# Only Mongo. The Job never opens Redis, and waiting on a Redis this
# chart may not even deploy would block an upgrade for no reason.
initContainers:
- name: wait-for-mongo
image: busybox:1.36
command:
- sh
- -c
- |
until nc -z {{ .Release.Name }}-mongo {{ .Values.mongo.port }}; do
echo "waiting for mongo..."; sleep 2;
done
{{- end }}
containers:
- name: migrate
image: "{{ .Values.server.image.repository }}:{{ .Values.server.image.tag }}"
env:
{{- include "vantage.server.env" . | nindent 12 }}
- name: VANTAGE_MIGRATE_ONLY
value: "true"
{{- end }}
@@ -1,3 +1,4 @@
{{- if .Values.mongo.enabled }}
{{- if .Values.mongo.persistence.enabled }}
apiVersion: v1
kind: PersistentVolumeClaim
@@ -88,3 +89,4 @@ spec:
ports:
- port: {{ .Values.mongo.port }}
targetPort: {{ .Values.mongo.port }}
{{- end }}
@@ -1,3 +1,4 @@
{{- if .Values.redis.enabled }}
{{- if .Values.redis.persistence.enabled }}
apiVersion: v1
kind: PersistentVolumeClaim
@@ -88,3 +89,4 @@ spec:
ports:
- port: {{ .Values.redis.port }}
targetPort: {{ .Values.redis.port }}
{{- end }}
+79 -28
View File
@@ -25,7 +25,19 @@ metadata:
{{- include "vantage.labels" . | nindent 4 }}
app.kubernetes.io/component: server
spec:
replicas: 1
{{- $replicas := int .Values.server.replicaCount }}
replicas: {{ $replicas }}
{{- if and .Values.server.persistence.enabled (eq .Values.server.persistence.accessMode "ReadWriteOnce") }}
# A ReadWriteOnce volume cannot be mounted by a second pod at all, and cannot
# be handed to a new pod while the old one still holds it. Persistence is off
# by default now that nothing writes to it; if it is on, replicas are capped
# at one and updates go through Recreate.
{{- if gt $replicas 1 }}
{{- fail "server.persistence.enabled with a ReadWriteOnce volume cannot be combined with server.replicaCount > 1. Nothing in the server writes to that volume any more (workflow logs live in MongoDB); set server.persistence.enabled=false, or use a ReadWriteMany accessMode if you are keeping it for another reason." }}
{{- end }}
strategy:
type: Recreate
{{- end }}
selector:
matchLabels:
app.kubernetes.io/instance: {{ .Release.Name }}
@@ -40,9 +52,13 @@ spec:
imagePullSecrets:
{{- toYaml .Values.imagePullSecrets | nindent 8 }}
{{- end }}
# Wait for redis & mongo to be reachable, approximating compose's
# `depends_on: condition: service_healthy`
{{- if or .Values.redis.enabled .Values.mongo.enabled }}
# Wait for the dependencies this chart deploys to be reachable,
# approximating compose's `depends_on: condition: service_healthy`. An
# external Redis or Mongo is assumed to be up already — waiting on one
# would only turn someone else's outage into a stuck pod.
initContainers:
{{- if .Values.redis.enabled }}
- name: wait-for-redis
image: busybox:1.36
command:
@@ -52,6 +68,8 @@ spec:
until nc -z {{ .Release.Name }}-redis {{ .Values.redis.port }}; do
echo "waiting for redis..."; sleep 2;
done
{{- end }}
{{- if .Values.mongo.enabled }}
- name: wait-for-mongo
image: busybox:1.36
command:
@@ -61,6 +79,8 @@ spec:
until nc -z {{ .Release.Name }}-mongo {{ .Values.mongo.port }}; do
echo "waiting for mongo..."; sleep 2;
done
{{- end }}
{{- end }}
containers:
- name: server
image: "{{ .Values.server.image.repository }}:{{ .Values.server.image.tag }}"
@@ -68,40 +88,71 @@ spec:
- containerPort: {{ .Values.server.service.httpPort }}
- containerPort: {{ .Values.server.service.grpcPort }}
env:
- name: MONGO_URI
value: {{ tpl .Values.server.env.mongoUri . | quote }}
- name: REDIS_ADDR
value: "{{ .Release.Name }}-redis:{{ .Values.redis.port }}"
- name: GRPC_HOST
value: {{ .Values.server.env.grpcHost | quote }}
- name: GRPC_PORT
value: {{ .Values.server.service.grpcPort | quote }}
- name: HTTP_PORT
value: {{ .Values.server.service.httpPort | quote }}
- name: KEY_ENCRYPTION_KEY
value: {{ .Values.server.env.keyEncryptionKey | quote }}
- name: VANTAGE_WORKFLOW_LOG_DIR
value: {{ .Values.server.env.vantageWorkflowLogDir | quote }}
- name: GUACD_ADDR
value: "{{ .Release.Name }}-guacd:{{ .Values.guacd.service.port }}"
- name: APP_ROOT_LABEL
value: {{ .Values.server.env.appRootLabel | quote }}
- name: PROXY_ADVERTISE_HOST
value: {{ .Values.server.env.proxyAdvertiseHost | quote }}
- name: PROXY_LISTEN_HOST
value: {{ .Values.server.env.proxyListenHost | quote }}
{{- include "vantage.server.env" . | nindent 12 }}
{{- if .Values.server.migrationJob.enabled }}
# Schema setup ran in the pre-upgrade Job. Pods that repeated it
# would race each other, and the rename migration is not a race
# that tolerates a loser.
- name: VANTAGE_SKIP_MIGRATIONS
value: "true"
{{- end }}
# Liveness never touches Mongo or Redis: restarting every pod cannot
# fix a database outage, and each restart drops every agent command
# stream and console session it was carrying. Readiness does check
# both, so a pod that cannot serve leaves the Service and stays up.
startupProbe:
httpGet:
path: /healthz
port: {{ .Values.server.service.httpPort }}
periodSeconds: 5
# Generous: without the migration Job this pod runs every migration
# before it listens, and the rename has a ten-minute budget.
failureThreshold: 150
livenessProbe:
httpGet:
path: /healthz
port: {{ .Values.server.service.httpPort }}
periodSeconds: 20
failureThreshold: 3
readinessProbe:
httpGet:
path: /readyz
port: {{ .Values.server.service.httpPort }}
periodSeconds: 10
failureThreshold: 3
{{- if .Values.server.persistence.enabled }}
# Nothing in the server writes here any more — workflow logs moved to
# MongoDB so that every replica can read and write them. The mount
# remains only so an operator upgrading from a file-log release can
# still reach the old files before turning persistence off.
volumeMounts:
- name: server-data
mountPath: /data
volumes:
- name: server-data
{{- if .Values.server.persistence.enabled }}
persistentVolumeClaim:
claimName: {{ .Release.Name }}-server-data
{{- else }}
emptyDir: {}
{{- end }}
---
{{- if gt (int .Values.server.replicaCount) 1 }}
apiVersion: policy/v1
kind: PodDisruptionBudget
metadata:
name: {{ .Release.Name }}-server
labels:
{{- include "vantage.labels" . | nindent 4 }}
app.kubernetes.io/component: server
spec:
# Agents reconnect on their own, but a drain that took every replica at once
# would disconnect every agent in the fleet simultaneously and stall every
# workflow run in flight.
minAvailable: 1
selector:
matchLabels:
app.kubernetes.io/instance: {{ .Release.Name }}
app.kubernetes.io/component: server
---
{{- end }}
apiVersion: v1
kind: Service
metadata:
+24 -1
View File
@@ -6,7 +6,9 @@ metadata:
{{- include "vantage.labels" . | nindent 4 }}
app.kubernetes.io/component: web
spec:
replicas: 1
# web holds no per-process state: sessions live in Redis and every request is
# proxied to the server. It is the one component here that scales freely.
replicas: {{ .Values.web.replicaCount }}
selector:
matchLabels:
app.kubernetes.io/instance: {{ .Release.Name }}
@@ -39,6 +41,27 @@ spec:
env:
- name: API_URL
value: {{ tpl .Values.web.env.apiUrl . | quote }}
# /healthz is served by this Next process; /api is rewritten to the
# server, so a probe there would report the backend's health and keep
# passing while this pod was wedged.
startupProbe:
httpGet:
path: /healthz
port: {{ .Values.web.service.port }}
periodSeconds: 3
failureThreshold: 20
livenessProbe:
httpGet:
path: /healthz
port: {{ .Values.web.service.port }}
periodSeconds: 20
failureThreshold: 3
readinessProbe:
httpGet:
path: /healthz
port: {{ .Values.web.service.port }}
periodSeconds: 10
failureThreshold: 3
---
apiVersion: v1
kind: Service
+39 -2
View File
@@ -1,6 +1,8 @@
# Default values for the vantage chart.
redis:
enabled: true
addr: ""
image:
repository: redis
tag: "8"
@@ -10,8 +12,15 @@ redis:
storageClass: ""
accessMode: ReadWriteOnce
port: 6379
auth:
username: ""
password: ""
existingSecret: ""
usernameKey: username
passwordKey: password
mongo:
enabled: true
image:
repository: mongo
tag: "7"
@@ -31,6 +40,11 @@ guacd:
port: 4822
server:
replicaCount: 1
migrationJob:
enabled: true
backoffLimit: 0
activeDeadlineSeconds: 900
image:
repository: gitea.hostxtra.co.uk/mrhid6/vantage/server
tag: latest
@@ -42,18 +56,18 @@ server:
mongoUri: "mongodb://{{ .Release.Name }}-mongo:27017/vantage"
grpcHost: "{{ .Release.Name }}-server:9090"
keyEncryptionKey: ""
vantageWorkflowLogDir: ""
appRootLabel: vantage
proxyAdvertiseHost: "{{ .Release.Name }}-server"
proxyListenHost: "0.0.0.0"
persistence:
enabled: true
enabled: false
size: 1Gi
storageClass: ""
accessMode: ReadWriteOnce
hostPath: /data
web:
replicaCount: 1
image:
repository: gitea.hostxtra.co.uk/mrhid6/vantage/web
tag: latest
@@ -63,4 +77,27 @@ web:
env:
apiUrl: "http://{{ .Release.Name }}-server:8080"
ingress:
enabled: false
className: traefik
entrypoint: websecure
annotations: {}
web:
host: ""
extraHosts: []
api:
enabled: false
paths:
- /api
- /auth
grpc:
enabled: true
host: ""
annotations: {}
tls:
enabled: true
secretName: ""
grpcSecretName: ""
certResolver: ""
imagePullSecrets: []
-1
View File
@@ -18,4 +18,3 @@ KEY_ENCRYPTION_KEY=
MONGO_URI=mongodb://mongo:27017/vantage
# Where workflow run logs are written inside the server container.
# VANTAGE_WORKFLOW_LOG_DIR=/data/workflow-logs
-1
View File
@@ -45,7 +45,6 @@ services:
GRPC_PORT: "9090"
HTTP_PORT: "8080"
KEY_ENCRYPTION_KEY: ${KEY_ENCRYPTION_KEY:-}
VANTAGE_WORKFLOW_LOG_DIR: ${VANTAGE_WORKFLOW_LOG_DIR:-}
GUACD_ADDR: guacd:4822
PROXY_ADVERTISE_HOST: server
depends_on:
@@ -50,7 +50,6 @@ GRPC_HOST=vantage.example.com:9090
KEY_ENCRYPTION_KEY=
# Optional: where workflow run logs are written inside the container.
VANTAGE_WORKFLOW_LOG_DIR=/data/workflow-logs
```
Generate the encryption key:
@@ -13,14 +13,16 @@ it is absent.
| -------------------------- | --------------- | --------------------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------ |
| `GRPC_HOST` | **yes** | | The `host:port` agents dial. Boot fails without it. There is deliberately no fallback to the web host: that would hand every agent a port that does not speak gRPC |
| `MONGO_URI` | no | `mongodb://localhost:27017` | The database name is taken from the URI path, falling back to `vantage`. There is no separate `MONGO_DB` |
| `REDIS_ADDR` | no | `localhost:6379` | Sessions only |
| `REDIS_ADDR` | no | `localhost:6379` | Sessions, and the bus that routes agent commands between server replicas. Every replica must point at the **same** Redis |
| `REDIS_USERNAME` | no | | Redis 6+ ACL user. Leave empty against a legacy `requirepass` instance, which authenticates with the password alone |
| `REDIS_PASSWORD` | no | | Leave empty for an unauthenticated Redis. Both of these exist so an install can use a managed Redis rather than the bundled one |
| `KEY_ENCRYPTION_KEY` | yes in practice | | 64 hex characters (32 bytes) for AES-256-GCM. Required for private keys, vault secrets, OIDC client secrets and console credentials |
| `GITEA_HOST` | yes | `gitea.example.com` | Used to build the install scripts and agent download URLs. The default is a placeholder that will not resolve |
| `GUACD_ADDR` | no | `guacd:4822` | The [browser console](../vantage/browser-console.md) daemon |
| `PROXY_ADVERTISE_HOST` | no | `server` | The hostname **guacd** uses to reach the control plane's console relay. Wrong here and every console session fails at connect with guacd unable to resolve the relay |
| `PROXY_LISTEN_HOST` | no | `0.0.0.0` | Interface the ephemeral relay listeners bind. Narrow it only if guacd shares a known interface |
| `POD_IP` | no | | Kubernetes only, set by the Helm chart from the downward API. Overrides `PROXY_ADVERTISE_HOST`, because a console relay belongs to one replica and a Service address names all of them |
| `APP_ROOT_LABEL` | no | `vantage` | The app root label for the host and session organisation guard |
| `VANTAGE_WORKFLOW_LOG_DIR` | no | | Where workflow run logs are written |
:::danger `KEY_ENCRYPTION_KEY` has no recovery path
It encrypts SSH private keys, vault secrets, OIDC client secrets and console
+87 -15
View File
@@ -4,10 +4,12 @@ import (
"context"
"log"
"os"
"strings"
"time"
"gitea.hostxtra.co.uk/mrhid6/vantage/server/internal/api"
"gitea.hostxtra.co.uk/mrhid6/vantage/server/internal/auth"
"gitea.hostxtra.co.uk/mrhid6/vantage/server/internal/bus"
"gitea.hostxtra.co.uk/mrhid6/vantage/server/internal/db"
grpcserver "gitea.hostxtra.co.uk/mrhid6/vantage/server/internal/grpc"
"gitea.hostxtra.co.uk/mrhid6/vantage/server/internal/monitorsched"
@@ -18,7 +20,21 @@ import (
func main() {
mongoURI := getEnv("MONGO_URI", "mongodb://localhost:27017")
if os.Getenv("GRPC_HOST") == "" {
// Two flags, so the schema work can be lifted out of the serving pods.
//
// Under Docker Compose neither is set and nothing changes: one process
// migrates and then serves. Under Kubernetes with more than one replica
// that is unsafe — every pod would run MigrateOrgToInstance at once, and
// renaming collections while a sibling reads them is not a race anyone
// wins. The chart therefore runs a pre-upgrade Job with MIGRATE_ONLY and
// starts the Deployment with SKIP_MIGRATIONS.
migrateOnly := boolEnv("VANTAGE_MIGRATE_ONLY")
skipMigrations := boolEnv("VANTAGE_SKIP_MIGRATIONS")
// Not required in migrate-only mode: that process never serves gRPC, and
// demanding it would put an agent-facing address in a Job that has no
// business knowing one.
if !migrateOnly && os.Getenv("GRPC_HOST") == "" {
log.Fatal("GRPC_HOST is required (host:port agents dial for gRPC)")
}
@@ -28,6 +44,24 @@ func main() {
}
log.Println("connected to MongoDB")
if migrateOnly || !skipMigrations {
runSchemaSetup()
} else {
log.Println("VANTAGE_SKIP_MIGRATIONS set: assuming migrations ran elsewhere")
}
if migrateOnly {
log.Println("VANTAGE_MIGRATE_ONLY set: schema setup complete, exiting")
return
}
serve()
}
// runSchemaSetup performs every write that must happen exactly once before the
// application serves: migrations, index builders and default step seeding. It
// is fatal on anything that would leave the schema half-moved.
func runSchemaSetup() {
// Migrations 0001 to 0003 still speak the pre-rename shape (orgs, org_id),
// so they must run before 0004 renames everything underneath them.
if err := services.RunMigrations(); err != nil {
@@ -87,24 +121,27 @@ func main() {
}
}
services.StartLogSweeper()
services.StartAuditSweeper()
}
func serve() {
redisAddr := getEnv("REDIS_ADDR", "localhost:6379")
if err := auth.InitRedis(redisAddr); err != nil {
redisUser := os.Getenv("REDIS_USERNAME")
redisPass := os.Getenv("REDIS_PASSWORD")
if err := auth.InitRedis(redisAddr, redisUser, redisPass); err != nil {
log.Fatalf("failed to connect to Redis: %v", err)
}
log.Println("connected to Redis")
go func() {
ticker := time.NewTicker(2 * time.Minute)
defer ticker.Stop()
for range ticker.C {
if err := services.MarkOfflineServers(); err != nil {
log.Printf("mark offline error: %v", err)
}
}
}()
// The bus carries agent commands and step results between replicas. It is
// not optional even on a single-replica deployment: dispatch takes the same
// path either way, so the code exercised in production is the code
// exercised everywhere.
if err := bus.Init(redisAddr, redisUser, redisPass); err != nil {
log.Fatalf("failed to connect the message bus: %v", err)
}
log.Printf("message bus ready as node %s", bus.NodeID())
ctx := context.Background()
go func() {
if err := grpcserver.StartGRPC(9090); err != nil {
@@ -112,9 +149,33 @@ func main() {
}
}()
monitorsched.Start(context.Background())
// Everything below runs on exactly one replica at a time.
//
// These are cluster-singleton jobs, not per-pod work: N replicas would mean
// every monitor check firing N times, every incident notification delivered
// to the customer N times, every retention sweep deleting concurrently, and
// N reapers racing to purge the same instance. They share one lock rather
// than holding four, because they are one role — housekeeping — and
// splitting them would only spread that role across pods for no benefit.
bus.RunAsLeader(ctx, "housekeeping", func(jobCtx context.Context) {
services.StartLogSweeper(jobCtx)
services.StartAuditSweeper(jobCtx)
services.StartReaper(jobCtx)
monitorsched.Start(jobCtx)
services.StartReaper(context.Background())
ticker := time.NewTicker(2 * time.Minute)
defer ticker.Stop()
for {
select {
case <-jobCtx.Done():
return
case <-ticker.C:
if err := services.MarkOfflineServers(); err != nil {
log.Printf("mark offline error: %v", err)
}
}
}
})
r := gin.New()
r.Use(gin.Recovery())
@@ -147,3 +208,14 @@ func getEnv(key, fallback string) string {
}
return fallback
}
// boolEnv reads a flag env var. Anything other than a recognised truthy value
// is false, so a typo leaves the safe default (migrate here, serve here) rather
// than silently skipping schema setup.
func boolEnv(key string) bool {
switch strings.ToLower(strings.TrimSpace(os.Getenv(key))) {
case "1", "true", "yes", "on":
return true
}
return false
}
+3
View File
@@ -20,6 +20,9 @@ func actorFromCtx(c *gin.Context) string {
}
func RegisterRoutes(r *gin.Engine) {
r.GET("/healthz", handleHealthz)
r.GET("/readyz", handleReadyz)
r.GET("/install", handleInstallScript)
r.GET("/install.ps1", handleInstallScriptWindows)
r.GET("/update", handleUpdateScript)
+61
View File
@@ -0,0 +1,61 @@
package api
import (
"context"
"net/http"
"time"
"gitea.hostxtra.co.uk/mrhid6/vantage/server/internal/auth"
"gitea.hostxtra.co.uk/mrhid6/vantage/server/internal/db"
"github.com/gin-gonic/gin"
)
// The two probes answer different questions on purpose.
//
// /healthz is liveness: the process is up and serving. It touches nothing
// external, because a Mongo outage must not make Kubernetes restart every
// server pod — a restart loop cannot fix someone else's database, and it
// destroys every open command stream and console session on the way.
//
// /readyz is readiness: this pod can serve a request end to end, which needs
// both Mongo and Redis. A failing readiness probe pulls the pod out of the
// Service and leaves it running, which is the behaviour that matters during a
// dependency blip.
//
// Both sit outside /api, so neither the session middleware nor the licence
// gate applies. Neither reveals anything beyond up or down.
const probeTimeout = 2 * time.Second
func handleHealthz(c *gin.Context) {
c.JSON(http.StatusOK, gin.H{"status": "ok"})
}
func handleReadyz(c *gin.Context) {
ctx, cancel := context.WithTimeout(c.Request.Context(), probeTimeout)
defer cancel()
checks := gin.H{"mongo": "ok", "redis": "ok"}
ready := true
if db.Client == nil {
checks["mongo"] = "not initialised"
ready = false
} else if err := db.Client.Ping(ctx, nil); err != nil {
checks["mongo"] = "unreachable"
ready = false
}
if err := auth.PingRedis(ctx); err != nil {
checks["redis"] = "unreachable"
ready = false
}
status := http.StatusOK
state := "ok"
if !ready {
status = http.StatusServiceUnavailable
state = "unready"
}
c.JSON(status, gin.H{"status": state, "checks": checks})
}
+41 -25
View File
@@ -5,7 +5,6 @@ import (
"fmt"
"io"
"net/http"
"os"
"regexp"
"strconv"
"strings"
@@ -50,23 +49,42 @@ func getServerRunLog(c *gin.Context) {
c.JSON(http.StatusBadRequest, gin.H{"error": "invalid id"})
return
}
path := services.ServerRunLogPath(runID, serverID)
b, err := os.ReadFile(path)
if err != nil {
if !services.HasServerRunLog(runID, serverID) {
c.JSON(http.StatusNotFound, gin.H{"error": "no logs"})
return
}
c.Data(http.StatusOK, "text/plain; charset=utf-8", b)
c.Header("Content-Type", "text/plain; charset=utf-8")
c.Status(http.StatusOK)
// Streamed in pages rather than read whole. A log capped at 200k lines is
// tens of megabytes, and holding that in memory per concurrent download is
// how one curious user takes a pod down.
var after int64
for {
lines, last, err := services.ReadServerRunLog(runID, serverID, after, logPageSize)
if err != nil || len(lines) == 0 {
return
}
for _, l := range lines {
_, _ = c.Writer.WriteString(l)
_, _ = c.Writer.WriteString("\n")
}
c.Writer.Flush()
after = last
}
}
// logPageSize bounds one read of the log store. Large enough that a normal log
// is one or two queries, small enough that no single response buffers much.
const logPageSize = 2000
func streamServerRunLog(c *gin.Context) {
runID, serverID := c.Param("runId"), c.Param("serverId")
if !uuidLike.MatchString(runID) || !uuidLike.MatchString(serverID) {
c.JSON(http.StatusBadRequest, gin.H{"error": "invalid id"})
return
}
path := services.ServerRunLogPath(runID, serverID)
c.Writer.Header().Set("Content-Type", "text/event-stream")
c.Writer.Header().Set("Cache-Control", "no-cache")
c.Writer.Header().Set("Connection", "keep-alive")
@@ -78,29 +96,27 @@ func streamServerRunLog(c *gin.Context) {
return
}
var offset int64
// The cursor is a sequence number now, not a byte offset: lines come from
// the log store, which any pod can read, rather than from a file only this
// one has.
var after int64
sendNew := func() {
f, err := os.Open(path)
if err != nil {
return
}
defer f.Close()
if _, err := f.Seek(offset, 0); err != nil {
return
}
buf := make([]byte, 32*1024)
for {
n, _ := f.Read(buf)
if n <= 0 {
break
lines, last, err := services.ReadServerRunLog(runID, serverID, after, logPageSize)
if err != nil || len(lines) == 0 {
return
}
offset += int64(n)
for _, line := range splitSSE(buf[:n]) {
_, _ = c.Writer.WriteString("data: " + line + "\n")
after = last
for _, line := range lines {
for _, part := range splitSSE([]byte(line)) {
_, _ = c.Writer.WriteString("data: " + part + "\n")
}
_, _ = c.Writer.WriteString("\n")
}
_, _ = c.Writer.WriteString("\n")
flusher.Flush()
if len(lines) < logPageSize {
return
}
}
}
+21 -2
View File
@@ -5,6 +5,7 @@ import (
"crypto/rand"
"encoding/hex"
"encoding/json"
"errors"
"time"
"github.com/redis/go-redis/v9"
@@ -25,13 +26,31 @@ type Session struct {
var rdb *redis.Client
func InitRedis(addr string) error {
rdb = redis.NewClient(&redis.Options{Addr: addr})
// InitRedis connects the session store.
//
// Username and password may both be empty for an unauthenticated instance. For
// a legacy `requirepass` Redis, pass the password with an empty username —
// go-redis then sends AUTH with one argument instead of two.
func InitRedis(addr, username, password string) error {
rdb = redis.NewClient(&redis.Options{
Addr: addr,
Username: username,
Password: password,
})
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
return rdb.Ping(ctx).Err()
}
// PingRedis reports whether the session store is reachable. Readiness depends
// on it: a pod that cannot reach Redis can authenticate nobody.
func PingRedis(ctx context.Context) error {
if rdb == nil {
return errors.New("redis not initialised")
}
return rdb.Ping(ctx).Err()
}
func randomHex(n int) (string, error) {
b := make([]byte, n)
if _, err := rand.Read(b); err != nil {
+220
View File
@@ -0,0 +1,220 @@
// Package bus is the control plane's inter-process message bus.
//
// It exists because an agent's CommandStream lands on exactly one server
// process. With a single replica that process is also the one handling every
// REST request, so an in-memory map was enough. With several replicas it is
// not: the pod asked to run a workflow step is almost never the pod holding
// that agent's stream, and a map cannot reach across the gap.
//
// Redis is already a hard dependency (sessions), so the bus adds no new
// infrastructure. It carries three things:
//
// presence which pod, if any, currently holds an agent's stream
// commands a request/ack exchange delivering a ServerCommand to that pod
// results step results travelling back to the pod driving the run
//
// Everything here is deliberately best-effort delivery with an explicit ack
// rather than a queue. A command whose owner pod died between the presence
// check and the publish must fail loudly and immediately — the caller answers
// 503 and the operator retries — not sit in a queue waiting for a stream that
// no longer exists.
package bus
import (
"context"
"crypto/rand"
"encoding/hex"
"encoding/json"
"errors"
"fmt"
"os"
"time"
"github.com/redis/go-redis/v9"
)
// ErrNoResponder means nothing was subscribed to the channel, or the subscriber
// did not ack within the timeout. Either way the command did not reach an agent.
var ErrNoResponder = errors.New("no responder on channel")
var (
rdb *redis.Client
nodeID string
)
// Init connects the bus. It takes its own client rather than sharing the
// session store's: a subscription occupies its connection for as long as it
// lives, and every agent stream on this pod holds one, so they must not come
// out of the pool that ordinary session reads depend on.
func Init(addr, username, password string) error {
rdb = redis.NewClient(&redis.Options{
Addr: addr,
Username: username,
Password: password,
})
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
if err := rdb.Ping(ctx).Err(); err != nil {
return err
}
// Hostname is the pod name under Kubernetes, which makes a log line naming
// a node directly actionable. The random suffix keeps two processes on one
// host (or a reused pod name) from claiming each other's locks.
host, _ := os.Hostname()
if host == "" {
host = "server"
}
b := make([]byte, 4)
_, _ = rand.Read(b)
nodeID = host + "-" + hex.EncodeToString(b)
return nil
}
// NodeID identifies this process on the bus. Stable for the process lifetime.
func NodeID() string { return nodeID }
// Client exposes the underlying Redis client for the leader election and
// presence helpers in this package. It is nil before Init.
func Client() *redis.Client { return rdb }
// Channel names. Every key and channel is prefixed so a Redis shared with the
// session store (which uses km:) stays legible.
const (
prefix = "vantage:"
// CommandChannel carries envelopes to whichever pod holds serverID's stream.
CommandChannel = prefix + "cmd:"
// ackChannel carries the owner pod's answer back to the requesting pod.
ackChannel = prefix + "ack:"
// ResultChannel carries a StepResult back to the pod driving the run.
ResultChannel = prefix + "res:"
// ProxyEndChannel carries a console relay's terminal reason back to the pod
// serving the WebSocket, which is the pod that has to write the audit event.
ProxyEndChannel = prefix + "proxyend:"
// PresenceKey records which node holds an agent's command stream.
PresenceKey = prefix + "agent:"
// leaderKey records the holder of a named singleton job.
leaderKey = prefix + "leader:"
)
// Publish sends v, JSON-encoded, to channel. It reports how many subscribers
// received it, which is the only signal Redis pub/sub gives that anyone was
// listening.
func Publish(ctx context.Context, channel string, v any) (int64, error) {
b, err := json.Marshal(v)
if err != nil {
return 0, err
}
return rdb.Publish(ctx, channel, b).Result()
}
// Subscribe returns a channel of raw payloads and a function that unsubscribes.
// It blocks until Redis confirms the subscription, so a caller may publish
// immediately afterwards without racing its own subscriber.
func Subscribe(ctx context.Context, channel string) (<-chan []byte, func(), error) {
ps := rdb.Subscribe(ctx, channel)
if _, err := ps.Receive(ctx); err != nil {
_ = ps.Close()
return nil, nil, err
}
out := make(chan []byte, 64)
go func() {
defer close(out)
for msg := range ps.Channel() {
select {
case out <- []byte(msg.Payload):
default:
// A subscriber too slow to keep up would otherwise stall every
// other subscriber sharing this connection. Dropping is correct
// here: commands are acked, and a dropped ack fails the caller
// loudly rather than hanging it.
}
}
}()
return out, func() { _ = ps.Close() }, nil
}
// Request publishes req to channel and waits for a single reply on replyTo.
//
// The reply subscription is established before the publish, so a responder that
// answers instantly cannot beat the subscriber into place. A publish that
// reaches no subscriber fails immediately with ErrNoResponder rather than
// burning the whole timeout: nobody is going to answer.
func Request(ctx context.Context, channel, replyTo string, req any, timeout time.Duration) ([]byte, error) {
waitCtx, cancel := context.WithTimeout(ctx, timeout)
defer cancel()
replies, unsub, err := Subscribe(waitCtx, replyTo)
if err != nil {
return nil, fmt.Errorf("subscribe reply channel: %w", err)
}
defer unsub()
n, err := Publish(waitCtx, channel, req)
if err != nil {
return nil, err
}
if n == 0 {
return nil, ErrNoResponder
}
select {
case b, ok := <-replies:
if !ok {
return nil, ErrNoResponder
}
return b, nil
case <-waitCtx.Done():
return nil, ErrNoResponder
}
}
// Reply answers a Request on its reply channel.
func Reply(ctx context.Context, replyTo string, v any) error {
_, err := Publish(ctx, replyTo, v)
return err
}
// AckChannelFor names the reply channel for a command. The command ID is
// already unique per dispatch, so it needs no further qualification.
func AckChannelFor(commandID string) string { return ackChannel + commandID }
// SetPresence claims serverID for this node for ttl. Called repeatedly by the
// pod holding the stream; the TTL is what bounds how long a crashed pod keeps
// claiming an agent it can no longer reach.
func SetPresence(ctx context.Context, serverID string, ttl time.Duration) error {
return rdb.Set(ctx, PresenceKey+serverID, nodeID, ttl).Err()
}
// ClearPresence releases serverID, but only if this node still holds it. A
// blind DEL would let a pod whose stream had already been re-established
// elsewhere delete the new owner's claim on its way out.
func ClearPresence(ctx context.Context, serverID string) error {
return releaseIfOwner.Run(ctx, rdb, []string{PresenceKey + serverID}, nodeID).Err()
}
// PresenceHolder returns the node holding serverID's stream, or "" if none does.
func PresenceHolder(ctx context.Context, serverID string) string {
v, err := rdb.Get(ctx, PresenceKey+serverID).Result()
if err != nil {
return ""
}
return v
}
// IsConnected reports whether any pod currently holds serverID's stream.
func IsConnected(ctx context.Context, serverID string) bool {
n, err := rdb.Exists(ctx, PresenceKey+serverID).Result()
return err == nil && n > 0
}
var releaseIfOwner = redis.NewScript(`
if redis.call("GET", KEYS[1]) == ARGV[1] then
return redis.call("DEL", KEYS[1])
end
return 0
`)
+100
View File
@@ -0,0 +1,100 @@
package bus
import (
"context"
"log"
"time"
"github.com/redis/go-redis/v9"
)
const (
leaderTTL = 30 * time.Second
leaderRenew = 10 * time.Second
)
// RunAsLeader runs job on exactly one process at a time, cluster-wide.
//
// The background work this guards is not merely wasteful when duplicated. Every
// replica running the monitor scheduler means every check fires N times, every
// incident notification is delivered N times to the customer, and every hourly
// rollup is written N times. The Free-instance reaper is worse: it deletes
// whole instances, and two processes deleting the same one concurrently is not
// a race anyone wins.
//
// Redis rather than a Kubernetes Lease so that Docker Compose, which has no
// API server, takes the identical code path — one implementation to reason
// about, not two.
//
// job is given a context cancelled the moment leadership is lost, and must
// return when it is cancelled. Losing the lock (a paused process, a Redis
// blip) is treated as fatal to that run of the job: the successor may already
// have started, and two schedulers overlapping is the exact thing being
// prevented.
func RunAsLeader(ctx context.Context, name string, job func(context.Context)) {
go func() {
key := leaderKey + name
for {
if ctx.Err() != nil {
return
}
ok, err := rdb.SetNX(ctx, key, nodeID, leaderTTL).Result()
if err != nil {
log.Printf("leader %s: acquire failed: %v", name, err)
}
if !ok {
select {
case <-ctx.Done():
return
case <-time.After(leaderRenew):
}
continue
}
log.Printf("leader %s: acquired by %s", name, nodeID)
jobCtx, cancel := context.WithCancel(ctx)
go job(jobCtx)
holdLeadership(ctx, key, name)
cancel()
log.Printf("leader %s: released by %s", name, nodeID)
// Give up the key on a clean shutdown so a successor takes over in
// milliseconds rather than waiting out the TTL.
if ctx.Err() != nil {
relCtx, relCancel := context.WithTimeout(context.Background(), 2*time.Second)
_ = releaseIfOwner.Run(relCtx, rdb, []string{key}, nodeID).Err()
relCancel()
return
}
}
}()
}
// holdLeadership renews the lock until it is lost or ctx ends.
func holdLeadership(ctx context.Context, key, name string) {
t := time.NewTicker(leaderRenew)
defer t.Stop()
for {
select {
case <-ctx.Done():
return
case <-t.C:
n, err := renewIfOwner.Run(ctx, rdb, []string{key}, nodeID, int(leaderTTL/time.Millisecond)).Int()
if err != nil && err != redis.Nil {
log.Printf("leader %s: renew failed, standing down: %v", name, err)
return
}
if n == 0 {
log.Printf("leader %s: lock lost, standing down", name)
return
}
}
}
}
var renewIfOwner = redis.NewScript(`
if redis.call("GET", KEYS[1]) == ARGV[1] then
return redis.call("PEXPIRE", KEYS[1], ARGV[2])
end
return 0
`)
+7 -2
View File
@@ -178,8 +178,13 @@ func (s *vantageServer) CommandStream(stream pb.Vantage_CommandStreamServer) err
log.Printf("update last seen %s: %v", srv.ServerID, err)
}
ch := services.Dispatcher.Connect(srv.ServerID)
defer services.Dispatcher.Disconnect(srv.ServerID)
// Serve claims this agent's presence on the bus and subscribes this pod to
// its command channel, so a dispatch issued by any other replica arrives
// here. The teardown releases both — an agent that reconnects to a
// different pod must not leave this one advertising a stream it no longer
// has.
ch, release := services.Dispatcher.Serve(stream.Context(), srv.ServerID)
defer release()
log.Printf("agent %s connected command stream", srv.ServerID)
defer log.Printf("agent %s disconnected command stream", srv.ServerID)
+24
View File
@@ -37,10 +37,24 @@ type Session struct {
reason string
conn net.Conn
closing bool
onEnd func(string)
rendezvous *time.Timer
}
// OnEnd registers a callback invoked exactly once, when the session finishes
// tearing down, with the recorded reason (empty if it ended cleanly).
//
// It exists because the process that observes a relay failing is not
// necessarily the process that has to record it: with several replicas the
// listener lives on the pod holding the agent's stream, while the audit event
// belongs to the pod serving the browser's WebSocket.
func (s *Session) OnEnd(fn func(string)) {
s.mu.Lock()
s.onEnd = fn
s.mu.Unlock()
}
// NewSession binds an ephemeral port on listenHost. allowed is the set of IPs
// permitted to connect; an empty set allows any, which is the degraded case
// when guacd's host could not be resolved.
@@ -112,6 +126,16 @@ func (s *Session) Close(reason string) {
if conn != nil {
_ = conn.Close()
}
// Read after teardown, not before: relay()'s goroutines set the reason
// as they unwind, and reporting one recorded earlier than that would
// name the symptom rather than the cause.
s.mu.Lock()
fn, final := s.onEnd, s.reason
s.mu.Unlock()
if fn != nil {
fn(final)
}
})
}
+11 -3
View File
@@ -15,13 +15,21 @@ import (
// An immediate pass then daily, following StartLogSweeper's shape. Daily rather
// than hourly because the unit of retention is a day: sweeping twenty-four times
// to delete the same nothing is load without a purpose.
func StartAuditSweeper() {
// It takes a context because it runs under leader election: several replicas
// all trimming the same audit logs is duplicated deletion of customer data, and
// the loop must stop the moment this process stops being the leader.
func StartAuditSweeper(ctx context.Context) {
go func() {
sweepAuditLogs()
t := time.NewTicker(24 * time.Hour)
defer t.Stop()
for range t.C {
sweepAuditLogs()
for {
select {
case <-ctx.Done():
return
case <-t.C:
sweepAuditLogs()
}
}
}()
}
+169 -34
View File
@@ -1,12 +1,17 @@
package services
import (
"context"
"encoding/json"
"errors"
"fmt"
"log"
"net"
"os"
"sync"
"time"
"gitea.hostxtra.co.uk/mrhid6/vantage/server/internal/bus"
"gitea.hostxtra.co.uk/mrhid6/vantage/server/internal/grpc/pb"
"gitea.hostxtra.co.uk/mrhid6/vantage/server/internal/proxy"
)
@@ -16,23 +21,51 @@ import (
// so this is fatal rather than a degraded mode.
var ErrAgentOffline = errors.New("agent is not connected")
// ConsoleProxy is a pending relay: a bound listener guacd can dial and a
// dispatched command telling the agent to meet it.
// A console session spans two processes once there is more than one replica.
//
// The browser's WebSocket lands on an arbitrary pod. The agent's ProxyStream
// lands on the pod holding that agent's command stream. The relay listener has
// to be on the latter — that is the only process that can match an incoming
// ProxyStream to a waiting listener — while guacd is dialled from the former.
//
// So the WebSocket's pod asks, over the bus, for a relay to be bound on the
// agent's pod, and gets back an address to hand to guacd. That address is the
// owner pod's own, which is why it must resolve to a single pod (POD_IP under
// Kubernetes) rather than to the Service, which would send guacd to a pod
// holding no listener roughly (n-1)/n of the time.
//
// Teardown needs no message of its own. When the browser goes away guac closes
// its connection to the relay, the relay sees the read end, and the session
// closes itself — the same path a single-process deployment always took. Only
// the *reason* has to cross back, because the pod that writes the audit event
// is not the pod that observed the failure.
// How long Close waits for the relay's terminal reason to arrive before giving
// up. The audit event is written immediately afterwards; a reason that has not
// crossed the bus in this long is not going to improve the record by being
// waited for longer.
const proxyEndGrace = 2 * time.Second
// ConsoleProxy is a relay as seen by the pod serving the WebSocket.
type ConsoleProxy struct {
ProxyID string
Host string
Port int
session *proxy.Session
serverID string
mu sync.Mutex
reason string
closed bool
stop func()
ended chan struct{}
}
func (c *ConsoleProxy) Close() {
proxy.Default.Remove(c.ProxyID)
c.session.Close("")
// proxyEnd is the terminal event a relay's owner pod publishes.
type proxyEnd struct {
Reason string `json:"reason"`
}
func (c *ConsoleProxy) Reason() string { return c.session.Reason() }
func proxyListenHost() string {
if v := os.Getenv("PROXY_LISTEN_HOST"); v != "" {
return v
@@ -40,7 +73,14 @@ func proxyListenHost() string {
return "0.0.0.0"
}
// proxyAdvertiseHost is the address guacd will dial to reach a relay bound by
// *this* process. POD_IP wins over the configured value because with several
// replicas the configured value names the Service, and a Service cannot address
// one pod. The chart sets POD_IP from the downward API.
func proxyAdvertiseHost() string {
if v := os.Getenv("POD_IP"); v != "" {
return v
}
if v := os.Getenv("PROXY_ADVERTISE_HOST"); v != "" {
return v
}
@@ -74,17 +114,55 @@ func guacdHosts(addr string) []string {
return ips
}
// DispatchOpenProxy tells the agent to dial its own loopback on port and relay
// it back under proxyID.
func DispatchOpenProxy(serverID, proxyID string, port uint32) error {
return Dispatcher.dispatch(serverID, &pb.ServerCommand{
CommandId: proxyID,
OpenProxy: &pb.OpenProxyCmd{ProxyId: proxyID, Port: port},
})
// localRelay is a listener bound by this process on behalf of a remote request.
type localRelay struct {
proxyID string
host string
port int
session *proxy.Session
}
// OpenConsoleProxy binds a relay listener, registers it, and asks the agent to
// connect. The caller must Close the result.
// openLocalRelay binds a listener here and registers it, so the agent's
// ProxyStream — which will arrive at this process — can be matched to it.
// Called on the owner pod, from the dispatch handler.
func openLocalRelay(instanceID, serverID, proxyID string) (*localRelay, error) {
sess, err := proxy.NewSession(proxyListenHost(), guacdHosts(guacdAddr()))
if err != nil {
return nil, err
}
sess.OnEnd(func(reason string) {
ctx, cancel := context.WithTimeout(context.Background(), dispatchAckTimeout)
defer cancel()
if _, err := bus.Publish(ctx, bus.ProxyEndChannel+proxyID, proxyEnd{Reason: reason}); err != nil {
log.Printf("proxy: publish end for %s: %v", proxyID, err)
}
})
proxy.Default.Add(&proxy.Entry{
ProxyID: proxyID,
InstanceID: instanceID,
ServerID: serverID,
Session: sess,
})
return &localRelay{
proxyID: proxyID,
host: proxyAdvertiseHost(),
port: sess.Port(),
session: sess,
}, nil
}
// abandon tears down a relay that was bound but whose command never reached the
// agent, so the listener does not sit out its rendezvous timeout for nothing.
func (r *localRelay) abandon() {
proxy.Default.Remove(r.proxyID)
r.session.Close("dispatch_failed")
}
// OpenConsoleProxy asks the pod holding serverID's stream to bind a relay and
// tell the agent to meet it. The caller must Close the result.
func OpenConsoleProxy(instanceID, serverID string, targetPort int) (*ConsoleProxy, error) {
if !Dispatcher.IsConnected(serverID) {
return nil, ErrAgentOffline
@@ -95,28 +173,85 @@ func OpenConsoleProxy(instanceID, serverID string, targetPort int) (*ConsoleProx
return nil, fmt.Errorf("generate proxy id: %w", err)
}
sess, err := proxy.NewSession(proxyListenHost(), guacdHosts(guacdAddr()))
// Subscribed before the relay is asked for: a relay that fails immediately
// (the agent never claims it, the dial is refused) publishes its reason at
// once, and that reason is the whole content of the audit event.
cp := &ConsoleProxy{ProxyID: proxyID, serverID: serverID, ended: make(chan struct{})}
endCtx, endCancel := context.WithCancel(context.Background())
ends, unsub, err := bus.Subscribe(endCtx, bus.ProxyEndChannel+proxyID)
if err != nil {
return nil, err
endCancel()
return nil, fmt.Errorf("subscribe relay end: %w", err)
}
cp.stop = func() {
endCancel()
unsub()
}
go cp.watchEnd(ends)
proxy.Default.Add(&proxy.Entry{
ProxyID: proxyID,
InstanceID: instanceID,
ServerID: serverID,
Session: sess,
ack, err := Dispatcher.send(CommandEnvelope{
ServerID: serverID,
Command: &pb.ServerCommand{
CommandId: proxyID,
OpenProxy: &pb.OpenProxyCmd{ProxyId: proxyID, Port: uint32(targetPort)},
},
Proxy: &ProxyRelayRequest{InstanceID: instanceID, ProxyID: proxyID},
})
if err := DispatchOpenProxy(serverID, proxyID, uint32(targetPort)); err != nil {
proxy.Default.Remove(proxyID)
sess.Close("dispatch_failed")
if err != nil {
cp.stop()
if errors.Is(err, ErrAgentNotConnected) {
return nil, ErrAgentOffline
}
return nil, err
}
if ack.ProxyHost == "" || ack.ProxyPort == 0 {
cp.stop()
return nil, fmt.Errorf("relay opened without an address")
}
return &ConsoleProxy{
ProxyID: proxyID,
Host: proxyAdvertiseHost(),
Port: sess.Port(),
session: sess,
}, nil
cp.Host = ack.ProxyHost
cp.Port = ack.ProxyPort
return cp, nil
}
func (c *ConsoleProxy) watchEnd(ends <-chan []byte) {
defer close(c.ended)
b, ok := <-ends
if !ok {
return
}
var e proxyEnd
if err := json.Unmarshal(b, &e); err != nil {
return
}
c.mu.Lock()
if c.reason == "" {
c.reason = e.Reason
}
c.mu.Unlock()
}
// Close waits briefly for the relay's terminal reason, then releases the
// subscription. It is safe to call more than once.
func (c *ConsoleProxy) Close() {
c.mu.Lock()
if c.closed {
c.mu.Unlock()
return
}
c.closed = true
c.mu.Unlock()
select {
case <-c.ended:
case <-time.After(proxyEndGrace):
}
c.stop()
}
// Reason reports why the relay ended, empty if it ended cleanly.
func (c *ConsoleProxy) Reason() string {
c.mu.Lock()
defer c.mu.Unlock()
return c.reason
}
+224 -34
View File
@@ -1,63 +1,253 @@
package services
import (
"context"
"encoding/json"
"errors"
"fmt"
"log"
"net/http"
"strings"
"sync"
"time"
"gitea.hostxtra.co.uk/mrhid6/vantage/server/internal/bus"
"gitea.hostxtra.co.uk/mrhid6/vantage/server/internal/grpc/pb"
"github.com/google/uuid"
)
type commandDispatcher struct {
mu sync.RWMutex
channels map[string]chan *pb.ServerCommand
// Commands travel over the bus even when the sender is also the owner.
//
// An agent's CommandStream terminates on exactly one process, and with several
// replicas that is almost never the process handling the REST request that
// wants to talk to it. Publishing unconditionally — rather than checking for a
// local stream first and falling back — means one code path, exercised on every
// deployment including the single-replica ones, instead of a rare cross-pod
// path that only fails in production.
const (
// How long a caller waits for the owning pod to acknowledge. Generous
// enough to cross a loaded cluster, short enough that a REST handler
// answering 503 does not look hung.
dispatchAckTimeout = 5 * time.Second
// Presence must outlive a renew or two, or a momentarily slow pod would
// look offline and its agent would be declared unreachable.
presenceTTL = 30 * time.Second
presenceRenew = 10 * time.Second
)
// CommandEnvelope is what actually crosses the bus. It is the command plus the
// small amount of context the owning pod needs to act on it locally.
type CommandEnvelope struct {
ServerID string `json:"server_id"`
Command *pb.ServerCommand `json:"command"`
ReplyTo string `json:"reply_to"`
Log *LogRequest `json:"log,omitempty"`
Proxy *ProxyRelayRequest `json:"proxy,omitempty"`
}
var Dispatcher = &commandDispatcher{
channels: make(map[string]chan *pb.ServerCommand),
// LogRequest asks the owner pod to open a step log before it dispatches.
//
// Step output arrives on the owner pod's gRPC stream, so that is where it is
// masked and written. Shipping raw output back to the run's pod first would put
// unmasked bytes on the bus for no gain.
type LogRequest struct {
RunID string `json:"run_id"`
ServerID string `json:"server_id"`
Mask []string `json:"mask,omitempty"`
}
func (d *commandDispatcher) Connect(serverID string) chan *pb.ServerCommand {
ch := make(chan *pb.ServerCommand, 16)
d.mu.Lock()
d.channels[serverID] = ch
d.mu.Unlock()
return ch
// ProxyRelayRequest asks the owner pod to bind a console relay listener and
// register it before dispatching OpenProxyCmd.
//
// The listener has to live on the owner pod: the agent's ProxyStream arrives
// there, and only there can it be matched to a waiting listener. The pod
// serving the browser's WebSocket learns the address from the ack and hands
// that to guacd.
type ProxyRelayRequest struct {
InstanceID string `json:"instance_id"`
ProxyID string `json:"proxy_id"`
}
func (d *commandDispatcher) Disconnect(serverID string) {
d.mu.Lock()
delete(d.channels, serverID)
d.mu.Unlock()
// CommandAck is the owner pod's answer.
type CommandAck struct {
OK bool `json:"ok"`
Error string `json:"error,omitempty"`
Node string `json:"node,omitempty"`
ProxyHost string `json:"proxy_host,omitempty"`
ProxyPort int `json:"proxy_port,omitempty"`
}
type commandDispatcher struct{}
var Dispatcher = &commandDispatcher{}
// ErrAgentNotConnected is returned when no pod holds the agent's stream.
var ErrAgentNotConnected = errors.New("agent is not connected")
// Serve subscribes this pod to serverID's command channel and claims presence
// for it, returning the channel the gRPC handler should send from and a
// teardown function.
//
// send is the local delivery: it is called on the bus goroutine and must not
// block for long, which is why it pushes onto a buffered channel rather than
// writing to the gRPC stream directly.
func (d *commandDispatcher) Serve(ctx context.Context, serverID string) (<-chan *pb.ServerCommand, func()) {
out := make(chan *pb.ServerCommand, 16)
envelopes, unsub, err := bus.Subscribe(ctx, bus.CommandChannel+serverID)
if err != nil {
// Without a subscription this pod cannot receive commands for the
// agent. Claiming presence anyway would advertise a stream nobody can
// reach, which is worse than the agent appearing offline.
log.Printf("dispatch: subscribe for %s failed, commands will not reach it: %v", serverID, err)
close(out)
return out, func() {}
}
runCtx, cancel := context.WithCancel(ctx)
if err := bus.SetPresence(runCtx, serverID, presenceTTL); err != nil {
log.Printf("dispatch: claim presence for %s: %v", serverID, err)
}
go func() {
t := time.NewTicker(presenceRenew)
defer t.Stop()
for {
select {
case <-runCtx.Done():
return
case <-t.C:
if err := bus.SetPresence(runCtx, serverID, presenceTTL); err != nil {
log.Printf("dispatch: renew presence for %s: %v", serverID, err)
}
}
}
}()
go func() {
for {
select {
case <-runCtx.Done():
return
case raw, ok := <-envelopes:
if !ok {
return
}
d.handleEnvelope(runCtx, raw, out)
}
}
}()
return out, func() {
cancel()
unsub()
relCtx, relCancel := context.WithTimeout(context.Background(), 2*time.Second)
if err := bus.ClearPresence(relCtx, serverID); err != nil {
log.Printf("dispatch: release presence for %s: %v", serverID, err)
}
relCancel()
}
}
// handleEnvelope performs the owner-pod side of a dispatch: any local setup the
// command needs, then queueing it for the stream, then the ack.
func (d *commandDispatcher) handleEnvelope(ctx context.Context, raw []byte, out chan *pb.ServerCommand) {
var env CommandEnvelope
if err := json.Unmarshal(raw, &env); err != nil {
log.Printf("dispatch: undecodable envelope: %v", err)
return
}
if env.Command == nil || env.ReplyTo == "" {
return
}
ack := CommandAck{OK: true, Node: bus.NodeID()}
if env.Log != nil {
if err := StepLogs.Open(env.Command.CommandId, env.Log.RunID, env.Log.ServerID, env.Log.Mask); err != nil {
log.Printf("dispatch: open step log for %s: %v", env.Command.CommandId, err)
}
}
var relay *localRelay
if env.Proxy != nil {
r, err := openLocalRelay(env.Proxy.InstanceID, env.ServerID, env.Proxy.ProxyID)
if err != nil {
ack = CommandAck{OK: false, Error: err.Error(), Node: bus.NodeID()}
} else {
relay = r
ack.ProxyHost = r.host
ack.ProxyPort = r.port
}
}
if ack.OK {
select {
case out <- env.Command:
default:
ack = CommandAck{OK: false, Error: "command queue full", Node: bus.NodeID()}
if relay != nil {
relay.abandon()
}
}
}
if err := bus.Reply(ctx, env.ReplyTo, ack); err != nil {
log.Printf("dispatch: reply on %s: %v", env.ReplyTo, err)
}
}
// IsConnected reports whether any pod holds the agent's command stream.
func (d *commandDispatcher) IsConnected(serverID string) bool {
d.mu.RLock()
_, ok := d.channels[serverID]
d.mu.RUnlock()
return ok
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
defer cancel()
return bus.IsConnected(ctx, serverID)
}
func (d *commandDispatcher) dispatch(serverID string, cmd *pb.ServerCommand) error {
d.mu.RLock()
ch, ok := d.channels[serverID]
d.mu.RUnlock()
if !ok {
return fmt.Errorf("agent for server %s is not connected", serverID)
}
select {
case ch <- cmd:
return nil
default:
return fmt.Errorf("command queue full for server %s", serverID)
}
_, err := d.send(CommandEnvelope{ServerID: serverID, Command: cmd})
return err
}
func DispatchRunStep(serverID, commandID string, cmd *pb.RunStepCmd) error {
return Dispatcher.dispatch(serverID, &pb.ServerCommand{CommandId: commandID, RunStep: cmd})
// send publishes the envelope and waits for the owning pod's ack. A missing
// responder and a refusing responder are both errors: in neither case did the
// command reach the agent.
func (d *commandDispatcher) send(env CommandEnvelope) (CommandAck, error) {
if env.Command == nil || env.Command.CommandId == "" {
return CommandAck{}, fmt.Errorf("command id is required")
}
env.ReplyTo = bus.AckChannelFor(env.Command.CommandId)
ctx, cancel := context.WithTimeout(context.Background(), dispatchAckTimeout)
defer cancel()
raw, err := bus.Request(ctx, bus.CommandChannel+env.ServerID, env.ReplyTo, env, dispatchAckTimeout)
if err != nil {
if errors.Is(err, bus.ErrNoResponder) {
return CommandAck{}, fmt.Errorf("%w: %s", ErrAgentNotConnected, env.ServerID)
}
return CommandAck{}, err
}
var ack CommandAck
if err := json.Unmarshal(raw, &ack); err != nil {
return CommandAck{}, fmt.Errorf("undecodable ack: %w", err)
}
if !ack.OK {
return ack, fmt.Errorf("agent for server %s: %s", env.ServerID, ack.Error)
}
return ack, nil
}
// DispatchRunStep sends a step to the agent and asks the owning pod to capture
// its output. mask holds the secret values that must not reach the log.
func DispatchRunStep(serverID, commandID, runID string, mask []string, cmd *pb.RunStepCmd) error {
_, err := Dispatcher.send(CommandEnvelope{
ServerID: serverID,
Command: &pb.ServerCommand{CommandId: commandID, RunStep: cmd},
Log: &LogRequest{RunID: runID, ServerID: serverID, Mask: mask},
})
return err
}
func DispatchCleanupWorkspace(serverID, workspaceID string) {
+285 -88
View File
@@ -2,9 +2,8 @@ package services
import (
"bytes"
"context"
"log"
"os"
"path/filepath"
"strings"
"sync"
"time"
@@ -12,47 +11,130 @@ import (
"gitea.hostxtra.co.uk/mrhid6/vantage/server/internal/db"
"go.mongodb.org/mongo-driver/v2/bson"
"go.mongodb.org/mongo-driver/v2/mongo"
"go.mongodb.org/mongo-driver/v2/mongo/options"
)
func WorkflowLogDir() string {
dir := os.Getenv("VANTAGE_WORKFLOW_LOG_DIR")
if dir == "" {
dir = filepath.Join("data", "workflow-logs")
}
_ = os.MkdirAll(dir, 0700)
return dir
}
// Workflow run logs live in MongoDB, not on disk.
//
// They used to be a file per (run, server) under a volume. That is the cheapest
// possible writer, and it is wrong the moment there is more than one server
// process: step output arrives on whichever pod holds the agent's stream, the
// run's markers are written by whichever pod started the run, and the browser
// asks for the log through whichever pod the load balancer picked. Three pods,
// one file, one local disk — two of them see an empty log.
//
// Mongo makes every pod an equal reader and writer, which is the property that
// matters. It costs writes on the hot path, so the writer batches (see
// stepLogWriter) and both caps below exist to keep a runaway step from turning
// a workflow into a database incident.
const (
// A single line longer than this is truncated. Base64 blobs and minified
// output are the usual cause; nobody reads column 9000 of a log line.
maxLogLineBytes = 8 * 1024
func ServerRunLogPath(runID, serverID string) string {
return filepath.Join(WorkflowLogDir(), runID, serverID+".log")
// A single (run, server) log stops accepting lines here, with one final
// marker saying so. 200k lines is far past what anyone reads and still
// only a few tens of MB. Without a cap, `yes` in a step fills the database.
maxLogLines = 200_000
// The writer flushes on whichever comes first. Batching is what keeps a
// chatty step to a handful of writes a second instead of one per line.
logFlushLines = 128
logFlushInterval = 250 * time.Millisecond
)
const logLinesCol = "workflow_log_lines"
const logSeqCol = "workflow_log_seq"
func logCtx() (context.Context, context.CancelFunc) {
return context.WithTimeout(context.Background(), 10*time.Second)
}
func logTS() string {
return time.Now().UTC().Format("2006-01-02T15:04:05.000") + "Z"
}
func AppendMarker(runID, serverID, text string) (int64, error) {
path := ServerRunLogPath(runID, serverID)
if err := os.MkdirAll(filepath.Dir(path), 0700); err != nil {
return 0, err
func seqID(runID, serverID string) string { return runID + "/" + serverID }
// reserveSeq atomically claims n consecutive sequence numbers for a
// (run, server) log and returns the first.
//
// The counter is a document rather than a per-process integer because two pods
// write the same log concurrently: the run's pod emits markers while the
// agent's pod emits step output. Ordering between them is only meaningful if
// they draw from the same counter.
func reserveSeq(ctx context.Context, runID, serverID string, n int) (int64, error) {
var doc struct {
Seq int64 `bson:"seq"`
}
f, err := os.OpenFile(path, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0600)
err := db.Col(logSeqCol).FindOneAndUpdate(ctx,
bson.M{"_id": seqID(runID, serverID)},
bson.M{"$inc": bson.M{"seq": int64(n)}},
options.FindOneAndUpdate().SetUpsert(true).SetReturnDocument(options.After),
).Decode(&doc)
if err != nil {
return 0, err
}
defer f.Close()
off, _ := f.Seek(0, 2)
if _, err := f.WriteString("[" + logTS() + "] " + text + "\n"); err != nil {
return off, err
}
return off, nil
return doc.Seq - int64(n), nil
}
type logLine struct {
RunID string `bson:"run_id"`
ServerID string `bson:"server_id"`
Seq int64 `bson:"seq"`
At time.Time `bson:"at"`
Line string `bson:"line"`
}
// writeLines reserves a block of sequence numbers and inserts the batch.
func writeLines(ctx context.Context, runID, serverID string, lines []string) (int64, error) {
if len(lines) == 0 {
return 0, nil
}
start, err := reserveSeq(ctx, runID, serverID, len(lines))
if err != nil {
return 0, err
}
now := time.Now().UTC()
docs := make([]any, 0, len(lines))
for i, l := range lines {
docs = append(docs, logLine{
RunID: runID,
ServerID: serverID,
Seq: start + int64(i),
At: now,
Line: l,
})
}
// Unordered: one rejected document must not discard the rest of the batch.
_, err = db.Col(logLinesCol).InsertMany(ctx, docs, options.InsertMany().SetOrdered(false))
return start, err
}
// AppendMarker writes one control line (step banners, retries, run outcome) and
// returns its sequence number, which is what a StepRun's log_offset records so
// the UI can scroll to where a step began.
func AppendMarker(runID, serverID, text string) (int64, error) {
ctx, cancel := logCtx()
defer cancel()
return writeLines(ctx, runID, serverID, []string{"[" + logTS() + "] " + text})
}
// stepLogWriter accumulates one command's output, splits it into lines, masks
// secrets and flushes batches to Mongo.
type stepLogWriter struct {
mu sync.Mutex
f *os.File
mu sync.Mutex
runID string
serverID string
secrets []string
carry []byte
secrets []string
pending []string
written int
capped bool
stop chan struct{}
done chan struct{}
}
type stepLogRegistry struct {
@@ -60,18 +142,30 @@ type stepLogRegistry struct {
writers map[string]*stepLogWriter
}
// The registry stays process-local, and correctly so: a step's output arrives
// on the pod holding that agent's stream, and that is the same pod the
// dispatch envelope asked to open the writer. Nothing here crosses pods —
// only the lines it produces do, by virtue of landing in Mongo.
var StepLogs = &stepLogRegistry{writers: make(map[string]*stepLogWriter)}
func (r *stepLogRegistry) Open(commandID, path string, secrets []string) error {
if err := os.MkdirAll(filepath.Dir(path), 0700); err != nil {
return err
func (r *stepLogRegistry) Open(commandID, runID, serverID string, secrets []string) error {
w := &stepLogWriter{
runID: runID,
serverID: serverID,
secrets: secrets,
stop: make(chan struct{}),
done: make(chan struct{}),
}
f, err := os.OpenFile(path, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0600)
if err != nil {
return err
}
w := &stepLogWriter{f: f, secrets: secrets}
go w.flushLoop()
r.mu.Lock()
if old := r.writers[commandID]; old != nil {
// A retry reuses the command ID. Close the previous attempt's writer
// rather than leaking its flush goroutine.
r.mu.Unlock()
old.close()
r.mu.Lock()
}
r.writers[commandID] = w
r.mu.Unlock()
return nil
@@ -89,24 +183,22 @@ func (r *stepLogRegistry) Append(commandID string, data []byte) {
return
}
w.mu.Lock()
defer w.mu.Unlock()
buf := append(w.carry, data...)
for {
i := bytes.IndexByte(buf, '\n')
if i < 0 {
break
}
w.writeLine(buf[:i])
w.queueLine(buf[:i])
buf = buf[i+1:]
}
w.carry = append([]byte{}, buf...)
}
full := len(w.pending) >= logFlushLines
w.mu.Unlock()
func (w *stepLogWriter) writeLine(line []byte) {
masked := maskBytes(line, w.secrets)
_, _ = w.f.WriteString("[" + logTS() + "] ")
_, _ = w.f.Write(masked)
_, _ = w.f.WriteString("\n")
if full {
w.flush()
}
}
func (r *stepLogRegistry) Close(commandID string) {
@@ -117,13 +209,72 @@ func (r *stepLogRegistry) Close(commandID string) {
if w == nil {
return
}
w.close()
}
// queueLine masks, truncates and enqueues a line. Caller holds w.mu.
func (w *stepLogWriter) queueLine(line []byte) {
if w.capped {
return
}
if w.written+len(w.pending) >= maxLogLines {
w.capped = true
w.pending = append(w.pending, "["+logTS()+"] [vantage] log truncated: this step exceeded the per-server line limit")
return
}
masked := maskBytes(line, w.secrets)
if len(masked) > maxLogLineBytes {
masked = append(masked[:maxLogLineBytes], []byte(" [truncated]")...)
}
w.pending = append(w.pending, "["+logTS()+"] "+string(masked))
}
func (w *stepLogWriter) flush() {
w.mu.Lock()
if len(w.pending) == 0 {
w.mu.Unlock()
return
}
batch := w.pending
w.pending = nil
w.written += len(batch)
runID, serverID := w.runID, w.serverID
w.mu.Unlock()
ctx, cancel := logCtx()
defer cancel()
if _, err := writeLines(ctx, runID, serverID, batch); err != nil {
// Dropped rather than retried. A log line is not worth stalling the
// step it describes, and a Mongo that cannot take writes has larger
// problems than a missing line.
log.Printf("step log: write %d line(s) for run %s: %v", len(batch), runID, err)
}
}
func (w *stepLogWriter) flushLoop() {
defer close(w.done)
t := time.NewTicker(logFlushInterval)
defer t.Stop()
for {
select {
case <-w.stop:
return
case <-t.C:
w.flush()
}
}
}
func (w *stepLogWriter) close() {
close(w.stop)
<-w.done
w.mu.Lock()
defer w.mu.Unlock()
if len(w.carry) > 0 {
w.writeLine(w.carry)
w.queueLine(w.carry)
w.carry = nil
}
_ = w.f.Close()
w.mu.Unlock()
w.flush()
}
func maskBytes(b []byte, secrets []string) []byte {
@@ -137,86 +288,132 @@ func maskBytes(b []byte, secrets []string) []byte {
return []byte(s)
}
func StartLogSweeper() {
// ReadServerRunLog returns the log for one server in a run, from sequence
// number after onwards, along with the highest sequence returned.
//
// Callers page with it rather than fetching everything: the live stream asks
// repeatedly for what is new, and the plain-text endpoint walks the whole log
// in chunks so a very large one is never held in memory whole.
func ReadServerRunLog(runID, serverID string, after int64, limit int) ([]string, int64, error) {
ctx, cancel := logCtx()
defer cancel()
cur, err := db.Col(logLinesCol).Find(ctx,
bson.M{"run_id": runID, "server_id": serverID, "seq": bson.M{"$gt": after}},
options.Find().SetSort(bson.D{{Key: "seq", Value: 1}}).SetLimit(int64(limit)),
)
if err != nil {
return nil, after, err
}
defer cur.Close(ctx)
lines := make([]string, 0, 64)
last := after
for cur.Next(ctx) {
var l logLine
if err := cur.Decode(&l); err != nil {
return lines, last, err
}
lines = append(lines, l.Line)
last = l.Seq
}
return lines, last, cur.Err()
}
// HasServerRunLog reports whether any line exists, so a log read can answer 404
// rather than 200 with an empty body.
func HasServerRunLog(runID, serverID string) bool {
ctx, cancel := logCtx()
defer cancel()
err := db.Col(logLinesCol).FindOne(ctx, bson.M{"run_id": runID, "server_id": serverID}).Err()
return err == nil
}
func StartLogSweeper(ctx context.Context) {
go func() {
sweepLogs()
t := time.NewTicker(time.Hour)
defer t.Stop()
for range t.C {
sweepLogs()
for {
select {
case <-ctx.Done():
return
case <-t.C:
sweepLogs()
}
}
}()
}
// sweepLogs deletes the lines of finished runs past their instance's retention.
//
// It walks workflow_runs rather than the log collection: retention is a
// property of the run (its instance, its finish time), and a run row is the
// only place both are recorded.
func sweepLogs() {
base := WorkflowLogDir()
entries, err := os.ReadDir(base)
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Minute)
defer cancel()
cur, err := db.Col("workflow_runs").Find(ctx,
bson.M{"finished_at": bson.M{"$ne": nil}},
options.Find().SetProjection(bson.M{"run_id": 1, "instance_id": 1, "finished_at": 1}),
)
if err != nil {
log.Printf("log sweep: list runs: %v", err)
return
}
defer cur.Close(ctx)
cache := map[string]int{}
now := time.Now()
for _, e := range entries {
if !e.IsDir() {
continue
for cur.Next(ctx) {
var run struct {
RunID string `bson:"run_id"`
InstanceID string `bson:"instance_id"`
FinishedAt *time.Time `bson:"finished_at"`
}
runID := e.Name()
dir := filepath.Join(base, runID)
instanceID, finishedAt, found, err := runRetentionInfo(runID)
if err != nil {
log.Printf("log sweep: retention lookup failed for run %s: %v", runID, err)
continue
}
if found && finishedAt == nil {
if err := cur.Decode(&run); err != nil || run.FinishedAt == nil {
continue
}
days, ok := cache[instanceID]
days, ok := cache[run.InstanceID]
if !ok {
days = defaultRetentionDays
if instanceID != "" {
if v, err := GetWorkflowLogRetentionDays(instanceID); err == nil {
if run.InstanceID != "" {
if v, err := GetWorkflowLogRetentionDays(run.InstanceID); err == nil {
days = v
}
}
cache[instanceID] = days
cache[run.InstanceID] = days
}
if days <= 0 {
continue
}
cutoff := now.AddDate(0, 0, -days)
if found {
if finishedAt.Before(cutoff) {
_ = os.RemoveAll(dir)
}
if !run.FinishedAt.Before(now.AddDate(0, 0, -days)) {
continue
}
if fi, e := os.Stat(dir); e == nil && fi.ModTime().Before(cutoff) {
_ = os.RemoveAll(dir)
if _, err := db.Col(logLinesCol).DeleteMany(ctx, bson.M{"run_id": run.RunID}); err != nil {
log.Printf("log sweep: delete lines for run %s: %v", run.RunID, err)
continue
}
_, _ = db.Col(logSeqCol).DeleteMany(ctx, bson.M{"_id": bson.M{"$regex": "^" + run.RunID + "/"}})
}
}
const defaultRetentionDays = 30
func runRetentionInfo(runID string) (string, *time.Time, bool, error) {
ctx, cancel := wfCtx()
// EnsureLogIndexes builds the indexes the log store depends on. The compound
// index is not an optimisation: every read is a range scan over it, and without
// it a log read collection-scans every line in the database.
func EnsureLogIndexes() error {
ctx, cancel := logCtx()
defer cancel()
var run struct {
InstanceID string `bson:"instance_id"`
FinishedAt *time.Time `bson:"finished_at"`
if _, err := db.Col(logLinesCol).Indexes().CreateOne(ctx, mongo.IndexModel{
Keys: bson.D{{Key: "run_id", Value: 1}, {Key: "server_id", Value: 1}, {Key: "seq", Value: 1}},
}); err != nil {
return err
}
err := db.Col("workflow_runs").FindOne(ctx, bson.M{"run_id": runID}).Decode(&run)
if err == mongo.ErrNoDocuments {
return "", nil, false, nil
}
if err != nil {
return "", nil, false, err
}
return run.InstanceID, run.FinishedAt, true, nil
return nil
}
+61 -29
View File
@@ -1,43 +1,75 @@
package services
import (
"sync"
"context"
"encoding/json"
"log"
"gitea.hostxtra.co.uk/mrhid6/vantage/server/internal/bus"
"gitea.hostxtra.co.uk/mrhid6/vantage/server/internal/grpc/pb"
)
type stepResultRegistry struct {
mu sync.Mutex
pending map[string]chan *pb.StepResult
}
var StepResults = &stepResultRegistry{pending: make(map[string]chan *pb.StepResult)}
func (r *stepResultRegistry) Await(commandID string) <-chan *pb.StepResult {
ch := make(chan *pb.StepResult, 1)
r.mu.Lock()
r.pending[commandID] = ch
r.mu.Unlock()
return ch
}
func (r *stepResultRegistry) Cancel(commandID string) {
r.mu.Lock()
delete(r.pending, commandID)
r.mu.Unlock()
// Step results travel back over the bus for the same reason commands travel out
// over it: the pod driving a workflow run and the pod holding the agent's
// stream are two different processes, and a map in one of them cannot be read
// by the other.
//
// Await subscribes before the command is dispatched (see dispatchAndWait),
// which is what stops a fast agent from answering into a channel nobody is
// listening on yet.
type stepResultRegistry struct{}
var StepResults = &stepResultRegistry{}
// Await subscribes to a command's result channel. The returned cancel function
// must be called once the caller is done, whether a result arrived or not —
// it is what releases the Redis subscription.
func (r *stepResultRegistry) Await(commandID string) (<-chan *pb.StepResult, func()) {
out := make(chan *pb.StepResult, 1)
ctx, cancel := context.WithCancel(context.Background())
raw, unsub, err := bus.Subscribe(ctx, bus.ResultChannel+commandID)
if err != nil {
log.Printf("step results: subscribe for %s: %v", commandID, err)
cancel()
close(out)
return out, func() {}
}
go func() {
defer close(out)
select {
case <-ctx.Done():
return
case b, ok := <-raw:
if !ok {
return
}
var res pb.StepResult
if err := json.Unmarshal(b, &res); err != nil {
log.Printf("step results: undecodable result for %s: %v", commandID, err)
return
}
out <- &res
}
}()
return out, func() {
cancel()
unsub()
}
}
// Deliver publishes a result received from an agent. Called on the pod holding
// that agent's stream, which is not usually the pod waiting for it.
func (r *stepResultRegistry) Deliver(res *pb.StepResult) {
if res == nil {
if res == nil || res.CommandId == "" {
return
}
r.mu.Lock()
ch, ok := r.pending[res.CommandId]
if ok {
delete(r.pending, res.CommandId)
}
r.mu.Unlock()
if ok {
ch <- res
ctx, cancel := context.WithTimeout(context.Background(), dispatchAckTimeout)
defer cancel()
if _, err := bus.Publish(ctx, bus.ResultChannel+res.CommandId, res); err != nil {
log.Printf("step results: publish for %s: %v", res.CommandId, err)
}
}
+17 -10
View File
@@ -240,7 +240,6 @@ func runServer(instanceID, runID string, srvIdx int, steps []models.ResolvedStep
marker := fmt.Sprintf("===== step %d/%d: %s (%s) =====", step.Order+1, len(steps), step.Name, step.Interpreter)
offset, _ := AppendMarker(runID, serverID, marker)
logPath := ServerRunLogPath(runID, serverID)
secretsSlice := secretValues(secretVals)
commandID := uuid.New().String()
@@ -250,15 +249,17 @@ func runServer(instanceID, runID string, srvIdx int, steps []models.ResolvedStep
_, _ = AppendMarker(runID, serverID, fmt.Sprintf("retry %d/%d after failure", attempts-1, maxAttempts-1))
}
_ = StepLogs.Open(commandID, logPath, secretsSlice)
res = dispatchAndWait(serverID, commandID, &pb.RunStepCmd{
// The log writer is opened by the pod that owns this agent's
// stream, not here: that is where the output arrives, and the mask
// list travels with the dispatch so it is applied before anything
// is stored.
res = dispatchAndWait(serverID, commandID, runID, secretsSlice, &pb.RunStepCmd{
Interpreter: step.Interpreter,
Script: step.Script,
Env: cmdEnv,
TimeoutSeconds: 0,
WorkspaceId: runID,
})
StepLogs.Close(commandID)
if res != nil && res.ExitCode == 0 {
break
}
@@ -322,10 +323,14 @@ func runServer(instanceID, runID string, srvIdx int, steps []models.ResolvedStep
})
}
func dispatchAndWait(serverID, commandID string, cmd *pb.RunStepCmd) *pb.StepResult {
ch := StepResults.Await(commandID)
if err := DispatchRunStep(serverID, commandID, cmd); err != nil {
StepResults.Cancel(commandID)
// dispatchAndWait subscribes to the result before dispatching, because the two
// happen on different pods and an agent that answers quickly would otherwise
// publish into a channel this process had not yet joined.
func dispatchAndWait(serverID, commandID, runID string, mask []string, cmd *pb.RunStepCmd) *pb.StepResult {
ch, done := StepResults.Await(commandID)
defer done()
if err := DispatchRunStep(serverID, commandID, runID, mask, cmd); err != nil {
return &pb.StepResult{ExitCode: 1, Stderr: "[vantage] dispatch failed: " + err.Error()}
}
wait := time.Duration(cmd.TimeoutSeconds)*time.Second + stepDispatchGrace
@@ -333,10 +338,12 @@ func dispatchAndWait(serverID, commandID string, cmd *pb.RunStepCmd) *pb.StepRes
wait = 30*time.Minute + stepDispatchGrace
}
select {
case res := <-ch:
case res, ok := <-ch:
if !ok || res == nil {
return &pb.StepResult{ExitCode: 1, Stderr: "[vantage] lost the result channel before the agent answered"}
}
return res
case <-time.After(wait):
StepResults.Cancel(commandID)
return &pb.StepResult{ExitCode: 124, Stderr: "[vantage] timed out waiting for agent result"}
}
}
+3
View File
@@ -60,6 +60,9 @@ func EnsureWorkflowIndexes() error {
}); err != nil {
return err
}
if err := EnsureLogIndexes(); err != nil {
return err
}
_, err := db.Col("workflow_runs").Indexes().CreateOne(ctx, mongo.IndexModel{
Keys: bson.D{{Key: "run_id", Value: 1}}, Options: options.Index().SetUnique(true),
})
+15
View File
@@ -0,0 +1,15 @@
import { NextResponse } from "next/server";
// Liveness and readiness for the Next server itself. Deliberately at /healthz
// and not /api/healthz: next.config.ts rewrites the whole of /api to the Go
// server, so a probe there would report the backend's health instead of this
// process's — and would keep passing while this pod was wedged.
//
// It answers without touching the backend on purpose. web is stateless; a
// backend outage must not take every web replica out of its Service as well,
// which would turn one failure into two.
export const dynamic = "force-dynamic";
export function GET() {
return NextResponse.json({ status: "ok" });
}