diff --git a/internal/ai/chat/agentic.go b/internal/ai/chat/agentic.go index ad0600fac..79921ad7f 100644 --- a/internal/ai/chat/agentic.go +++ b/internal/ai/chat/agentic.go @@ -871,14 +871,18 @@ func (a *AgenticLoop) executeWithTools(ctx context.Context, sessionID string, me Str("session_id", sessionID). Msg("[AgenticLoop] Starting turn") - // Build the request with dynamic system prompt (includes current mode) - systemPrompt := a.getSystemPrompt() + // Build the request with dynamic system prompt (includes current mode). + // Mark the frozen base plus mode text as cacheable so a prompt-caching + // provider reuses it across turns without caching the per-turn time or + // accumulated knowledge that follow it. + stablePrompt, systemPrompt := a.systemPromptParts() req := providers.ChatRequest{ - Messages: providerMessages, - System: systemPrompt, - Tools: tools, - ExecutionID: a.executionID, - StreamIdleTimeout: streamIdleTimeout, + Messages: providerMessages, + System: systemPrompt, + SystemCacheablePrefix: stablePrompt, + Tools: tools, + ExecutionID: a.executionID, + StreamIdleTimeout: streamIdleTimeout, } if isPatrolDetectionExecution(a.currentExecutionProfile()) && patrolFindingsReadCompleted { req.Tools = withoutProviderTool(req.Tools, agentcapabilities.PatrolGetFindingsToolName) diff --git a/internal/ai/chat/agentic_prompt.go b/internal/ai/chat/agentic_prompt.go index 02bac7579..cd68f394c 100644 --- a/internal/ai/chat/agentic_prompt.go +++ b/internal/ai/chat/agentic_prompt.go @@ -13,6 +13,16 @@ import ( // service start, so anything that must stay fresh per turn (mode, current time) // is appended here rather than baked into baseSystemPrompt. func (a *AgenticLoop) getSystemPrompt() string { + _, prompt := a.systemPromptParts() + return prompt +} + +// systemPromptParts returns the full system prompt together with its stable +// leading prefix. The prefix (frozen base plus mode text) is identical across +// the turns of a run, so a provider with prompt caching can reuse it; the +// per-turn time and accumulated knowledge that follow it change every turn and +// must stay outside the cached block. +func (a *AgenticLoop) systemPromptParts() (stablePrefix, full string) { a.mu.Lock() isAutonomous := a.autonomousMode profile := a.executionProfile @@ -68,12 +78,13 @@ Treat this as the current date and time. Answer "what time is it" / "what's the questions directly from this value — do not run a command or ask for a target host just to report the current time.`, time.Now().Format("Mon, 02 Jan 2006 15:04:05 MST")) - prompt := a.baseSystemPrompt + modeContext + currentTime + stablePrefix = a.baseSystemPrompt + modeContext + prompt := stablePrefix + currentTime // Append accumulated knowledge facts to system prompt if ka := a.knowledgeAccumulator; ka != nil && ka.Len() > 0 { prompt += "\n\n" + ka.Render() } - return prompt + return stablePrefix, prompt } diff --git a/internal/ai/chat/agentic_test.go b/internal/ai/chat/agentic_test.go index 066179e90..58a0b3523 100644 --- a/internal/ai/chat/agentic_test.go +++ b/internal/ai/chat/agentic_test.go @@ -461,6 +461,30 @@ func TestAgenticLoop_SystemPromptIncludesCurrentTime(t *testing.T) { } } +func TestAgenticLoop_SystemPromptPartsSeparateStablePrefix(t *testing.T) { + mockProvider := &MockProvider{} + executor := tools.NewPulseToolExecutor(tools.ExecutorConfig{}) + loop := NewAgenticLoop(mockProvider, executor, "BASE PROMPT BODY") + + stable, full := loop.systemPromptParts() + + if !strings.HasPrefix(full, stable) { + t.Fatalf("stable prefix must lead the full prompt, stable=%q full=%q", stable, full) + } + if !strings.Contains(stable, "BASE PROMPT BODY") || !strings.Contains(stable, "EXECUTION MODE:") { + t.Fatalf("stable prefix must carry the frozen base and mode text, got %q", stable) + } + if strings.Contains(stable, "CURRENT TIME:") { + t.Fatalf("per-turn time must stay outside the cacheable prefix, got %q", stable) + } + if !strings.Contains(full, "CURRENT TIME:") { + t.Fatalf("full prompt must still carry CURRENT TIME, got %q", full) + } + if loop.getSystemPrompt() != full { + t.Fatal("getSystemPrompt must return the full prompt from systemPromptParts") + } +} + func TestAgenticLoop_AnswerQuestion(t *testing.T) { mockProvider := &MockProvider{} executor := tools.NewPulseToolExecutor(tools.ExecutorConfig{}) diff --git a/internal/ai/providers/anthropic.go b/internal/ai/providers/anthropic.go index fb69c75fd..b4ba51400 100644 --- a/internal/ai/providers/anthropic.go +++ b/internal/ai/providers/anthropic.go @@ -67,7 +67,7 @@ type anthropicRequest struct { Model string `json:"model"` Messages []anthropicMessage `json:"messages"` MaxTokens int `json:"max_tokens"` - System string `json:"system,omitempty"` + System interface{} `json:"system,omitempty"` Temperature float64 `json:"temperature,omitempty"` Tools []anthropicTool `json:"tools,omitempty"` ToolChoice *anthropicToolChoice `json:"tool_choice,omitempty"` @@ -114,6 +114,48 @@ type anthropicCacheControl struct { Type string `json:"type"` // "ephemeral" } +// anthropicSystemBlock is a system-prompt content block. Anthropic accepts the +// system prompt as either a plain string or an array of content blocks; the +// array form is required to place a prompt-caching breakpoint on the stable +// leading portion of the prompt. +type anthropicSystemBlock struct { + Type string `json:"type"` // "text" + Text string `json:"text"` + CacheControl *anthropicCacheControl `json:"cache_control,omitempty"` +} + +// buildAnthropicSystem serializes the system prompt for an Anthropic request. +// When cacheablePrefix is a non-empty exact prefix of system, the stable prefix +// is returned as a content block with a cache breakpoint and the volatile +// remainder as a second block. Anthropic caches content in tools -> system -> +// messages order, so this breakpoint also covers the tool definitions and the +// redundant per-tool breakpoint can be dropped. Otherwise the plain string form +// is returned unchanged. +func buildAnthropicSystem(system, cacheablePrefix string) interface{} { + if system == "" { + return nil + } + if cacheablePrefix == "" || !strings.HasPrefix(system, cacheablePrefix) { + return system + } + blocks := []anthropicSystemBlock{{ + Type: "text", + Text: cacheablePrefix, + CacheControl: &anthropicCacheControl{Type: "ephemeral"}, + }} + if remainder := system[len(cacheablePrefix):]; remainder != "" { + blocks = append(blocks, anthropicSystemBlock{Type: "text", Text: remainder}) + } + return blocks +} + +// anthropicSystemIsCached reports whether buildAnthropicSystem returned the +// block form, meaning a system breakpoint already covers the tool definitions. +func anthropicSystemIsCached(system interface{}) bool { + _, ok := system.([]anthropicSystemBlock) + return ok +} + // anthropicResponse is the response from the Anthropic API type anthropicResponse struct { ID string `json:"id"` @@ -245,11 +287,12 @@ func (c *AnthropicClient) Chat(ctx context.Context, req ChatRequest) (*ChatRespo maxTokens = 4096 } + systemValue := buildAnthropicSystem(req.System, req.SystemCacheablePrefix) anthropicReq := anthropicRequest{ Model: model, Messages: messages, MaxTokens: maxTokens, - System: req.System, + System: systemValue, } if req.Temperature > 0 { @@ -279,8 +322,12 @@ func (c *AnthropicClient) Chat(ctx context.Context, req ChatRequest) (*ChatRespo } } // Mark the last tool with cache_control so Anthropic caches all tool - // definitions (and everything before them) on subsequent turns. - anthropicReq.Tools[len(anthropicReq.Tools)-1].CacheControl = &anthropicCacheControl{Type: "ephemeral"} + // definitions (and everything before them) on subsequent turns. A + // system cache breakpoint already covers the tool definitions, so it is + // omitted then to avoid a redundant cache write. + if !anthropicSystemIsCached(systemValue) { + anthropicReq.Tools[len(anthropicReq.Tools)-1].CacheControl = &anthropicCacheControl{Type: "ephemeral"} + } } // Add tool_choice only for explicit overrides. Nil keeps Anthropic's default @@ -417,6 +464,15 @@ func (c *AnthropicClient) Chat(ctx context.Context, req ChatRequest) (*ChatRespo logEvent = logEvent. Int("cache_creation_tokens", anthropicResp.Usage.CacheCreationInputTokens). Int("cache_read_tokens", anthropicResp.Usage.CacheReadInputTokens) + // Prompt-cache effectiveness is otherwise invisible in production + // because the parsed-response event is Debug. Surface the counters at + // Info so operators can confirm cache reads and writes from real runs. + log.Info(). + Int("input_tokens", anthropicResp.Usage.InputTokens). + Int("cache_creation_tokens", anthropicResp.Usage.CacheCreationInputTokens). + Int("cache_read_tokens", anthropicResp.Usage.CacheReadInputTokens). + Str("model", anthropicResp.Model). + Msg("anthropic prompt cache usage") } logEvent.Msg("anthropic response parsed") @@ -542,7 +598,7 @@ type anthropicStreamRequest struct { Model string `json:"model"` Messages []anthropicMessage `json:"messages"` MaxTokens int `json:"max_tokens"` - System string `json:"system,omitempty"` + System interface{} `json:"system,omitempty"` Temperature float64 `json:"temperature,omitempty"` Tools []anthropicTool `json:"tools,omitempty"` ToolChoice *anthropicToolChoice `json:"tool_choice,omitempty"` @@ -587,11 +643,12 @@ func (c *AnthropicClient) ChatStream(ctx context.Context, req ChatRequest, callb maxTokens = 4096 } + systemValue := buildAnthropicSystem(req.System, req.SystemCacheablePrefix) anthropicReq := anthropicStreamRequest{ Model: model, Messages: messages, MaxTokens: maxTokens, - System: req.System, + System: systemValue, Stream: true, } @@ -620,8 +677,11 @@ func (c *AnthropicClient) ChatStream(ctx context.Context, req ChatRequest, callb } } } - // Mark the last tool with cache_control for prompt caching (same as non-streaming). - anthropicReq.Tools[len(anthropicReq.Tools)-1].CacheControl = &anthropicCacheControl{Type: "ephemeral"} + // Mark the last tool with cache_control for prompt caching (same as + // non-streaming), unless a system breakpoint already covers the tools. + if !anthropicSystemIsCached(systemValue) { + anthropicReq.Tools[len(anthropicReq.Tools)-1].CacheControl = &anthropicCacheControl{Type: "ephemeral"} + } } // Add tool_choice only for explicit overrides, same as non-streaming. diff --git a/internal/ai/providers/anthropic_test.go b/internal/ai/providers/anthropic_test.go index 12ff434f0..4879e5993 100644 --- a/internal/ai/providers/anthropic_test.go +++ b/internal/ai/providers/anthropic_test.go @@ -734,3 +734,171 @@ func TestAnthropicClient_SupportsThinking(t *testing.T) { t.Fatal("expected SupportsThinking to be false") } } + +func TestBuildAnthropicSystem(t *testing.T) { + t.Run("plain string when no cacheable prefix", func(t *testing.T) { + got := buildAnthropicSystem("base+mode+time", "") + if s, ok := got.(string); !ok || s != "base+mode+time" { + t.Fatalf("buildAnthropicSystem = %#v, want plain string", got) + } + if anthropicSystemIsCached(got) { + t.Fatal("plain system must not be reported as cached") + } + }) + + t.Run("non-prefix falls back to plain string", func(t *testing.T) { + got := buildAnthropicSystem("base+mode+time", "different") + if _, ok := got.(string); !ok { + t.Fatalf("buildAnthropicSystem = %#v, want plain string fallback", got) + } + }) + + t.Run("splits stable prefix with a cache breakpoint", func(t *testing.T) { + got := buildAnthropicSystem("base+mode+time", "base+mode") + blocks, ok := got.([]anthropicSystemBlock) + if !ok || len(blocks) != 2 { + t.Fatalf("buildAnthropicSystem = %#v, want two blocks", got) + } + if blocks[0].Text != "base+mode" || blocks[0].CacheControl == nil || blocks[0].CacheControl.Type != "ephemeral" { + t.Fatalf("unexpected stable block: %+v", blocks[0]) + } + if blocks[1].Text != "+time" || blocks[1].CacheControl != nil { + t.Fatalf("unexpected volatile block: %+v", blocks[1]) + } + if !anthropicSystemIsCached(got) { + t.Fatal("block system must be reported as cached") + } + }) +} + +func TestAnthropicClient_Chat_CacheableSystemPrefix(t *testing.T) { + var got anthropicRequest + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if err := json.NewDecoder(r.Body).Decode(&got); err != nil { + t.Fatalf("decode request: %v", err) + } + _ = json.NewEncoder(w).Encode(anthropicResponse{ + ID: "msg_1", + Type: "message", + Role: "assistant", + Model: "claude-3-5-sonnet", + StopReason: "end_turn", + Content: []anthropicContent{{Type: "text", Text: "ok"}}, + Usage: anthropicUsage{InputTokens: 1, OutputTokens: 1}, + }) + })) + defer server.Close() + + client := NewAnthropicClientWithBaseURL("test-key", "claude-3-5-sonnet", server.URL, 0) + if _, err := client.Chat(context.Background(), ChatRequest{ + System: "stable promptCURRENT TIME: now", + SystemCacheablePrefix: "stable prompt", + Messages: []Message{{Role: "user", Content: "Hi"}}, + Tools: []Tool{{Name: "get_time", Description: "d", InputSchema: map[string]any{"type": "object"}}}, + }); err != nil { + t.Fatalf("Chat: %v", err) + } + + blocks, ok := got.System.([]interface{}) + if !ok || len(blocks) != 2 { + t.Fatalf("system = %#v, want two content blocks", got.System) + } + first, _ := blocks[0].(map[string]interface{}) + if first["text"] != "stable prompt" { + t.Fatalf("first block = %#v", first) + } + cc, _ := first["cache_control"].(map[string]interface{}) + if cc["type"] != "ephemeral" { + t.Fatalf("first block cache_control = %#v", first["cache_control"]) + } + second, _ := blocks[1].(map[string]interface{}) + if second["text"] != "CURRENT TIME: now" { + t.Fatalf("second block = %#v", second) + } + if _, present := second["cache_control"]; present { + t.Fatalf("volatile block must not carry a breakpoint: %#v", second) + } + if got.Tools[0].CacheControl != nil { + t.Fatalf("tool breakpoint must be dropped when the system block is cached: %+v", got.Tools[0]) + } +} + +func TestAnthropicClient_Chat_KeepsToolBreakpointWithoutCacheableSystem(t *testing.T) { + var got anthropicRequest + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if err := json.NewDecoder(r.Body).Decode(&got); err != nil { + t.Fatalf("decode request: %v", err) + } + _ = json.NewEncoder(w).Encode(anthropicResponse{ + ID: "msg_1", + Type: "message", + Role: "assistant", + Model: "claude-3-5-sonnet", + StopReason: "end_turn", + Content: []anthropicContent{{Type: "text", Text: "ok"}}, + Usage: anthropicUsage{InputTokens: 1, OutputTokens: 1}, + }) + })) + defer server.Close() + + client := NewAnthropicClientWithBaseURL("test-key", "claude-3-5-sonnet", server.URL, 0) + if _, err := client.Chat(context.Background(), ChatRequest{ + System: "plain prompt", + Messages: []Message{{Role: "user", Content: "Hi"}}, + Tools: []Tool{{Name: "get_time", Description: "d", InputSchema: map[string]any{"type": "object"}}}, + }); err != nil { + t.Fatalf("Chat: %v", err) + } + if _, ok := got.System.(string); !ok { + t.Fatalf("system = %#v, want plain string", got.System) + } + if got.Tools[0].CacheControl == nil { + t.Fatal("tool breakpoint must remain when the system prompt is not cached") + } +} + +func TestAnthropicClient_ChatStream_CacheableSystemPrefix(t *testing.T) { + var got anthropicStreamRequest + stream := []string{ + `{"type":"message_start","message":{"usage":{"input_tokens":5}}}`, + `{"type":"content_block_start","content_block":{"type":"text"}}`, + `{"type":"content_block_delta","delta":{"type":"text_delta","text":"ok"}}`, + `{"type":"content_block_stop"}`, + `{"type":"message_delta","delta":{"stop_reason":"end_turn"},"usage":{"output_tokens":1}}`, + `{"type":"message_stop"}`, + } + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if err := json.NewDecoder(r.Body).Decode(&got); err != nil { + t.Fatalf("decode request: %v", err) + } + w.Header().Set("Content-Type", "text/event-stream") + w.WriteHeader(http.StatusOK) + for _, event := range stream { + _, _ = w.Write([]byte("data: " + event + "\n\n")) + w.(http.Flusher).Flush() + } + })) + defer server.Close() + + client := NewAnthropicClientWithBaseURL("test-key", "claude-3-5-sonnet", server.URL, 0) + if err := client.ChatStream(context.Background(), ChatRequest{ + System: "stableCURRENT", + SystemCacheablePrefix: "stable", + Messages: []Message{{Role: "user", Content: "Hi"}}, + Tools: []Tool{{Name: "get_time", Description: "d", InputSchema: map[string]any{"type": "object"}}}, + }, func(StreamEvent) {}); err != nil { + t.Fatalf("ChatStream: %v", err) + } + + blocks, ok := got.System.([]interface{}) + if !ok || len(blocks) != 2 { + t.Fatalf("system = %#v, want two content blocks", got.System) + } + first, _ := blocks[0].(map[string]interface{}) + if cc, _ := first["cache_control"].(map[string]interface{}); cc["type"] != "ephemeral" { + t.Fatalf("first block cache_control = %#v", first["cache_control"]) + } + if got.Tools[0].CacheControl != nil { + t.Fatalf("tool breakpoint must be dropped when the system block is cached: %+v", got.Tools[0]) + } +} diff --git a/internal/ai/providers/provider.go b/internal/ai/providers/provider.go index ac3596741..86c9d7e13 100644 --- a/internal/ai/providers/provider.go +++ b/internal/ai/providers/provider.go @@ -107,11 +107,18 @@ type ChatRequest struct { // tokens. Ollama otherwise loads models at its server default (typically // 4096) regardless of the model's trained window and silently truncates // large prompts (#1624). Providers without such a control ignore it. - MinContextTokens int `json:"-"` - System string `json:"system,omitempty"` // System prompt (Anthropic style) - Tools []Tool `json:"tools,omitempty"` // Available tools - ToolChoice *ToolChoice `json:"tool_choice,omitempty"` // nil = model-owned automatic selection; none = text-only safety brake; required = provider-native forced tool use where supported - StreamIdleTimeout time.Duration `json:"-"` // Optional use-case-specific inter-chunk stall allowance; provider request payloads must not serialize it. + MinContextTokens int `json:"-"` + System string `json:"system,omitempty"` // System prompt (Anthropic style) + // SystemCacheablePrefix marks the stable leading portion of System. When + // set and it is an exact prefix of System, providers with prompt-caching + // support (Anthropic) place a cache breakpoint after this prefix so the + // frozen base plus mode text is reused across turns while per-turn content + // appended after it stays outside the cached block. Providers without that + // control ignore it and send System unchanged. + SystemCacheablePrefix string `json:"-"` + Tools []Tool `json:"tools,omitempty"` // Available tools + ToolChoice *ToolChoice `json:"tool_choice,omitempty"` // nil = model-owned automatic selection; none = text-only safety brake; required = provider-native forced tool use where supported + StreamIdleTimeout time.Duration `json:"-"` // Optional use-case-specific inter-chunk stall allowance; provider request payloads must not serialize it. } func (r ChatRequest) NormalizeCollections() ChatRequest {