refactor(agent-pod): merge hub into pi container, drop kubectl-exec spawn step

agent-console complained about the two-step workflow: kubectl exec a
script to spawn an agent, separate terminal to watch it. Root cause
was hub being a separate Go sidecar with no access to the pi binary
or its config, so it could only watch, never trigger.

hub.js now runs inside the pi container itself (same filesystem, same
PATH) and exposes POST /run, which spawns `pi -p --mode json` directly
and streams it over the same /console WebSocket everyone else is
already watching. agent-console gained a `run <agent> <prompt>`
subcommand that POSTs then watches -- no separate exec step, no
separate terminal. Drops the golang:1.25-alpine sidecar container and
agent-run-configmap.yaml entirely (both superseded).

ConfigMaps regenerated via local-harness/hub/sync-to-homelab.sh
instead of hand-copied, so hub.js/config.json/settings.json each have
exactly one source of truth.
This commit is contained in:
Story Crater Bot
2026-08-16 10:24:03 -07:00
parent 83efe8ef9e
commit 12b260f35a
5 changed files with 128 additions and 314 deletions
+114 -204
View File
@@ -1,209 +1,119 @@
apiVersion: v1
data:
hub.js: |
#!/usr/bin/env node
// agent-hub: lives inside the pi container (not a sidecar) so it can spawn
// `pi -p --mode json` directly. One persistent in-cluster service --
// POST /run to trigger an agent, GET /console (WebSocket) to watch every
// concurrent run live as it happens, relaying pi's own session protocol
// verbatim (same event shape Claude Code sessions use).
const http = require("node:http");
const crypto = require("node:crypto");
const { spawn } = require("node:child_process");
const readline = require("node:readline");
const { WebSocketServer } = require("ws");
const PORT = process.env.HUB_PORT || 9090;
const sessions = new Map(); // id -> {id, agent, status, events, startedAt, endedAt}
const viewers = new Set(); // WebSocket connections watching /console
function broadcast(type, session) {
const msg = JSON.stringify({ type, session });
for (const ws of viewers) {
if (ws.readyState === ws.OPEN) ws.send(msg);
}
}
function startSession(agent) {
const id = crypto.randomUUID();
const session = { id, agent, status: "running", events: [], startedAt: new Date().toISOString() };
sessions.set(id, session);
broadcast("start", session);
return session;
}
function addEvent(session, rawLine) {
session.events.push(JSON.parse(rawLine));
broadcast("event", session);
}
function endSession(session, status) {
session.status = status;
session.endedAt = new Date().toISOString();
broadcast("end", session);
}
// Spawns `pi -p --mode json <prompt>` and relays it live. Same shape as the
// old agent-run.js wrapper, just in-process instead of a separate exec.
function runAgent(agent, prompt, extraArgs = []) {
const session = startSession(agent);
const child = spawn("pi", ["-p", "--mode", "json", ...extraArgs, prompt], {
stdio: ["ignore", "pipe", "pipe"],
});
const rl = readline.createInterface({ input: child.stdout });
rl.on("line", (line) => {
if (!line.trim()) return;
try {
addEvent(session, line);
} catch {
// non-JSON stdout noise, ignore
}
});
child.stderr.on("data", (chunk) => process.stderr.write(chunk));
child.on("close", (code) => endSession(session, code === 0 ? "done" : "error"));
return session;
}
const server = http.createServer((req, res) => {
const url = new URL(req.url, "http://localhost");
if (url.pathname === "/healthz") {
res.writeHead(200).end();
return;
}
if (url.pathname === "/sessions" && req.method === "GET") {
res.writeHead(200, { "Content-Type": "application/json" });
res.end(JSON.stringify([...sessions.values()]));
return;
}
if (url.pathname === "/run" && req.method === "POST") {
let body = "";
req.on("data", (chunk) => (body += chunk));
req.on("end", () => {
try {
const { agent, prompt, provider, model } = JSON.parse(body);
if (!agent || !prompt) throw new Error("agent and prompt are required");
const extraArgs = [];
if (provider) extraArgs.push("--provider", provider);
if (model) extraArgs.push("--model", model);
const session = runAgent(agent, prompt, extraArgs);
res.writeHead(200, { "Content-Type": "application/json" });
res.end(JSON.stringify({ id: session.id }));
} catch (err) {
res.writeHead(400, { "Content-Type": "application/json" });
res.end(JSON.stringify({ error: err.message }));
}
});
return;
}
res.writeHead(404).end();
});
const wss = new WebSocketServer({ server, path: "/console" });
wss.on("connection", (ws) => {
for (const session of sessions.values()) {
ws.send(JSON.stringify({ type: "snapshot", session }));
}
viewers.add(ws);
ws.on("close", () => viewers.delete(ws));
});
server.listen(PORT, () => console.log(`agent-hub listening on :${PORT}`));
kind: ConfigMap
metadata:
name: hub-src
namespace: agent-pod
data:
go.mod: |
module agent-hub
go 1.25.0
require github.com/gorilla/websocket v1.5.3
go.sum: |
github.com/gorilla/websocket v1.5.3 h1:saDtZ6Pbx/0u+bgYQ3q96pZgCzfhKXGPqt7kZ72aNNg=
github.com/gorilla/websocket v1.5.3/go.mod h1:YR8l580nyteQvAITg2hZ9XVh4b55+EU/adAjf1fMHhE=
main.go: |
// agent-hub sits inside agent-pod alongside pi. Spawned agent runs report
// their lifecycle over localhost HTTP, relaying pi's own `--mode json`
// session protocol (session/agent_start/turn_start/message_start/
// message_update/message_end/turn_end/agent_end -- the same shape Claude
// Code sessions use) verbatim. agent-console (outside the pod, via kubectl
// port-forward) watches all concurrent agents live over one WebSocket and
// can render each session as a real transcript, not a log tail.
package main
import (
"encoding/json"
"log"
"net/http"
"os"
"sync"
"time"
"github.com/gorilla/websocket"
)
type Session struct {
ID string `json:"id"`
Agent string `json:"agent"`
Status string `json:"status"` // running | done | error
Events []json.RawMessage `json:"events"` // raw pi --mode json events, in order
StartedAt time.Time `json:"startedAt"`
EndedAt time.Time `json:"endedAt,omitempty"`
}
type Notification struct {
Type string `json:"type"` // snapshot | start | event | end
Session Session `json:"session"`
}
type Hub struct {
mu sync.Mutex
sessions map[string]*Session
viewers map[*websocket.Conn]struct{}
}
func newHub() *Hub {
return &Hub{
sessions: make(map[string]*Session),
viewers: make(map[*websocket.Conn]struct{}),
}
}
func (h *Hub) broadcast(n Notification) {
h.mu.Lock()
defer h.mu.Unlock()
for conn := range h.viewers {
if err := conn.WriteJSON(n); err != nil {
conn.Close()
delete(h.viewers, conn)
}
}
}
func (h *Hub) start(id, agent string) {
h.mu.Lock()
s := &Session{ID: id, Agent: agent, Status: "running", StartedAt: time.Now()}
h.sessions[id] = s
h.mu.Unlock()
h.broadcast(Notification{Type: "start", Session: *s})
}
func (h *Hub) event(id string, raw json.RawMessage) {
h.mu.Lock()
s, ok := h.sessions[id]
if !ok {
h.mu.Unlock()
return
}
s.Events = append(s.Events, raw)
snapshot := *s
h.mu.Unlock()
h.broadcast(Notification{Type: "event", Session: snapshot})
}
func (h *Hub) end(id, status string) {
h.mu.Lock()
s, ok := h.sessions[id]
if !ok {
h.mu.Unlock()
return
}
s.Status = status
s.EndedAt = time.Now()
snapshot := *s
h.mu.Unlock()
h.broadcast(Notification{Type: "end", Session: snapshot})
}
func (h *Hub) snapshot() []Session {
h.mu.Lock()
defer h.mu.Unlock()
out := make([]Session, 0, len(h.sessions))
for _, s := range h.sessions {
out = append(out, *s)
}
return out
}
var upgrader = websocket.Upgrader{
CheckOrigin: func(r *http.Request) bool { return true },
}
func main() {
addr := os.Getenv("HUB_ADDR")
if addr == "" {
addr = ":9090"
}
h := newHub()
http.HandleFunc("/healthz", func(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusOK)
})
// Plain JSON snapshot for curl/debugging -- the WebSocket below is the
// live path agent-console actually uses.
http.HandleFunc("/sessions", func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
json.NewEncoder(w).Encode(h.snapshot())
})
// Agent side: fire-and-forget POSTs, no persistent connection needed
// since each pi run is short-lived.
http.HandleFunc("/agent/start", func(w http.ResponseWriter, r *http.Request) {
var body struct{ ID, Agent string }
if err := json.NewDecoder(r.Body).Decode(&body); err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
h.start(body.ID, body.Agent)
w.WriteHeader(http.StatusOK)
})
http.HandleFunc("/agent/event", func(w http.ResponseWriter, r *http.Request) {
var body struct {
ID string
Event json.RawMessage
}
if err := json.NewDecoder(r.Body).Decode(&body); err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
h.event(body.ID, body.Event)
w.WriteHeader(http.StatusOK)
})
http.HandleFunc("/agent/end", func(w http.ResponseWriter, r *http.Request) {
var body struct{ ID, Status string }
if err := json.NewDecoder(r.Body).Decode(&body); err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
h.end(body.ID, body.Status)
w.WriteHeader(http.StatusOK)
})
// Console side: connect, get a snapshot of everything running, then
// stream live start/event/end notifications as they happen.
http.HandleFunc("/console", func(w http.ResponseWriter, r *http.Request) {
conn, err := upgrader.Upgrade(w, r, nil)
if err != nil {
return
}
h.mu.Lock()
for _, s := range h.sessions {
conn.WriteJSON(Notification{Type: "snapshot", Session: *s})
}
h.viewers[conn] = struct{}{}
h.mu.Unlock()
// Drain reads to detect disconnect; console never sends anything.
go func() {
defer func() {
h.mu.Lock()
delete(h.viewers, conn)
h.mu.Unlock()
conn.Close()
}()
for {
if _, _, err := conn.ReadMessage(); err != nil {
return
}
}
}()
})
log.Printf("agent-hub listening on %s", addr)
log.Fatal(http.ListenAndServe(addr, nil))
}