// Package cluster implements the Coven — the firewall-isolated Ley-Lines that bind // independent Sanctums into a single distributed forge. // // Topology: // // Master Sanctum — holds the Grimoire and the Tablet. // Coven-Worker — provides "Mana" (CPU/RAM) to the Cauldron for // sharded Grid-Casts (e.g., a full glibc rebuild). // // The join protocol admits new nodes after firewall-level validation. // A new node registers with the Master, begins pulling Essences from the // Tomb until it reaches parity with the rest of the Coven. // // Fester Integration: // // When a Fester master URL is configured, the Coven delegates distributed // build scheduling, node telemetry, and build dispatch to the Fester cluster // controller via its HTTP + WebSocket API. Sorcery-go handles security // (eBPF warding, tomb protection, essence verification) while Fester handles // the distributed execution brain. // // Fester API docs: https://git.dcos.net/dcosnet/fester package cluster import ( "bytes" "context" "encoding/json" "fmt" "io" "net/http" "net/url" "strings" "sync" "time" "github.com/gorilla/websocket" "dcos.net/sorcery-go/pkg/cas" ) // Node is one Sanctum in the Coven. type Node struct { ID string Arch string Role string // "master" or "worker" Address string // host:port JoinTime time.Time ComputePower int // "Mana" — drives HPC scheduling // Fester-derived fields (populated when Fester integration is active). CPU float64 `json:"cpu,omitempty"` // 0..1 Memory float64 `json:"memory,omitempty"` // 0..1 ActiveBuilds int `json:"active_builds,omitempty"` Temperature float64 `json:"temperature,omitempty"` // Celsius MaxJobs int `json:"max_jobs,omitempty"` Policy string `json:"policy,omitempty"` // preferred/avoid/neutral Status string `json:"status,omitempty"` // online/offline/draining } // Coven is the cluster manager. type Coven struct { mu sync.RWMutex Self *Node Peers map[string]*Node // Fester integration — when set, delegates scheduling and telemetry // to a Fester master instance. Nil means standalone mode. Fester *FesterClient } // NewCoven bootstraps a Master. func NewCoven(selfID, arch, address string) *Coven { return &Coven{ Self: &Node{ ID: selfID, Arch: arch, Role: "master", Address: address, JoinTime: time.Now(), }, Peers: make(map[string]*Node), } } // NewCovenWithFester bootstraps a Master connected to a Fester cluster. func NewCovenWithFester(selfID, arch, address string, festerURL string) *Coven { c := NewCoven(selfID, arch, address) if festerURL != "" { c.Fester = NewFesterClient(festerURL) } return c } // Register admits a new Worker into the Coven. func (c *Coven) Register(n *Node) error { c.mu.Lock() defer c.mu.Unlock() if n.Role == "" { n.Role = "worker" } n.JoinTime = time.Now() c.Peers[n.ID] = n return nil } // List returns every known Sanctum (including self). func (c *Coven) List() []*Node { c.mu.RLock() defer c.mu.RUnlock() out := []*Node{c.Self} for _, p := range c.Peers { out = append(out, p) } return out } // Drain migrates active builds off a node so it can be safely taken offline // for maintenance. In Fester mode, this sets the node's policy to "avoid". func (c *Coven) Drain(nodeID string) error { // If Fester is active, delegate the drain. if c.Fester != nil { return c.Fester.SetNodePolicy(nodeID, "avoid") } c.mu.Lock() defer c.mu.Unlock() if _, ok := c.Peers[nodeID]; !ok { return fmt.Errorf("cluster: unknown node %s", nodeID) } c.Peers[nodeID].ComputePower = 0 return nil } // Pulse is the heartbeat broadcast — every node reports its current load. // The Cockpit "Grid Heatmap" renders this in real time. type Pulse struct { NodeID string CPU float64 // 0..1 Memory float64 // 0..1 ActiveBuilds int Temperature float64 MaxJobs int Status string } // PulseSnapshot returns the latest heartbeat from every node. // In Fester mode, this queries the Fester /api/nodes endpoint for live // telemetry. In standalone mode, returns zeros (no probe agent). func (c *Coven) PulseSnapshot() []Pulse { // If Fester is active, fetch live telemetry. if c.Fester != nil { pulses, err := c.Fester.GetNodePulses() if err == nil && len(pulses) > 0 { return pulses } // Fall through to local state on error. } c.mu.RLock() defer c.mu.RUnlock() out := []Pulse{{NodeID: c.Self.ID}} for _, p := range c.Peers { out = append(out, Pulse{ NodeID: p.ID, CPU: p.CPU, Memory: p.Memory, ActiveBuilds: p.ActiveBuilds, Temperature: p.Temperature, }) } return out } // --------------------------------------------------------------------------- // Fester HTTP API Client // --------------------------------------------------------------------------- // FesterClient connects to a Fester master's HTTP API for node telemetry, // build scheduling, distributed execution, and shared artifact caching. // // It maps Fester's REST endpoints to the Coven's cluster abstractions: // // Fester Endpoint → Coven Method // GET /api/nodes → GetNodes, GetNodePulses // GET /api/nodes/runtimes → GetRuntimes // POST /api/nodes/{name}/policy → SetNodePolicy // POST /api/nodes/{name}/probe → ProbeNode // POST /api/build → SubmitBuild // GET /api/builds → ListBuilds // POST /api/builds/{id}/cancel → CancelBuild // GET /api/metrics/json → GetMetrics // GET /api/cause/explain/{n} → CauseExplain // WS /ws → WatchEvents // PUT /api/cas/{sha256} → CAS.PushFile, CAS.PushArtifact // HEAD /api/cas/{sha256} → CAS.CheckArtifact // GET /api/cas/{sha256} → CAS.RetrieveArtifact // GET /api/cas/stats → CAS.Stats // // Reference: https://git.dcos.net/dcosnet/fester type FesterClient struct { baseURL string http *http.Client userAgent string CAS *cas.Client // shared content-addressable store } // NewFesterClient creates a client for a Fester master. The stack operates behind a firewall (OPNsense/IPFire); no transport-layer encryption is used. func NewFesterClient(baseURL string) *FesterClient { return &FesterClient{ baseURL: strings.TrimRight(baseURL, "/"), http: &http.Client{Timeout: 15 * time.Second}, userAgent: "sorcery-go/cluster (AGPL-3.0)", CAS: cas.NewClient(baseURL), } } // --------------------------------------------------------------------------- // Fester API types (mirror Fester's JSON schemas) // --------------------------------------------------------------------------- // festerNode is the JSON representation from GET /api/nodes. type festerNode struct { Name string `json:"name"` Host string `json:"host"` Arch string `json:"arch"` Runtime string `json:"runtime"` MaxJobs int `json:"max_jobs"` ActiveJobs int `json:"active_jobs"` CPU float64 `json:"cpu_load"` Memory float64 `json:"memory_load"` Temperature float64 `json:"heat"` Policy string `json:"policy"` Status string `json:"state"` LastProbe string `json:"last_seen"` ProbeError string `json:"probe_error"` Container string `json:"container,omitempty"` VM string `json:"vm,omitempty"` Firecracker *festerNodeFC `json:"firecracker,omitempty"` } // festerNodeFC is the firecracker config block from a node. type festerNodeFC struct { Kernel string `json:"kernel"` Rootfs string `json:"rootfs"` SSHPort int `json:"ssh_port"` SSHKey string `json:"ssh_key"` } // festerBuild is the JSON representation from GET /api/builds. type festerBuild struct { ID string `json:"id"` Target string `json:"target"` Node string `json:"node"` Status string `json:"status"` StartedAt string `json:"started_at"` EndedAt string `json:"ended_at"` ExitCode int `json:"exit_code"` } // festerMetrics is the JSON from GET /api/metrics/json. type festerMetrics struct { Nodes []festerNodeMetric `json:"nodes"` Builds festerBuildMetrics `json:"builds"` Timestamp string `json:"timestamp"` } type festerNodeMetric struct { Name string `json:"name"` CPU float64 `json:"cpu"` Memory float64 `json:"memory"` Load1 float64 `json:"load1"` Jobs int `json:"jobs"` MaxJobs int `json:"max_jobs"` } type festerBuildMetrics struct { Active int `json:"active"` Total int `json:"total"` Failed int `json:"failed"` Success int `json:"success"` Cached int `json:"cached"` } // festerBuildRequest is the POST body for /api/build. type festerBuildRequest struct { Target string `json:"target"` Node string `json:"node,omitempty"` // empty = let Fester pick Cmd string `json:"cmd"` Dir string `json:"dir"` Watch bool `json:"watch,omitempty"` } // festerCauseNode is the JSON from GET /api/cause/explain/{node}. type festerCauseNode struct { Node string `json:"node"` Reason string `json:"reason"` Children []string `json:"children"` Events []festerCauseEvent `json:"events"` } type festerCauseEvent struct { Timestamp string `json:"timestamp"` Type string `json:"type"` Message string `json:"message"` Node string `json:"node"` } // --------------------------------------------------------------------------- // Node management // --------------------------------------------------------------------------- // GetNodes fetches the full node list from Fester's /api/nodes and returns // them as Coven Node structs. Fields like Runtime, Container, and // Firecracker config are mapped for sorcery-go's runtime selection. func (fc *FesterClient) GetNodes() ([]*Node, error) { // Fester returns {"master": {...}, "nodes": [...]} var resp struct { Nodes []festerNode `json:"nodes"` } if err := fc.get("/api/nodes", &resp); err != nil { return nil, fmt.Errorf("fester: get nodes: %w", err) } nodes := make([]*Node, 0, len(resp.Nodes)) for _, fn := range resp.Nodes { n := &Node{ ID: fn.Name, Address: fn.Host, Arch: fn.Arch, ComputePower: fn.MaxJobs, CPU: fn.CPU, Memory: fn.Memory, ActiveBuilds: fn.ActiveJobs, Temperature: fn.Temperature, MaxJobs: fn.MaxJobs, Policy: fn.Policy, Status: fn.Status, } nodes = append(nodes, n) } return nodes, nil } // GetNodePulses converts Fester's /api/nodes response to Pulse structs // for the Cockpit Grid Heatmap. func (fc *FesterClient) GetNodePulses() ([]Pulse, error) { nodes, err := fc.GetNodes() if err != nil { return nil, err } pulses := make([]Pulse, 0, len(nodes)) for _, n := range nodes { pulses = append(pulses, Pulse{ NodeID: n.ID, CPU: n.CPU, Memory: n.Memory, ActiveBuilds: n.ActiveBuilds, Temperature: n.Temperature, MaxJobs: n.MaxJobs, Status: n.Status, }) } return pulses, nil } // SetNodePolicy sets a node's scheduling policy via POST /api/nodes/{name}/policy. func (fc *FesterClient) SetNodePolicy(nodeName, policy string) error { body := map[string]string{"policy": policy} return fc.post(fmt.Sprintf("/api/nodes/%s/policy", url.PathEscape(nodeName)), body, nil) } // ProbeNode triggers a manual probe of a node via POST /api/nodes/{name}/probe. func (fc *FesterClient) ProbeNode(nodeName string) error { return fc.post(fmt.Sprintf("/api/nodes/%s/probe", url.PathEscape(nodeName)), nil, nil) } // --------------------------------------------------------------------------- // Build management // --------------------------------------------------------------------------- // SubmitBuild submits a build to Fester via POST /api/build. func (fc *FesterClient) SubmitBuild(ctx context.Context, target, cmd, dir string, preferredNode string) (*festerBuild, error) { req := festerBuildRequest{ Target: target, Node: preferredNode, Cmd: cmd, Dir: dir, } var build festerBuild if err := fc.postWithContext(ctx, "/api/build", req, &build); err != nil { return nil, fmt.Errorf("fester: submit build: %w", err) } return &build, nil } // ListBuilds fetches the build history from Fester's /api/builds. func (fc *FesterClient) ListBuilds() ([]festerBuild, error) { var builds []festerBuild if err := fc.get("/api/builds", &builds); err != nil { return nil, fmt.Errorf("fester: list builds: %w", err) } return builds, nil } // CancelBuild cancels a running build via POST /api/builds/{id}/cancel. func (fc *FesterClient) CancelBuild(buildID string) error { return fc.post(fmt.Sprintf("/api/builds/%s/cancel", url.PathEscape(buildID)), nil, nil) } // --------------------------------------------------------------------------- // Observability // --------------------------------------------------------------------------- // GetMetrics fetches the full metrics snapshot from Fester's /api/metrics/json. func (fc *FesterClient) GetMetrics() (*festerMetrics, error) { var m festerMetrics if err := fc.get("/api/metrics/json", &m); err != nil { return nil, fmt.Errorf("fester: get metrics: %w", err) } return &m, nil } // RuntimeInfo describes a single runtime's availability. type RuntimeInfo struct { Available bool `json:"available"` Version string `json:"version,omitempty"` Note string `json:"note"` } // RuntimesResponse is the response from GET /api/nodes/runtimes. type RuntimesResponse struct { DefaultRuntime string `json:"default_runtime"` Runtimes map[string]RuntimeInfo `json:"runtimes"` } // GetRuntimes queries Fester's /api/nodes/runtimes endpoint to discover // which runtimes (host, lxc, podman, firecracker, libvirt, tmux) are // available. Sorcery-go uses this to align its own runtime selection // with what Fester can actually execute. func (fc *FesterClient) GetRuntimes() (*RuntimesResponse, error) { var resp RuntimesResponse if err := fc.get("/api/nodes/runtimes", &resp); err != nil { return nil, fmt.Errorf("fester: get runtimes: %w", err) } return &resp, nil } // CauseExplain fetches the causal explanation for a node from // Fester's /api/cause/explain/{node}. func (fc *FesterClient) CauseExplain(nodeName string) (*festerCauseNode, error) { var cause festerCauseNode if err := fc.get(fmt.Sprintf("/api/cause/explain/%s", url.PathEscape(nodeName)), &cause); err != nil { return nil, fmt.Errorf("fester: cause explain %s: %w", nodeName, err) } return &cause, nil } // --------------------------------------------------------------------------- // Event streaming (WebSocket) // --------------------------------------------------------------------------- // FesterEvent represents a real-time event from Fester's WebSocket stream. // This mirrors Fester's event schema from its EventBus. type FesterEvent struct { Type string `json:"type"` // build_started, build_completed, node_probe, scheduler_decision, etc. Timestamp time.Time `json:"timestamp"` Node string `json:"node"` BuildID string `json:"build_id"` Message string `json:"message"` Data json.RawMessage `json:"data,omitempty"` } // WatchEvents connects to Fester's WebSocket at /ws and streams events. // The handler is called for each event. Blocks until ctx is cancelled or // an error occurs. // // This uses raw HTTP upgrade since we need to support the same WebSocket // protocol as Fester's frontend. The gorilla/websocket dependency is // already in go.mod for the Cockpit. func (fc *FesterClient) WatchEvents(ctx context.Context, handler func(FesterEvent)) error { wsURL := fc.baseURL wsURL = strings.Replace(wsURL, "http://", "ws://", 1) wsURL = strings.Replace(wsURL, "https://", "wss://", 1) wsURL += "/ws" wsDialer := &websocket.Dialer{} wsConn, _, err := wsDialer.DialContext(ctx, wsURL, nil) if err != nil { return fmt.Errorf("fester: ws connect: %w", err) } defer wsConn.Close() for { select { case <-ctx.Done(): return ctx.Err() default: } _, msg, err := wsConn.ReadMessage() if err != nil { if websocket.IsCloseError(err, websocket.CloseNormalClosure) { return nil } return fmt.Errorf("fester: ws read: %w", err) } var event FesterEvent if err := json.Unmarshal(msg, &event); err != nil { // Skip malformed events. continue } handler(event) } } // --------------------------------------------------------------------------- // HTTP helpers // --------------------------------------------------------------------------- func (fc *FesterClient) get(path string, v interface{}) error { req, err := http.NewRequest("GET", fc.baseURL+path, nil) if err != nil { return err } req.Header.Set("User-Agent", fc.userAgent) req.Header.Set("Accept", "application/json") resp, err := fc.http.Do(req) if err != nil { return err } defer resp.Body.Close() if resp.StatusCode != http.StatusOK { body, _ := io.ReadAll(io.LimitReader(resp.Body, 1024)) return fmt.Errorf("fester: %s %s → %d: %s", "GET", path, resp.StatusCode, string(body)) } if v != nil { return json.NewDecoder(resp.Body).Decode(v) } return nil } func (fc *FesterClient) post(path string, body interface{}, v interface{}) error { return fc.postWithContext(context.Background(), path, body, v) } func (fc *FesterClient) postWithContext(ctx context.Context, path string, body interface{}, v interface{}) error { var reqBody io.Reader if body != nil { data, err := json.Marshal(body) if err != nil { return fmt.Errorf("fester: marshal body: %w", err) } reqBody = bytes.NewReader(data) } req, err := http.NewRequestWithContext(ctx, "POST", fc.baseURL+path, reqBody) if err != nil { return err } req.Header.Set("User-Agent", fc.userAgent) req.Header.Set("Content-Type", "application/json") req.Header.Set("Accept", "application/json") resp, err := fc.http.Do(req) if err != nil { return err } defer resp.Body.Close() if resp.StatusCode >= 400 { respBody, _ := io.ReadAll(io.LimitReader(resp.Body, 1024)) return fmt.Errorf("fester: POST %s → %d: %s", path, resp.StatusCode, string(respBody)) } if v != nil { return json.NewDecoder(resp.Body).Decode(v) } return nil }