diff --git a/README.md b/README.md index 26e1290..d259f07 100644 --- a/README.md +++ b/README.md @@ -281,7 +281,11 @@ When enabled, PicoClaw exposes: - `GET /api/v1/status` — Node status - `POST /api/v1/command` — Receive commands from fleet -And periodically sends heartbeats (including gene stats) to the configured fleet manager. +The edge reporter sends: +- registration payloads to `/fleet/register` on startup +- heartbeats to `/fleet/heartbeat` on the configured interval, including gene stats when available +- status snapshots to `/fleet/status` +- event reports to `/fleet/events` via `ReportEvent` ## Skills (6 built-in) diff --git a/pkg/edge/reporter.go b/pkg/edge/reporter.go index 48f06a8..627cf8c 100644 --- a/pkg/edge/reporter.go +++ b/pkg/edge/reporter.go @@ -6,6 +6,8 @@ import ( "fmt" "log" "net/http" + "strings" + "sync" "time" ) @@ -19,14 +21,19 @@ type GeneStatsProvider interface { // Reporter periodically sends heartbeat and event reports // to the upstream Fleet Manager (L2 NanoClaw or L3 MoltClaw). type Reporter struct { - config Config - client *http.Client - stopCh chan struct{} - geneProvider GeneStatsProvider + config Config + client *http.Client + stopCh chan struct{} + stopOnce sync.Once + geneProvider GeneStatsProvider } // NewReporter creates a new Edge Reporter. func NewReporter(cfg Config) *Reporter { + if cfg.Cloud.HeartbeatInterval <= 0 { + cfg.Cloud.HeartbeatInterval = DefaultConfig().Cloud.HeartbeatInterval + } + return &Reporter{ config: cfg, client: &http.Client{Timeout: 10 * time.Second}, @@ -44,11 +51,14 @@ func (r *Reporter) StartHeartbeat() { // Send initial registration r.sendRegistration() + r.sendHeartbeat() + r.sendStatus("online", nil) for { select { case <-ticker.C: r.sendHeartbeat() + r.sendStatus("online", nil) case <-r.stopCh: log.Println("[edge] heartbeat reporter stopped") return @@ -58,7 +68,9 @@ func (r *Reporter) StartHeartbeat() { // Stop halts the heartbeat reporter. func (r *Reporter) Stop() { - close(r.stopCh) + r.stopOnce.Do(func() { + close(r.stopCh) + }) } // SetGeneProvider attaches a gene stats provider to include @@ -80,6 +92,10 @@ func (r *Reporter) PublishGene(geneData map[string]interface{}) error { // ReportEvent sends a one-off event to the Fleet Manager. func (r *Reporter) ReportEvent(eventType string, data map[string]any) error { + if data == nil { + data = map[string]any{} + } + payload := map[string]any{ "node_id": r.config.NodeID, "type": eventType, @@ -89,6 +105,18 @@ func (r *Reporter) ReportEvent(eventType string, data map[string]any) error { return r.post("/fleet/events", payload) } +// ReportStatus sends the current node status to the Fleet Manager. +func (r *Reporter) ReportStatus(status string, data map[string]any) error { + if status == "" { + status = "online" + } + if data == nil { + data = map[string]any{} + } + + return r.sendStatus(status, data) +} + func (r *Reporter) sendRegistration() { payload := map[string]any{ "node_id": r.config.NodeID, @@ -121,13 +149,28 @@ func (r *Reporter) sendHeartbeat() { } } +func (r *Reporter) sendStatus(status string, data map[string]any) error { + if data == nil { + data = map[string]any{} + } + + payload := map[string]any{ + "node_id": r.config.NodeID, + "status": status, + "data": data, + "timestamp": time.Now().UTC().Format(time.RFC3339), + } + + return r.post("/fleet/status", payload) +} + func (r *Reporter) post(path string, payload any) error { body, err := json.Marshal(payload) if err != nil { return fmt.Errorf("marshal: %w", err) } - url := r.config.Cloud.Endpoint + path + url := strings.TrimRight(r.config.Cloud.Endpoint, "/") + path req, err := http.NewRequest("POST", url, bytes.NewReader(body)) if err != nil { return fmt.Errorf("request: %w", err) @@ -147,4 +190,4 @@ func (r *Reporter) post(path string, payload any) error { return fmt.Errorf("post %s: status %d", path, resp.StatusCode) } return nil -} \ No newline at end of file +} diff --git a/pkg/edge/reporter_test.go b/pkg/edge/reporter_test.go new file mode 100644 index 0000000..d5adef6 --- /dev/null +++ b/pkg/edge/reporter_test.go @@ -0,0 +1,130 @@ +package edge + +import ( + "encoding/json" + "net/http" + "net/http/httptest" + "testing" + "time" +) + +type reporterRequest struct { + Path string + Authorization string + Payload map[string]any +} + +type fakeGeneStats struct{} + +func (fakeGeneStats) GetStats() map[string]interface{} { + return map[string]interface{}{ + "genes": float64(3), + } +} + +func (fakeGeneStats) GetHighConfidenceGenes(minConfidence float64, minVerifiedBy int) []map[string]interface{} { + return []map[string]interface{}{} +} + +func TestReporterSendsRegistrationHeartbeatStatusAndEvents(t *testing.T) { + var requests []reporterRequest + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodPost { + t.Fatalf("method = %s, want POST", r.Method) + } + + var payload map[string]any + if err := json.NewDecoder(r.Body).Decode(&payload); err != nil { + t.Fatalf("decode payload: %v", err) + } + requests = append(requests, reporterRequest{ + Path: r.URL.Path, + Authorization: r.Header.Get("Authorization"), + Payload: payload, + }) + w.WriteHeader(http.StatusAccepted) + })) + defer server.Close() + + reporter := NewReporter(Config{ + NodeID: "edge-01", + NodeName: "Pond Edge 01", + Cloud: CloudConfig{ + Endpoint: server.URL + "/", + HeartbeatInterval: time.Second, + Token: "fleet-token", + }, + }) + reporter.SetGeneProvider(fakeGeneStats{}) + + reporter.sendRegistration() + reporter.sendHeartbeat() + if err := reporter.ReportStatus("degraded", map[string]any{"reason": "low_battery"}); err != nil { + t.Fatalf("ReportStatus: %v", err) + } + if err := reporter.ReportEvent("sensor_alert", map[string]any{"sensor": "do"}); err != nil { + t.Fatalf("ReportEvent: %v", err) + } + + wantPaths := []string{"/fleet/register", "/fleet/heartbeat", "/fleet/status", "/fleet/events"} + if len(requests) != len(wantPaths) { + t.Fatalf("requests = %d, want %d", len(requests), len(wantPaths)) + } + for i, want := range wantPaths { + if requests[i].Path != want { + t.Fatalf("request %d path = %s, want %s", i, requests[i].Path, want) + } + if requests[i].Authorization != "Bearer fleet-token" { + t.Fatalf("request %d auth = %q, want bearer token", i, requests[i].Authorization) + } + if requests[i].Payload["node_id"] != "edge-01" { + t.Fatalf("request %d node_id = %v, want edge-01", i, requests[i].Payload["node_id"]) + } + if _, ok := requests[i].Payload["timestamp"]; i > 0 && !ok { + t.Fatalf("request %d missing timestamp", i) + } + } + + if requests[0].Payload["node_name"] != "Pond Edge 01" { + t.Fatalf("registration node_name = %v, want Pond Edge 01", requests[0].Payload["node_name"]) + } + if _, ok := requests[1].Payload["gene_stats"]; !ok { + t.Fatal("heartbeat missing gene_stats") + } + if requests[2].Payload["status"] != "degraded" { + t.Fatalf("status = %v, want degraded", requests[2].Payload["status"]) + } + if requests[3].Payload["type"] != "sensor_alert" { + t.Fatalf("event type = %v, want sensor_alert", requests[3].Payload["type"]) + } +} + +func TestReporterReturnsHTTPError(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + http.Error(w, "nope", http.StatusBadGateway) + })) + defer server.Close() + + reporter := NewReporter(Config{ + NodeID: "edge-01", + Cloud: CloudConfig{ + Endpoint: server.URL, + HeartbeatInterval: time.Second, + }, + }) + + if err := reporter.ReportEvent("sensor_alert", nil); err == nil { + t.Fatal("ReportEvent error = nil, want HTTP status error") + } +} + +func TestReporterUsesDefaultHeartbeatIntervalAndIdempotentStop(t *testing.T) { + reporter := NewReporter(Config{}) + + if reporter.config.Cloud.HeartbeatInterval != DefaultConfig().Cloud.HeartbeatInterval { + t.Fatalf("interval = %s, want default %s", reporter.config.Cloud.HeartbeatInterval, DefaultConfig().Cloud.HeartbeatInterval) + } + + reporter.Stop() + reporter.Stop() +}