Files
felhom.eu/hub/internal/api/handler.go
T

2073 lines
84 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package api
import (
"bytes"
"context"
"crypto/subtle"
"database/sql"
"encoding/base64"
"encoding/json"
"fmt"
"io"
"log"
"net/http"
"strings"
"time"
"gitea.dooplex.hu/admin/felhom-hub/internal/assets"
"gitea.dooplex.hu/admin/felhom-hub/internal/claim"
"gitea.dooplex.hu/admin/felhom-hub/internal/configgen"
"gitea.dooplex.hu/admin/felhom-hub/internal/mailrelay"
"gitea.dooplex.hu/admin/felhom-hub/internal/notify"
"gitea.dooplex.hu/admin/felhom-hub/internal/store"
)
// ConfigTemplateProvider returns the controller.yaml template for config generation.
type ConfigTemplateProvider interface {
Template() string
}
// LatestVersionProvider returns the latest controller image version known to the hub's registry
// checker (e.g. "0.86.0"), or "" if unknown. Satisfied by *web.VersionChecker; nil when the
// registry checker is disabled. Used to advertise latest on the controller report ACK.
type LatestVersionProvider interface {
LatestVersion() string
}
// Handler handles API endpoints for report ingest and customer queries.
type Handler struct {
store *store.Store
apiKey string
resendAPIKey string
fromEmail string
logger *log.Logger
httpClient *http.Client
templateProvider ConfigTemplateProvider
dispatcher *notify.Dispatcher
assetsMgr *assets.Manager
latestVersion LatestVersionProvider
// App-email passthrough (POST /api/v1/mail). nil sender = endpoint returns 503.
mailSender mailrelay.Sender
mailLimiter *mailRateLimiter
mailFromAllow map[string]bool
// S1 offsite connectivity: the wgsync reconciler seam (internal/api/wg.go). nil = peer-sync
// disabled — mutations still persist, responses carry sync:"disabled".
wgSyncer WGSyncer
// claimEngine is the customer-claim code engine (v0.50.0). nil = claim arc disabled: no codes
// issued, no claim field in ACKs/configs — pre-arc behavior exactly.
claimEngine *claim.Engine
// wgRegisteredHook (v0.51.0, DR-tier-by-default scenario A) fires after a host's FIRST WG
// peer registration — main.go wires it to the web server's PBSDRAutoProvision so a DR-ON
// customer's pbs_dr descriptor lands hands-free (WG registers → provision → the agent's next
// desired-state tick). nil = no cascade hook (pre-v0.51.0 behavior). Runs in a detached
// goroutine; must never delay or fail the registration response.
wgRegisteredHook func(ctx context.Context, customerID string)
}
// SetClaimEngine wires the customer-claim code engine (nil-safe everywhere it is used).
func (h *Handler) SetClaimEngine(e *claim.Engine) {
h.claimEngine = e
}
// SetWGRegisteredHook wires the post-WG-registration cascade hook (v0.51.0; nil-safe).
func (h *Handler) SetWGRegisteredHook(f func(ctx context.Context, customerID string)) {
h.wgRegisteredHook = f
}
// SetLatestVersionProvider wires the registry version checker so the controller report ACK can
// advertise the latest available version (Phase 2). nil-safe (no latest_version field emitted).
func (h *Handler) SetLatestVersionProvider(p LatestVersionProvider) {
h.latestVersion = p
}
// New creates a new API handler.
func New(store *store.Store, apiKey, resendAPIKey, fromEmail string, templateProvider ConfigTemplateProvider, logger *log.Logger) *Handler {
return &Handler{
store: store,
apiKey: apiKey,
resendAPIKey: resendAPIKey,
fromEmail: fromEmail,
logger: logger,
httpClient: &http.Client{Timeout: 10 * time.Second},
templateProvider: templateProvider,
}
}
// SetDispatcher sets the notification dispatcher for event-triggered emails.
func (h *Handler) SetDispatcher(d *notify.Dispatcher) {
h.dispatcher = d
}
// SetAssetManager sets the asset manager for serving app assets to controllers.
func (h *Handler) SetAssetManager(am *assets.Manager) {
h.assetsMgr = am
}
// checkAuth verifies the Bearer token against the global API key or a per-customer API key.
// Returns true if authorized.
func (h *Handler) checkAuth(r *http.Request) bool {
_, _, ok := h.checkAuthCustomer(r)
return ok
}
// checkAuthCustomer verifies the Bearer token and returns the authenticated customer identity.
// For per-customer keys: returns (customerID, false, true).
// For global key: returns ("", true, true) — caller must allow any customer_id.
// On failure: returns ("", false, false).
func (h *Handler) checkAuthCustomer(r *http.Request) (customerID string, isGlobal bool, ok bool) {
auth := r.Header.Get("Authorization")
if !strings.HasPrefix(auth, "Bearer ") {
return "", false, false
}
token := strings.TrimPrefix(auth, "Bearer ")
// Check global key first
if h.apiKey != "" && subtle.ConstantTimeCompare([]byte(token), []byte(h.apiKey)) == 1 {
return "", true, true
}
// Check per-customer key
cfg, err := h.store.GetCustomerConfigByAPIKey(token)
if err != nil || cfg == nil {
return "", false, false
}
return cfg.CustomerID, false, true
}
// checkAuthHost resolves a Bearer token to a HOST identity (the agent's auth
// path). It is a sibling of checkAuthCustomer — the controller path is unchanged.
// - global key -> ("", "", true, true) caller trusts body.host_id
// - per-host key -> (hostID, customerID, false, true)
// - failure -> ("", "", false, false)
func (h *Handler) checkAuthHost(r *http.Request) (hostID, customerID string, isGlobal, ok bool) {
auth := r.Header.Get("Authorization")
if !strings.HasPrefix(auth, "Bearer ") {
return "", "", false, false
}
token := strings.TrimPrefix(auth, "Bearer ")
// Global key first (same constant-time compare as checkAuthCustomer).
if h.apiKey != "" && subtle.ConstantTimeCompare([]byte(token), []byte(h.apiKey)) == 1 {
return "", "", true, true
}
host, err := h.store.GetHostByAPIKey(token)
if err != nil || host == nil {
return "", "", false, false
}
return host.HostID, host.CustomerID, false, true
}
// ServeHTTP routes API requests.
func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
path := strings.TrimPrefix(r.URL.Path, "/api/v1")
switch {
case r.Method == http.MethodPost && path == "/report":
h.handleReport(w, r)
case r.Method == http.MethodPost && path == "/host-report":
h.handleHostReport(w, r)
case r.Method == http.MethodPost && path == "/host-enroll":
h.handleHostEnroll(w, r)
case r.Method == http.MethodPost && path == "/admin/hosts":
h.handleAdminCreateHost(w, r)
case r.Method == http.MethodPut && strings.HasPrefix(path, "/hosts/") && strings.HasSuffix(path, "/escrow"):
hostID := strings.TrimSuffix(strings.TrimPrefix(path, "/hosts/"), "/escrow")
h.handleHostEscrowPut(w, r, hostID)
// G1 break-glass: day-0 vaults the root@pam console credential (self-scoped host key); the
// operator retrieves it via the /admin/ path (global key only).
case r.Method == http.MethodPut && strings.HasPrefix(path, "/hosts/") && strings.HasSuffix(path, "/recovery-credential"):
hostID := strings.TrimSuffix(strings.TrimPrefix(path, "/hosts/"), "/recovery-credential")
h.handleHostRecoveryCredentialPut(w, r, hostID)
case r.Method == http.MethodGet && strings.HasPrefix(path, "/admin/hosts/") && strings.HasSuffix(path, "/recovery-credential"):
hostID := strings.TrimSuffix(strings.TrimPrefix(path, "/admin/hosts/"), "/recovery-credential")
h.handleAdminGetRecoveryCredential(w, r, hostID)
// DR capstone (slice 10D). Recovery-mode toggle (global key); re-enroll + restore-directive
// (gated on recovery mode — no old key needed, the box is lost).
case r.Method == http.MethodPut && strings.HasPrefix(path, "/admin/hosts/") && strings.HasSuffix(path, "/recovery-mode"):
hostID := strings.TrimSuffix(strings.TrimPrefix(path, "/admin/hosts/"), "/recovery-mode")
h.handleSetRecoveryMode(w, r, hostID)
case r.Method == http.MethodDelete && strings.HasPrefix(path, "/admin/hosts/") && strings.HasSuffix(path, "/recovery-mode"):
hostID := strings.TrimSuffix(strings.TrimPrefix(path, "/admin/hosts/"), "/recovery-mode")
h.handleClearRecoveryMode(w, r, hostID)
case r.Method == http.MethodPost && strings.HasPrefix(path, "/hosts/") && strings.HasSuffix(path, "/re-enroll"):
hostID := strings.TrimSuffix(strings.TrimPrefix(path, "/hosts/"), "/re-enroll")
h.handleReEnroll(w, r, hostID)
case r.Method == http.MethodGet && strings.HasPrefix(path, "/hosts/") && strings.HasSuffix(path, "/restore-directive"):
hostID := strings.TrimSuffix(strings.TrimPrefix(path, "/hosts/"), "/restore-directive")
h.handleGetRestoreDirective(w, r, hostID)
// S2 offsite connectivity: box-facing WG pubkey registration (per-host key, self-scoped).
case r.Method == http.MethodPost && strings.HasPrefix(path, "/hosts/") && strings.HasSuffix(path, "/wg"):
hostID := strings.TrimSuffix(strings.TrimPrefix(path, "/hosts/"), "/wg")
h.handleRegisterHostWG(w, r, hostID)
// PBS DR tier (SLICE 1): the agent's consume-once fetch of its PBS token secret (api/pbsdr.go).
case r.Method == http.MethodPost && strings.HasPrefix(path, "/hosts/") && strings.HasSuffix(path, "/pbs/consume-token"):
hostID := strings.TrimSuffix(strings.TrimPrefix(path, "/hosts/"), "/pbs/consume-token")
h.handleConsumePBSToken(w, r, hostID)
// Desired-state serving (slice 10A) — per-host-key, self-scoped (a host reads only its own).
case r.Method == http.MethodGet && strings.HasPrefix(path, "/hosts/") && strings.HasSuffix(path, "/desired-state"):
hostID := strings.TrimSuffix(strings.TrimPrefix(path, "/hosts/"), "/desired-state")
h.handleGetDesiredState(w, r, hostID)
case r.Method == http.MethodGet && strings.HasPrefix(path, "/hosts/") && strings.HasSuffix(path, "/jobs"):
hostID := strings.TrimSuffix(strings.TrimPrefix(path, "/hosts/"), "/jobs")
h.handleGetJobs(w, r, hostID)
// Job completion (slice 10B) — per-host-key, self-scoped: DELETE /hosts/{id}/jobs/{job_id}.
case r.Method == http.MethodDelete && strings.HasPrefix(path, "/hosts/") && strings.Contains(path, "/jobs/"):
rest := strings.TrimPrefix(path, "/hosts/")
if i := strings.Index(rest, "/jobs/"); i > 0 {
h.handleDeleteJob(w, r, rest[:i], rest[i+len("/jobs/"):])
} else {
http.NotFound(w, r)
}
// Admin-set (slice 10A) — global/operator key only; bumps the generation.
case r.Method == http.MethodPut && strings.HasPrefix(path, "/admin/hosts/") && strings.HasSuffix(path, "/desired-state"):
hostID := strings.TrimSuffix(strings.TrimPrefix(path, "/admin/hosts/"), "/desired-state")
h.handleAdminSetDesiredState(w, r, hostID)
case r.Method == http.MethodPost && strings.HasPrefix(path, "/admin/hosts/") && strings.HasSuffix(path, "/jobs"):
hostID := strings.TrimSuffix(strings.TrimPrefix(path, "/admin/hosts/"), "/jobs")
h.handleAdminEnqueueJob(w, r, hostID)
// S1 offsite connectivity — WG endpoint record + peer registry (global key only, api/wg.go).
// DELETE carries the pubkey in the body: base64 '/'+'+' keep pubkeys out of URL paths.
case r.Method == http.MethodPut && path == "/admin/wg/endpoint":
h.handleAdminSetWGEndpoint(w, r)
case r.Method == http.MethodGet && path == "/admin/wg/endpoint":
h.handleAdminGetWGEndpoint(w, r)
case r.Method == http.MethodPost && path == "/admin/wg/peers":
h.handleAdminAddWGPeer(w, r)
case r.Method == http.MethodDelete && path == "/admin/wg/peers":
h.handleAdminDeleteWGPeer(w, r)
case r.Method == http.MethodGet && path == "/admin/wg/peers":
h.handleAdminListWGPeers(w, r)
// H1: the fleet operator OOB peer (register/rotate at an explicit /32) + read-back.
case r.Method == http.MethodPut && path == "/admin/wg/operator-peer":
h.handleAdminSetOperatorPeer(w, r)
case r.Method == http.MethodGet && path == "/admin/wg/operator-peer":
h.handleAdminGetOperatorPeer(w, r)
case r.Method == http.MethodPost && path == "/claim/reset-request":
h.handleClaimResetRequest(w, r)
case r.Method == http.MethodPost && path == "/event":
h.handleEvent(w, r)
case r.Method == http.MethodPost && path == "/mail":
h.handleMail(w, r)
case r.Method == http.MethodPost && path == "/notify":
h.handleNotify(w, r)
case r.Method == http.MethodPost && path == "/preferences":
h.handleSavePreferences(w, r)
case r.Method == http.MethodGet && path == "/customers":
h.handleCustomers(w, r)
case r.Method == http.MethodGet && strings.HasPrefix(path, "/customers/"):
parts := strings.Split(strings.TrimPrefix(path, "/customers/"), "/")
customerID := parts[0]
if len(parts) > 1 && parts[1] == "history" {
h.handleCustomerHistory(w, r, customerID)
} else {
h.handleCustomer(w, r, customerID)
}
case r.Method == http.MethodGet && strings.HasPrefix(path, "/recovery/"):
customerID := strings.TrimPrefix(path, "/recovery/")
h.handleRecovery(w, r, customerID)
case r.Method == http.MethodGet && strings.HasPrefix(path, "/config/"):
customerID := strings.TrimPrefix(path, "/config/")
h.handleConfigRetrieve(w, r, customerID)
case r.Method == http.MethodPost && strings.HasPrefix(path, "/offsite/consume-password/"):
customerID := strings.TrimPrefix(path, "/offsite/consume-password/")
h.handleOffsiteConsumePassword(w, r, customerID)
case r.Method == http.MethodGet && strings.HasPrefix(path, "/artifacts/"):
customerID := strings.TrimPrefix(path, "/artifacts/")
h.handleArtifactManifest(w, r, customerID)
case r.Method == http.MethodGet && path == "/assets/manifest":
h.handleAssetsManifest(w, r)
case r.Method == http.MethodGet && strings.HasPrefix(path, "/assets/file/"):
filename := strings.TrimPrefix(path, "/assets/file/")
h.handleAssetFile(w, r, filename)
default:
http.NotFound(w, r)
}
}
func (h *Handler) handleReport(w http.ResponseWriter, r *http.Request) {
authCustomerID, isGlobal, ok := h.checkAuthCustomer(r)
if !ok {
http.Error(w, "Unauthorized", http.StatusUnauthorized)
return
}
body, err := io.ReadAll(io.LimitReader(r.Body, 1<<20)) // 1MB limit
if err != nil {
http.Error(w, "Bad request", http.StatusBadRequest)
return
}
// Extract customer_id from JSON
var payload struct {
CustomerID string `json:"customer_id"`
}
if err := json.Unmarshal(body, &payload); err != nil || payload.CustomerID == "" {
http.Error(w, "Invalid payload: customer_id required", http.StatusBadRequest)
return
}
// Validate customer_id matches authenticated customer (unless global key)
if !isGlobal && authCustomerID != payload.CustomerID {
http.Error(w, "Forbidden: customer_id mismatch", http.StatusForbidden)
return
}
if err := h.store.SaveReport(payload.CustomerID, body); err != nil {
h.logger.Printf("[ERROR] Failed to save report from %s: %v", payload.CustomerID, err)
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
// Parse and save app telemetry (backward-compatible — old controllers won't have this field)
var telemetryPayload struct {
AppTelemetry []store.AppTelemetryRecord `json:"app_telemetry"`
}
if err := json.Unmarshal(body, &telemetryPayload); err == nil && len(telemetryPayload.AppTelemetry) > 0 {
if err := h.store.SaveAppTelemetry(payload.CustomerID, time.Now(), telemetryPayload.AppTelemetry); err != nil {
h.logger.Printf("[WARN] Failed to save app telemetry for %s: %v", payload.CustomerID, err)
}
}
// On-demand log tails (v0.43.0) — the controller ships these on the cycle after the ACK
// requested them. Storing one clears its pending request (consume-once, in SaveAppLogTail)
// so the next ACK stops advertising it. Backward-compatible: old controllers never send this.
var tailPayload struct {
LogTails []struct {
App string `json:"app"`
CollectedAt time.Time `json:"collected_at"`
Lines []string `json:"lines"`
} `json:"log_tails"`
}
if err := json.Unmarshal(body, &tailPayload); err == nil {
for _, lt := range tailPayload.LogTails {
if lt.App == "" {
continue
}
if err := h.store.SaveAppLogTail(payload.CustomerID, lt.App, lt.CollectedAt, lt.Lines); err != nil {
h.logger.Printf("[WARN] Failed to save log tail %s/%s: %v", payload.CustomerID, lt.App, err)
} else {
h.logger.Printf("[INFO] Log tail received for %s/%s (%d lines)", payload.CustomerID, lt.App, len(lt.Lines))
}
}
}
// Controller self-log tail (v0.46.0) — the controller's OWN debug ring, shipped on the
// cycle after the ACK's controller_log_requested. SaveLogBundle runs the secret gate
// (a hit stores a BLOCKED flag row, nothing else) and clears the pending request
// (consume-once). Backward-compatible: old controllers never send this.
var selfTailPayload struct {
ControllerLogTail *struct {
CollectedAt time.Time `json:"collected_at"`
Lines []string `json:"lines"`
} `json:"controller_log_tail"`
}
if err := json.Unmarshal(body, &selfTailPayload); err == nil && selfTailPayload.ControllerLogTail != nil {
lt := selfTailPayload.ControllerLogTail
blocked, berr := h.store.SaveLogBundle(payload.CustomerID, store.LogBundleComponentController, lt.CollectedAt, lt.Lines)
switch {
case berr != nil:
h.logger.Printf("[WARN] Failed to save controller log bundle for %s: %v", payload.CustomerID, berr)
case blocked:
h.logger.Printf("[WARN] controller log bundle for %s BLOCKED: possible secret in log content — nothing stored", payload.CustomerID)
default:
h.logger.Printf("[INFO] controller log bundle received for %s (%d lines)", payload.CustomerID, len(lt.Lines))
}
}
// DR recipe — persist the controller's secret-free customer/apps half (preserving any host half).
// Backward-compatible (old controllers won't have this field); a failure must not drop the report.
var drPayload struct {
DRRecipe json.RawMessage `json:"dr_recipe"`
}
if err := json.Unmarshal(body, &drPayload); err == nil && len(drPayload.DRRecipe) > 0 {
var ver drRecipeVersionOnly
_ = json.Unmarshal(drPayload.DRRecipe, &ver)
if err := h.store.SaveDRRecipeAppHalf(payload.CustomerID, ver.RecipeVersion, drPayload.DRRecipe); err != nil {
h.logger.Printf("[WARN] Failed to save DR-recipe app-half for %s: %v", payload.CustomerID, err)
} else {
h.logger.Printf("[INFO] DR-recipe app-half stored for customer %s (v%d)", payload.CustomerID, ver.RecipeVersion)
}
}
h.logger.Printf("[INFO] Received report from %s (%d bytes)", payload.CustomerID, len(body))
// Build response with optional customer_blocked flag
resp := map[string]interface{}{"status": "ok"}
if custCfg, err := h.store.GetCustomerConfig(payload.CustomerID); err == nil && custCfg != nil {
if custCfg.Status == "blocked" {
resp["customer_blocked"] = true
}
// Config-refresh (v0.26.0): advertise the per-customer config_version. The controller compares
// it against its last-applied version and, on a change, re-pulls controller.yaml + self-restarts
// (pull-based config delivery — the hub never connects into the box). Only emitted for
// config-managed customers (a report-only box without a config row gets no field and is unaffected).
resp["config_version"] = custCfg.ConfigVersion
// Customer-claim arc (v0.50.0, F-4): ensure a claim code exists for every reporting managed
// customer (idempotent — the live-box entry point; Day-0 boxes get theirs at config retrieve),
// ingest the controller's reported claimed flag (set-only — a wiped settings.json can never
// un-claim), and serve the ACTIVE code hash + generation in the ACK. The hash is bcrypt (non-
// reversible) — safe to serve on the authenticated report channel; the plaintext code exists
// only in the customer's mailbox.
if h.claimEngine != nil {
cs, cerr := h.claimEngine.EnsureIssued(custCfg)
if cerr != nil {
h.logger.Printf("[WARN] claim issue for %s on report: %v", payload.CustomerID, cerr)
}
var claimedPayload struct {
Claimed *bool `json:"claimed"`
}
if err := json.Unmarshal(body, &claimedPayload); err == nil &&
claimedPayload.Claimed != nil && *claimedPayload.Claimed {
if err := h.claimEngine.MarkClaimed(custCfg); err != nil {
h.logger.Printf("[WARN] claim mark-claimed for %s: %v", payload.CustomerID, err)
} else if cs != nil && cs.ClaimedAt == nil {
cs, _ = h.store.GetClaim(payload.CustomerID) // refresh for the ACK below
}
}
if cs != nil {
resp["claim"] = map[string]interface{}{
"code_hash": cs.CodeHash,
"generation": cs.Generation,
"issued_at": cs.IssuedAt.UTC().Format(time.RFC3339),
}
}
}
}
// SLICE 3 — escrow status for the hub-verified auto-confirm: the controller flips its offbox
// EscrowState pending→escrowed ONLY when sha256(its local repo password) matches restic_pw_sha256
// (blob-presence alone must never confirm — a stale blob may not cover the current key). The hash is
// non-reversible (256-bit random secret) — safe to serve; omitted entirely when no escrow row exists.
if es, err := h.store.GetEscrowStatusForCustomer(payload.CustomerID); err == nil && es != nil {
resp["escrow"] = es
}
// v0.43.0 — pending log-tail requests (same additive ACK-flag pattern as escrow): the
// controller collects the named apps' tails and ships them on its NEXT report; the field
// is omitted when nothing is pending. The hub never connects into the box.
if apps, err := h.store.GetPendingLogTailRequests(payload.CustomerID); err == nil && len(apps) > 0 {
resp["log_tail_requests"] = apps
}
// v0.46.0 — pending CONTROLLER self-log pull (same additive ACK-flag pattern): the
// controller ships controller_log_tail on its NEXT report; omitted when nothing pending.
if pending, err := h.store.PendingLogBundleRequest(payload.CustomerID, store.LogBundleComponentController); err == nil && pending {
resp["controller_log_requested"] = true
}
// Phase 2 managed updates: advertise the effective controller-version FLOOR (per-customer override
// else global default) and the latest available version. The controller compares its current
// version against the floor and auto-updates when below it (latest stays the customer's opt-in
// "update to latest" button — NOT the auto-target). Both fields are omitted when empty, so an old
// controller that ignores them, or a hub with no floor configured, behaves exactly as before.
// Part D: the per-box MinAgent conditional floor — HOLD the controller-version floor for a box
// whose host agent is below the current golden's required MinAgent (never push a controller past
// the agent it depends on). A held box gets NO directive (behaves as if no floor) but is flagged
// on the dashboard, never silently stale.
if fd := h.store.ResolveManagedFloor(payload.CustomerID); fd.Floor != "" {
resp["min_controller_version"] = fd.Floor
} else if fd.Held {
h.logger.Printf("[INFO] managed floor HELD for %s: agent %q < MinAgent %s (controller floor withheld)",
payload.CustomerID, fd.AgentVersion, fd.MinAgent)
}
if h.latestVersion != nil {
if latest := h.latestVersion.LatestVersion(); latest != "" {
resp["latest_version"] = latest
}
}
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusOK)
json.NewEncoder(w).Encode(resp)
}
// defaultHostPollSeconds is the cadence the hub hands every agent this slice (no
// per-host override UI yet — that is a later slice).
const defaultHostPollSeconds = 900
// maxHostReportBytes bounds a host-report body. Larger than the controller path's
// 1 MiB because host reports carry the full guest list + (later) storage/backup
// arrays. We read one byte past it and reject explicitly (413) rather than letting
// LimitReader silently truncate — a truncated-but-valid JSON would otherwise be
// accepted as a partial report, dropping guests from the mirror.
const maxHostReportBytes = 4 << 20 // 4 MiB
// hostReportPayload is the subset of the agent host-report (slice-3 contract,
// §3 / agent spec §4) the hub needs for denorm + guest reality. The remaining fields
// (backups/restore_tests/pbs_snapshots/audit_tail) are ignored, so an empty or absent
// collection is accepted without error.
//
// storage_targets (slice 5) is now parsed: the agent populates it, and the hub accepts
// + persists it. Persistence is the full report_json row (which carries the targets
// verbatim) plus the denorm counts below — the RICH manifest schema (desired class/role/
// policy/creds) is hub-owned and lands in slice 10; this slice only mirrors what the agent
// observes.
type hostReportPayload struct {
HostID string `json:"host_id"`
AgentVersion string `json:"agent_version"`
Host struct {
CPUPercent float64 `json:"cpu_percent"`
MemoryPercent float64 `json:"memory_percent"`
DiskPercent float64 `json:"disk_percent"`
} `json:"host"`
Guests []struct {
VMID int `json:"vmid"`
Name string `json:"name"`
Status string `json:"status"`
ControllerVersion string `json:"controller_version"`
} `json:"guests"`
StorageTargets []hostStorageTarget `json:"storage_targets"`
Backups []hostBackup `json:"backups"` // slice 6
RestoreTests []hostRestoreTest `json:"restore_tests"` // slice 6
PBSSnapshots []hostPBSSnapshot `json:"pbs_snapshots"` // slice 6 Phase B
Cloudflared struct {
Status string `json:"status"`
} `json:"cloudflared"`
// DR recipe — the agent's storage/guest/PBS half (secret-free). RawMessage = stored verbatim,
// ignore-unknown (forward-compat). Persisted to dr_recipe, assembled with the controller half.
DRRecipe json.RawMessage `json:"dr_recipe"`
// LogTail (v0.46.0) — the agent's on-demand debug-ring tail, present only on the
// heartbeat right after the envelope's log_tail_requested (agent ≥ 0.83.0).
LogTail *struct {
CollectedAt time.Time `json:"collected_at"`
Lines []string `json:"lines"`
} `json:"log_tail"`
}
// drRecipeVersionOnly extracts just recipe_version from a half's JSON (ignore-unknown). 0 if absent.
type drRecipeVersionOnly struct {
RecipeVersion int `json:"recipe_version"`
}
// hostPBSSnapshot mirrors the agent's hub.PBSSnapshot wire contract (slice 6 Phase B). The
// hub persists it via report_json and surfaces a FAILED verify prominently (the loudest
// offsite-DR signal — same treatment as a failed restore-test).
type hostPBSSnapshot struct {
Namespace string `json:"namespace"`
BackupType string `json:"backup_type"`
BackupID string `json:"backup_id"`
BackupTime string `json:"backup_time"`
SizeBytes int64 `json:"size_bytes"`
Owner string `json:"owner"`
Protected bool `json:"protected"`
Encrypted bool `json:"encrypted"`
VerifyState string `json:"verify_state"`
VerifyUPID string `json:"verify_upid,omitempty"`
}
// hostBackup / hostRestoreTest mirror the agent's hub.Backup / hub.RestoreTest wire
// contract field-for-field (slice 6, doc 03 §8). DUPLICATED contract — the golden stays
// byte-identical with felhom-agent's copy and the key-set tests guard drift. The hub
// persists these via report_json (no new columns this slice) and surfaces a FAILED
// restore-test prominently (the loudest DR signal). The rich backup policy is slice 10.
type hostBackup struct {
TargetID string `json:"target_id"`
VMID int `json:"vmid"`
Archive string `json:"archive"`
Mode string `json:"mode"`
CrashConsistent bool `json:"crash_consistent"`
SizeBytes int64 `json:"size_bytes"`
Success bool `json:"success"`
Error string `json:"error,omitempty"`
StartedAt string `json:"started_at"`
DurationSeconds float64 `json:"duration_seconds"`
UncoveredVolumes []string `json:"uncovered_volumes"`
}
type hostRestoreTest struct {
SourceArchive string `json:"source_archive"`
SourceTier string `json:"source_tier"`
ScratchVMID int `json:"scratch_vmid"`
Pass bool `json:"pass"`
Verified string `json:"verified"`
Error string `json:"error,omitempty"`
TestedAt string `json:"tested_at"`
DurationSeconds float64 `json:"duration_seconds"`
// Warnings are the guest-start task's warning line(s) on a PASS (e.g. the systemd-nesting
// advisory). The verdict is liveness-only, so a passed restore-test can carry warnings.
Warnings []string `json:"warnings,omitempty"`
// WarningsRecognized is true iff every warning is the known-benign anchor. Absent ⇒ false,
// which is the SAFE default: the hub then treats it as an unrecognized warning (the louder
// path), so a missing flag can only over-notice, never hide a real warning.
WarningsRecognized bool `json:"warnings_recognized,omitempty"`
}
// hostStorageTarget mirrors the agent's hub.StorageTarget wire contract field-for-field.
// It is a DUPLICATED contract (no shared types module yet); testdata/host-report.golden.json
// must stay byte-identical with felhom-agent's copy and the key-set test guards drift.
// The hub does not act on these yet beyond persisting + counting them (slice 10 adds the
// authoritative manifest), but mirroring the full shape keeps the cross-repo contract honest.
type hostStorageTarget struct {
Name string `json:"name"`
Type string `json:"type"`
DurableID string `json:"durable_id"`
State string `json:"state"`
Reachable bool `json:"reachable"`
TotalBytes int64 `json:"total_bytes"`
UsedBytes int64 `json:"used_bytes"`
AvailBytes int64 `json:"avail_bytes"`
UsedFraction float64 `json:"used_fraction"`
Content string `json:"content"`
MountPath string `json:"mount_path"`
BackingDevice string `json:"backing_device"`
ClassHint string `json:"class_hint"`
Role string `json:"role"`
ThinPool *struct {
DataUsedFraction float64 `json:"data_used_fraction"`
MetadataUsedFraction *float64 `json:"metadata_used_fraction"`
} `json:"thin_pool,omitempty"`
Smart struct {
Health string `json:"health"`
TemperatureC *int `json:"temperature_c"`
PowerOnHours *int `json:"power_on_hours"`
ReallocatedSectors *int `json:"reallocated_sectors"`
PendingSectors *int `json:"pending_sectors"`
OfflineUncorrectable *int `json:"offline_uncorrectable"`
CriticalWarning *int `json:"critical_warning"`
MediaErrors *int `json:"media_errors"`
PercentageUsed *int `json:"percentage_used"`
} `json:"smart"`
}
// handleHostReport ingests the agent's host-report (the heartbeat) and returns the
// control envelope (agent spec §5).
func (h *Handler) handleHostReport(w http.ResponseWriter, r *http.Request) {
hostID, custID, isGlobal, ok := h.checkAuthHost(r)
if !ok {
http.Error(w, "Unauthorized", http.StatusUnauthorized)
return
}
body, err := io.ReadAll(io.LimitReader(r.Body, maxHostReportBytes+1))
if err != nil {
http.Error(w, "Bad request", http.StatusBadRequest)
return
}
if len(body) > maxHostReportBytes {
http.Error(w, "Payload too large", http.StatusRequestEntityTooLarge)
return
}
var rep hostReportPayload
if err := json.Unmarshal(body, &rep); err != nil || rep.HostID == "" {
http.Error(w, "Invalid payload: host_id required", http.StatusBadRequest)
return
}
if isGlobal {
// Global-key bootstrap: trust body.host_id but require the host to exist
// (it must be minted first) and resolve its customer from the row.
host, err := h.store.GetHost(rep.HostID)
if err != nil {
h.logger.Printf("[ERROR] host lookup failed for %s: %v", rep.HostID, err)
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
if host == nil {
http.Error(w, "Unknown host_id (mint via /admin/hosts first)", http.StatusBadRequest)
return
}
hostID, custID = rep.HostID, host.CustomerID
} else if rep.HostID != hostID {
http.Error(w, "Forbidden: host_id mismatch", http.StatusForbidden)
return
}
running := 0
for _, g := range rep.Guests {
if g.Status == "running" {
running++
}
}
denorm := store.HostReportDenorm{
AgentVersion: rep.AgentVersion,
CPUPercent: rep.Host.CPUPercent,
MemoryPercent: rep.Host.MemoryPercent,
DiskPercent: rep.Host.DiskPercent,
GuestTotal: len(rep.Guests),
GuestRunning: running,
CloudflaredStatus: rep.Cloudflared.Status,
}
if err := h.store.SaveHostReport(hostID, custID, body, denorm); err != nil {
h.logger.Printf("[ERROR] Failed to save host-report from %s: %v", hostID, err)
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
for _, g := range rep.Guests {
status := g.Status
if status == "" {
status = "unknown"
}
guest := &store.Guest{
GuestID: store.GuestID(hostID, g.VMID),
CustomerID: custID,
HostID: hostID,
VMID: g.VMID,
DisplayName: g.Name,
Status: status,
ControllerVersion: g.ControllerVersion,
}
if err := h.store.UpsertGuestFromReport(guest); err != nil {
// A guest upsert failure must not drop the whole report (liveness).
h.logger.Printf("[WARN] Failed to upsert guest %s: %v", guest.GuestID, err)
}
}
// storage_targets (slice 5): persisted as part of report_json above. Count + surface
// disconnected ones in the log (the slice-10 manifest will reconcile them; for now the
// signal is the visibility — a disconnected target is the storage analog of host-down).
disconnected := 0
for _, st := range rep.StorageTargets {
if st.State == "disconnected" {
disconnected++
}
}
if disconnected > 0 {
h.logger.Printf("[WARN] host %s reports %d disconnected storage target(s) of %d",
hostID, disconnected, len(rep.StorageTargets))
}
// restore_tests (slice 6): a FAILED self-restore-test is the loudest DR signal there is
// — surface it prominently. A PASS that carried start warnings (e.g. the systemd-nesting
// advisory) is surfaced too: INFO when every warning is recognized-benign, escalated to
// WARN when an UNRECOGNIZED warning stood out (as loud as a failed PBS verify is for
// backups), so a real restore warning can't hide behind a green pass. A backup whose
// vzdump failed is also worth a warning.
for _, rt := range rep.RestoreTests {
switch {
case !rt.Pass:
h.logger.Printf("[WARN] host %s restore-test FAILED: archive=%s tier=%s scratch=%d err=%q",
hostID, rt.SourceArchive, rt.SourceTier, rt.ScratchVMID, rt.Error)
case len(rt.Warnings) == 0:
// clean pass — nothing to surface here (counted in the summary line below).
case rt.WarningsRecognized:
h.logger.Printf("[INFO] host %s restore-test passed WITH WARNINGS (recognized): archive=%s tier=%s warnings=%v",
hostID, rt.SourceArchive, rt.SourceTier, rt.Warnings)
default:
h.logger.Printf("[WARN] host %s restore-test passed WITH UNRECOGNIZED WARNINGS: archive=%s tier=%s warnings=%v",
hostID, rt.SourceArchive, rt.SourceTier, rt.Warnings)
}
}
for _, bk := range rep.Backups {
if !bk.Success {
h.logger.Printf("[WARN] host %s backup FAILED: target=%s vmid=%d err=%q",
hostID, bk.TargetID, bk.VMID, bk.Error)
}
}
// pbs_snapshots (slice 6 Phase B): a FAILED PBS verify is the loudest offsite-DR signal.
for _, ps := range rep.PBSSnapshots {
if ps.VerifyState == "failed" {
h.logger.Printf("[WARN] host %s PBS verify FAILED: %s/%s ns=%s owner=%s",
hostID, ps.BackupType, ps.BackupID, ps.Namespace, ps.Owner)
}
}
h.logger.Printf("[INFO] host-report from %s (%d guests, %d storage targets, %d backups, %d restore-tests, %d pbs-snapshots, %d bytes)",
hostID, len(rep.Guests), len(rep.StorageTargets), len(rep.Backups), len(rep.RestoreTests), len(rep.PBSSnapshots), len(body))
// Agent log tail (v0.46.0) — the debug-ring bundle a prior envelope requested.
// SaveLogBundle runs the secret gate + clears the pending request (consume-once).
// A failure must NOT drop the heartbeat; just warn.
if rep.LogTail != nil {
blocked, berr := h.store.SaveLogBundle(hostID, store.LogBundleComponentAgent, rep.LogTail.CollectedAt, rep.LogTail.Lines)
switch {
case berr != nil:
h.logger.Printf("[WARN] Failed to save agent log bundle for %s: %v", hostID, berr)
case blocked:
h.logger.Printf("[WARN] agent log bundle for %s BLOCKED: possible secret in log content — nothing stored", hostID)
default:
h.logger.Printf("[INFO] agent log bundle received for %s (%d lines)", hostID, len(rep.LogTail.Lines))
}
}
// DR recipe — persist the agent's secret-free storage/guest/PBS half (preserving any app half).
// A failure here must NOT drop the heartbeat (the report already saved); just warn.
if len(rep.DRRecipe) > 0 && custID != "" {
var ver drRecipeVersionOnly
_ = json.Unmarshal(rep.DRRecipe, &ver)
if err := h.store.SaveDRRecipeHostHalf(custID, hostID, ver.RecipeVersion, rep.DRRecipe); err != nil {
h.logger.Printf("[WARN] Failed to save DR-recipe host-half for customer %s (host %s): %v", custID, hostID, err)
} else {
h.logger.Printf("[INFO] DR-recipe host-half stored for customer %s (host %s, v%d)", custID, hostID, ver.RecipeVersion)
}
}
blocked := false
if cc, err := h.store.GetCustomerConfig(custID); err == nil && cc != nil && cc.Status == "blocked" {
blocked = true
}
// Control envelope (slice 10A): the cheap change-notification. desired_generation is the
// host's current generation (the agent re-fetches the full desired-state only when it
// advances past its cached one); has_signed_ops flags a non-empty signed-jobs queue (the
// agent fetches/executes them in 10B). Both degrade safely to their slice-4 defaults on a
// store error — a heartbeat must never fail on the control channel.
var desiredGen int64
if host, err := h.store.GetHost(hostID); err == nil && host != nil {
desiredGen = host.DesiredGeneration
}
hasSignedOps := false
if n, err := h.store.CountSignedJobs(hostID); err == nil && n > 0 {
hasSignedOps = true
}
resp := map[string]interface{}{
"status": "ok",
"poll_interval_seconds": defaultHostPollSeconds,
"blocked": blocked,
"desired_generation": desiredGen,
"has_signed_ops": hasSignedOps,
}
// v0.46.0 — pending agent log pull: the NEXT heartbeat carries log_tail (agent ≥
// 0.83.0; older agents ignore the flag and the request stays visibly pending).
if pending, err := h.store.PendingLogBundleRequest(hostID, store.LogBundleComponentAgent); err == nil && pending {
resp["log_tail_requested"] = true
}
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusOK)
json.NewEncoder(w).Encode(resp)
}
// handleAdminCreateHost mints a host identity (host_id + per-host api_key).
//
// PROVISIONAL (slice-3 bootstrap): global-key only, so the demo agent can
// authenticate before enrollment (slices 78) exists. Enrollment will mint host
// identity + pin signing keys; this endpoint should be removed/locked down then
// (tracked under doc 05 §11 auth-tightening at cutover).
func (h *Handler) handleAdminCreateHost(w http.ResponseWriter, r *http.Request) {
_, _, isGlobal, ok := h.checkAuthHost(r)
if !ok || !isGlobal {
http.Error(w, "Forbidden: global key required", http.StatusForbidden)
return
}
body, err := io.ReadAll(io.LimitReader(r.Body, 1<<20))
if err != nil {
http.Error(w, "Bad request", http.StatusBadRequest)
return
}
var req struct {
CustomerID string `json:"customer_id"`
HostID string `json:"host_id"`
DisplayName string `json:"display_name"`
}
if err := json.Unmarshal(body, &req); err != nil || req.CustomerID == "" {
http.Error(w, "Invalid payload: customer_id required", http.StatusBadRequest)
return
}
cc, err := h.store.GetCustomerConfig(req.CustomerID)
if err != nil {
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
if cc == nil {
http.Error(w, "Unknown customer_id", http.StatusBadRequest)
return
}
hostID := req.HostID
if hostID == "" {
sfx, err := configgen.RandomHex(3) // 6 hex chars — human-legible for the demo
if err != nil {
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
hostID = req.CustomerID + "-" + sfx
}
apiKey, err := configgen.RandomHex(32)
if err != nil {
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
if err := h.store.UpsertHost(&store.Host{HostID: hostID, CustomerID: req.CustomerID, APIKey: apiKey}); err != nil {
h.logger.Printf("[ERROR] Failed to mint host for %s: %v", req.CustomerID, err)
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
h.logger.Printf("[INFO] provisional host mint: %s (customer %s)", hostID, req.CustomerID)
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusCreated)
json.NewEncoder(w).Encode(map[string]string{"host_id": hostID, "api_key": apiKey})
}
// handleHostEnroll is the passphrase-authed, mint-once-reuse host enrollment for Day-0
// (option C, SPIKE-day0-firstboot-handshake-2026-06-26). It is the sibling of the
// global-key handleAdminCreateHost: the operator/host-bootstrap script carries ONLY the
// customer's retrieval passphrase (no global key in the field deploy path), POSTs the
// customer_id, and gets back the host credential — minted on first call, REUSED byte-for-
// byte on every subsequent call (so re-running the bootstrap never orphans a live agent's
// key). Auth (passphrase) is checked BEFORE any mint — a bad-auth call never writes a row.
// The proven GET /config/{id} controller pull and POST /admin/hosts escape hatch are
// untouched.
func (h *Handler) handleHostEnroll(w http.ResponseWriter, r *http.Request) {
body, err := io.ReadAll(io.LimitReader(r.Body, 1<<20))
if err != nil {
http.Error(w, "Bad request", http.StatusBadRequest)
return
}
var req struct {
CustomerID string `json:"customer_id"`
}
if err := json.Unmarshal(body, &req); err != nil || req.CustomerID == "" {
http.Error(w, "Invalid payload: customer_id required", http.StatusBadRequest)
return
}
// Passphrase auth — mirrors handleConfigRetrieve exactly (header, 404-then-401 order,
// constant-time compare). Happens BEFORE any mint.
password := r.Header.Get("X-Retrieval-Password")
if password == "" {
http.Error(w, "Unauthorized: X-Retrieval-Password header required", http.StatusUnauthorized)
return
}
cc, err := h.store.GetCustomerConfig(req.CustomerID)
if err != nil {
h.logger.Printf("[ERROR] host-enroll: customer lookup failed for %s: %v", req.CustomerID, err)
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
if cc == nil {
http.Error(w, "Not found", http.StatusNotFound)
return
}
if subtle.ConstantTimeCompare([]byte(password), []byte(cc.RetrievalPassword)) != 1 {
http.Error(w, "Unauthorized: invalid password", http.StatusUnauthorized)
return
}
// Mint-once-reuse: an existing host for this customer is returned as-is (idempotent).
existing, err := h.store.GetHostByCustomer(req.CustomerID)
if err != nil {
h.logger.Printf("[ERROR] host-enroll: host lookup failed for %s: %v", req.CustomerID, err)
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
if existing != nil {
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusOK)
json.NewEncoder(w).Encode(map[string]string{"host_id": existing.HostID, "api_key": existing.APIKey})
return
}
// First enroll: mint (mirrors handleAdminCreateHost's mint block).
sfx, err := configgen.RandomHex(3) // 6 hex chars — host_id suffix
if err != nil {
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
hostID := req.CustomerID + "-" + sfx
apiKey, err := configgen.RandomHex(32)
if err != nil {
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
if err := h.store.UpsertHost(&store.Host{HostID: hostID, CustomerID: req.CustomerID, APIKey: apiKey}); err != nil {
h.logger.Printf("[ERROR] host-enroll: failed to mint host for %s: %v", req.CustomerID, err)
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
h.logger.Printf("[INFO] host enrolled: %s (customer %s)", hostID, req.CustomerID)
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusCreated)
json.NewEncoder(w).Encode(map[string]string{"host_id": hostID, "api_key": apiKey})
}
// escrowUploadRequest is the agent→hub wire shape for the OPAQUE PBS recovery-code escrow blob
// (slice 7, doc 03 §8a). It MUST stay in lockstep with the agent's emit struct
// (felhom-agent cmd/felhom-agent escrowUploadRequest). The hub stores the bytes and NEVER decrypts
// them (it has no recovery code).
type escrowUploadRequest struct {
BlobB64 string `json:"blob_b64"` // base64 of the opaque R-wrapped blob (ciphertext)
KeyFingerprint string `json:"key_fingerprint"` // for operator display only
Posture string `json:"posture"` // e.g. "zero_knowledge"
CreatedAt string `json:"created_at"` // RFC3339
// Slice 10D.1 — optional DR bundle, stored alongside the K-escrow (both opaque/non-secret).
IdentityBlobB64 string `json:"identity_blob_b64,omitempty"` // age-wrapped {tunnel_token, pbs_token}
DirectiveJSON json.RawMessage `json:"directive,omitempty"` // non-secret directive (pbs repo/ns, expected fp, tunnel id)
// SLICE 3 — sha256 hex of the restic repo password sealed in the identity blob (non-reversible hash
// of a 256-bit random secret — safe to store/serve; present only when a staged password was folded in).
ResticPwSHA256 string `json:"restic_pw_sha256,omitempty"`
}
// handleHostEscrowPut stores a host's opaque escrow blob (doc 03 §8a). Authed with the PER-HOST key
// (a host may only write its own escrow; the global operator key is also accepted). The hub keeps
// the ciphertext and never opens it. Last-write-wins (rotation). No serving this slice (slice 10).
func (h *Handler) handleHostEscrowPut(w http.ResponseWriter, r *http.Request, pathHostID string) {
authHostID, _, isGlobal, ok := h.checkAuthHost(r)
if !ok {
http.Error(w, "Unauthorized", http.StatusUnauthorized)
return
}
if pathHostID == "" {
http.Error(w, "Missing host_id", http.StatusBadRequest)
return
}
// A per-host key may only write ITS OWN escrow; the global key may write any.
if !isGlobal && authHostID != pathHostID {
http.Error(w, "Forbidden: host_id mismatch", http.StatusForbidden)
return
}
body, err := io.ReadAll(io.LimitReader(r.Body, 1<<20)) // 1 MB cap; the blob is ~hundreds of bytes
if err != nil {
http.Error(w, "Bad request", http.StatusBadRequest)
return
}
var req escrowUploadRequest
if err := json.Unmarshal(body, &req); err != nil || req.BlobB64 == "" {
http.Error(w, "Invalid payload: blob_b64 required", http.StatusBadRequest)
return
}
blob, err := base64.StdEncoding.DecodeString(req.BlobB64)
if err != nil || len(blob) == 0 {
http.Error(w, "Invalid payload: blob_b64 not valid base64", http.StatusBadRequest)
return
}
createdAt := req.CreatedAt
if createdAt == "" {
createdAt = time.Now().UTC().Format(time.RFC3339)
}
// Store the OPAQUE bytes. No decrypt path exists — the hub cannot open this.
if err := h.store.SaveHostEscrow(pathHostID, blob, req.KeyFingerprint, req.Posture, createdAt, req.ResticPwSHA256); err != nil {
h.logger.Printf("[ERROR] Failed to store escrow for host %s: %v", pathHostID, err)
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
// Slice 10D.1: optionally store the IDENTITY escrow blob + the non-secret DR directive alongside
// the K-escrow (both opaque / non-secret — no usable secret hub-side). Additive: a slice-7
// upload without these is unchanged.
if req.IdentityBlobB64 != "" {
idBlob, derr := base64.StdEncoding.DecodeString(req.IdentityBlobB64)
if derr != nil || len(idBlob) == 0 {
http.Error(w, "Invalid payload: identity_blob_b64 not valid base64", http.StatusBadRequest)
return
}
directive := req.DirectiveJSON
if len(directive) == 0 || !json.Valid(directive) {
directive = json.RawMessage("{}")
}
if err := h.store.SaveHostDRBundle(pathHostID, idBlob, string(directive)); err != nil {
h.logger.Printf("[ERROR] Failed to store DR bundle for host %s: %v", pathHostID, err)
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
h.logger.Printf("[INFO] stored DR bundle for host %s (identity %d bytes + directive)", pathHostID, len(idBlob))
}
h.logger.Printf("[INFO] stored opaque escrow blob for host %s (%d bytes, posture=%s, fp=%s)",
pathHostID, len(blob), req.Posture, req.KeyFingerprint)
w.WriteHeader(http.StatusOK)
w.Write([]byte(`{"status":"ok"}`))
}
// handleHostRecoveryCredentialPut vaults a host's break-glass root@pam console credential (TASK G1).
// SELF-SCOPED (a host key writes only its own; global may write any) — day-0 posts it with the
// host api_key. The secret is stored at rest and NEVER logged (only the username + a length are
// logged). This is the human fallback for when both the sshd path AND the agent-independent
// auto-heal have failed: the operator retrieves it to reach the PVE web console (pveproxy — a
// failure domain distinct from sshd).
func (h *Handler) handleHostRecoveryCredentialPut(w http.ResponseWriter, r *http.Request, pathHostID string) {
authHostID, _, isGlobal, ok := h.checkAuthHost(r)
if !ok {
http.Error(w, "Unauthorized", http.StatusUnauthorized)
return
}
if pathHostID == "" {
http.Error(w, "Missing host_id", http.StatusBadRequest)
return
}
if !isGlobal && authHostID != pathHostID {
http.Error(w, "Forbidden: host_id mismatch", http.StatusForbidden)
return
}
body, err := io.ReadAll(io.LimitReader(r.Body, 1<<16)) // 64 KiB cap; a username+password is tiny
if err != nil {
http.Error(w, "Bad request", http.StatusBadRequest)
return
}
var req struct {
Username string `json:"username"`
Password string `json:"password"`
}
if err := json.Unmarshal(body, &req); err != nil || req.Username == "" || req.Password == "" {
http.Error(w, "Invalid payload: username + password required", http.StatusBadRequest)
return
}
// The host must exist (mint-first) — a per-host key already proves it; the global path re-checks.
if isGlobal {
host, herr := h.store.GetHost(pathHostID)
if herr != nil || host == nil {
http.Error(w, "Unknown host_id", http.StatusBadRequest)
return
}
}
if err := h.store.SaveHostRecoveryCredential(pathHostID, req.Username, req.Password); err != nil {
h.logger.Printf("[ERROR] Failed to vault recovery credential for host %s: %v", pathHostID, err)
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
// SECRET DISCIPLINE: log the username + a length only — NEVER the password.
h.logger.Printf("[INFO] vaulted break-glass recovery credential for host %s (user=%s, secret %d chars)",
pathHostID, req.Username, len(req.Password))
w.WriteHeader(http.StatusOK)
w.Write([]byte(`{"status":"ok"}`))
}
// handleAdminGetRecoveryCredential returns a host's vaulted break-glass credential to the OPERATOR
// (global key only — a per-host key must NOT read its own console password back out). This is the
// authenticated retrieval path the break-glass runbook uses. The response body carries the secret by
// necessity; it is never written to the hub log.
func (h *Handler) handleAdminGetRecoveryCredential(w http.ResponseWriter, r *http.Request, pathHostID string) {
_, _, isGlobal, ok := h.checkAuthHost(r)
if !ok || !isGlobal {
http.Error(w, "Unauthorized", http.StatusUnauthorized) // operator/global key ONLY
return
}
if pathHostID == "" {
http.Error(w, "Missing host_id", http.StatusBadRequest)
return
}
cred, err := h.store.GetHostRecoveryCredential(pathHostID)
if err != nil {
h.logger.Printf("[ERROR] Failed to read recovery credential for host %s: %v", pathHostID, err)
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
if cred == nil {
http.Error(w, "No recovery credential vaulted for this host", http.StatusNotFound)
return
}
h.logger.Printf("[INFO] operator retrieved break-glass recovery credential for host %s (user=%s)", pathHostID, cred.Username)
resp, _ := json.Marshal(map[string]string{
"host_id": cred.HostID,
"username": cred.Username,
"password": cred.Secret,
"set_at": cred.SetAt.UTC().Format(time.RFC3339),
})
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusOK)
w.Write(resp)
}
// handleGetDesiredState serves a host its authoritative desired-state (slice 10A). Per-host key,
// SELF-SCOPED: a host reads ONLY its own (the global operator key may read any). The agent fetches
// this when the heartbeat envelope's desired_generation has advanced past its cached one. The
// response carries the generation the state corresponds to, so the agent caches it atomically.
func (h *Handler) handleGetDesiredState(w http.ResponseWriter, r *http.Request, pathHostID string) {
authHostID, _, isGlobal, ok := h.checkAuthHost(r)
if !ok {
http.Error(w, "Unauthorized", http.StatusUnauthorized)
return
}
if pathHostID == "" {
http.Error(w, "Missing host_id", http.StatusBadRequest)
return
}
if !isGlobal && authHostID != pathHostID {
http.Error(w, "Forbidden: host_id mismatch", http.StatusForbidden)
return
}
host, err := h.store.GetHost(pathHostID)
if err != nil {
h.logger.Printf("[ERROR] desired-state lookup for %s: %v", pathHostID, err)
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
if host == nil {
http.Error(w, "Unknown host_id", http.StatusNotFound)
return
}
desired := host.DesiredJSON
if strings.TrimSpace(desired) == "" {
desired = "{}"
}
// S2: merge the hub-OWNED wireguard block at read time (no peer → pass-through unchanged;
// the stored operator blob is never modified). See api/wg.go mergeWireguard.
desired = h.mergeWireguard(pathHostID, desired)
resp := map[string]interface{}{
"generation": host.DesiredGeneration,
"desired_state": json.RawMessage(desired), // opaque to the hub — agent owns the schema
}
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusOK)
json.NewEncoder(w).Encode(resp)
}
// handleGetJobs serves a host its pending signed-op blobs (slice 10A). Per-host key, SELF-SCOPED.
// The blobs are OPAQUE (the hub never forged or opened them); the agent verifies + executes them
// in 10B. 10A only serves the queue.
func (h *Handler) handleGetJobs(w http.ResponseWriter, r *http.Request, pathHostID string) {
authHostID, _, isGlobal, ok := h.checkAuthHost(r)
if !ok {
http.Error(w, "Unauthorized", http.StatusUnauthorized)
return
}
if pathHostID == "" {
http.Error(w, "Missing host_id", http.StatusBadRequest)
return
}
if !isGlobal && authHostID != pathHostID {
http.Error(w, "Forbidden: host_id mismatch", http.StatusForbidden)
return
}
jobs, err := h.store.GetSignedJobs(pathHostID)
if err != nil {
h.logger.Printf("[ERROR] jobs lookup for %s: %v", pathHostID, err)
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
out := make([]map[string]string, 0, len(jobs))
for _, j := range jobs {
out = append(out, map[string]string{
"job_id": j.JobID,
"blob_b64": base64.StdEncoding.EncodeToString(j.Blob),
"created_at": j.CreatedAt,
})
}
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusOK)
json.NewEncoder(w).Encode(map[string]interface{}{"jobs": out})
}
// handleDeleteJob clears a processed job from a host's queue (slice 10B). Per-host key,
// SELF-SCOPED (a host clears only its own jobs; the global key may clear any). Idempotent.
func (h *Handler) handleDeleteJob(w http.ResponseWriter, r *http.Request, pathHostID, jobID string) {
authHostID, _, isGlobal, ok := h.checkAuthHost(r)
if !ok {
http.Error(w, "Unauthorized", http.StatusUnauthorized)
return
}
if pathHostID == "" || jobID == "" {
http.Error(w, "Missing host_id or job_id", http.StatusBadRequest)
return
}
if !isGlobal && authHostID != pathHostID {
http.Error(w, "Forbidden: host_id mismatch", http.StatusForbidden)
return
}
if err := h.store.DeleteSignedJob(pathHostID, jobID); err != nil {
h.logger.Printf("[ERROR] delete job %s for %s: %v", jobID, pathHostID, err)
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
h.logger.Printf("[INFO] host %s cleared signed-op job %s (executed or rejected)", pathHostID, jobID)
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusOK)
w.Write([]byte(`{"status":"ok"}`))
}
// handleAdminSetDesiredState sets a host's desired-state (slice 10A). GLOBAL/operator key ONLY —
// a per-host key cannot author its own intent. The body is the desired-state JSON (opaque to the
// hub: it stores + serves bytes, never validates/interprets the schema — the agent/CLI owns it).
// Writing BUMPS desired_generation so the next heartbeat signals the agent to re-fetch.
func (h *Handler) handleAdminSetDesiredState(w http.ResponseWriter, r *http.Request, pathHostID string) {
_, _, isGlobal, ok := h.checkAuthHost(r)
if !ok || !isGlobal {
http.Error(w, "Forbidden: global key required", http.StatusForbidden)
return
}
if pathHostID == "" {
http.Error(w, "Missing host_id", http.StatusBadRequest)
return
}
body, err := io.ReadAll(io.LimitReader(r.Body, 1<<20))
if err != nil {
http.Error(w, "Bad request", http.StatusBadRequest)
return
}
// Validate it is well-formed JSON (the hub does not interpret the schema, but a malformed
// blob would break the agent's parse — reject it at the door).
if !json.Valid(body) {
http.Error(w, "Invalid payload: body must be JSON", http.StatusBadRequest)
return
}
// S2: the wireguard block is HUB-owned, merged at read time — an operator copy-paste of a
// served desired-state must never write it back into the stored blob (it would go stale and
// shadow the live assignment). Reject at the door.
var topKeys map[string]json.RawMessage
if err := json.Unmarshal(body, &topKeys); err == nil {
if _, has := topKeys["wireguard"]; has {
http.Error(w, "wireguard is hub-owned; register via POST /hosts/{id}/wg", http.StatusBadRequest)
return
}
}
gen, err := h.store.SetHostDesired(pathHostID, body)
if err == sql.ErrNoRows {
http.Error(w, "Unknown host_id", http.StatusNotFound)
return
}
if err != nil {
h.logger.Printf("[ERROR] set desired-state for %s: %v", pathHostID, err)
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
h.logger.Printf("[INFO] admin-set desired-state for host %s (generation now %d, %d bytes)", pathHostID, gen, len(body))
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusOK)
json.NewEncoder(w).Encode(map[string]interface{}{"status": "ok", "generation": gen})
}
// handleAdminEnqueueJob appends an opaque signed-op blob to a host's queue (slice 10A). GLOBAL key
// ONLY. The blob is pre-signed off-hub (the hub holds no signing key); the hub stores it verbatim.
// This is the minimal operator path to seed the queue so HasSignedOps/serving are exercisable; the
// rich operator/signing UX is later. Execution is 10B.
func (h *Handler) handleAdminEnqueueJob(w http.ResponseWriter, r *http.Request, pathHostID string) {
_, _, isGlobal, ok := h.checkAuthHost(r)
if !ok || !isGlobal {
http.Error(w, "Forbidden: global key required", http.StatusForbidden)
return
}
if pathHostID == "" {
http.Error(w, "Missing host_id", http.StatusBadRequest)
return
}
host, err := h.store.GetHost(pathHostID)
if err != nil || host == nil {
http.Error(w, "Unknown host_id", http.StatusNotFound)
return
}
body, err := io.ReadAll(io.LimitReader(r.Body, 1<<20))
if err != nil {
http.Error(w, "Bad request", http.StatusBadRequest)
return
}
var req struct {
JobID string `json:"job_id"`
BlobB64 string `json:"blob_b64"`
}
if err := json.Unmarshal(body, &req); err != nil || req.BlobB64 == "" {
http.Error(w, "Invalid payload: blob_b64 required", http.StatusBadRequest)
return
}
blob, err := base64.StdEncoding.DecodeString(req.BlobB64)
if err != nil || len(blob) == 0 {
http.Error(w, "Invalid payload: blob_b64 not valid base64", http.StatusBadRequest)
return
}
if req.JobID == "" {
req.JobID, _ = configgen.RandomHex(8)
}
if err := h.store.EnqueueSignedJob(pathHostID, req.JobID, blob); err != nil {
h.logger.Printf("[ERROR] enqueue job for %s: %v", pathHostID, err)
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
h.logger.Printf("[INFO] enqueued signed-op job %s for host %s (%d bytes)", req.JobID, pathHostID, len(blob))
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusCreated)
json.NewEncoder(w).Encode(map[string]interface{}{"status": "ok", "job_id": req.JobID})
}
// handleClaimResetRequest is the controller-forwarded "Elfelejtett jelszó" (v0.50.0): the box
// asks the hub to email a fresh reset code to the REGISTERED customer address — the requester
// never chooses the destination. Auth: the customer's own report Bearer key (self-scoped).
// The response is deliberately neutral 200 on every authorized outcome (cap reached, email
// failure) — the customer-facing message is always "ha az e-mail cím regisztrálva van…"; the
// real outcome goes to the operator log + notification_log.
func (h *Handler) handleClaimResetRequest(w http.ResponseWriter, r *http.Request) {
authCustomerID, isGlobal, ok := h.checkAuthCustomer(r)
if !ok {
http.Error(w, "Unauthorized", http.StatusUnauthorized)
return
}
body, err := io.ReadAll(io.LimitReader(r.Body, 4096))
if err != nil {
http.Error(w, "Bad request", http.StatusBadRequest)
return
}
var payload struct {
CustomerID string `json:"customer_id"`
}
if err := json.Unmarshal(body, &payload); err != nil || payload.CustomerID == "" {
http.Error(w, "Invalid payload: customer_id required", http.StatusBadRequest)
return
}
if !isGlobal && authCustomerID != payload.CustomerID {
http.Error(w, "Forbidden: customer_id mismatch", http.StatusForbidden)
return
}
if h.claimEngine == nil {
http.Error(w, "Claim engine not available", http.StatusServiceUnavailable)
return
}
cfg, err := h.store.GetCustomerConfig(payload.CustomerID)
if err != nil || cfg == nil {
http.Error(w, "Not found", http.StatusNotFound)
return
}
if err := h.claimEngine.RequestReset(cfg); err != nil {
// Neutral to the box; loud to the operator (cap reached / send failure / no email).
h.logger.Printf("[WARN] claim reset-request for %s: %v", payload.CustomerID, err)
} else {
h.logger.Printf("[INFO] claim reset-request for %s: reset code emailed to the registered address", payload.CustomerID)
}
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusOK)
w.Write([]byte(`{"status":"ok"}`))
}
// allowedEventTypes lists all valid event_type values the Hub accepts.
var allowedEventTypes = map[string]bool{
// Controller-pushed events
"controller_started": true,
"claim_lockout": true, // v0.50.0 — claim/reset code brute-force lockout tripped
"controller_updated": true,
"backup_completed": true,
"backup_failed": true,
"db_dump_completed": true,
"db_dump_failed": true,
"backup_integrity_ok": true,
"backup_integrity_failed": true,
"crossdrive_completed": true,
"crossdrive_failed": true,
"storage_disconnected": true,
"storage_reconnected": true,
"disk_warning": true,
"disk_critical": true,
"health_degraded": true,
"health_critical": true,
"health_recovered": true,
"app_deployed": true,
"app_removed": true,
"app_start_failed": true, // controller fix-3 (CAMPAIGN-3): a deployed app is not running
"disaster_recovery_started": true,
"disaster_recovery_completed": true,
// Controller→agent channel health (controller v0.90.0) — operator-only (not customer toggles)
"agent_channel_pin_mismatch": true,
"agent_channel_unauthorized": true,
"agent_channel_unreachable": true,
"agent_channel_timeout": true,
"agent_channel_misconfigured": true,
"agent_channel_construction_error": true,
"agent_channel_unknown": true,
"agent_channel_recovered": true,
// Hub-generated events
"node_stale": true,
"node_down": true,
"node_recovered": true,
// Hub-generated host-domain events (v0.7.0, slice 3)
"host_stale": true,
"host_down": true,
"host_recovered": true,
// Hub-generated host root-fs disk-pressure (v0.23.0) — distinct from the controller's GUEST disk_*
"host_disk_warning": true,
"host_disk_critical": true,
"storage_fill_warning": true,
"storage_fill_critical": true,
"expected_backup_missed": true,
"expected_dbdump_missed": true,
// Special
"test": true,
}
// handleEvent processes structured events from controllers (new endpoint, replaces /notify for updated controllers).
func (h *Handler) handleEvent(w http.ResponseWriter, r *http.Request) {
authCustomerID, isGlobal, ok := h.checkAuthCustomer(r)
if !ok {
http.Error(w, "Unauthorized", http.StatusUnauthorized)
return
}
body, err := io.ReadAll(io.LimitReader(r.Body, 1<<20))
if err != nil {
http.Error(w, "Bad request", http.StatusBadRequest)
return
}
var payload struct {
CustomerID string `json:"customer_id"`
EventType string `json:"event_type"`
Severity string `json:"severity"`
Message string `json:"message"`
Details json.RawMessage `json:"details"`
}
if err := json.Unmarshal(body, &payload); err != nil {
http.Error(w, "Invalid JSON", http.StatusBadRequest)
return
}
if payload.CustomerID == "" || payload.EventType == "" {
http.Error(w, "customer_id and event_type are required", http.StatusBadRequest)
return
}
// Validate customer_id matches authenticated customer (unless global key)
if !isGlobal && authCustomerID != payload.CustomerID {
http.Error(w, "Forbidden: customer_id mismatch", http.StatusForbidden)
return
}
// Validate event_type
if !allowedEventTypes[payload.EventType] {
http.Error(w, fmt.Sprintf("Invalid event_type: %s", payload.EventType), http.StatusBadRequest)
return
}
// Validate/default severity (exact-match lowercase; unknown values coerce to info)
switch payload.Severity {
case "info", "warning", "error", "critical":
default:
payload.Severity = "info"
}
// Store details as JSON string
detailsStr := "{}"
if len(payload.Details) > 0 && string(payload.Details) != "null" {
detailsStr = string(payload.Details)
}
_, err = h.store.SaveEvent(payload.CustomerID, payload.EventType, payload.Severity, payload.Message, detailsStr, "controller")
if err != nil {
h.logger.Printf("[ERROR] Failed to save event from %s: %v", payload.CustomerID, err)
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
h.logger.Printf("[INFO] Event from %s: %s (%s) — %s", payload.CustomerID, payload.EventType, payload.Severity, payload.Message)
// Dispatch notifications (non-blocking)
if h.dispatcher != nil {
go h.dispatcher.ProcessEvent(payload.CustomerID, payload.EventType, payload.Severity, payload.Message, detailsStr, "controller")
}
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusOK)
w.Write([]byte(`{"ok":true}`))
}
func (h *Handler) handleCustomers(w http.ResponseWriter, r *http.Request) {
customers, err := h.store.GetCustomers()
if err != nil {
h.logger.Printf("[ERROR] Failed to get customers: %v", err)
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
type customerJSON struct {
ID string `json:"id"`
Name string `json:"name"`
ControllerVersion string `json:"controller_version"`
ControllerURL string `json:"controller_url,omitempty"`
HealthStatus string `json:"health_status"`
LastSeen time.Time `json:"last_seen"`
CPUPercent float64 `json:"cpu_percent"`
MemoryPercent float64 `json:"memory_percent"`
ContainerTotal int `json:"container_total"`
ContainerRunning int `json:"container_running"`
BackupLastSnapshot *time.Time `json:"backup_last_snapshot"`
}
result := make([]customerJSON, 0, len(customers))
for _, c := range customers {
result = append(result, customerJSON{
ID: c.CustomerID,
Name: c.CustomerName,
ControllerVersion: c.ControllerVersion,
ControllerURL: c.ControllerURL,
HealthStatus: c.HealthStatus,
LastSeen: c.ReceivedAt,
CPUPercent: c.CPUPercent,
MemoryPercent: c.MemoryPercent,
ContainerTotal: c.ContainerTotal,
ContainerRunning: c.ContainerRunning,
BackupLastSnapshot: c.BackupLastSnapshot,
})
}
w.Header().Set("Content-Type", "application/json")
json.NewEncoder(w).Encode(result)
}
func (h *Handler) handleCustomer(w http.ResponseWriter, r *http.Request, customerID string) {
customer, err := h.store.GetCustomer(customerID)
if err != nil {
h.logger.Printf("[ERROR] Failed to get customer %s: %v", customerID, err)
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
if customer == nil {
http.NotFound(w, r)
return
}
w.Header().Set("Content-Type", "application/json")
// Return the full report JSON directly
w.Write([]byte(customer.ReportJSON))
}
func (h *Handler) handleCustomerHistory(w http.ResponseWriter, r *http.Request, customerID string) {
period := r.URL.Query().Get("period")
var since time.Duration
switch period {
case "7d":
since = 7 * 24 * time.Hour
case "30d":
since = 30 * 24 * time.Hour
default:
since = 24 * time.Hour
}
history, err := h.store.GetCustomerHistory(customerID, since)
if err != nil {
h.logger.Printf("[ERROR] Failed to get history for %s: %v", customerID, err)
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
type historyEntry struct {
ReceivedAt time.Time `json:"received_at"`
HealthStatus string `json:"health_status"`
CPUPercent float64 `json:"cpu_percent"`
MemoryPercent float64 `json:"memory_percent"`
}
result := make([]historyEntry, 0, len(history))
for _, h := range history {
result = append(result, historyEntry{
ReceivedAt: h.ReceivedAt,
HealthStatus: h.HealthStatus,
CPUPercent: h.CPUPercent,
MemoryPercent: h.MemoryPercent,
})
}
w.Header().Set("Content-Type", "application/json")
json.NewEncoder(w).Encode(result)
}
// handleNotify processes notification events from customer controllers.
func (h *Handler) handleNotify(w http.ResponseWriter, r *http.Request) {
if !h.checkAuth(r) {
http.Error(w, "Unauthorized", http.StatusUnauthorized)
return
}
body, err := io.ReadAll(io.LimitReader(r.Body, 1<<20))
if err != nil {
http.Error(w, "Bad request", http.StatusBadRequest)
return
}
var payload struct {
CustomerID string `json:"customer_id"`
EventType string `json:"event_type"`
Severity string `json:"severity"`
Message string `json:"message"`
Details string `json:"details"`
}
if err := json.Unmarshal(body, &payload); err != nil || payload.CustomerID == "" || payload.EventType == "" {
http.Error(w, "Invalid payload: customer_id and event_type required", http.StatusBadRequest)
return
}
h.logger.Printf("[INFO] Notification from %s: %s (%s) — %s", payload.CustomerID, payload.EventType, payload.Severity, payload.Message)
// Check if customer is blocked
if h.store.IsCustomerBlocked(payload.CustomerID) {
h.logger.Printf("[INFO] Notification suppressed for blocked customer %s", payload.CustomerID)
h.store.LogNotification(payload.CustomerID, payload.EventType, payload.Severity, payload.Message, "skipped", "customer blocked", "customer")
w.WriteHeader(http.StatusOK)
w.Write([]byte(`{"status":"ok","sent":false,"reason":"blocked"}`))
return
}
// Look up customer notification preferences
prefs, err := h.store.GetNotificationPrefs(payload.CustomerID)
if err != nil {
h.logger.Printf("[ERROR] Failed to get notification prefs for %s: %v", payload.CustomerID, err)
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
// Check if customer has email configured and event type is enabled
if prefs == nil || prefs.Email == "" {
h.logger.Printf("[INFO] No email configured for %s, skipping notification", payload.CustomerID)
h.store.LogNotification(payload.CustomerID, payload.EventType, payload.Severity, payload.Message, "skipped", "no email configured", "customer")
w.WriteHeader(http.StatusOK)
w.Write([]byte(`{"status":"ok","sent":false,"reason":"no_email"}`))
return
}
// Check if event type is in the enabled list (test events always pass)
eventEnabled := payload.EventType == "test"
for _, e := range prefs.EnabledEvents {
if e == payload.EventType {
eventEnabled = true
break
}
}
if !eventEnabled {
h.logger.Printf("[INFO] Event %s not enabled for %s, skipping", payload.EventType, payload.CustomerID)
h.store.LogNotification(payload.CustomerID, payload.EventType, payload.Severity, payload.Message, "skipped", "event not enabled", "customer")
w.WriteHeader(http.StatusOK)
w.Write([]byte(`{"status":"ok","sent":false,"reason":"event_disabled"}`))
return
}
// Send email via Resend API
if h.resendAPIKey == "" {
h.logger.Printf("[WARN] Resend API key not configured, cannot send notification email")
h.store.LogNotification(payload.CustomerID, payload.EventType, payload.Severity, payload.Message, "skipped", "resend api key not configured", "customer")
w.WriteHeader(http.StatusOK)
w.Write([]byte(`{"status":"ok","sent":false,"reason":"no_api_key"}`))
return
}
subject, emailBody := formatNotificationEmail(payload.CustomerID, payload.EventType, payload.Severity, payload.Message, payload.Details)
sendErr := h.sendResendEmail(prefs.Email, subject, emailBody)
if sendErr != nil {
h.logger.Printf("[ERROR] Failed to send notification email to %s: %v", prefs.Email, sendErr)
h.store.LogNotification(payload.CustomerID, payload.EventType, payload.Severity, payload.Message, "failed", sendErr.Error(), "customer")
http.Error(w, "Failed to send email", http.StatusInternalServerError)
return
}
h.logger.Printf("[INFO] Notification email sent to %s for %s/%s", prefs.Email, payload.CustomerID, payload.EventType)
h.store.LogNotification(payload.CustomerID, payload.EventType, payload.Severity, payload.Message, "sent", "", "customer")
w.WriteHeader(http.StatusOK)
w.Write([]byte(`{"status":"ok","sent":true}`))
}
// handleSavePreferences stores notification preferences pushed from a customer controller.
func (h *Handler) handleSavePreferences(w http.ResponseWriter, r *http.Request) {
if !h.checkAuth(r) {
http.Error(w, "Unauthorized", http.StatusUnauthorized)
return
}
body, err := io.ReadAll(io.LimitReader(r.Body, 1<<20))
if err != nil {
http.Error(w, "Bad request", http.StatusBadRequest)
return
}
var payload struct {
CustomerID string `json:"customer_id"`
Email string `json:"email"`
EnabledEvents []string `json:"enabled_events"`
CooldownHours int `json:"cooldown_hours"`
}
if err := json.Unmarshal(body, &payload); err != nil || payload.CustomerID == "" {
http.Error(w, "Invalid payload: customer_id required", http.StatusBadRequest)
return
}
if err := h.store.SaveNotificationPrefs(payload.CustomerID, payload.Email, payload.EnabledEvents, payload.CooldownHours); err != nil {
h.logger.Printf("[ERROR] Failed to save notification prefs for %s: %v", payload.CustomerID, err)
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
h.logger.Printf("[INFO] Notification preferences updated for %s: email=%s, events=%v", payload.CustomerID, payload.Email, payload.EnabledEvents)
w.WriteHeader(http.StatusOK)
w.Write([]byte(`{"status":"ok"}`))
}
// handleRecovery returns the generated controller.yaml for disaster recovery.
// Auth: X-Retrieval-Password header (same as config retrieval).
//
// The infra-backup payload was retired (Phase-1, 2026-06-16): it pushed plaintext
// customer secrets to the hub (a zero-knowledge violation) and had been dead since
// slice 8C. DR config now comes from the generated controller.yaml here; the data
// bytes come from the agent's PBS whole-CT snapshot. A secret-free DR recipe is the
// later DR slice's job.
func (h *Handler) handleRecovery(w http.ResponseWriter, r *http.Request, customerID string) {
if customerID == "" {
http.Error(w, "Missing customer_id", http.StatusBadRequest)
return
}
password := r.Header.Get("X-Retrieval-Password")
if password == "" {
http.Error(w, "Unauthorized: X-Retrieval-Password header required", http.StatusUnauthorized)
return
}
cfg, err := h.store.GetCustomerConfig(customerID)
if err != nil {
h.logger.Printf("[ERROR] Recovery: failed to get customer config for %s: %v", customerID, err)
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
if cfg == nil {
http.Error(w, "Not found", http.StatusNotFound)
return
}
if subtle.ConstantTimeCompare([]byte(password), []byte(cfg.RetrievalPassword)) != 1 {
http.Error(w, "Unauthorized: invalid password", http.StatusUnauthorized)
return
}
// Generate controller.yaml. The claim state is baked read-only (no issue on the DR path — a
// recovered box whose settings.json is gone re-gates on the EXISTING hash; the reset flow
// covers a customer who lost the password with the box).
var configYAML string
if h.templateProvider != nil {
var claimState *store.ClaimState
if h.claimEngine != nil {
claimState, _ = h.store.GetClaim(customerID)
}
yamlOutput, err := configgen.Generate(h.templateProvider.Template(), cfg, claimState)
if err != nil {
h.logger.Printf("[ERROR] Recovery: failed to generate config for %s: %v", customerID, err)
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
configYAML = yamlOutput
}
// infra_backup retired: the response keeps has_infra_backup=false so any old client
// degrades gracefully to the config_yaml-only path.
resp := struct {
CustomerID string `json:"customer_id"`
ConfigYAML string `json:"config_yaml"`
HasInfraBackup bool `json:"has_infra_backup"`
}{
CustomerID: customerID,
ConfigYAML: configYAML,
HasInfraBackup: false,
}
h.logger.Printf("[INFO] Recovery data downloaded for customer %s (config only; infra-backup retired)", customerID)
w.Header().Set("Content-Type", "application/json")
json.NewEncoder(w).Encode(resp)
}
// handleConfigRetrieve returns a generated controller.yaml for a customer.
// Auth: X-Retrieval-Password header (not Bearer token).
func (h *Handler) handleConfigRetrieve(w http.ResponseWriter, r *http.Request, customerID string) {
if customerID == "" {
http.Error(w, "Missing customer_id", http.StatusBadRequest)
return
}
password := r.Header.Get("X-Retrieval-Password")
if password == "" {
http.Error(w, "Unauthorized: X-Retrieval-Password header required", http.StatusUnauthorized)
return
}
cfg, err := h.store.GetCustomerConfig(customerID)
if err != nil {
h.logger.Printf("[ERROR] Failed to get customer config for %s: %v", customerID, err)
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
if cfg == nil {
http.Error(w, "Not found", http.StatusNotFound)
return
}
// Constant-time comparison to prevent timing attacks
if subtle.ConstantTimeCompare([]byte(password), []byte(cfg.RetrievalPassword)) != 1 {
http.Error(w, "Unauthorized: invalid password", http.StatusUnauthorized)
return
}
if h.templateProvider == nil {
http.Error(w, "Config generation not available", http.StatusServiceUnavailable)
return
}
// Customer-claim arc (v0.50.0): the REAL config pull (Day-0 installer / controller refresh) is
// the Day-0 claim entry point — issue + email the first code here (idempotent), and bake the
// active hash into the generated web.claim_code_hash so the box is gated from FIRST boot.
// (The operator-UI preview deliberately does NOT issue — it only bakes an existing hash.)
var claimState *store.ClaimState
if h.claimEngine != nil {
var cerr error
claimState, cerr = h.claimEngine.EnsureIssued(cfg)
if cerr != nil {
// Loud but non-fatal: the config is still served; if a hash was stored the gate is armed
// and the operator resends the email from the customer page.
h.logger.Printf("[WARN] claim issue for %s on config retrieve: %v", customerID, cerr)
}
}
yamlOutput, err := configgen.Generate(h.templateProvider.Template(), cfg, claimState)
if err != nil {
h.logger.Printf("[ERROR] Failed to generate config for %s: %v", customerID, err)
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
h.logger.Printf("[INFO] Config downloaded for customer %s", customerID)
w.Header().Set("Content-Type", "text/yaml; charset=utf-8")
w.Write([]byte(yamlOutput))
}
// artifactManifestResponse is the wire shape the host-bootstrap script consumes: the operator-vouched
// current agent binary + golden archive (version + sha256 each). The script fetches each artifact from
// Gitea with the config-retrieve git token and verifies its sha256 against THESE values before install
// — so the hub is the checksum trust root, a different root than Gitea (which only stores the bytes).
type artifactManifestResponse struct {
Agent artifactEntry `json:"agent"`
Golden artifactEntry `json:"golden"`
}
type artifactEntry struct {
Version string `json:"version"`
SHA256 string `json:"sha256"`
}
// handleArtifactManifest serves the current artifact set for a customer. Auth mirrors
// handleConfigRetrieve EXACTLY (X-Retrieval-Password header, 404-then-401 order, constant-time
// compare) so a script that can pull the controller.yaml can pull the manifest with the same secret.
// v0.16.0 returns the GLOBAL current set for every customer (per-customer pinning is a future hook).
// An unset manifest returns empty fields (not an error) — the script falls back to the local golden
// and fails clearly on a missing binary.
func (h *Handler) handleArtifactManifest(w http.ResponseWriter, r *http.Request, customerID string) {
if customerID == "" {
http.Error(w, "Missing customer_id", http.StatusBadRequest)
return
}
password := r.Header.Get("X-Retrieval-Password")
if password == "" {
http.Error(w, "Unauthorized: X-Retrieval-Password header required", http.StatusUnauthorized)
return
}
cfg, err := h.store.GetCustomerConfig(customerID)
if err != nil {
h.logger.Printf("[ERROR] artifacts: customer lookup failed for %s: %v", customerID, err)
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
if cfg == nil {
http.Error(w, "Not found", http.StatusNotFound)
return
}
if subtle.ConstantTimeCompare([]byte(password), []byte(cfg.RetrievalPassword)) != 1 {
http.Error(w, "Unauthorized: invalid password", http.StatusUnauthorized)
return
}
m := h.store.GetArtifactManifest()
resp := artifactManifestResponse{
Agent: artifactEntry{Version: m.AgentVersion, SHA256: m.AgentSHA256},
Golden: artifactEntry{Version: m.GoldenVersion, SHA256: m.GoldenSHA256},
}
h.logger.Printf("[INFO] Artifact manifest served for customer %s (agent=%s golden=%s)", customerID, m.AgentVersion, m.GoldenVersion)
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusOK)
json.NewEncoder(w).Encode(resp)
}
// sendResendEmail sends an email via the Resend HTTP API.
func (h *Handler) sendResendEmail(to, subject, textBody string) error {
payload := map[string]interface{}{
"from": h.fromEmail,
"to": []string{to},
"subject": subject,
"text": textBody,
}
jsonData, err := json.Marshal(payload)
if err != nil {
return fmt.Errorf("marshaling email payload: %w", err)
}
req, err := http.NewRequest("POST", "https://api.resend.com/emails", bytes.NewReader(jsonData))
if err != nil {
return fmt.Errorf("creating request: %w", err)
}
req.Header.Set("Authorization", "Bearer "+h.resendAPIKey)
req.Header.Set("Content-Type", "application/json")
resp, err := h.httpClient.Do(req)
if err != nil {
return fmt.Errorf("sending request: %w", err)
}
defer resp.Body.Close()
if resp.StatusCode >= 400 {
respBody, _ := io.ReadAll(io.LimitReader(resp.Body, 1024))
return fmt.Errorf("resend API returned %d: %s", resp.StatusCode, string(respBody))
}
return nil
}
// formatNotificationEmail creates a Hungarian email subject and body.
func formatNotificationEmail(customerID, eventType, severity, message, details string) (string, string) {
severityLabel := map[string]string{
"info": "Információ",
"warning": "Figyelmeztetés",
"error": "Hiba",
"critical": "Kritikus",
}
label := severityLabel[severity]
if label == "" {
label = severity
}
subject := fmt.Sprintf("[Felhom] %s: %s", label, message)
now := time.Now().Format("2006-01-02 15:04")
emailText := fmt.Sprintf(`Kedves Ügyfél!
A Felhom rendszered a következő figyelmeztetést jelezte:
%s
Részletek:
- Szerver: %s
- Időpont: %s
- Szint: %s
- Típus: %s`, message, customerID, now, label, eventType)
if details != "" {
emailText += fmt.Sprintf("\n- Megjegyzés: %s", details)
}
emailText += `
Ha kérdésed van, vedd fel a kapcsolatot az üzemeltetővel.
Üdvözlettel,
Felhom.eu monitoring`
return subject, emailText
}
// --- Asset endpoints ---
func (h *Handler) handleAssetsManifest(w http.ResponseWriter, r *http.Request) {
if !h.checkAuth(r) {
http.Error(w, "Unauthorized", http.StatusUnauthorized)
return
}
if h.assetsMgr == nil {
http.Error(w, "Assets not configured", http.StatusServiceUnavailable)
return
}
data, err := h.assetsMgr.MarshalManifestJSON()
if err != nil {
http.Error(w, "Internal error", http.StatusInternalServerError)
return
}
w.Header().Set("Content-Type", "application/json")
w.Write(data)
}
func (h *Handler) handleAssetFile(w http.ResponseWriter, r *http.Request, filename string) {
if !h.checkAuth(r) {
http.Error(w, "Unauthorized", http.StatusUnauthorized)
return
}
if h.assetsMgr == nil {
http.Error(w, "Assets not configured", http.StatusServiceUnavailable)
return
}
h.assetsMgr.ServeFile(w, r, filename)
}