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 7–8) 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) } // v0.52.0 (take-two F-15): serve the ACTIVE code state in the response — same shape and same // bcrypt-only guarantee as the report ACK — so the box accepts the emailed code the moment it // lands instead of waiting for the next ACK (~15 min). Served on every authorized outcome: on // a cap-reached refusal it is the unrotated row (a controller-side no-op by generation). resp := map[string]interface{}{"status": "ok"} if cs, err := h.store.GetClaim(payload.CustomerID); err == nil && cs != nil { resp["claim"] = map[string]interface{}{ "code_hash": cs.CodeHash, "generation": cs.Generation, "issued_at": cs.IssuedAt.UTC().Format(time.RFC3339), } } w.Header().Set("Content-Type", "application/json") w.WriteHeader(http.StatusOK) json.NewEncoder(w).Encode(resp) } // 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, // controller v0.134.1 — enlarged offsite push refused by the quota gate (warning; the controller's // dynamic Hungarian message is customer-grade — deliberately NO customerMessages entry, which would // discard the numbers (templates.go:129 priority)). "offbox_enlarge_blocked": 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) }