Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 5 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down
57 changes: 50 additions & 7 deletions pkg/edge/reporter.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,8 @@ import (
"fmt"
"log"
"net/http"
"strings"
"sync"
"time"
)

Expand All @@ -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},
Expand All @@ -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
Expand All @@ -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
Expand All @@ -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,
Expand All @@ -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,
Expand Down Expand Up @@ -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)
Expand All @@ -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
}
}
130 changes: 130 additions & 0 deletions pkg/edge/reporter_test.go
Original file line number Diff line number Diff line change
@@ -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()
}