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
9 changes: 8 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -281,7 +281,14 @@ 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.
And periodically reports to the configured fleet manager:
- `POST /fleet/register` - node id, name, capabilities, version, and timestamp
- `POST /fleet/heartbeat` - online status plus gene stats when the Gene Evolution engine is enabled
- `POST /fleet/status` - structured status snapshots for operational state
- `POST /fleet/events` - one-off edge events
- `POST /fleet/genes/publish` - high-confidence gene sharing

If `cloud_token` is set, reporter requests include `Authorization: Bearer <token>`.

## Skills (6 built-in)

Expand Down
95 changes: 72 additions & 23 deletions pkg/edge/reporter.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,8 +4,11 @@ import (
"bytes"
"encoding/json"
"fmt"
"io"
"log"
"net/http"
"strings"
"sync"
"time"
)

Expand All @@ -19,10 +22,12 @@ 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
now func() time.Time
}

// NewReporter creates a new Edge Reporter.
Expand All @@ -31,24 +36,43 @@ func NewReporter(cfg Config) *Reporter {
config: cfg,
client: &http.Client{Timeout: 10 * time.Second},
stopCh: make(chan struct{}),
now: time.Now,
}
}

// SetHTTPClient overrides the reporter HTTP client. It is primarily useful for
// tests and embedders that need custom transport behavior.
func (r *Reporter) SetHTTPClient(client *http.Client) {
if client != nil {
r.client = client
}
}

// StartHeartbeat begins the periodic heartbeat loop.
func (r *Reporter) StartHeartbeat() {
ticker := time.NewTicker(r.config.Cloud.HeartbeatInterval)
interval := r.config.Cloud.HeartbeatInterval
if interval <= 0 {
interval = 30 * time.Second
}
ticker := time.NewTicker(interval)
defer ticker.Stop()

log.Printf("[edge] heartbeat reporter started (interval: %s, endpoint: %s)",
r.config.Cloud.HeartbeatInterval, r.config.Cloud.Endpoint)
interval, r.config.Cloud.Endpoint)

// Send initial registration
r.sendRegistration()
if err := r.Register(); err != nil {
log.Printf("[edge] registration failed: %v", err)
} else {
log.Printf("[edge] registered with Fleet Manager as %s", r.config.NodeID)
}

for {
select {
case <-ticker.C:
r.sendHeartbeat()
if err := r.SendHeartbeat(); err != nil {
log.Printf("[edge] heartbeat failed: %v", err)
}
case <-r.stopCh:
log.Println("[edge] heartbeat reporter stopped")
return
Expand All @@ -58,7 +82,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 @@ -73,52 +99,67 @@ func (r *Reporter) PublishGene(geneData map[string]interface{}) error {
payload := map[string]any{
"node_id": r.config.NodeID,
"gene": geneData,
"timestamp": time.Now().UTC().Format(time.RFC3339),
"timestamp": r.timestamp(),
}
return r.post("/fleet/genes/publish", payload)
}

// ReportEvent sends a one-off event to the Fleet Manager.
func (r *Reporter) ReportEvent(eventType string, data map[string]any) error {
if strings.TrimSpace(eventType) == "" {
return fmt.Errorf("event type is required")
}
payload := map[string]any{
"node_id": r.config.NodeID,
"type": eventType,
"data": data,
"timestamp": time.Now().UTC().Format(time.RFC3339),
"timestamp": r.timestamp(),
}
return r.post("/fleet/events", payload)
}

func (r *Reporter) sendRegistration() {
// Register sends an initial node registration report to the Fleet Manager.
func (r *Reporter) Register() error {
payload := map[string]any{
"node_id": r.config.NodeID,
"node_name": r.config.NodeName,
"type": "picclaw",
"capabilities": []string{"sensor", "exec", "cron", "alert"},
"version": "0.1.0",
"timestamp": r.timestamp(),
}
if err := r.post("/fleet/register", payload); err != nil {
log.Printf("[edge] registration failed: %v", err)
} else {
log.Printf("[edge] registered with Fleet Manager as %s", r.config.NodeID)
}
return r.post("/fleet/register", payload)
}

func (r *Reporter) sendHeartbeat() {
// SendHeartbeat reports that this edge node is online.
func (r *Reporter) SendHeartbeat() error {
payload := map[string]any{
"node_id": r.config.NodeID,
"status": "online",
"timestamp": time.Now().UTC().Format(time.RFC3339),
"timestamp": r.timestamp(),
}

// Include gene evolution stats if provider is available
if r.geneProvider != nil {
payload["gene_stats"] = r.geneProvider.GetStats()
}

if err := r.post("/fleet/heartbeat", payload); err != nil {
log.Printf("[edge] heartbeat failed: %v", err)
return r.post("/fleet/heartbeat", payload)
}

// ReportStatus sends a structured status snapshot to the Fleet Manager.
func (r *Reporter) ReportStatus(status string, data map[string]any) error {
if strings.TrimSpace(status) == "" {
status = "unknown"
}
payload := map[string]any{
"node_id": r.config.NodeID,
"node_name": r.config.NodeName,
"status": status,
"data": data,
"timestamp": r.timestamp(),
}
return r.post("/fleet/status", payload)
}

func (r *Reporter) post(path string, payload any) error {
Expand All @@ -127,7 +168,7 @@ func (r *Reporter) post(path string, payload any) error {
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 @@ -144,7 +185,15 @@ func (r *Reporter) post(path string, payload any) error {
defer resp.Body.Close()

if resp.StatusCode >= 300 {
detail, _ := io.ReadAll(io.LimitReader(resp.Body, 512))
if len(bytes.TrimSpace(detail)) > 0 {
return fmt.Errorf("post %s: status %d: %s", path, resp.StatusCode, strings.TrimSpace(string(detail)))
}
return fmt.Errorf("post %s: status %d", path, resp.StatusCode)
}
return nil
}
}

func (r *Reporter) timestamp() string {
return r.now().UTC().Format(time.RFC3339)
}
Loading