diff --git a/cllama b/cllama index 4125d52..5519a86 160000 --- a/cllama +++ b/cllama @@ -1 +1 @@ -Subproject commit 4125d5272df17032fd9d98c98822182aa26746d6 +Subproject commit 5519a865b21575b9eb4cc94695e5a3bea013c740 diff --git a/cmd/claw-wall/channel_memory.go b/cmd/claw-wall/channel_memory.go index f7b4c8a..d318e97 100644 --- a/cmd/claw-wall/channel_memory.go +++ b/cmd/claw-wall/channel_memory.go @@ -22,6 +22,7 @@ const ( type channelMemoryClient struct { ingestURL string + digestURL string token string client *http.Client } @@ -52,35 +53,139 @@ type channelMemoryIngestSource struct { GuildID string `json:"guild_id,omitempty"` } +type channelMemoryDigestRequest struct { + SourceKind string `json:"source_kind,omitempty"` + ChannelIDs []string `json:"channel_ids,omitempty"` + Since string `json:"since,omitempty"` + Budget channelMemoryDigestBudget `json:"budget,omitempty"` +} + +type channelMemoryDigestBudget struct { + MaxBlocks int `json:"max_blocks,omitempty"` +} + +type channelMemoryDigestResponse struct { + Status string `json:"status"` + GeneratedAt string `json:"generated_at"` + Coverage channelMemoryDigestCoverage `json:"coverage"` + Blocks []channelMemoryDigestBlock `json:"blocks"` + Cost channelMemoryDigestCost `json:"cost"` +} + +type channelMemoryDigestCoverage struct { + From string `json:"from,omitempty"` + To string `json:"to,omitempty"` + SourceMessages int `json:"source_messages"` + DigestMessages int `json:"digest_messages"` + RawRecentMessages int `json:"raw_recent_messages"` + Gaps []channelMemoryCoverageGap `json:"gaps,omitempty"` +} + +type channelMemoryCoverageGap struct { + ID int64 `json:"id,omitempty"` + ChannelID string `json:"channel_id"` + From string `json:"from"` + To string `json:"to"` + Reason string `json:"reason,omitempty"` + CreatedAt string `json:"created_at,omitempty"` +} + +type channelMemoryDigestCost struct { + DeterministicOnly bool `json:"deterministic_only"` + LLMCallsToday int `json:"llm_calls_today"` +} + +type channelMemoryDigestBlock struct { + ID int64 `json:"id,omitempty"` + Kind string `json:"kind"` + EventType string `json:"event_type,omitempty"` + Text string `json:"text"` + SourceChannel string `json:"source_channel"` + SourceMessages []string `json:"source_messages"` + CoveredContentHashes []string `json:"covered_content_hashes,omitempty"` + SourceWindow channelMemorySourceWindow `json:"source_window"` + Sparse bool `json:"sparse"` + Score float64 `json:"score"` + GeneratedAt string `json:"generated_at"` + Stale bool `json:"stale,omitempty"` + Dirty bool `json:"dirty,omitempty"` + Processor string `json:"processor"` +} + +type channelMemorySourceWindow struct { + From string `json:"from"` + To string `json:"to"` +} + func newChannelMemoryClient(rawURL, token string, timeout time.Duration) (*channelMemoryClient, error) { - rawURL = strings.TrimSpace(rawURL) - if rawURL == "" { + return newChannelMemoryClientWithDigest(rawURL, "", token, timeout) +} + +func newChannelMemoryClientWithDigest(rawIngestURL, rawDigestURL, token string, timeout time.Duration) (*channelMemoryClient, error) { + rawIngestURL = strings.TrimSpace(rawIngestURL) + rawDigestURL = strings.TrimSpace(rawDigestURL) + if rawIngestURL == "" && rawDigestURL == "" { return nil, nil } - parsed, err := url.Parse(rawURL) - if err != nil { - return nil, fmt.Errorf("parse channel-memory ingest URL: %w", err) + if rawDigestURL == "" { + rawDigestURL = deriveChannelMemoryDigestURL(rawIngestURL) } - if parsed.Scheme != "http" && parsed.Scheme != "https" { - return nil, fmt.Errorf("channel-memory ingest URL must use http or https") + if rawIngestURL != "" { + if err := validateChannelMemoryURL(rawIngestURL, "ingest"); err != nil { + return nil, err + } } - if strings.TrimSpace(parsed.Host) == "" { - return nil, fmt.Errorf("channel-memory ingest URL must include a host") + if rawDigestURL != "" { + if err := validateChannelMemoryURL(rawDigestURL, "digest"); err != nil { + return nil, err + } } if timeout <= 0 { timeout = 2 * time.Second } return &channelMemoryClient{ - ingestURL: rawURL, + ingestURL: rawIngestURL, + digestURL: rawDigestURL, token: strings.TrimSpace(token), client: &http.Client{Timeout: timeout}, }, nil } +func validateChannelMemoryURL(rawURL, name string) error { + parsed, err := url.Parse(rawURL) + if err != nil { + return fmt.Errorf("parse channel-memory %s URL: %w", name, err) + } + if parsed.Scheme != "http" && parsed.Scheme != "https" { + return fmt.Errorf("channel-memory %s URL must use http or https", name) + } + if strings.TrimSpace(parsed.Host) == "" { + return fmt.Errorf("channel-memory %s URL must include a host", name) + } + return nil +} + +func deriveChannelMemoryDigestURL(rawIngestURL string) string { + parsed, err := url.Parse(strings.TrimSpace(rawIngestURL)) + if err != nil { + return "" + } + path := strings.TrimRight(parsed.Path, "/") + if strings.HasSuffix(path, "/ingest") { + parsed.Path = strings.TrimSuffix(path, "/ingest") + "/digest" + return parsed.String() + } + return "" +} + func (c *channelMemoryClient) enabled() bool { return c != nil && strings.TrimSpace(c.ingestURL) != "" } +func (c *channelMemoryClient) digestEnabled() bool { + return c != nil && strings.TrimSpace(c.digestURL) != "" +} + func (c *channelMemoryClient) ingestMessages(ctx context.Context, messages []wallMessage) (int, error) { if !c.enabled() || len(messages) == 0 { return 0, nil @@ -126,6 +231,94 @@ func (c *channelMemoryClient) ingestMessage(ctx context.Context, msg wallMessage return nil } +func (c *channelMemoryClient) digest(ctx context.Context, req channelMemoryDigestRequest) (*channelMemoryDigestResponse, error) { + if !c.digestEnabled() { + return nil, nil + } + body, err := json.Marshal(req) + if err != nil { + return nil, err + } + httpReq, err := http.NewRequestWithContext(ctx, http.MethodPost, c.digestURL, bytes.NewReader(body)) + if err != nil { + return nil, err + } + httpReq.Header.Set("Content-Type", "application/json") + if c.token != "" { + httpReq.Header.Set("Authorization", "Bearer "+c.token) + } + + resp, err := c.client.Do(httpReq) + if err != nil { + return nil, err + } + defer resp.Body.Close() + respBody, err := io.ReadAll(io.LimitReader(resp.Body, 1<<20)) + if err != nil { + return nil, err + } + if resp.StatusCode < 200 || resp.StatusCode >= 300 { + return nil, fmt.Errorf("channel-memory digest returned %s: %s", resp.Status, strings.TrimSpace(string(respBody))) + } + var parsed channelMemoryDigestResponse + if err := json.Unmarshal(respBody, &parsed); err != nil { + return nil, fmt.Errorf("decode channel-memory digest: %w", err) + } + return &parsed, nil +} + +func (c *channelMemoryClient) fetchAwarenessDigest(ctx context.Context, channelIDs []string, since time.Duration) channelAwarenessDigest { + digest := channelAwarenessDigest{Requested: true, Status: "unavailable"} + if !c.digestEnabled() { + return digest + } + req := channelMemoryDigestRequest{ + SourceKind: channelMemorySourceKind, + ChannelIDs: normalizeChannelIDs(channelIDs), + Budget: channelMemoryDigestBudget{ + MaxBlocks: defaultDigestMaxBlocks, + }, + } + if since > 0 { + req.Since = since.String() + } + resp, err := c.digest(ctx, req) + if err != nil || resp == nil { + return digest + } + status := normalizeDigestStatus(resp.Status) + if status == "ok" && digestBlocksStale(resp.Blocks) { + status = "stale" + } + digest.Status = status + digest.GeneratedAt = strings.TrimSpace(resp.GeneratedAt) + digest.SourceMessages = resp.Coverage.SourceMessages + digest.DigestMessages = resp.Coverage.DigestMessages + digest.RawRecentMessages = resp.Coverage.RawRecentMessages + digest.CoverageGaps = len(resp.Coverage.Gaps) + digest.DeterministicOnly = resp.Cost.DeterministicOnly + digest.Blocks = append([]channelMemoryDigestBlock(nil), resp.Blocks...) + return digest +} + +func normalizeDigestStatus(status string) string { + switch strings.TrimSpace(status) { + case "ok", "stale", "unavailable", "coverage_gap": + return strings.TrimSpace(status) + default: + return "unavailable" + } +} + +func digestBlocksStale(blocks []channelMemoryDigestBlock) bool { + for _, block := range blocks { + if block.Stale || block.Dirty { + return true + } + } + return false +} + func channelMemoryPayloadForMessage(msg wallMessage) channelMemoryIngestRequest { scope := "channel:" + strings.TrimSpace(msg.ChannelID) return channelMemoryIngestRequest{ diff --git a/cmd/claw-wall/main.go b/cmd/claw-wall/main.go index 1f0bf22..1f127dd 100644 --- a/cmd/claw-wall/main.go +++ b/cmd/claw-wall/main.go @@ -27,6 +27,7 @@ type config struct { ToolToken string AgentChannelsPath string ChannelMemoryIngestURL string + ChannelMemoryDigestURL string ChannelMemoryToken string ChannelMemoryTimeout time.Duration } @@ -58,7 +59,7 @@ func run(args []string) error { if err != nil { return fmt.Errorf("claw-wall: parse CLAW_WALL_TOKENS: %w", err) } - channelMemory, err := newChannelMemoryClient(cfg.ChannelMemoryIngestURL, cfg.ChannelMemoryToken, cfg.ChannelMemoryTimeout) + channelMemory, err := newChannelMemoryClientWithDigest(cfg.ChannelMemoryIngestURL, cfg.ChannelMemoryDigestURL, cfg.ChannelMemoryToken, cfg.ChannelMemoryTimeout) if err != nil { return fmt.Errorf("claw-wall: configure channel-memory: %w", err) } @@ -168,6 +169,7 @@ func loadConfig() (config, error) { ToolToken: strings.TrimSpace(os.Getenv("CLAW_WALL_TOOL_TOKEN")), AgentChannelsPath: envOr("CLAW_WALL_AGENT_CHANNELS_FILE", "/etc/claw-wall/agent-channels.json"), ChannelMemoryIngestURL: strings.TrimSpace(os.Getenv("CLAW_WALL_CHANNEL_MEMORY_INGEST_URL")), + ChannelMemoryDigestURL: strings.TrimSpace(os.Getenv("CLAW_WALL_CHANNEL_MEMORY_DIGEST_URL")), ChannelMemoryToken: strings.TrimSpace(os.Getenv("CLAW_WALL_CHANNEL_MEMORY_TOKEN")), ChannelMemoryTimeout: channelMemoryTimeout, }, nil diff --git a/cmd/claw-wall/main_test.go b/cmd/claw-wall/main_test.go index 1de3d93..3fe58b2 100644 --- a/cmd/claw-wall/main_test.go +++ b/cmd/claw-wall/main_test.go @@ -294,6 +294,217 @@ func TestChannelAwarenessHeaderReportsBackfillStatus(t *testing.T) { } } +func TestChannelAwarenessHandlerServesDigestWhenAvailable(t *testing.T) { + store := newConversationStore(50) + store.now = func() time.Time { return time.Unix(112, 0) } + for i := 0; i < 12; i++ { + store.merge("chan-1", []wallMessage{{ + ID: fmt.Sprintf("%d", 100+i), + Author: "agent-status", + Content: fmt.Sprintf("HEARTBEAT_OK message-%03d", i), + Timestamp: time.Unix(int64(100+i), 0), + }}) + } + + memory := newDigestChannelMemoryServer(t, channelMemoryDigestResponse{ + Status: "ok", + GeneratedAt: "2026-05-21T20:00:00Z", + Coverage: channelMemoryDigestCoverage{ + SourceMessages: 12, + DigestMessages: 1, + }, + Blocks: []channelMemoryDigestBlock{{ + Kind: "telemetry_count", + Text: "[00:01-00:02] runtime/status noise elided: 2 messages.", + SourceChannel: "chan-1", + SourceMessages: []string{"100", "101"}, + SourceWindow: channelMemorySourceWindow{ + From: "1970-01-01T00:01:40Z", + To: "1970-01-01T00:01:41Z", + }, + Sparse: true, + Score: 0.25, + GeneratedAt: "2026-05-21T20:00:00Z", + Processor: "deterministic", + }}, + Cost: channelMemoryDigestCost{DeterministicOnly: true}, + }) + defer memory.Close() + client, err := newChannelMemoryClientWithDigest("", memory.URL+"/digest", "memory-token", time.Second) + if err != nil { + t.Fatalf("newChannelMemoryClientWithDigest: %v", err) + } + + server := httptest.NewServer(newHandler(store, handlerConfig{ + channelMemory: client, + toolToken: "tool-token", + agentChannels: map[string]map[string]struct{}{ + "trader-0": {"chan-1": {}}, + }, + })) + defer server.Close() + + resp, err := http.Get(server.URL + "/channel-awareness?channels=chan-1&since=24h&limit=12&max_chars=4096&context_kind=raw_window%2Bdigest") + if err != nil { + t.Fatalf("GET: %v", err) + } + defer resp.Body.Close() + body, err := io.ReadAll(resp.Body) + if err != nil { + t.Fatalf("read response: %v", err) + } + text := string(body) + if !strings.Contains(text, "kind=raw_window+digest") || !strings.Contains(text, "digest=ok") || !strings.Contains(text, "deterministic_only=true") { + t.Fatalf("expected digest awareness header, got %q", text) + } + if !strings.Contains(text, "digest_blocks=1") || !strings.Contains(text, "digest_source_messages=12") { + t.Fatalf("expected digest counts, got %q", text) + } + if !strings.Contains(text, "[digest kind=telemetry_count source_channel=chan-1 source_messages=100,101") { + t.Fatalf("expected digest block provenance, got %q", text) + } + if strings.Contains(text, "source=chan-1/100") || !strings.Contains(text, "source=chan-1/111") { + t.Fatalf("expected digest mode to keep only recent raw messages, got %q", text) + } + gotReqs := memory.requests() + if len(gotReqs) != 1 || gotReqs[0].SourceKind != channelMemorySourceKind || gotReqs[0].Since != "24h0m0s" || gotReqs[0].Budget.MaxBlocks != defaultDigestMaxBlocks { + t.Fatalf("unexpected digest request: %+v", gotReqs) + } + + req, err := http.NewRequest(http.MethodPost, server.URL+"/get_channel_messages", strings.NewReader(`{"channels":["chan-1"],"message_ids":["100","101"]}`)) + if err != nil { + t.Fatalf("build source retrieval request: %v", err) + } + req.Header.Set("Authorization", "Bearer tool-token") + req.Header.Set("X-Claw-ID", "trader-0") + retrievalResp, err := http.DefaultClient.Do(req) + if err != nil { + t.Fatalf("source retrieval: %v", err) + } + defer retrievalResp.Body.Close() + if retrievalResp.StatusCode != http.StatusOK { + body, _ := io.ReadAll(retrievalResp.Body) + t.Fatalf("expected source retrieval 200, got %d: %s", retrievalResp.StatusCode, string(body)) + } + var exact retrievalResult + if err := json.NewDecoder(retrievalResp.Body).Decode(&exact); err != nil { + t.Fatalf("decode source retrieval: %v", err) + } + if exact.Status != "ok" || len(exact.Messages) != 2 || exact.Messages[0].ID != "100" || exact.Messages[1].ID != "101" { + t.Fatalf("digest source provenance did not round-trip to exact messages: %+v", exact) + } +} + +func TestChannelAwarenessDigestUnavailableFallsBackToRawWindow(t *testing.T) { + store := newConversationStore(50) + store.now = func() time.Time { return time.Unix(212, 0) } + for i := 0; i < 12; i++ { + store.merge("chan-1", []wallMessage{{ + ID: fmt.Sprintf("%d", 200+i), + Author: "alice", + Content: fmt.Sprintf("raw-message-%03d", i), + Timestamp: time.Unix(int64(200+i), 0), + }}) + } + + memory := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + http.Error(w, "digest unavailable", http.StatusBadGateway) + })) + defer memory.Close() + client, err := newChannelMemoryClientWithDigest("", memory.URL+"/digest", "", time.Second) + if err != nil { + t.Fatalf("newChannelMemoryClientWithDigest: %v", err) + } + + server := httptest.NewServer(newHandler(store, handlerConfig{channelMemory: client})) + defer server.Close() + + resp, err := http.Get(server.URL + "/channel-awareness?channels=chan-1&since=24h&limit=12&max_chars=4096&context_kind=raw_window%2Bdigest") + if err != nil { + t.Fatalf("GET: %v", err) + } + defer resp.Body.Close() + body, err := io.ReadAll(resp.Body) + if err != nil { + t.Fatalf("read response: %v", err) + } + text := string(body) + if !strings.Contains(text, "digest=unavailable") || !strings.Contains(text, "raw-message-000") || !strings.Contains(text, "raw-message-011") { + t.Fatalf("expected full raw fallback on digest failure, got %q", text) + } +} + +func TestChannelAwarenessDigestStaleAndCoverageGapStatuses(t *testing.T) { + store := newConversationStore(50) + store.merge("chan-1", []wallMessage{ + {ID: "300", Author: "alice", Content: "raw fallback", Timestamp: time.Unix(300, 0)}, + }) + + for _, tc := range []struct { + name string + response channelMemoryDigestResponse + want string + }{ + { + name: "stale", + response: channelMemoryDigestResponse{ + Status: "ok", + Blocks: []channelMemoryDigestBlock{{ + Kind: "raw_excerpt", + Text: "stale digest", + SourceChannel: "chan-1", + SourceMessages: []string{"299"}, + Stale: true, + Processor: "deterministic", + }}, + Cost: channelMemoryDigestCost{DeterministicOnly: true}, + }, + want: "digest=stale", + }, + { + name: "coverage-gap", + response: channelMemoryDigestResponse{ + Status: "coverage_gap", + Coverage: channelMemoryDigestCoverage{ + Gaps: []channelMemoryCoverageGap{{ChannelID: "chan-1", From: "2026-05-21T19:00:00Z", To: "2026-05-21T19:15:00Z"}}, + }, + Blocks: []channelMemoryDigestBlock{{ + Kind: "raw_excerpt", + Text: "partial digest", + SourceChannel: "chan-1", + SourceMessages: []string{"299"}, + Processor: "deterministic", + }}, + Cost: channelMemoryDigestCost{DeterministicOnly: true}, + }, + want: "digest=coverage_gap", + }, + } { + t.Run(tc.name, func(t *testing.T) { + memory := newDigestChannelMemoryServer(t, tc.response) + defer memory.Close() + client, err := newChannelMemoryClientWithDigest("", memory.URL+"/digest", "", time.Second) + if err != nil { + t.Fatalf("newChannelMemoryClientWithDigest: %v", err) + } + server := httptest.NewServer(newHandler(store, handlerConfig{channelMemory: client})) + defer server.Close() + resp, err := http.Get(server.URL + "/channel-awareness?channels=chan-1&context_kind=raw_window%2Bdigest") + if err != nil { + t.Fatalf("GET: %v", err) + } + defer resp.Body.Close() + body, err := io.ReadAll(resp.Body) + if err != nil { + t.Fatalf("read response: %v", err) + } + if !strings.Contains(string(body), tc.want) { + t.Fatalf("expected %s in response, got %q", tc.want, string(body)) + } + }) + } +} + func TestToolSearchRequiresAuthAndChannelAllowlist(t *testing.T) { store := newConversationStore(50) store.merge("chan-1", []wallMessage{ @@ -530,6 +741,7 @@ func TestLoadConfigDefaultsPollIntervalToThirtySeconds(t *testing.T) { t.Setenv("CLAW_WALL_BACKFILL_MAX_PAGES", "") t.Setenv("CLAW_WALL_DISCORD_BASE_URL", "") t.Setenv("CLAW_WALL_CHANNEL_MEMORY_INGEST_URL", "") + t.Setenv("CLAW_WALL_CHANNEL_MEMORY_DIGEST_URL", "") t.Setenv("CLAW_WALL_CHANNEL_MEMORY_TOKEN", "") t.Setenv("CLAW_WALL_CHANNEL_MEMORY_TIMEOUT", "") t.Setenv("CLAW_WALL_TOKENS", "chan-1:token-a") @@ -553,7 +765,7 @@ func TestLoadConfigDefaultsPollIntervalToThirtySeconds(t *testing.T) { if cfg.DiscordBaseURL != "" { t.Fatalf("expected empty discord base url by default, got %q", cfg.DiscordBaseURL) } - if cfg.ChannelMemoryIngestURL != "" || cfg.ChannelMemoryToken != "" || cfg.ChannelMemoryTimeout != 2*time.Second { + if cfg.ChannelMemoryIngestURL != "" || cfg.ChannelMemoryDigestURL != "" || cfg.ChannelMemoryToken != "" || cfg.ChannelMemoryTimeout != 2*time.Second { t.Fatalf("unexpected channel-memory defaults: %+v", cfg) } } @@ -565,6 +777,7 @@ func TestLoadConfigReadsDiscordBaseURLOverride(t *testing.T) { t.Setenv("CLAW_WALL_RETENTION", "") t.Setenv("CLAW_WALL_BACKFILL_MAX_PAGES", "") t.Setenv("CLAW_WALL_DISCORD_BASE_URL", "http://fake-discord:9000/api/v10") + t.Setenv("CLAW_WALL_CHANNEL_MEMORY_DIGEST_URL", "") t.Setenv("CLAW_WALL_CHANNEL_MEMORY_TIMEOUT", "") t.Setenv("CLAW_WALL_TOKENS", "chan-1:token-a") @@ -585,6 +798,7 @@ func TestLoadConfigReadsChannelMemoryConfig(t *testing.T) { t.Setenv("CLAW_WALL_BACKFILL_MAX_PAGES", "") t.Setenv("CLAW_WALL_DISCORD_BASE_URL", "") t.Setenv("CLAW_WALL_CHANNEL_MEMORY_INGEST_URL", "http://channel-memory:8080/ingest") + t.Setenv("CLAW_WALL_CHANNEL_MEMORY_DIGEST_URL", "http://channel-memory:8080/digest") t.Setenv("CLAW_WALL_CHANNEL_MEMORY_TOKEN", "memory-token") t.Setenv("CLAW_WALL_CHANNEL_MEMORY_TIMEOUT", "750ms") t.Setenv("CLAW_WALL_TOKENS", "chan-1:token-a") @@ -593,7 +807,7 @@ func TestLoadConfigReadsChannelMemoryConfig(t *testing.T) { if err != nil { t.Fatalf("loadConfig: %v", err) } - if cfg.ChannelMemoryIngestURL != "http://channel-memory:8080/ingest" || cfg.ChannelMemoryToken != "memory-token" || cfg.ChannelMemoryTimeout != 750*time.Millisecond { + if cfg.ChannelMemoryIngestURL != "http://channel-memory:8080/ingest" || cfg.ChannelMemoryDigestURL != "http://channel-memory:8080/digest" || cfg.ChannelMemoryToken != "memory-token" || cfg.ChannelMemoryTimeout != 750*time.Millisecond { t.Fatalf("unexpected channel-memory config: %+v", cfg) } } @@ -604,6 +818,8 @@ func TestLoadConfigUsesPollIntervalOverride(t *testing.T) { t.Setenv("CLAW_WALL_POLL_INTERVAL", "42") t.Setenv("CLAW_WALL_RETENTION", "6h") t.Setenv("CLAW_WALL_BACKFILL_MAX_PAGES", "7") + t.Setenv("CLAW_WALL_CHANNEL_MEMORY_INGEST_URL", "") + t.Setenv("CLAW_WALL_CHANNEL_MEMORY_DIGEST_URL", "") t.Setenv("CLAW_WALL_CHANNEL_MEMORY_TIMEOUT", "") t.Setenv("CLAW_WALL_TOKENS", "chan-1:token-a") @@ -1183,6 +1399,44 @@ func (s *recordingChannelMemoryServer) requests() []channelMemoryIngestRequest { return out } +type digestChannelMemoryServer struct { + *httptest.Server + mu sync.Mutex + received []channelMemoryDigestRequest +} + +func newDigestChannelMemoryServer(t *testing.T, response channelMemoryDigestResponse) *digestChannelMemoryServer { + t.Helper() + srv := &digestChannelMemoryServer{} + srv.Server = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodPost || r.URL.Path != "/digest" { + http.Error(w, "unexpected request", http.StatusNotFound) + return + } + var payload channelMemoryDigestRequest + if err := json.NewDecoder(r.Body).Decode(&payload); err != nil { + http.Error(w, "invalid json", http.StatusBadRequest) + return + } + srv.mu.Lock() + srv.received = append(srv.received, payload) + srv.mu.Unlock() + w.Header().Set("Content-Type", "application/json") + if err := json.NewEncoder(w).Encode(response); err != nil { + http.Error(w, "encode digest response", http.StatusInternalServerError) + } + })) + return srv +} + +func (s *digestChannelMemoryServer) requests() []channelMemoryDigestRequest { + s.mu.Lock() + defer s.mu.Unlock() + out := make([]channelMemoryDigestRequest, len(s.received)) + copy(out, s.received) + return out +} + func makeDiscordMessages(start time.Time, count int, step time.Duration) []discordAPIMessage { messages := make([]discordAPIMessage, 0, count) for i := 0; i < count; i++ { diff --git a/cmd/claw-wall/store.go b/cmd/claw-wall/store.go index 3d3e3fb..592dd0b 100644 --- a/cmd/claw-wall/store.go +++ b/cmd/claw-wall/store.go @@ -13,11 +13,13 @@ import ( ) const ( - defaultResponseLimit = 20 - backgroundContextSize = 10 - defaultTailLimit = 40 - defaultAwarenessLimit = 60 - defaultTailMaxChars = 32 * 1024 + defaultResponseLimit = 20 + backgroundContextSize = 10 + defaultTailLimit = 40 + defaultAwarenessLimit = 60 + defaultTailMaxChars = 32 * 1024 + defaultDigestRawLimit = 10 + defaultDigestMaxBlocks = 32 ) const ( @@ -133,6 +135,18 @@ type channelMemoryReplayResponse struct { Pushed int `json:"pushed"` } +type channelAwarenessDigest struct { + Requested bool + Status string + GeneratedAt string + SourceMessages int + DigestMessages int + RawRecentMessages int + CoverageGaps int + DeterministicOnly bool + Blocks []channelMemoryDigestBlock +} + func newConversationStore(limit int, retentions ...time.Duration) *conversationStore { var retention time.Duration if len(retentions) > 0 { @@ -778,15 +792,24 @@ func newHandler(store *conversationStore, cfgs ...handlerConfig) http.Handler { } maxChars = parsed } + digest := channelAwarenessDigest{} + contextKind := channelAwarenessContextKindFromRequest(r) + if contextKind == "raw_window+digest" { + digest = cfg.channelMemory.fetchAwarenessDigest(r.Context(), channelIDs, since) + } + rawLimit := limit + if contextKind == "raw_window+digest" && len(digest.Blocks) > 0 && defaultDigestRawLimit < rawLimit { + rawLimit = defaultDigestRawLimit + } result := store.tail(tailRequest{ ChannelIDs: channelIDs, Since: since, - Limit: limit, + Limit: rawLimit, MaxChars: maxChars, Now: store.currentTime(), }) w.Header().Set("Content-Type", "text/plain; charset=utf-8") - _, _ = w.Write([]byte(formatChannelAwareness(result, since))) + _, _ = w.Write([]byte(formatChannelAwareness(result, since, contextKind, digest))) }) mux.HandleFunc("/search_channel_context", func(w http.ResponseWriter, r *http.Request) { if !authorizeToolRequest(w, r, cfg) { @@ -976,9 +999,28 @@ func formatTailContext(result tailResult, since time.Duration, kind string) stri return b.String() } -func formatChannelAwareness(result tailResult, since time.Duration) string { +func channelAwarenessContextKindFromRequest(r *http.Request) string { + switch strings.TrimSpace(r.URL.Query().Get("context_kind")) { + case "raw_window+digest", "raw_window digest": + return "raw_window+digest" + default: + return "raw_window" + } +} + +func formatChannelAwareness(result tailResult, since time.Duration, contextKind string, digest channelAwarenessDigest) string { + if contextKind != "raw_window+digest" { + contextKind = "raw_window" + } var b strings.Builder - fmt.Fprintf(&b, "[channel-awareness] kind=raw_window since=%s channels=%s messages=%d available=%d omitted=%d retained=%d/since-%s digest=unavailable", + rawBody := formatWallMessages(result.Messages) + digestBody := formatDigestBlocks(digest.Blocks) + digestStatus := "unavailable" + if digest.Requested { + digestStatus = normalizeDigestStatus(digest.Status) + } + fmt.Fprintf(&b, "[channel-awareness] kind=%s since=%s channels=%s messages=%d available=%d omitted=%d retained=%d/since-%s digest=%s", + contextKind, formatDurationForHeader(since), strings.Join(result.ChannelIDs, ","), len(result.Messages), @@ -986,7 +1028,20 @@ func formatChannelAwareness(result tailResult, since time.Duration) string { result.Omitted, result.Available, formatDurationForHeader(since), + digestStatus, ) + fmt.Fprintf(&b, " raw_bytes=%d digest_bytes=%d", len(rawBody), len(digestBody)) + if contextKind == "raw_window+digest" { + fmt.Fprintf(&b, " digest_blocks=%d digest_source_messages=%d coverage_gaps=%d deterministic_only=%t", + len(digest.Blocks), + digest.SourceMessages, + digest.CoverageGaps, + digest.DeterministicOnly, + ) + if digest.GeneratedAt != "" { + fmt.Fprintf(&b, " digest_generated_at=%s", digest.GeneratedAt) + } + } if result.CapReason != "" { fmt.Fprintf(&b, " cap=%s", result.CapReason) } @@ -999,13 +1054,71 @@ func formatChannelAwareness(result tailResult, since time.Duration) string { fmt.Fprintf(&b, "[omitted %d older retained messages due to %s; newest retained messages follow]\n", result.Omitted, result.CapReason) } - if len(result.Messages) > 0 { + if digestBody != "" { + b.WriteByte('\n') + b.WriteString("[digest]\n") + b.WriteString(digestBody) b.WriteByte('\n') - b.WriteString(formatWallMessages(result.Messages)) + } + if rawBody != "" { + b.WriteByte('\n') + b.WriteString("[raw recent]\n") + b.WriteString(rawBody) } return b.String() } +func formatDigestBlocks(blocks []channelMemoryDigestBlock) string { + var b strings.Builder + for _, block := range blocks { + text := strings.TrimSpace(block.Text) + if text == "" { + continue + } + if b.Len() > 0 { + b.WriteByte('\n') + } + fmt.Fprintf(&b, "[digest kind=%s source_channel=%s source_messages=%s window=%s..%s sparse=%t score=%.2f processor=%s] %s", + headerToken(block.Kind), + headerToken(block.SourceChannel), + headerToken(strings.Join(trimNonEmptyStrings(block.SourceMessages), ",")), + headerToken(block.SourceWindow.From), + headerToken(block.SourceWindow.To), + block.Sparse, + block.Score, + headerToken(block.Processor), + text, + ) + } + return b.String() +} + +func trimNonEmptyStrings(values []string) []string { + out := make([]string, 0, len(values)) + seen := make(map[string]struct{}, len(values)) + for _, value := range values { + value = strings.TrimSpace(value) + if value == "" { + continue + } + if _, ok := seen[value]; ok { + continue + } + seen[value] = struct{}{} + out = append(out, value) + } + return out +} + +func headerToken(value string) string { + value = strings.TrimSpace(value) + if value == "" { + return "none" + } + replacer := strings.NewReplacer(" ", "_", "\t", "_", "\n", "_", "\r", "_") + return replacer.Replace(value) +} + func parseDurationQuery(rawSince, rawWindow string) (time.Duration, error) { raw := strings.TrimSpace(rawSince) if raw == "" { diff --git a/cmd/claw/compose_up.go b/cmd/claw/compose_up.go index 982d6dc..1df8f32 100644 --- a/cmd/claw/compose_up.go +++ b/cmd/claw/compose_up.go @@ -56,9 +56,11 @@ const ( conversationWallBackfillPages = "25" conversationWallInternalPort = "8080" conversationWallDockerfile = "dockerfiles/claw-wall/Dockerfile" + conversationWallDiscordBaseEnv = "CLAW_WALL_DISCORD_BASE_URL" conversationWallToolTokenEnv = "CLAW_WALL_TOOL_TOKEN" conversationWallAllowlistPath = "/etc/claw-wall/agent-channels.json" conversationWallMemoryIngestEnv = "CLAW_WALL_CHANNEL_MEMORY_INGEST_URL" + conversationWallMemoryDigestEnv = "CLAW_WALL_CHANNEL_MEMORY_DIGEST_URL" conversationWallMemoryTimeout = "2s" clawInternalNetworkName = "claw-internal" historyReplayAuthService = "cllama-history" @@ -1813,7 +1815,7 @@ func injectConversationWall(p *pod.Pod, resolvedClaws map[string]*driver.Resolve } settings := conversationWallChannelContext(p, svc) awarenessSettings := conversationWallChannelAwarenessContext(p, svc) - svc.Claw.Feeds = appendConversationWallAwarenessFeed(svc.Claw.Feeds, channelIDs, awarenessSettings) + svc.Claw.Feeds = appendConversationWallAwarenessFeed(svc.Claw.Feeds, channelIDs, awarenessSettings, p.ChannelMemory != nil) svc.Claw.Feeds = appendConversationWallFeed(svc.Claw.Feeds, channelIDs, settings) svc.Claw.Tools = appendConversationWallToolPolicy(svc.Claw.Tools) } @@ -1825,6 +1827,9 @@ func injectConversationWall(p *pod.Pod, resolvedClaws map[string]*driver.Resolve "CLAW_WALL_RETENTION": envOrDefault("CLAW_WALL_RETENTION", conversationWallRetention), "CLAW_WALL_BACKFILL_MAX_PAGES": envOrDefault("CLAW_WALL_BACKFILL_MAX_PAGES", conversationWallBackfillPages), } + if discordBaseURL := strings.TrimSpace(os.Getenv(conversationWallDiscordBaseEnv)); discordBaseURL != "" { + wallEnv[conversationWallDiscordBaseEnv] = discordBaseURL + } if err := configureConversationWallChannelMemory(p, wallEnv); err != nil { return err } @@ -1868,6 +1873,7 @@ func configureConversationWallChannelMemory(p *pod.Pod, wallEnv map[string]strin return fmt.Errorf("channel-memory service %q: %w", serviceName, err) } wallEnv[conversationWallMemoryIngestEnv] = buildFeedURL(baseURL, "/ingest") + wallEnv[conversationWallMemoryDigestEnv] = buildFeedURL(baseURL, "/digest") wallEnv["CLAW_WALL_CHANNEL_MEMORY_TIMEOUT"] = envOrDefault("CLAW_WALL_CHANNEL_MEMORY_TIMEOUT", conversationWallMemoryTimeout) targetSvc := p.Services[serviceName] if targetSvc.Environment == nil { @@ -2007,13 +2013,19 @@ func appendConversationWallFeed(feeds []pod.FeedEntry, channelIDs []string, sett }) } -func appendConversationWallAwarenessFeed(feeds []pod.FeedEntry, channelIDs []string, settings conversationWallContextSettings) []pod.FeedEntry { +func appendConversationWallAwarenessFeed(feeds []pod.FeedEntry, channelIDs []string, settings conversationWallContextSettings, channelMemory bool) []pod.FeedEntry { + contextKind := "raw_window" + if channelMemory { + contextKind = "raw_window+digest" + } + escapedContextKind := strings.ReplaceAll(contextKind, "+", "%2B") path := fmt.Sprintf( - "/channel-awareness?channels=%s&since=%s&limit=%d&max_chars=%d&context_kind=raw_window", + "/channel-awareness?channels=%s&since=%s&limit=%d&max_chars=%d&context_kind=%s", strings.Join(channelIDs, ","), settings.Since, settings.Limit, settings.MaxChars, + escapedContextKind, ) for _, feed := range feeds { if feed.Name == conversationWallAwarenessName && feed.Source == conversationWallServiceName && feed.Path == path { diff --git a/cmd/claw/compose_up_test.go b/cmd/claw/compose_up_test.go index d9a61c5..a2cd40e 100644 --- a/cmd/claw/compose_up_test.go +++ b/cmd/claw/compose_up_test.go @@ -3880,6 +3880,7 @@ func testFeedByName(t *testing.T, feeds []pod.FeedEntry, name string) pod.FeedEn } func TestInjectConversationWallAddsServiceAndFeed(t *testing.T) { + t.Setenv(conversationWallDiscordBaseEnv, "http://fake-discord:8090") t.Setenv("CLAW_WALL_POLL_INTERVAL", "") t.Setenv("CLAW_WALL_RETENTION", "") t.Setenv("CLAW_WALL_BACKFILL_MAX_PAGES", "") @@ -3929,6 +3930,9 @@ func TestInjectConversationWallAddsServiceAndFeed(t *testing.T) { if wall.Environment["CLAW_WALL_BACKFILL_MAX_PAGES"] != conversationWallBackfillPages { t.Fatalf("unexpected CLAW_WALL_BACKFILL_MAX_PAGES: %q", wall.Environment["CLAW_WALL_BACKFILL_MAX_PAGES"]) } + if wall.Environment[conversationWallDiscordBaseEnv] != "http://fake-discord:8090" { + t.Fatalf("unexpected %s: %q", conversationWallDiscordBaseEnv, wall.Environment[conversationWallDiscordBaseEnv]) + } traderFeeds := p.Services["trader"].Claw.Feeds if len(traderFeeds) != 2 { @@ -4060,6 +4064,9 @@ func TestInjectConversationWallWiresChannelMemory(t *testing.T) { if wall.Environment[conversationWallMemoryIngestEnv] != "http://channel-memory:8080/ingest" { t.Fatalf("unexpected channel-memory ingest URL: %q", wall.Environment[conversationWallMemoryIngestEnv]) } + if wall.Environment[conversationWallMemoryDigestEnv] != "http://channel-memory:8080/digest" { + t.Fatalf("unexpected channel-memory digest URL: %q", wall.Environment[conversationWallMemoryDigestEnv]) + } if wall.Environment["CLAW_WALL_CHANNEL_MEMORY_TIMEOUT"] != conversationWallMemoryTimeout { t.Fatalf("unexpected channel-memory timeout: %q", wall.Environment["CLAW_WALL_CHANNEL_MEMORY_TIMEOUT"]) } @@ -4069,6 +4076,10 @@ func TestInjectConversationWallWiresChannelMemory(t *testing.T) { if p.Services["channel-memory"].Environment["CHANNEL_MEMORY_TOKEN"] != wall.Environment["CLAW_WALL_CHANNEL_MEMORY_TOKEN"] { t.Fatalf("expected channel-memory token to match claw-wall token") } + awarenessFeed := testFeedByName(t, p.Services["trader"].Claw.Feeds, conversationWallAwarenessName) + if awarenessFeed.Path != "/channel-awareness?channels=chan-1&since=24h&limit=60&max_chars=32768&context_kind=raw_window%2Bdigest" { + t.Fatalf("unexpected digest-backed awareness feed path: %q", awarenessFeed.Path) + } networks, ok := p.Services["channel-memory"].Compose["networks"].([]string) if !ok || len(networks) != 1 || networks[0] != clawInternalNetworkName { t.Fatalf("expected channel-memory on %s, got %#v", clawInternalNetworkName, p.Services["channel-memory"].Compose["networks"]) diff --git a/cmd/claw/skill_data/SKILL.md b/cmd/claw/skill_data/SKILL.md index 14882a1..7c6c23d 100644 --- a/cmd/claw/skill_data/SKILL.md +++ b/cmd/claw/skill_data/SKILL.md @@ -389,7 +389,7 @@ The proxy sits between agents and LLM providers. Agents get bearer tokens, proxy ### claw-wall sidecar -Auto-injected by `claw up` when any cllama-enabled service has Discord channel IDs. Polls Discord channels and serves the recent channel transcript to agents through `channel-context` tail feeds; legacy unread-mailbox cursor paging remains available as `mode=delta`. On startup, wall backfills Discord history before its first forward poll up to `CLAW_WALL_RETENTION` (default `24h`) and `CLAW_WALL_BACKFILL_MAX_PAGES` (default `25`), while `CLAW_WALL_LIMIT` is a per-channel safety cap (default `5000`). Configure the generated tail window with pod or service `x-claw.context.channel` (`since`, `limit`, `max-chars`, `buffer`). Since `v0.15.0` channel-consuming services also get a default-on `channel-awareness` feed (uncursored 24h raw window; `x-claw.context.channel.max-chars` tunes both feeds together, default 32 KB) plus two cllama-mediated retrieval tools - `search_channel_context` and `get_channel_messages` - auto-subscribed via a compiler-owned claw-wall descriptor. Feed headers include `backfill_status`; `partial` or `rate_limited` means the backing window did not fully satisfy the requested horizon. Calls are gated by a generated per-agent channel allowlist, claw-wall service-token auth, and forwarded `X-Claw-ID`. The service name `claw-wall` is reserved - declaring it in `claw-pod.yml` is a hard error. +Auto-injected by `claw up` when any cllama-enabled service has Discord channel IDs. Polls Discord channels and serves the recent channel transcript to agents through `channel-context` tail feeds; legacy unread-mailbox cursor paging remains available as `mode=delta`. On startup, wall backfills Discord history before its first forward poll up to `CLAW_WALL_RETENTION` (default `24h`) and `CLAW_WALL_BACKFILL_MAX_PAGES` (default `25`), while `CLAW_WALL_LIMIT` is a per-channel safety cap (default `5000`). Configure the generated tail window with pod or service `x-claw.context.channel` (`since`, `limit`, `max-chars`, `buffer`). Since `v0.15.0` channel-consuming services also get a default-on `channel-awareness` feed plus two cllama-mediated retrieval tools - `search_channel_context` and `get_channel_messages` - auto-subscribed via a compiler-owned claw-wall descriptor. Without `x-claw.channel-memory.service`, `channel-awareness` is an uncursored raw window. With channel-memory wired, `claw up` changes that feed to `context_kind=raw_window+digest`; claw-wall fetches digest blocks from channel-memory, emits a compact `[digest]` section with source provenance, and keeps a bounded `[raw recent]` tail. Feed headers include `backfill_status`, `digest`, `raw_bytes`, `digest_bytes`, `digest_blocks`, `coverage_gaps`, and `deterministic_only`; `partial` or `rate_limited` backfill means the backing window did not fully satisfy the requested horizon. Calls are gated by a generated per-agent channel allowlist, claw-wall service-token auth, and forwarded `X-Claw-ID`. The service name `claw-wall` is reserved - declaring it in `claw-pod.yml` is a hard error. ### cllama feed injection budgets diff --git a/cmd/claw/spike_channel_digest_test.go b/cmd/claw/spike_channel_digest_test.go new file mode 100644 index 0000000..3908410 --- /dev/null +++ b/cmd/claw/spike_channel_digest_test.go @@ -0,0 +1,457 @@ +//go:build spike + +package main + +import ( + "context" + "encoding/json" + "fmt" + "os" + "os/exec" + "path/filepath" + "runtime" + "strconv" + "strings" + "testing" + "time" +) + +// TestSpikeChannelDigestGeneratedPod runs a generated pod with fake Discord, +// real claw-wall, real channel-memory, real cllama, and a fake upstream +// provider. It proves #267's generated-compose seam and captures the +// provider-visible raw_window+digest block that cllama injects. +// +// Requires: Docker. Does NOT require Discord or provider credentials. +// Run with: go test -tags spike -v -run TestSpikeChannelDigestGeneratedPod -timeout 10m ./cmd/claw/... +func TestSpikeChannelDigestGeneratedPod(t *testing.T) { + if _, err := exec.LookPath("docker"); err != nil { + t.Skip("docker not on PATH - skipping") + } + + _, thisFile, _, ok := runtime.Caller(0) + if !ok { + t.Fatal("runtime.Caller failed") + } + repoRoot, err := filepath.Abs(filepath.Join(filepath.Dir(thisFile), "..", "..")) + if err != nil { + t.Fatalf("resolve repo root: %v", err) + } + + workDir := t.TempDir() + generatedPath := filepath.Join(workDir, "compose.generated.yml") + projectName := "channel-digest-spike" + networkName := projectName + "_claw-internal" + spikeCleanupProject(projectName, generatedPath) + t.Cleanup(func() { + spikeCleanupProject(projectName, generatedPath) + }) + t.Cleanup(func() { + if !t.Failed() { + return + } + out, _ := exec.Command("docker", "compose", "-f", generatedPath, "logs", "--tail", "200").CombinedOutput() + t.Logf("compose logs:\n%s", out) + }) + + channelID := "267000000000000001" + agentImage := "channel-digest-agent:spike-267" + channelMemoryImage := "channel-memory:spike-267" + pythonImage := "python:3.12-alpine" + + // The generated pod points at the normal infra image refs, so overwrite + // those tags with images built from the worktree under test. + spikeEnsurePulledImage(t, pythonImage) + spikeBuildImage(t, repoRoot, resolveConversationWallImageRef(), conversationWallDockerfile) + spikeEnsureRepoInfraImages(t, repoRoot, infraComponentClawdash) + spikeEnsureCllamaPassthroughImage(t, repoRoot) + spikeBuildImage(t, repoRoot, channelMemoryImage, "examples/channel-memory/Dockerfile") + + rollcallDir := filepath.Join(repoRoot, "examples", "rollcall") + spikeBuildImage(t, rollcallDir, "nullclaw:latest", "Dockerfile.nullclaw-base") + spikeWriteFile(t, filepath.Join(workDir, "AGENTS.md"), "# Digest Spike Agent\n\nUse runtime channel context.") + spikeWriteFile(t, filepath.Join(workDir, "Clawfile"), `FROM nullclaw:latest + +CLAW_TYPE nullclaw +AGENT AGENTS.md +MODEL primary openai/gpt-4o +HANDLE discord +`) + spikeBuildImage(t, workDir, agentImage, "Clawfile") + + capturesDir := filepath.Join(workDir, "captures") + if err := os.MkdirAll(capturesDir, 0o777); err != nil { + t.Fatalf("create captures dir: %v", err) + } + spikeWriteFile(t, filepath.Join(workDir, "fake_discord.py"), fakeDiscordDigestSpikeScript()) + spikeWriteFile(t, filepath.Join(workDir, "fake_provider.py"), fakeProviderDigestSpikeScript()) + spikeWriteFile(t, filepath.Join(workDir, "claw-pod.yml"), fmt.Sprintf(`name: channel-digest-spike + +x-claw: + pod: channel-digest-spike + channel-memory: + service: channel-memory + +services: + fake-discord: + image: %s + command: ["python", "/app/fake_discord.py"] + environment: + CHANNEL_ID: "%s" + MESSAGE_COUNT: "30" + volumes: + - ./fake_discord.py:/app/fake_discord.py:ro + expose: + - "8090" + networks: + - claw-internal + + fake-provider: + image: %s + command: ["python", "/app/fake_provider.py"] + volumes: + - ./fake_provider.py:/app/fake_provider.py:ro + - ./captures:/captures:rw + expose: + - "8080" + networks: + - claw-internal + + channel-memory: + image: %s + expose: + - "8080" + + digest-agent: + image: %s + environment: + DISCORD_BOT_TOKEN: fake-token + x-claw: + agent: AGENTS.md + cllama: passthrough + cllama-env: + OPENAI_API_KEY: sk-fake + OPENAI_BASE_URL: http://fake-provider:8080/v1 + CLLAMA_FEED_MAX_RESPONSE_BYTES: "200000" + CLLAMA_FEED_MAX_TOTAL_BYTES: "400000" + models: + primary: openai/gpt-4o + handles: + discord: + id: digest-bot + guilds: + - id: guild-1 + channels: + - id: "%s" +`, pythonImage, channelID, pythonImage, channelMemoryImage, agentImage, channelID)) + + t.Setenv(conversationWallDiscordBaseEnv, "http://fake-discord:8090") + t.Setenv("CLAW_WALL_RETENTION", "24h") + t.Setenv("CLAW_WALL_BACKFILL_MAX_PAGES", "2") + t.Setenv("CLAW_WALL_POLL_INTERVAL", "1") + t.Setenv("CLAW_WALL_CHANNEL_MEMORY_TIMEOUT", "5s") + t.Setenv("CLLAMA_UI_PORT", spikeFreePort(t)) + t.Setenv("CLAWDASH_ADDR", ":"+spikeFreePort(t)) + + prevDetach := composeUpDetach + composeUpDetach = true + defer func() { composeUpDetach = prevDetach }() + + if err := runComposeUp(filepath.Join(workDir, "claw-pod.yml")); err != nil { + t.Fatalf("runComposeUp: %v", err) + } + + composeText := spikeReadFile(t, generatedPath) + if !strings.Contains(composeText, conversationWallMemoryDigestEnv+": http://channel-memory:8080/digest") { + t.Fatalf("compose did not wire claw-wall digest URL:\n%s", composeText) + } + if !strings.Contains(composeText, conversationWallDiscordBaseEnv+": http://fake-discord:8090") { + t.Fatalf("compose did not forward fake Discord base URL:\n%s", composeText) + } + feedsText := spikeReadFile(t, filepath.Join(workDir, ".claw-runtime", "context", "digest-agent", "feeds.json")) + if !strings.Contains(feedsText, "context_kind=raw_window%2Bdigest") { + t.Fatalf("generated feeds.json did not request digest-backed awareness:\n%s", feedsText) + } + + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Minute) + defer cancel() + rawURL := "http://claw-wall:8080/channel-awareness?channels=" + channelID + "&since=24h&limit=60&max_chars=200000&context_kind=raw_window" + digestURL := "http://claw-wall:8080/channel-awareness?channels=" + channelID + "&since=24h&limit=60&max_chars=200000&context_kind=raw_window%2Bdigest" + + var ( + rawBody string + digestBody string + lastProbe string + ) + if err := waitForCondition(ctx, func() bool { + var err error + rawBody, err = spikeDockerProbe(networkName, pythonImage, probeGetURLScript, rawURL) + if err != nil { + lastProbe = err.Error() + return false + } + digestBody, err = spikeDockerProbe(networkName, pythonImage, probeGetURLScript, digestURL) + if err != nil { + lastProbe = err.Error() + return false + } + lastProbe = digestBody + messages, hasMessages := spikeHeaderInt(digestBody, "messages") + return hasMessages && + messages > 0 && + strings.Contains(digestBody, "digest=ok") && + !strings.Contains(digestBody, "digest_blocks=0") + }); err != nil { + t.Fatalf("digest-backed awareness never became available: %v\nlast probe:\n%s", err, lastProbe) + } + + rawBytes, ok := spikeHeaderInt(rawBody, "raw_bytes") + if !ok || rawBytes <= 0 { + t.Fatalf("raw awareness missing raw_bytes header:\n%s", rawBody) + } + digestRawBytes, ok := spikeHeaderInt(digestBody, "raw_bytes") + if !ok || digestRawBytes <= 0 { + t.Fatalf("digest awareness missing raw_bytes header:\n%s", digestBody) + } + digestBytes, ok := spikeHeaderInt(digestBody, "digest_bytes") + if !ok || digestBytes <= 0 { + t.Fatalf("digest awareness missing digest_bytes header:\n%s", digestBody) + } + if digestRawBytes >= rawBytes { + t.Fatalf("expected digest mode to reduce retained raw bytes, raw=%d digest-raw=%d\nraw:\n%s\ndigest:\n%s", rawBytes, digestRawBytes, rawBody, digestBody) + } + if len(digestBody) >= len(rawBody) { + t.Fatalf("expected digest-backed awareness body to be smaller than raw-only, raw=%d digest=%d\nraw:\n%s\ndigest:\n%s", len(rawBody), len(digestBody), rawBody, digestBody) + } + for _, want := range []string{"[channel-awareness] kind=raw_window+digest", "[digest]", "deterministic_only=true", "source_channel=" + channelID, "[raw recent]"} { + if !strings.Contains(digestBody, want) { + t.Fatalf("digest awareness missing %q:\n%s", want, digestBody) + } + } + + token := spikeReadAgentToken(t, filepath.Join(workDir, ".claw-runtime", "context", "digest-agent", "metadata.json")) + if out, err := spikeDockerProbe(networkName, pythonImage, probePostCllamaScript, token); err != nil { + t.Fatalf("cllama provider request failed: %v\n%s", err, out) + } else if !strings.Contains(out, "chatcmpl-spike") { + t.Fatalf("unexpected cllama response:\n%s", out) + } + + capturePath := filepath.Join(capturesDir, "latest.json") + if err := waitForCondition(ctx, func() bool { + _, err := os.Stat(capturePath) + return err == nil + }); err != nil { + t.Fatalf("fake provider did not write capture: %v", err) + } + captured := spikeReadFile(t, capturePath) + providerText := spikeOpenAIMessageText(t, captured) + for _, want := range []string{"BEGIN FEED: channel-awareness", "[channel-awareness] kind=raw_window+digest", "[digest]", "deterministic_only=true", "source_channel=" + channelID} { + if !strings.Contains(providerText, want) { + t.Fatalf("provider-visible request missing %q:\n%s", want, providerText) + } + } +} + +const probeGetURLScript = ` +import sys, urllib.request +with urllib.request.urlopen(sys.argv[1], timeout=5) as resp: + sys.stdout.write(resp.read().decode("utf-8")) +` + +const probePostCllamaScript = ` +import json, sys, urllib.request +payload = {"model":"openai/gpt-4o","messages":[{"role":"user","content":"Use the channel context."}]} +body = json.dumps(payload).encode("utf-8") +req = urllib.request.Request( + "http://cllama:8080/v1/chat/completions", + data=body, + headers={"Content-Type":"application/json","Authorization":"Bearer " + sys.argv[1]}, +) +with urllib.request.urlopen(req, timeout=15) as resp: + sys.stdout.write(resp.read().decode("utf-8")) +` + +func spikeDockerProbe(networkName, image, script string, args ...string) (string, error) { + cmdArgs := []string{"run", "--rm", "--network", networkName, image, "python", "-c", script} + cmdArgs = append(cmdArgs, args...) + cmd := exec.Command("docker", cmdArgs...) + out, err := cmd.CombinedOutput() + if err != nil { + return string(out), fmt.Errorf("%w: %s", err, strings.TrimSpace(string(out))) + } + return string(out), nil +} + +func spikeWriteFile(t *testing.T, path, content string) { + t.Helper() + if err := os.WriteFile(path, []byte(content), 0o644); err != nil { + t.Fatalf("write %s: %v", path, err) + } +} + +func spikeEnsurePulledImage(t *testing.T, image string) { + t.Helper() + if spikeImageExists(image) { + return + } + out, err := exec.Command("docker", "pull", image).CombinedOutput() + if err != nil { + t.Fatalf("docker pull %s: %v\n%s", image, err, out) + } +} + +func spikeReadAgentToken(t *testing.T, path string) string { + t.Helper() + var metadata struct { + Token string `json:"token"` + } + if err := json.Unmarshal([]byte(spikeReadFile(t, path)), &metadata); err != nil { + t.Fatalf("parse %s: %v", path, err) + } + if strings.TrimSpace(metadata.Token) == "" { + t.Fatalf("%s did not contain token", path) + } + return metadata.Token +} + +func spikeHeaderInt(body, key string) (int, bool) { + header, _, _ := strings.Cut(body, "\n") + for _, field := range strings.Fields(header) { + raw, ok := strings.CutPrefix(field, key+"=") + if !ok { + continue + } + value, err := strconv.Atoi(strings.TrimSpace(raw)) + return value, err == nil + } + return 0, false +} + +func spikeOpenAIMessageText(t *testing.T, raw string) string { + t.Helper() + var payload map[string]any + if err := json.Unmarshal([]byte(raw), &payload); err != nil { + t.Fatalf("parse captured provider request: %v\n%s", err, raw) + } + messages, _ := payload["messages"].([]any) + var b strings.Builder + for _, item := range messages { + msg, _ := item.(map[string]any) + if msg == nil { + continue + } + switch content := msg["content"].(type) { + case string: + b.WriteString(content) + b.WriteByte('\n') + case []any: + for _, block := range content { + m, _ := block.(map[string]any) + if text, _ := m["text"].(string); text != "" { + b.WriteString(text) + b.WriteByte('\n') + } + } + } + } + return b.String() +} + +func fakeDiscordDigestSpikeScript() string { + return `import json +import os +import urllib.parse +from datetime import datetime, timedelta, timezone +from http.server import BaseHTTPRequestHandler, HTTPServer + +channel_id = os.environ["CHANNEL_ID"] +count = int(os.environ.get("MESSAGE_COUNT", "30")) +base_id = 2670000000001000000 +now = datetime.now(timezone.utc).replace(microsecond=0) +messages = [] +for i in range(count): + ts = (now - timedelta(minutes=count - i)).isoformat().replace("+00:00", "Z") + content = "heartbeat_ok runtime status seq=%03d " % i + ("diagnostic-payload-%03d " % i) * 18 + messages.append({ + "id": str(base_id + i), + "content": content, + "timestamp": ts, + "author": {"id": "user-%03d" % i, "username": "ops-%02d" % (i % 3), "global_name": "Ops %02d" % (i % 3)}, + "channel_id": channel_id, + }) + +class Handler(BaseHTTPRequestHandler): + def do_GET(self): + parsed = urllib.parse.urlparse(self.path) + if parsed.path != "/channels/%s/messages" % channel_id: + self.send_response(404) + self.end_headers() + return + query = urllib.parse.parse_qs(parsed.query) + limit = min(100, max(1, int(query.get("limit", ["100"])[0]))) + after = query.get("after", [""])[0] + before = query.get("before", [""])[0] + selected = [] + for msg in messages: + if after and int(msg["id"]) <= int(after): + continue + if before and int(msg["id"]) >= int(before): + continue + selected.append(msg) + selected = selected[-limit:] + selected = list(reversed(selected)) + body = json.dumps(selected).encode("utf-8") + self.send_response(200) + self.send_header("Content-Type", "application/json") + self.send_header("Content-Length", str(len(body))) + self.end_headers() + self.wfile.write(body) + + def log_message(self, fmt, *args): + return + +HTTPServer(("0.0.0.0", 8090), Handler).serve_forever() +` +} + +func fakeProviderDigestSpikeScript() string { + return `import json +import os +from http.server import BaseHTTPRequestHandler, HTTPServer + +class Handler(BaseHTTPRequestHandler): + def do_GET(self): + if self.path == "/health": + body = b"ok" + self.send_response(200) + self.send_header("Content-Length", str(len(body))) + self.end_headers() + self.wfile.write(body) + return + self.send_response(404) + self.end_headers() + + def do_POST(self): + length = int(self.headers.get("Content-Length", "0")) + body = self.rfile.read(length) + os.makedirs("/captures", exist_ok=True) + with open("/captures/latest.json", "wb") as f: + f.write(body) + response = { + "id": "chatcmpl-spike", + "object": "chat.completion", + "choices": [{"index": 0, "message": {"role": "assistant", "content": "ok"}, "finish_reason": "stop"}], + "usage": {"prompt_tokens": 1, "completion_tokens": 1, "total_tokens": 2}, + } + encoded = json.dumps(response).encode("utf-8") + self.send_response(200) + self.send_header("Content-Type", "application/json") + self.send_header("Content-Length", str(len(encoded))) + self.end_headers() + self.wfile.write(encoded) + + def log_message(self, fmt, *args): + return + +HTTPServer(("0.0.0.0", 8080), Handler).serve_forever() +` +} diff --git a/examples/channel-memory/README.md b/examples/channel-memory/README.md index 8a756d6..ad66322 100644 --- a/examples/channel-memory/README.md +++ b/examples/channel-memory/README.md @@ -1,10 +1,11 @@ # Channel Memory Adapter -This is the first executable slice of the digest-backed channel-awareness design. -It is a standalone HTTP service; `claw-wall` does not call it yet. +This is the channel-memory adapter for digest-backed channel awareness. +When a pod declares `x-claw.channel-memory.service`, `claw up` wires `claw-wall` +to ingest retained channel messages and request digest blocks from this service. The adapter stores Discord-style channel messages in SQLite and keeps exact -source provenance for later digest-backed `channel-awareness` work. +source provenance for digest-backed `channel-awareness` work. ## Endpoints @@ -62,9 +63,10 @@ docker build -f examples/channel-memory/Dockerfile -t channel-memory:latest . ## Pod Wiring -Declare the adapter once at pod level. `claw up` injects the ingest URL and a -bearer token into `claw-wall`, and injects the matching `CHANNEL_MEMORY_TOKEN` -into the adapter service. +Declare the adapter once at pod level. `claw up` injects ingest and digest URLs +plus a bearer token into `claw-wall`, injects the matching +`CHANNEL_MEMORY_TOKEN` into the adapter service, and changes generated +`channel-awareness` feeds to request `context_kind=raw_window+digest`. ```yaml x-claw: diff --git a/internal/audit/event.go b/internal/audit/event.go index 2f75a69..75a60a5 100644 --- a/internal/audit/event.go +++ b/internal/audit/event.go @@ -35,6 +35,11 @@ type Event struct { Retained *int `json:"retained,omitempty"` Returned *int `json:"returned,omitempty"` Omitted *int `json:"omitted,omitempty"` + RawBytes *int `json:"raw_bytes,omitempty"` + DigestBytes *int `json:"digest_bytes,omitempty"` + DigestBlocks *int `json:"digest_blocks,omitempty"` + CoverageGaps *int `json:"coverage_gaps,omitempty"` + DeterministicOnly *bool `json:"deterministic_only,omitempty"` Status string `json:"status,omitempty"` // provider_pool event fields Provider string `json:"provider,omitempty"` diff --git a/internal/audit/normalize.go b/internal/audit/normalize.go index a0a7e08..7877c46 100644 --- a/internal/audit/normalize.go +++ b/internal/audit/normalize.go @@ -74,6 +74,21 @@ func NormalizeLine(line []byte) (*Event, error) { if value, ok := intField(raw, "omitted"); ok { event.Omitted = &value } + if value, ok := intField(raw, "raw_bytes"); ok { + event.RawBytes = &value + } + if value, ok := intField(raw, "digest_bytes"); ok { + event.DigestBytes = &value + } + if value, ok := intField(raw, "digest_blocks"); ok { + event.DigestBlocks = &value + } + if value, ok := intField(raw, "coverage_gaps"); ok { + event.CoverageGaps = &value + } + if value, ok := boolField(raw, "deterministic_only"); ok { + event.DeterministicOnly = &value + } if value, ok := float64Field(raw, "cost_usd"); ok { event.CostUSD = &value } diff --git a/internal/audit/normalize_test.go b/internal/audit/normalize_test.go index 4cf33a6..f2d5e34 100644 --- a/internal/audit/normalize_test.go +++ b/internal/audit/normalize_test.go @@ -81,12 +81,12 @@ func TestNormalizeLineParseFeedFetchEvent(t *testing.T) { } func TestNormalizeLineParseChannelContextOpEvent(t *testing.T) { - line := `{"ts":"2026-05-12T10:00:00Z","claw_id":"weston","type":"channel_context_op","kind":"raw_window","channels":["chan-1","chan-2"],"retained":60,"returned":40,"omitted":20,"source":"claw-wall","status":"ok","tool_name":"search_channel_context","status_code":200}` + line := `{"ts":"2026-05-12T10:00:00Z","claw_id":"weston","type":"channel_context_op","kind":"raw_window+digest","channels":["chan-1","chan-2"],"retained":60,"returned":40,"omitted":20,"raw_bytes":12000,"digest_bytes":800,"digest_blocks":3,"coverage_gaps":1,"deterministic_only":true,"source":"claw-wall","status":"coverage_gap","tool_name":"search_channel_context","status_code":200}` event, err := NormalizeLine([]byte(line)) if err != nil { t.Fatalf("unexpected error: %v", err) } - if event.Type != "channel_context_op" || event.ChannelKind != "raw_window" || event.SourceService != "claw-wall" || event.Status != "ok" { + if event.Type != "channel_context_op" || event.ChannelKind != "raw_window+digest" || event.SourceService != "claw-wall" || event.Status != "coverage_gap" { t.Fatalf("unexpected event: %+v", event) } if len(event.Channels) != 2 || event.Channels[0] != "chan-1" || event.Channels[1] != "chan-2" { @@ -95,6 +95,12 @@ func TestNormalizeLineParseChannelContextOpEvent(t *testing.T) { if event.Retained == nil || *event.Retained != 60 || event.Returned == nil || *event.Returned != 40 || event.Omitted == nil || *event.Omitted != 20 { t.Fatalf("unexpected counts: %+v", event) } + if event.RawBytes == nil || *event.RawBytes != 12000 || event.DigestBytes == nil || *event.DigestBytes != 800 || event.DigestBlocks == nil || *event.DigestBlocks != 3 || event.CoverageGaps == nil || *event.CoverageGaps != 1 { + t.Fatalf("unexpected digest counts: %+v", event) + } + if event.DeterministicOnly == nil || !*event.DeterministicOnly { + t.Fatalf("expected deterministic_only=true: %+v", event) + } if event.ToolName != "search_channel_context" { t.Fatalf("unexpected tool name: %+v", event) } diff --git a/skills/clawdapus/SKILL.md b/skills/clawdapus/SKILL.md index 14882a1..7c6c23d 100644 --- a/skills/clawdapus/SKILL.md +++ b/skills/clawdapus/SKILL.md @@ -389,7 +389,7 @@ The proxy sits between agents and LLM providers. Agents get bearer tokens, proxy ### claw-wall sidecar -Auto-injected by `claw up` when any cllama-enabled service has Discord channel IDs. Polls Discord channels and serves the recent channel transcript to agents through `channel-context` tail feeds; legacy unread-mailbox cursor paging remains available as `mode=delta`. On startup, wall backfills Discord history before its first forward poll up to `CLAW_WALL_RETENTION` (default `24h`) and `CLAW_WALL_BACKFILL_MAX_PAGES` (default `25`), while `CLAW_WALL_LIMIT` is a per-channel safety cap (default `5000`). Configure the generated tail window with pod or service `x-claw.context.channel` (`since`, `limit`, `max-chars`, `buffer`). Since `v0.15.0` channel-consuming services also get a default-on `channel-awareness` feed (uncursored 24h raw window; `x-claw.context.channel.max-chars` tunes both feeds together, default 32 KB) plus two cllama-mediated retrieval tools - `search_channel_context` and `get_channel_messages` - auto-subscribed via a compiler-owned claw-wall descriptor. Feed headers include `backfill_status`; `partial` or `rate_limited` means the backing window did not fully satisfy the requested horizon. Calls are gated by a generated per-agent channel allowlist, claw-wall service-token auth, and forwarded `X-Claw-ID`. The service name `claw-wall` is reserved - declaring it in `claw-pod.yml` is a hard error. +Auto-injected by `claw up` when any cllama-enabled service has Discord channel IDs. Polls Discord channels and serves the recent channel transcript to agents through `channel-context` tail feeds; legacy unread-mailbox cursor paging remains available as `mode=delta`. On startup, wall backfills Discord history before its first forward poll up to `CLAW_WALL_RETENTION` (default `24h`) and `CLAW_WALL_BACKFILL_MAX_PAGES` (default `25`), while `CLAW_WALL_LIMIT` is a per-channel safety cap (default `5000`). Configure the generated tail window with pod or service `x-claw.context.channel` (`since`, `limit`, `max-chars`, `buffer`). Since `v0.15.0` channel-consuming services also get a default-on `channel-awareness` feed plus two cllama-mediated retrieval tools - `search_channel_context` and `get_channel_messages` - auto-subscribed via a compiler-owned claw-wall descriptor. Without `x-claw.channel-memory.service`, `channel-awareness` is an uncursored raw window. With channel-memory wired, `claw up` changes that feed to `context_kind=raw_window+digest`; claw-wall fetches digest blocks from channel-memory, emits a compact `[digest]` section with source provenance, and keeps a bounded `[raw recent]` tail. Feed headers include `backfill_status`, `digest`, `raw_bytes`, `digest_bytes`, `digest_blocks`, `coverage_gaps`, and `deterministic_only`; `partial` or `rate_limited` backfill means the backing window did not fully satisfy the requested horizon. Calls are gated by a generated per-agent channel allowlist, claw-wall service-token auth, and forwarded `X-Claw-ID`. The service name `claw-wall` is reserved - declaring it in `claw-pod.yml` is a hard error. ### cllama feed injection budgets