From cfa1b0b4b76b82295ca19756b99572af4b8738c9 Mon Sep 17 00:00:00 2001 From: Kareem Elbahrawy Date: Mon, 20 Apr 2026 12:18:04 +0000 Subject: [PATCH] adapter: claude registry, hub ingest/handler, signal raw MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Wire the first real adapter end-to-end. POSTing a Claude hook body to /hook/claude now produces an Event on Hub.Events(). Core - Signal.Raw + Event.Raw plumbing (decision.go, hub.go, design.md). - Adapter struct, RegisterAdapter, Adapters, lookupAdapter, InstallConfig, InstallResult, AllAgents, ErrUnknownAgent (adapter.go). - Hub.Ingest, Hub.Handler (POST /hook/{agent}, 202/400/404/405/413/500), Hub.ServeHTTP, 1 MiB body cap (hub.go). Claude adapter - Self-registers from init(); MapHookEvent covers every row in design.md's Claude table. Install/uninstall stubbed pending next ticket. - Subagent identity follows the official schema: SessionID = agent_id, ParentSessionID = session_id (Claude fires Subagent* under the parent). Tests - adapters/claude/map_test.go: per-row mapping + timestamp + tolerance. - adapter_test.go: registry duplicate / empty / sorted snapshot / concurrent. - hub_ingest_test.go: E2E HTTP — happy path, 404/400/413/405, unknown event drop, PreCompact drop, subagent flow, 1000-concurrent ingest. - testdata/claude/*.json fixtures mirror real Claude payload shape. --- adapter.go | 101 +++++++- adapter_test.go | 128 ++++++++++ adapters/claude/adapter.go | 57 +++++ adapters/claude/map.go | 119 +++++++++ adapters/claude/map_test.go | 150 +++++++++++ decision.go | 6 + hub.go | 101 ++++++++ hub_ingest_test.go | 280 +++++++++++++++++++++ specs/design.md | 1 + testdata/claude/notification.json | 9 + testdata/claude/permission_request.json | 9 + testdata/claude/post_tool_use.json | 10 + testdata/claude/post_tool_use_failure.json | 12 + testdata/claude/pre_compact.json | 5 + testdata/claude/pre_tool_use_read.json | 9 + testdata/claude/session_end.json | 5 + testdata/claude/session_start.json | 8 + testdata/claude/stop.json | 5 + testdata/claude/subagent_start.json | 8 + testdata/claude/subagent_stop.json | 11 + testdata/claude/unknown_event.json | 5 + testdata/claude/user_prompt_submit.json | 5 + 22 files changed, 1042 insertions(+), 2 deletions(-) create mode 100644 adapter_test.go create mode 100644 adapters/claude/adapter.go create mode 100644 adapters/claude/map.go create mode 100644 adapters/claude/map_test.go create mode 100644 hub_ingest_test.go create mode 100644 testdata/claude/notification.json create mode 100644 testdata/claude/permission_request.json create mode 100644 testdata/claude/post_tool_use.json create mode 100644 testdata/claude/post_tool_use_failure.json create mode 100644 testdata/claude/pre_compact.json create mode 100644 testdata/claude/pre_tool_use_read.json create mode 100644 testdata/claude/session_end.json create mode 100644 testdata/claude/session_start.json create mode 100644 testdata/claude/stop.json create mode 100644 testdata/claude/subagent_start.json create mode 100644 testdata/claude/subagent_stop.json create mode 100644 testdata/claude/unknown_event.json create mode 100644 testdata/claude/user_prompt_submit.json diff --git a/adapter.go b/adapter.go index 6d6c4c4..4810fe7 100644 --- a/adapter.go +++ b/adapter.go @@ -1,4 +1,101 @@ package agentstatus -// Adapter type, RegisterAdapter, and the adapter registry will live here. -// Intentionally empty in the scaffolding ticket — see specs/design.md. +import ( + "errors" + "fmt" + "sort" + "sync" +) + +// Adapter is the per-agent extension point. Built-in adapters live under +// adapters/ and self-register from init(); third parties register the +// same way. +// +// MapHookEvent is required: it translates a single native hook payload into a +// Signal (or returns (nil, nil) to silently drop, e.g. for unknown event +// names or metadata-only events). +// +// InstallHooks and UninstallHooks may be nil if an adapter does not yet +// implement them; the orchestrator (deferred to a later ticket) treats nil as +// "skipped — not implemented". +type Adapter struct { + Name Agent + MapHookEvent func(event string, payload map[string]any) (*Signal, error) + InstallHooks func(cfg InstallConfig) (InstallResult, error) + UninstallHooks func(cfg InstallConfig) (InstallResult, error) +} + +// InstallConfig parameterizes hook installation. See specs/design.md. +type InstallConfig struct { + // Endpoint is the base URL the bridge POSTs hook payloads to (e.g. + // "http://localhost:9090/hook"). Adapters append / as needed. + Endpoint string + // Agents narrows install to a subset; empty means all registered agents. + Agents []Agent + // Project, when non-empty, targets a project-level config file instead of + // the user-level default. + Project string +} + +// InstallResult is one adapter's outcome from an install or uninstall pass. +type InstallResult struct { + Agent Agent + Installed bool + Skipped bool + Reason string + Path string +} + +// AllAgents enumerates the built-in agent identifiers, in a stable order +// suitable for InstallConfig.Agents. +var AllAgents = []Agent{Claude, Codex, OpenCode} + +var ( + registryMu sync.RWMutex + registry = map[Agent]Adapter{} +) + +// ErrUnknownAgent is returned by Hub.Ingest when no adapter is registered +// under the given Agent name. The HTTP handler maps this to 404. +var ErrUnknownAgent = errors.New("agentstatus: unknown agent") + +// RegisterAdapter adds an adapter to the package-level registry. It is +// goroutine-safe and intended to be called from init() in adapter +// subpackages. Returns an error if the name is empty or already registered. +func RegisterAdapter(a Adapter) error { + if a.Name == "" { + return errors.New("agentstatus: adapter name is empty") + } + if a.MapHookEvent == nil { + return fmt.Errorf("agentstatus: adapter %q has nil MapHookEvent", a.Name) + } + registryMu.Lock() + defer registryMu.Unlock() + if _, ok := registry[a.Name]; ok { + return fmt.Errorf("agentstatus: adapter %q already registered", a.Name) + } + registry[a.Name] = a + return nil +} + +// Adapters returns a snapshot of registered adapters, sorted by Name. The +// returned slice is owned by the caller; mutating it does not affect the +// registry. +func Adapters() []Adapter { + registryMu.RLock() + out := make([]Adapter, 0, len(registry)) + for _, a := range registry { + out = append(out, a) + } + registryMu.RUnlock() + sort.Slice(out, func(i, j int) bool { return out[i].Name < out[j].Name }) + return out +} + +// lookupAdapter is the internal accessor used by Hub.Ingest. +func lookupAdapter(name Agent) (Adapter, bool) { + registryMu.RLock() + defer registryMu.RUnlock() + a, ok := registry[name] + return a, ok +} diff --git a/adapter_test.go b/adapter_test.go new file mode 100644 index 0000000..957adbc --- /dev/null +++ b/adapter_test.go @@ -0,0 +1,128 @@ +package agentstatus + +import ( + "errors" + "strings" + "sync" + "testing" +) + +// withCleanRegistry runs fn with a freshly-empty registry, restoring whatever +// adapters were registered (e.g. by init() in blank-imported subpackages) +// when fn returns. Tests share the package-level registry, so this isolation +// is required. +func withCleanRegistry(t *testing.T, fn func()) { + t.Helper() + registryMu.Lock() + saved := registry + registry = map[Agent]Adapter{} + registryMu.Unlock() + + defer func() { + registryMu.Lock() + registry = saved + registryMu.Unlock() + }() + fn() +} + +func okMap(string, map[string]any) (*Signal, error) { return nil, nil } + +func TestRegisterAdapter_Duplicate(t *testing.T) { + withCleanRegistry(t, func() { + a := Adapter{Name: "fake", MapHookEvent: okMap} + if err := RegisterAdapter(a); err != nil { + t.Fatalf("first register: %v", err) + } + err := RegisterAdapter(a) + if err == nil { + t.Fatal("expected duplicate error") + } + if !strings.Contains(err.Error(), "already registered") { + t.Errorf("error text: %v", err) + } + }) +} + +func TestRegisterAdapter_EmptyName(t *testing.T) { + withCleanRegistry(t, func() { + err := RegisterAdapter(Adapter{MapHookEvent: okMap}) + if err == nil { + t.Fatal("expected empty-name error") + } + if !strings.Contains(err.Error(), "empty") { + t.Errorf("error text: %v", err) + } + }) +} + +func TestRegisterAdapter_NilMap(t *testing.T) { + withCleanRegistry(t, func() { + err := RegisterAdapter(Adapter{Name: "fake"}) + if err == nil { + t.Fatal("expected nil-map error") + } + }) +} + +func TestAdapters_SortedSnapshot(t *testing.T) { + withCleanRegistry(t, func() { + _ = RegisterAdapter(Adapter{Name: "zebra", MapHookEvent: okMap}) + _ = RegisterAdapter(Adapter{Name: "alpha", MapHookEvent: okMap}) + _ = RegisterAdapter(Adapter{Name: "mango", MapHookEvent: okMap}) + + got := Adapters() + if len(got) != 3 { + t.Fatalf("len: %d", len(got)) + } + want := []Agent{"alpha", "mango", "zebra"} + for i, a := range got { + if a.Name != want[i] { + t.Errorf("[%d]: got %q, want %q", i, a.Name, want[i]) + } + } + + // Mutating the returned slice must not affect the registry. + got[0] = Adapter{Name: "tampered"} + again := Adapters() + if again[0].Name != "alpha" { + t.Errorf("registry mutated via snapshot: %q", again[0].Name) + } + }) +} + +func TestRegisterAdapter_Concurrent(t *testing.T) { + withCleanRegistry(t, func() { + var wg sync.WaitGroup + for i := range 16 { + wg.Add(1) + go func(i int) { + defer wg.Done() + _ = RegisterAdapter(Adapter{ + Name: Agent("a-" + string(rune('a'+i))), + MapHookEvent: okMap, + }) + _ = Adapters() + }(i) + } + wg.Wait() + if got := len(Adapters()); got != 16 { + t.Errorf("registered: got %d, want 16", got) + } + }) +} + +func TestErrUnknownAgent_IsSentinel(t *testing.T) { + withCleanRegistry(t, func() { + h, err := NewHub(HubConfig{}) + if err != nil { + t.Fatalf("NewHub: %v", err) + } + t.Cleanup(func() { _ = h.Close() }) + + err = h.Ingest("nope", []byte(`{}`)) + if !errors.Is(err, ErrUnknownAgent) { + t.Fatalf("Ingest err: %v", err) + } + }) +} diff --git a/adapters/claude/adapter.go b/adapters/claude/adapter.go new file mode 100644 index 0000000..2886aa2 --- /dev/null +++ b/adapters/claude/adapter.go @@ -0,0 +1,57 @@ +package claude + +import ( + agentstatus "github.com/kareemaly/agentstatus" +) + +// Adapter is the registered Claude Code adapter. Imported for side effects +// from package init(). +// +// Caveats (per specs/design.md §"Known gaps"): +// +// - Auto-approved tools: when a user has pre-approved a tool, Claude does +// not fire PermissionRequest, so the "awaiting_input" status will be rarer +// than for users running with default permissions. +// - "Thinking" gap: between UserPromptSubmit and the next hook event no +// status signal fires. Status remains "working" (inferred from +// UserPromptSubmit) until PreToolUse or Stop. This is acceptable. +// - Subagent identity: per the Claude hooks schema, SubagentStart and +// SubagentStop fire under the parent session's `session_id` and carry the +// subagent's stable id in `agent_id`. We model the subagent as an +// independent session: emitted Event.SessionID = agent_id, and +// Event.ParentSessionID = parent's session_id. +var Adapter = agentstatus.Adapter{ + Name: agentstatus.Claude, + MapHookEvent: MapHookEvent, + InstallHooks: installHooks, + UninstallHooks: uninstallHooks, +} + +func init() { + if err := agentstatus.RegisterAdapter(Adapter); err != nil { + // Registry collisions during init mean a programming error in the + // importing binary (double-import of this package is impossible; a + // duplicate Name in another adapter is the only way). Panic so the + // binary fails fast at startup. + panic(err) + } +} + +// installHooks is a placeholder for the next ticket. It returns +// Skipped: true so the orchestrator can fan out without the Claude adapter +// claiming success. +func installHooks(_ agentstatus.InstallConfig) (agentstatus.InstallResult, error) { + return agentstatus.InstallResult{ + Agent: agentstatus.Claude, + Skipped: true, + Reason: "not yet implemented", + }, nil +} + +func uninstallHooks(_ agentstatus.InstallConfig) (agentstatus.InstallResult, error) { + return agentstatus.InstallResult{ + Agent: agentstatus.Claude, + Skipped: true, + Reason: "not yet implemented", + }, nil +} diff --git a/adapters/claude/map.go b/adapters/claude/map.go new file mode 100644 index 0000000..ca796cd --- /dev/null +++ b/adapters/claude/map.go @@ -0,0 +1,119 @@ +package claude + +import ( + "time" + + agentstatus "github.com/kareemaly/agentstatus" +) + +// MapHookEvent translates a single Claude Code hook payload into a Signal. +// +// Returning (nil, nil) means "drop silently" — used for unknown event names +// and for metadata-only events (PreCompact). MapHookEvent never returns an +// error today; the signature reserves the seam for adapters that need to +// surface parse failures separately from drops. +func MapHookEvent(event string, payload map[string]any) (*agentstatus.Signal, error) { + sessionID := getString(payload, "session_id") + at := getTime(payload) + + base := func(s *agentstatus.Status, activity bool) *agentstatus.Signal { + return &agentstatus.Signal{ + At: at, + Activity: activity, + Status: s, + SessionID: sessionID, + Raw: payload, + } + } + + withTool := func(s *agentstatus.Signal) *agentstatus.Signal { + s.Tool = getString(payload, "tool_name") + return s + } + + // subagentSession rebinds the session ids for SubagentStart/SubagentStop: + // per the Claude hooks schema the subagent's own identity is `agent_id`, + // while the parent session keeps emitting under its own `session_id`. + subagentSession := func(s *agentstatus.Signal) *agentstatus.Signal { + s.SessionID = getString(payload, "agent_id") + s.ParentSessionID = sessionID + return s + } + + switch event { + case "SessionStart": + s := agentstatus.StatusStarting + return base(&s, false), nil + + case "UserPromptSubmit": + return base(nil, true), nil + + case "PreToolUse": + return withTool(base(nil, true)), nil + + case "PostToolUse": + return base(nil, true), nil + + case "PostToolUseFailure": + s := agentstatus.StatusError + return withTool(base(&s, false)), nil + + case "Stop": + s := agentstatus.StatusIdle + return base(&s, false), nil + + case "Notification": + s := agentstatus.StatusAwaitingInput + return base(&s, false), nil + + case "PermissionRequest": + s := agentstatus.StatusAwaitingInput + return base(&s, false), nil + + case "SubagentStart": + s := agentstatus.StatusStarting + return subagentSession(base(&s, false)), nil + + case "SubagentStop": + s := agentstatus.StatusIdle + return subagentSession(base(&s, false)), nil + + case "SessionEnd": + s := agentstatus.StatusEnded + return base(&s, false), nil + + case "PreCompact": + // Metadata only; no status change. + return nil, nil + + default: + // Unknown event — log-and-drop. + return nil, nil + } +} + +func getString(m map[string]any, key string) string { + v, _ := m[key].(string) + return v +} + +// getTime resolves the wall-clock time for a Claude payload. Claude does not +// emit a timestamp on every hook today; we accept "timestamp" as either +// RFC3339 string or Unix-seconds number, and otherwise fall back to +// time.Now() so downstream consumers always see a non-zero At. +func getTime(m map[string]any) time.Time { + switch v := m["timestamp"].(type) { + case string: + if t, err := time.Parse(time.RFC3339Nano, v); err == nil { + return t + } + if t, err := time.Parse(time.RFC3339, v); err == nil { + return t + } + case float64: + secs := int64(v) + nsecs := int64((v - float64(secs)) * 1e9) + return time.Unix(secs, nsecs) + } + return time.Now() +} diff --git a/adapters/claude/map_test.go b/adapters/claude/map_test.go new file mode 100644 index 0000000..cc99f88 --- /dev/null +++ b/adapters/claude/map_test.go @@ -0,0 +1,150 @@ +package claude + +import ( + "encoding/json" + "os" + "path/filepath" + "testing" + + agentstatus "github.com/kareemaly/agentstatus" +) + +func loadFixture(t *testing.T, name string) map[string]any { + t.Helper() + path := filepath.Join("..", "..", "testdata", "claude", name) + b, err := os.ReadFile(path) + if err != nil { + t.Fatalf("read %s: %v", path, err) + } + var m map[string]any + if err := json.Unmarshal(b, &m); err != nil { + t.Fatalf("unmarshal %s: %v", path, err) + } + return m +} + +func TestMapHookEvent_AllRows(t *testing.T) { + t.Parallel() + + type want struct { + drop bool + status *agentstatus.Status + activity bool + tool string + parent string + sessionID string + } + + starting := agentstatus.StatusStarting + idle := agentstatus.StatusIdle + awaiting := agentstatus.StatusAwaitingInput + errSt := agentstatus.StatusError + ended := agentstatus.StatusEnded + + cases := []struct { + fixture string + event string + want want + }{ + {"session_start.json", "SessionStart", want{status: &starting, sessionID: "sess-1"}}, + {"user_prompt_submit.json", "UserPromptSubmit", want{activity: true, sessionID: "sess-1"}}, + {"pre_tool_use_read.json", "PreToolUse", want{activity: true, tool: "Read", sessionID: "sess-1"}}, + {"post_tool_use.json", "PostToolUse", want{activity: true, sessionID: "sess-1"}}, + {"post_tool_use_failure.json", "PostToolUseFailure", want{status: &errSt, tool: "Bash", sessionID: "sess-1"}}, + {"stop.json", "Stop", want{status: &idle, sessionID: "sess-1"}}, + {"notification.json", "Notification", want{status: &awaiting, sessionID: "sess-1"}}, + {"permission_request.json", "PermissionRequest", want{status: &awaiting, sessionID: "sess-1"}}, + {"subagent_start.json", "SubagentStart", want{status: &starting, sessionID: "agent-abc123", parent: "parent-1"}}, + {"subagent_stop.json", "SubagentStop", want{status: &idle, sessionID: "agent-abc123", parent: "parent-1"}}, + {"session_end.json", "SessionEnd", want{status: &ended, sessionID: "sess-1"}}, + {"pre_compact.json", "PreCompact", want{drop: true}}, + {"unknown_event.json", "NonExistent", want{drop: true}}, + } + + for _, tc := range cases { + t.Run(tc.event, func(t *testing.T) { + payload := loadFixture(t, tc.fixture) + sig, err := MapHookEvent(tc.event, payload) + if err != nil { + t.Fatalf("MapHookEvent: %v", err) + } + if tc.want.drop { + if sig != nil { + t.Fatalf("expected drop, got %+v", sig) + } + return + } + if sig == nil { + t.Fatal("expected signal, got nil") + } + if (sig.Status == nil) != (tc.want.status == nil) { + t.Fatalf("status presence: got %v want %v", sig.Status, tc.want.status) + } + if sig.Status != nil && *sig.Status != *tc.want.status { + t.Errorf("status: got %q want %q", *sig.Status, *tc.want.status) + } + if sig.Activity != tc.want.activity { + t.Errorf("activity: got %v want %v", sig.Activity, tc.want.activity) + } + if sig.Tool != tc.want.tool { + t.Errorf("tool: got %q want %q", sig.Tool, tc.want.tool) + } + if sig.SessionID != tc.want.sessionID { + t.Errorf("session: got %q want %q", sig.SessionID, tc.want.sessionID) + } + if sig.ParentSessionID != tc.want.parent { + t.Errorf("parent: got %q want %q", sig.ParentSessionID, tc.want.parent) + } + if sig.Raw == nil { + t.Error("Raw is nil") + } else if sig.Raw["hook_event_name"] != payload["hook_event_name"] { + t.Errorf("Raw mismatch: %v", sig.Raw) + } + if sig.At.IsZero() { + t.Error("At is zero") + } + }) + } +} + +func TestMapHookEvent_TimestampRFC3339(t *testing.T) { + payload := map[string]any{ + "hook_event_name": "Stop", + "session_id": "s1", + "timestamp": "2025-01-02T03:04:05Z", + } + sig, err := MapHookEvent("Stop", payload) + if err != nil || sig == nil { + t.Fatalf("map: %v %v", sig, err) + } + if sig.At.Year() != 2025 || sig.At.Month() != 1 || sig.At.Day() != 2 { + t.Errorf("At: got %v", sig.At) + } +} + +func TestMapHookEvent_TimestampNumeric(t *testing.T) { + payload := map[string]any{ + "hook_event_name": "Stop", + "session_id": "s1", + "timestamp": float64(1700000000), + } + sig, err := MapHookEvent("Stop", payload) + if err != nil || sig == nil { + t.Fatalf("map: %v %v", sig, err) + } + if sig.At.Unix() != 1700000000 { + t.Errorf("At unix: got %d", sig.At.Unix()) + } +} + +func TestMapHookEvent_MissingFieldsTolerated(t *testing.T) { + // Empty payload — must not panic, must not error. Status pointer still + // honored for the event. + sig, err := MapHookEvent("Stop", map[string]any{}) + if err != nil { + t.Fatalf("err: %v", err) + } + if sig == nil || sig.Status == nil || *sig.Status != agentstatus.StatusIdle { + t.Errorf("got %+v", sig) + } +} diff --git a/decision.go b/decision.go index 2f0e7ea..4da7368 100644 --- a/decision.go +++ b/decision.go @@ -7,6 +7,11 @@ import "time" // // Signal carries no Agent field: the Agent is known at the adapter/ingest seam // and attached to the emitted Event there. +// +// Raw is the original hook payload. Decide does not read or write it; the +// field flows around the decision machine into the emitted Event.Raw as a +// design-mandated escape hatch for consumers that need provider-specific +// fields not surfaced on Event. type Signal struct { At time.Time Activity bool @@ -15,6 +20,7 @@ type Signal struct { Work string SessionID string ParentSessionID string + Raw map[string]any } // Transition is the internal return value of Decide. It is emitted only when a diff --git a/hub.go b/hub.go index 319c445..1b397d5 100644 --- a/hub.go +++ b/hub.go @@ -1,14 +1,22 @@ package agentstatus import ( + "encoding/json" + "errors" + "fmt" "io" "log/slog" "maps" + "net/http" "sync" "github.com/kareemaly/agentstatus/internal/broadcast" ) +// MaxIngestBodyBytes caps the request body size accepted by Hub.Handler. +// Larger bodies receive 413 Request Entity Too Large. +const MaxIngestBodyBytes = 1 << 20 // 1 MiB + // HubConfig configures a Hub. Zero values are valid and produce the documented // defaults. type HubConfig struct { @@ -156,6 +164,99 @@ func (h *Hub) dispatchSignal(agent Agent, sig Signal) { Work: sig.Work, At: sig.At, Tags: tags, + Raw: sig.Raw, } h.bcast.Publish(ev) } + +// Ingest is the transport-agnostic entry point. Adapters are looked up by +// agent name; the payload is JSON-decoded into a map and handed to the +// adapter's MapHookEvent. A nil-signal return drops silently (unknown event +// or metadata-only event); a non-nil signal is dispatched through the +// decision machine. +// +// Ingest returning nil guarantees the signal was dispatched, not that any +// subscriber observed the resulting Event — slow subscribers may have +// dropped it per their own buffer policy. +func (h *Hub) Ingest(agent Agent, payload []byte) error { + adapter, ok := lookupAdapter(agent) + if !ok { + return fmt.Errorf("%w: %q", ErrUnknownAgent, agent) + } + + var m map[string]any + if err := json.Unmarshal(payload, &m); err != nil { + return fmt.Errorf("agentstatus: invalid JSON: %w", err) + } + + event, _ := m["hook_event_name"].(string) + sig, err := adapter.MapHookEvent(event, m) + if err != nil { + return fmt.Errorf("agentstatus: adapter %q map error: %w", agent, err) + } + if sig == nil { + return nil + } + h.dispatchSignal(agent, *sig) + return nil +} + +// Handler returns an http.Handler that accepts POST /hook/{agent} and +// forwards bodies into Ingest. The handler is safe for concurrent use. +// +// Status codes: +// - 202 Accepted on success +// - 400 Bad Request on malformed JSON or read error +// - 404 Not Found on unknown agent +// - 405 Method Not Allowed on non-POST (provided by net/http for the +// method-aware route) +// - 413 Request Entity Too Large when the body exceeds MaxIngestBodyBytes +// - 500 Internal Server Error for everything else (also routed through +// HubConfig.ErrorHandler) +func (h *Hub) Handler() http.Handler { + mux := http.NewServeMux() + mux.HandleFunc("POST /hook/{agent}", h.handleIngest) + return mux +} + +func (h *Hub) handleIngest(w http.ResponseWriter, r *http.Request) { + r.Body = http.MaxBytesReader(w, r.Body, MaxIngestBodyBytes) + body, err := io.ReadAll(r.Body) + if err != nil { + var maxErr *http.MaxBytesError + if errors.As(err, &maxErr) { + http.Error(w, "request body too large", http.StatusRequestEntityTooLarge) + return + } + http.Error(w, "read error", http.StatusBadRequest) + return + } + + agent := Agent(r.PathValue("agent")) + if err := h.Ingest(agent, body); err != nil { + switch { + case errors.Is(err, ErrUnknownAgent): + http.Error(w, "unknown agent", http.StatusNotFound) + case isJSONError(err): + http.Error(w, "invalid JSON", http.StatusBadRequest) + default: + h.errH(err) + http.Error(w, "internal error", http.StatusInternalServerError) + } + return + } + w.WriteHeader(http.StatusAccepted) +} + +func isJSONError(err error) bool { + var syntaxErr *json.SyntaxError + var typeErr *json.UnmarshalTypeError + return errors.As(err, &syntaxErr) || errors.As(err, &typeErr) +} + +// ServeHTTP is a convenience wrapper that mounts Handler at addr and blocks +// until http.ListenAndServe returns. Consumers needing graceful shutdown +// should use Handler() with their own *http.Server. +func (h *Hub) ServeHTTP(addr string) error { + return http.ListenAndServe(addr, h.Handler()) +} diff --git a/hub_ingest_test.go b/hub_ingest_test.go new file mode 100644 index 0000000..f31c526 --- /dev/null +++ b/hub_ingest_test.go @@ -0,0 +1,280 @@ +package agentstatus_test + +import ( + "bytes" + "fmt" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "sync" + "sync/atomic" + "testing" + "time" + + agentstatus "github.com/kareemaly/agentstatus" + _ "github.com/kareemaly/agentstatus/adapters/claude" +) + +func newServedHub(t *testing.T) (*agentstatus.Hub, *httptest.Server) { + t.Helper() + h, err := agentstatus.NewHub(agentstatus.HubConfig{}) + if err != nil { + t.Fatalf("NewHub: %v", err) + } + srv := httptest.NewServer(h.Handler()) + t.Cleanup(func() { + srv.Close() + _ = h.Close() + }) + return h, srv +} + +func loadPayload(t *testing.T, name string) []byte { + t.Helper() + b, err := os.ReadFile(filepath.Join("testdata", "claude", name)) + if err != nil { + t.Fatalf("read %s: %v", name, err) + } + return b +} + +func postFixture(t *testing.T, srv *httptest.Server, agent, fixture string) *http.Response { + t.Helper() + body := loadPayload(t, fixture) + resp, err := http.Post(srv.URL+"/hook/"+agent, "application/json", bytes.NewReader(body)) + if err != nil { + t.Fatalf("POST %s: %v", fixture, err) + } + return resp +} + +func recvOrFail(t *testing.T, ch <-chan agentstatus.Event, d time.Duration) agentstatus.Event { + t.Helper() + select { + case ev, ok := <-ch: + if !ok { + t.Fatal("channel closed") + } + return ev + case <-time.After(d): + t.Fatal("no event") + return agentstatus.Event{} + } +} + +func TestHTTP_HappyPath(t *testing.T) { + t.Parallel() + h, srv := newServedHub(t) + stream := h.Events() + + for _, fix := range []string{"session_start.json", "pre_tool_use_read.json", "stop.json"} { + resp := postFixture(t, srv, "claude", fix) + if resp.StatusCode != http.StatusAccepted { + t.Fatalf("%s: status %d", fix, resp.StatusCode) + } + _ = resp.Body.Close() + } + + ch := stream.Channel() + e1 := recvOrFail(t, ch, time.Second) + e2 := recvOrFail(t, ch, time.Second) + e3 := recvOrFail(t, ch, time.Second) + + if e1.Status != agentstatus.StatusStarting || e1.PrevStatus != "" { + t.Errorf("e1: %q←%q", e1.Status, e1.PrevStatus) + } + if e2.Status != agentstatus.StatusWorking || e2.PrevStatus != agentstatus.StatusStarting { + t.Errorf("e2: %q←%q", e2.Status, e2.PrevStatus) + } + if e2.Tool != "Read" { + t.Errorf("e2 tool: %q", e2.Tool) + } + if e3.Status != agentstatus.StatusIdle || e3.PrevStatus != agentstatus.StatusWorking { + t.Errorf("e3: %q←%q", e3.Status, e3.PrevStatus) + } + for _, e := range []agentstatus.Event{e1, e2, e3} { + if e.Agent != agentstatus.Claude { + t.Errorf("agent: %q", e.Agent) + } + if e.SessionID != "sess-1" { + t.Errorf("session: %q", e.SessionID) + } + if e.Raw == nil || e.Raw["session_id"] != "sess-1" { + t.Errorf("raw: %v", e.Raw) + } + } +} + +func TestHTTP_UnknownAgent(t *testing.T) { + t.Parallel() + _, srv := newServedHub(t) + resp, err := http.Post(srv.URL+"/hook/nope", "application/json", bytes.NewReader([]byte(`{}`))) + if err != nil { + t.Fatal(err) + } + defer func() { _ = resp.Body.Close() }() + if resp.StatusCode != http.StatusNotFound { + t.Errorf("status: %d", resp.StatusCode) + } +} + +func TestHTTP_MalformedJSON(t *testing.T) { + t.Parallel() + _, srv := newServedHub(t) + resp, err := http.Post(srv.URL+"/hook/claude", "application/json", bytes.NewReader([]byte(`{`))) + if err != nil { + t.Fatal(err) + } + defer func() { _ = resp.Body.Close() }() + if resp.StatusCode != http.StatusBadRequest { + t.Errorf("status: %d", resp.StatusCode) + } +} + +func TestHTTP_OversizedBody(t *testing.T) { + t.Parallel() + _, srv := newServedHub(t) + big := bytes.Repeat([]byte("a"), 2<<20) // 2 MiB + resp, err := http.Post(srv.URL+"/hook/claude", "application/json", bytes.NewReader(big)) + if err != nil { + t.Fatal(err) + } + defer func() { _ = resp.Body.Close() }() + if resp.StatusCode != http.StatusRequestEntityTooLarge { + t.Errorf("status: %d", resp.StatusCode) + } +} + +func TestHTTP_GetNotAllowed(t *testing.T) { + t.Parallel() + _, srv := newServedHub(t) + resp, err := http.Get(srv.URL + "/hook/claude") + if err != nil { + t.Fatal(err) + } + defer func() { _ = resp.Body.Close() }() + if resp.StatusCode != http.StatusMethodNotAllowed { + t.Errorf("status: %d", resp.StatusCode) + } +} + +func TestHTTP_UnknownEventNoEvent(t *testing.T) { + t.Parallel() + h, srv := newServedHub(t) + stream := h.Events() + + resp := postFixture(t, srv, "claude", "unknown_event.json") + if resp.StatusCode != http.StatusAccepted { + t.Fatalf("status: %d", resp.StatusCode) + } + _ = resp.Body.Close() + + select { + case ev := <-stream.Channel(): + t.Fatalf("unexpected event: %+v", ev) + case <-time.After(50 * time.Millisecond): + } +} + +func TestHTTP_PreCompactNoEvent(t *testing.T) { + t.Parallel() + h, srv := newServedHub(t) + stream := h.Events() + + resp := postFixture(t, srv, "claude", "pre_compact.json") + if resp.StatusCode != http.StatusAccepted { + t.Fatalf("status: %d", resp.StatusCode) + } + _ = resp.Body.Close() + + select { + case ev := <-stream.Channel(): + t.Fatalf("unexpected event: %+v", ev) + case <-time.After(50 * time.Millisecond): + } +} + +func TestHTTP_SubagentFlow(t *testing.T) { + t.Parallel() + h, srv := newServedHub(t) + stream := h.Events() + + for _, fix := range []string{"subagent_start.json", "subagent_stop.json"} { + resp := postFixture(t, srv, "claude", fix) + if resp.StatusCode != http.StatusAccepted { + t.Fatalf("%s: status %d", fix, resp.StatusCode) + } + _ = resp.Body.Close() + } + + ch := stream.Channel() + e1 := recvOrFail(t, ch, time.Second) + e2 := recvOrFail(t, ch, time.Second) + + for i, e := range []agentstatus.Event{e1, e2} { + if e.SessionID != "agent-abc123" { + t.Errorf("[%d] session: %q", i, e.SessionID) + } + if e.ParentSessionID != "parent-1" { + t.Errorf("[%d] parent: %q", i, e.ParentSessionID) + } + } + if e1.Status != agentstatus.StatusStarting { + t.Errorf("e1 status: %q", e1.Status) + } + if e2.Status != agentstatus.StatusIdle { + t.Errorf("e2 status: %q", e2.Status) + } +} + +func TestHTTP_ConcurrentIngest(t *testing.T) { + t.Parallel() + h, err := agentstatus.NewHub(agentstatus.HubConfig{BufferSize: 2048}) + if err != nil { + t.Fatalf("NewHub: %v", err) + } + srv := httptest.NewServer(h.Handler()) + t.Cleanup(func() { + srv.Close() + _ = h.Close() + }) + + stream := h.Events() + var got atomic.Int32 + done := make(chan struct{}) + go func() { + for range stream.Channel() { + if got.Add(1) == 1000 { + close(done) + return + } + } + }() + + const N = 1000 + var wg sync.WaitGroup + wg.Add(N) + for i := range N { + go func(i int) { + defer wg.Done() + body := fmt.Appendf(nil, `{"hook_event_name":"SessionStart","session_id":"s-%d"}`, i) + resp, err := http.Post(srv.URL+"/hook/claude", "application/json", bytes.NewReader(body)) + if err != nil { + t.Errorf("POST: %v", err) + return + } + _ = resp.Body.Close() + if resp.StatusCode != http.StatusAccepted { + t.Errorf("status: %d", resp.StatusCode) + } + }(i) + } + wg.Wait() + + select { + case <-done: + case <-time.After(5 * time.Second): + t.Fatalf("timed out: got %d/1000 events", got.Load()) + } +} diff --git a/specs/design.md b/specs/design.md index 6d8ec07..041d27d 100644 --- a/specs/design.md +++ b/specs/design.md @@ -180,6 +180,7 @@ type Signal struct { Work string SessionID string ParentSessionID string + Raw map[string]any // original hook payload; flows into Event.Raw } ``` diff --git a/testdata/claude/notification.json b/testdata/claude/notification.json new file mode 100644 index 0000000..c97afce --- /dev/null +++ b/testdata/claude/notification.json @@ -0,0 +1,9 @@ +{ + "hook_event_name": "Notification", + "session_id": "sess-1", + "transcript_path": "/Users/u/.claude/projects/repo/sess-1.jsonl", + "cwd": "/repo", + "message": "Claude needs your permission to use Bash", + "title": "Permission needed", + "notification_type": "permission_prompt" +} diff --git a/testdata/claude/permission_request.json b/testdata/claude/permission_request.json new file mode 100644 index 0000000..8fa7f42 --- /dev/null +++ b/testdata/claude/permission_request.json @@ -0,0 +1,9 @@ +{ + "hook_event_name": "PermissionRequest", + "session_id": "sess-1", + "transcript_path": "/Users/u/.claude/projects/repo/sess-1.jsonl", + "cwd": "/repo", + "permission_mode": "default", + "tool_name": "Bash", + "tool_input": { "command": "rm -rf node_modules" } +} diff --git a/testdata/claude/post_tool_use.json b/testdata/claude/post_tool_use.json new file mode 100644 index 0000000..81b9610 --- /dev/null +++ b/testdata/claude/post_tool_use.json @@ -0,0 +1,10 @@ +{ + "hook_event_name": "PostToolUse", + "session_id": "sess-1", + "transcript_path": "/Users/u/.claude/projects/repo/sess-1.jsonl", + "cwd": "/repo", + "permission_mode": "default", + "tool_name": "Read", + "tool_input": { "file_path": "/repo/main.go" }, + "tool_use_id": "toolu_01ABC" +} diff --git a/testdata/claude/post_tool_use_failure.json b/testdata/claude/post_tool_use_failure.json new file mode 100644 index 0000000..071b876 --- /dev/null +++ b/testdata/claude/post_tool_use_failure.json @@ -0,0 +1,12 @@ +{ + "hook_event_name": "PostToolUseFailure", + "session_id": "sess-1", + "transcript_path": "/Users/u/.claude/projects/repo/sess-1.jsonl", + "cwd": "/repo", + "permission_mode": "default", + "tool_name": "Bash", + "tool_input": { "command": "npm test" }, + "tool_use_id": "toolu_01ABC", + "error": "Command exited with non-zero status code 1", + "is_interrupt": false +} diff --git a/testdata/claude/pre_compact.json b/testdata/claude/pre_compact.json new file mode 100644 index 0000000..ca41747 --- /dev/null +++ b/testdata/claude/pre_compact.json @@ -0,0 +1,5 @@ +{ + "hook_event_name": "PreCompact", + "session_id": "sess-1", + "cwd": "/repo" +} diff --git a/testdata/claude/pre_tool_use_read.json b/testdata/claude/pre_tool_use_read.json new file mode 100644 index 0000000..e709502 --- /dev/null +++ b/testdata/claude/pre_tool_use_read.json @@ -0,0 +1,9 @@ +{ + "hook_event_name": "PreToolUse", + "session_id": "sess-1", + "transcript_path": "/Users/u/.claude/projects/repo/sess-1.jsonl", + "cwd": "/repo", + "permission_mode": "default", + "tool_name": "Read", + "tool_input": { "file_path": "/repo/main.go" } +} diff --git a/testdata/claude/session_end.json b/testdata/claude/session_end.json new file mode 100644 index 0000000..8909bd2 --- /dev/null +++ b/testdata/claude/session_end.json @@ -0,0 +1,5 @@ +{ + "hook_event_name": "SessionEnd", + "session_id": "sess-1", + "cwd": "/repo" +} diff --git a/testdata/claude/session_start.json b/testdata/claude/session_start.json new file mode 100644 index 0000000..625f095 --- /dev/null +++ b/testdata/claude/session_start.json @@ -0,0 +1,8 @@ +{ + "hook_event_name": "SessionStart", + "session_id": "sess-1", + "transcript_path": "/Users/u/.claude/projects/repo/sess-1.jsonl", + "cwd": "/repo", + "source": "startup", + "model": "claude-sonnet-4-6" +} diff --git a/testdata/claude/stop.json b/testdata/claude/stop.json new file mode 100644 index 0000000..6889f65 --- /dev/null +++ b/testdata/claude/stop.json @@ -0,0 +1,5 @@ +{ + "hook_event_name": "Stop", + "session_id": "sess-1", + "cwd": "/repo" +} diff --git a/testdata/claude/subagent_start.json b/testdata/claude/subagent_start.json new file mode 100644 index 0000000..91a8bbb --- /dev/null +++ b/testdata/claude/subagent_start.json @@ -0,0 +1,8 @@ +{ + "hook_event_name": "SubagentStart", + "session_id": "parent-1", + "transcript_path": "/Users/u/.claude/projects/repo/parent-1.jsonl", + "cwd": "/repo", + "agent_id": "agent-abc123", + "agent_type": "Explore" +} diff --git a/testdata/claude/subagent_stop.json b/testdata/claude/subagent_stop.json new file mode 100644 index 0000000..5bc318d --- /dev/null +++ b/testdata/claude/subagent_stop.json @@ -0,0 +1,11 @@ +{ + "hook_event_name": "SubagentStop", + "session_id": "parent-1", + "transcript_path": "/Users/u/.claude/projects/repo/parent-1.jsonl", + "cwd": "/repo", + "permission_mode": "default", + "stop_hook_active": false, + "agent_id": "agent-abc123", + "agent_type": "Explore", + "agent_transcript_path": "/Users/u/.claude/projects/repo/parent-1/subagents/agent-abc123.jsonl" +} diff --git a/testdata/claude/unknown_event.json b/testdata/claude/unknown_event.json new file mode 100644 index 0000000..ec8386a --- /dev/null +++ b/testdata/claude/unknown_event.json @@ -0,0 +1,5 @@ +{ + "hook_event_name": "NonExistent", + "session_id": "sess-1", + "cwd": "/repo" +} diff --git a/testdata/claude/user_prompt_submit.json b/testdata/claude/user_prompt_submit.json new file mode 100644 index 0000000..f23a8f1 --- /dev/null +++ b/testdata/claude/user_prompt_submit.json @@ -0,0 +1,5 @@ +{ + "hook_event_name": "UserPromptSubmit", + "session_id": "sess-1", + "cwd": "/repo" +}