From 07284c090ce0e7831d27327b88fdea39a52b0a13 Mon Sep 17 00:00:00 2001 From: Yanpeng Wang Date: Thu, 27 Aug 2026 19:01:41 +0800 Subject: [PATCH 1/2] refactor: remove inference geo from agent models --- docs/api/agents.md | 10 +- docs/api/sessions.md | 21 ++-- docs/capabilities.md | 2 +- docs/guides/multi-agent.md | 2 - docs/provenance.md | 10 ++ internal/agentruntime/agentcore.go | 18 ++-- internal/agentruntime/multiagent.go | 7 +- internal/agentruntime/multiagent_test.go | 3 +- internal/app/agent_service.go | 12 --- internal/app/agent_service_test.go | 64 ------------ internal/controlplane/session_service.go | 6 -- internal/domain/agent.go | 13 +-- internal/domain/budget.go | 36 +++---- internal/domain/budget_test.go | 10 +- internal/domain/usage.go | 4 + internal/httpapi/agents_test.go | 23 +++++ internal/httpapi/dto.go | 13 +-- internal/httpapi/openapi.yaml | 7 -- internal/httpapi/sdk_test.go | 97 ------------------- internal/httpapi/testhelpers_test.go | 6 -- internal/model/anthropic.go | 11 ++- internal/model/anthropic_test.go | 26 ++--- internal/model/client.go | 15 ++- internal/model/context.go | 15 ++- internal/pg/advisor.go | 9 +- internal/pg/advisor_test.go | 27 ++---- internal/pg/budget.go | 8 +- .../migrations/00031_model_request_usage.sql | 9 +- internal/pg/pgstore/models.go | 1 - internal/temporal/activities.go | 9 +- internal/temporal/agent_workflow.go | 1 - internal/temporal/agent_workflow_test.go | 8 +- internal/temporal/agent_workflow_tools.go | 2 - internal/temporal/types.go | 1 - 34 files changed, 153 insertions(+), 353 deletions(-) diff --git a/docs/api/agents.md b/docs/api/agents.md index a429b0c..e52c406 100644 --- a/docs/api/agents.md +++ b/docs/api/agents.md @@ -29,13 +29,9 @@ Required fields: - `name`: non-empty string; - `model`: model ID string or an object with a non-empty `id`. -The object model form also preserves supported `effort`, `speed`, and -`inference_geo` values. +The object model form also preserves supported `effort` and `speed` values. `effort` accepts either a level string such as `"high"` or the tagged object `{"type":"high"}`; responses use the tagged object form. -An explicit non-empty `inference_geo` is forwarded on every working and outcome -grader request. On Agent update, `model` is whole-object replacement for this -field: omitting `inference_geo` clears a previous pin. Optional collection and metadata fields may be omitted or supplied with their documented array/object shape; explicit `null` is not a create-time default. @@ -62,10 +58,6 @@ Creating a Session expands those pins into the full immutable definitions returned in `session.agent.multiagent.agents`; child Threads will execute those Session-owned snapshots rather than re-resolving Agent resources. Archived, missing, duplicate, and nested coordinator references are rejected. -If the coordinator pins `model.inference_geo`, every independently referenced -Agent must pin the same value; if the coordinator leaves it unset, every member -must also leave it unset. A model change that would violate this invariant must -replace or clear the roster in the same Agent update. Ordinary roster entries can execute as persistent child Session Threads with independent context, events, usage, and Workflow state. See the [multi-agent guide](../guides/multi-agent.md) for an end-to-end example and diff --git a/docs/api/sessions.md b/docs/api/sessions.md index 06db8ec..871ad43 100644 --- a/docs/api/sessions.md +++ b/docs/api/sessions.md @@ -49,11 +49,9 @@ Session-local overrides: Overrides replace model, system, tools, MCP servers, or skills for this session only. They do not mutate or renumber the agent. A model override may change the -model ID, speed, or inference geography; effort remains an Agent-level setting -and a session override does not replace it. Overrides also apply to `self` -copies in a coordinator roster. Independently referenced Agents are unaffected, -so a geography override that would make the coordinator disagree with one of -those pinned Agents is rejected. +model ID or speed; effort remains an Agent-level setting and a session override +does not replace it. Overrides also apply to `self` copies in a coordinator +roster. Independently referenced Agents are unaffected. For a coordinator, `session.agent.multiagent.agents` expands the Agent resource's Version references into full immutable Agent definitions. The @@ -249,12 +247,13 @@ File and Memory Store Resource objects. Ordered `vault_ids` are resolved at creation; update-time vault replacement is rejected. `usage` aggregates provider-reported token, prompt-cache, Web Fetch, and Web Search counters across every Session Thread. `usage.list_cost` is calculated -from Mango's current built-in price catalog, the Web Search request rate, and -$0.08 per Session active hour, then rounded to the nearest cent for the public -monetary projection. Thread list cost excludes Session runtime. Accounting -remains exact internally, and model-request admission checks the shared ceiling -before every request; an already in-flight request may take the Session over its -limit. +from Mango's current built-in price catalog, provider-reported execution facts +that affect that catalog's rates, the Web Search request rate, and $0.08 per +Session active hour, then rounded to the nearest cent for the public monetary +projection. Provider routing is not Agent configuration. Thread list cost +excludes Session runtime. Accounting remains exact internally, and +model-request admission checks the shared ceiling before every request; an +already in-flight request may take the Session over its limit. Provider-reported tokens remain visible even when a response-level billing rule makes their list cost zero, such as an unbilled Claude Fable 5 refusal. diff --git a/docs/capabilities.md b/docs/capabilities.md index cc7e6bb..1953bef 100644 --- a/docs/capabilities.md +++ b/docs/capabilities.md @@ -36,7 +36,7 @@ and service test suites. | Capability | Status | Supported scope and important constraints | | --- | --- | --- | -| Agents and Versions | Supported | Create, get, list, update, immutable Version history, archive, filters, and pagination. Model ID, effort, speed, and `inference_geo` reach working and grader requests. | +| Agents and Versions | Supported | Create, get, list, update, immutable Version history, archive, filters, and pagination. Model ID, effort, and speed reach working and grader requests. Provider routing policy remains outside the Agent contract. | | Environments | Supported | Cloud and self-hosted lifecycle, package configuration, limited-network declarations, filters, and pagination. Package execution requires a capable sandbox; limited egress is currently enforced only by OpenSandbox. | | Sessions | Supported | Create from immutable Agent snapshots, get/list/update/archive/delete, metadata, filters, exact shared public-list-cost budgets, usage, timing, and resource projections. Deletion fences admission and durably releases the Workflow and sandbox. | | Events and client actions | Limited | System context, messages, thinking, tool events, confirmation/custom/self-hosted result barriers, outcomes, retries, interrupts, and the budget-boundary `session.usage`/`budget_reached` idle sequence are implemented. File-backed message documents are limited to bounded UTF-8 text; File-sourced images and File documents in tool results are not supported. | diff --git a/docs/guides/multi-agent.md b/docs/guides/multi-agent.md index f8ff17e..7231c07 100644 --- a/docs/guides/multi-agent.md +++ b/docs/guides/multi-agent.md @@ -24,8 +24,6 @@ an Advisor, and persistent follow-up, see tools. The deterministic local model is useful for platform smoke tests but does not make open-ended delegation decisions. -- Keep every roster Agent on the same `inference_geo` value, or leave it unset - everywhere. The examples use the local API at `http://localhost:8080`. diff --git a/docs/provenance.md b/docs/provenance.md index 07cc37e..5ae88e6 100644 --- a/docs/provenance.md +++ b/docs/provenance.md @@ -32,6 +32,16 @@ release is never an automatic roadmap. outbound endpoint requires. Tests that exercise Mango through an Anthropic SDK are optional research evidence; raw HTTP and OpenAPI tests define Mango's transport contract. +- Claude Managed Agents' agent-level `inference_geo` and the public + [Claude data-residency design](https://platform.claude.com/docs/en/manage-claude/data-residency) + prompted a focused review on 2026-08-27. Mango rejected request-time + geography from its Agent and Session model configuration: it is a hosted + provider routing policy, other model platforms express placement through + different endpoints or deployment resources, and Mango's replaceable model + boundary cannot enforce a portable meaning for it. Operators select and + govern the configured model endpoint outside the Agent contract. The current + Anthropic adapter reads a provider-reported response region only as an + internal list-cost input; it never sends a geography request field. ## Built-in Agent tools diff --git a/internal/agentruntime/agentcore.go b/internal/agentruntime/agentcore.go index 95202ab..b228f56 100644 --- a/internal/agentruntime/agentcore.go +++ b/internal/agentruntime/agentcore.go @@ -133,11 +133,10 @@ func (a *AgentCore) Run(ctx context.Context, req RunRequest, sink EventSink) (Ru messageID := a.ids.NewID(domain.PrefixEvent) started := false resp, err = a.client.CreateMessageStream(ctx, model.Request{ - Model: req.AgentSnapshot.Model.ID, - InferenceGeo: req.AgentSnapshot.Model.InferenceGeo, - System: system, - Messages: messages, - Tools: toolSchemas, + Model: req.AgentSnapshot.Model.ID, + System: system, + Messages: messages, + Tools: toolSchemas, }, func(index int, text string) { if !started { previewer.PreviewStart(messageID, domain.EvAgentMessage) @@ -172,11 +171,10 @@ func (a *AgentCore) Run(ctx context.Context, req RunRequest, sink EventSink) (Ru } } else { resp, err = a.client.CreateMessage(ctx, model.Request{ - Model: req.AgentSnapshot.Model.ID, - InferenceGeo: req.AgentSnapshot.Model.InferenceGeo, - System: system, - Messages: messages, - Tools: toolSchemas, + Model: req.AgentSnapshot.Model.ID, + System: system, + Messages: messages, + Tools: toolSchemas, }) if err != nil { return RunOutcome{}, err diff --git a/internal/agentruntime/multiagent.go b/internal/agentruntime/multiagent.go index bca9431..33b914b 100644 --- a/internal/agentruntime/multiagent.go +++ b/internal/agentruntime/multiagent.go @@ -143,10 +143,9 @@ func AdvisorRequest( return model.Request{}, fmt.Errorf("encode advisor context: %w", err) } return model.Request{ - Model: advisorModel, - InferenceGeo: executor.InferenceGeo, - System: advisorReviewerSystem, - MaxTokens: advisorMaxTokens, + Model: advisorModel, + System: advisorReviewerSystem, + MaxTokens: advisorMaxTokens, Messages: []domain.Message{{ Role: domain.RoleUser, Content: []domain.ContentBlock{{ diff --git a/internal/agentruntime/multiagent_test.go b/internal/agentruntime/multiagent_test.go index 4bff82f..6b51c67 100644 --- a/internal/agentruntime/multiagent_test.go +++ b/internal/agentruntime/multiagent_test.go @@ -26,7 +26,7 @@ func TestAdvisorRequestQuotesExecutorContextWithoutReplayingReasoning(t *testing require.Equal(t, false, schema.InputSchema["additionalProperties"]) executor := model.Request{ - Model: "executor-model", InferenceGeo: "us", System: "executor system", + Model: "executor-model", System: "executor system", Tools: []model.ToolSchema{{ Name: "read", Description: "Read a file.", InputSchema: map[string]any{"type": "object"}, @@ -49,7 +49,6 @@ func TestAdvisorRequestQuotesExecutorContextWithoutReplayingReasoning(t *testing ) require.NoError(t, err) require.Equal(t, "advisor-model", request.Model) - require.Equal(t, "us", request.InferenceGeo) require.Empty(t, request.Tools) require.Equal(t, 2048, request.MaxTokens) require.Len(t, request.Messages, 1) diff --git a/internal/app/agent_service.go b/internal/app/agent_service.go index e366d74..1c85ceb 100644 --- a/internal/app/agent_service.go +++ b/internal/app/agent_service.go @@ -149,13 +149,6 @@ func (s *AgentService) Update(ctx context.Context, id string, patch domain.Agent // version used to perform semantic no-op detection. next.Multiagent = next.Multiagent.RebindAgentVersion(cur.ID, cur.Version+1) } - if patch.Model != nil && patch.Multiagent == nil && - next.Model.InferenceGeo != cur.Model.InferenceGeo && - cur.Multiagent.HasExternalAgent(cur.ID) { - return domain.Agent{}, false, domain.Validation( - "model.inference_geo must match every independently referenced multiagent roster member", - ) - } if changed && next.Multiagent != nil && !next.Multiagent.IsResolved() { return domain.Agent{}, false, domain.Validation( "legacy multiagent configuration must be replaced before updating the Agent", @@ -306,11 +299,6 @@ func (s *AgentService) resolveMultiagent( return nil, domain.Validation("multiagent references are limited to one coordinator level") } } - if target.Model.InferenceGeo != ownerModel.InferenceGeo { - return nil, domain.Validation( - "model.inference_geo must match every multiagent roster member", - ) - } if _, duplicate := seen[target.ID]; duplicate { return nil, domain.Validation("multiagent.agents must reference distinct agents") } diff --git a/internal/app/agent_service_test.go b/internal/app/agent_service_test.go index 3157d97..06e5644 100644 --- a/internal/app/agent_service_test.go +++ b/internal/app/agent_service_test.go @@ -260,70 +260,6 @@ func TestAgentService_MultiagentPinsLatestAndRebindsSelf(t *testing.T) { } } -func TestAgentService_MultiagentInferenceGeoMustMatch(t *testing.T) { - s := newAgentService(t) - ctx := context.Background() - globalPeer, err := s.Create(ctx, domain.Agent{ - Name: "global peer", Model: domain.Model{ID: "m", InferenceGeo: "global"}, - }) - if err != nil { - t.Fatal(err) - } - usPeer, err := s.Create(ctx, domain.Agent{ - Name: "US peer", Model: domain.Model{ID: "m", InferenceGeo: "us"}, - }) - if err != nil { - t.Fatal(err) - } - roster := func(id string) *domain.Multiagent { - return &domain.Multiagent{Type: "coordinator", Agents: []domain.AgentReference{{ - Type: "agent", ID: id, - }}} - } - if _, err := s.Create(ctx, domain.Agent{ - Name: "mismatch", Model: domain.Model{ID: "m", InferenceGeo: "us"}, - Multiagent: roster(globalPeer.ID), - }); err == nil { - t.Fatal("created a coordinator whose inference_geo differs from its roster") - } - - coordinator, err := s.Create(ctx, domain.Agent{ - Name: "coordinator", Model: domain.Model{ID: "m", InferenceGeo: "global"}, - Multiagent: roster(globalPeer.ID), - }) - if err != nil { - t.Fatal(err) - } - usModel := domain.Model{ID: "m", InferenceGeo: "us"} - if _, err := s.Update(ctx, coordinator.ID, domain.AgentPatch{Model: &usModel}); err == nil { - t.Fatal("changed coordinator inference_geo without replacing its external roster") - } - updated, err := s.Update(ctx, coordinator.ID, domain.AgentPatch{ - Model: &usModel, Multiagent: &domain.NullableMultiagent{Value: roster(usPeer.ID)}, - }) - if err != nil { - t.Fatalf("replace model and roster atomically: %v", err) - } - if updated.Model.InferenceGeo != "us" || updated.Multiagent.Agents[0].ID != usPeer.ID { - t.Fatalf("updated coordinator = %#v", updated) - } - - selfOnly, err := s.Create(ctx, domain.Agent{ - Name: "self", Model: domain.Model{ID: "m", InferenceGeo: "global"}, - Multiagent: &domain.Multiagent{Type: "coordinator", Agents: []domain.AgentReference{{Type: "self"}}}, - }) - if err != nil { - t.Fatal(err) - } - selfOnly, err = s.Update(ctx, selfOnly.ID, domain.AgentPatch{Model: &usModel}) - if err != nil { - t.Fatalf("self copies should inherit the coordinator override: %v", err) - } - if selfOnly.Model.InferenceGeo != "us" || selfOnly.Multiagent.Agents[0].Version != 2 { - t.Fatalf("updated self coordinator = %#v", selfOnly) - } -} - func TestAgentService_MultiagentRejectsInvalidReferences(t *testing.T) { s := newAgentService(t) ctx := context.Background() diff --git a/internal/controlplane/session_service.go b/internal/controlplane/session_service.go index 282ef1a..151d62c 100644 --- a/internal/controlplane/session_service.go +++ b/internal/controlplane/session_service.go @@ -203,12 +203,6 @@ func (s *SessionService) Create( if input.Overrides != nil { snapshot = agent.WithOverrides(*input.Overrides) } - if snapshot.Model.InferenceGeo != agent.Model.InferenceGeo && - agent.Multiagent.HasExternalAgent(agent.ID) { - return domain.Session{}, domain.Validation( - "agent override model.inference_geo must match every independently referenced multiagent roster member", - ) - } snapshot.Skills, err = app.ResolveAgentSkillReferences( ctx, s.skillRef, diff --git a/internal/domain/agent.go b/internal/domain/agent.go index 49df403..fa50c32 100644 --- a/internal/domain/agent.go +++ b/internal/domain/agent.go @@ -9,10 +9,9 @@ import ( ) type Model struct { - ID string - Effort string - Speed string - InferenceGeo string + ID string + Effort string + Speed string // EffortExplicit and SpeedExplicit distinguish an explicit Agent setting // from the Mango defaults echoed in the resolved resource. The // Messages adapter uses this distinction to avoid sending preview fields to @@ -54,9 +53,6 @@ func ValidateModel(model Model) error { default: return Validation("model speed must be standard or fast") } - if model.InferenceGeo != "" && strings.TrimSpace(model.InferenceGeo) == "" { - return Validation("model inference_geo must be a non-empty string") - } return nil } @@ -434,9 +430,6 @@ func (a Agent) SessionSnapshotJSON() map[string]any { if a.Model.Speed != "" { model["speed"] = a.Model.Speed } - if a.Model.InferenceGeo != "" { - model["inference_geo"] = a.Model.InferenceGeo - } system, description := "", "" if a.System != nil { system = *a.System diff --git a/internal/domain/budget.go b/internal/domain/budget.go index 0492332..10bdfb1 100644 --- a/internal/domain/budget.go +++ b/internal/domain/budget.go @@ -112,10 +112,10 @@ func RuntimeListCostNanoUSD(activeSeconds float64) int64 { } type modelListPrice struct { - inputPerToken int64 - outputPerToken int64 - geoSurcharge bool - fastEligible bool + inputPerToken int64 + outputPerToken int64 + usRegionSurcharge bool + fastEligible bool } var datedModelSuffix = regexp.MustCompile(`^(.*?)(?:-[0-9]{8})?$`) @@ -151,18 +151,18 @@ func ModelUsageListCostNanoUSDAt( // Fast mode for currently supported Opus models is $10/$50 per MTok. inputRate, outputRate = 10_000, 50_000 } - geoNumerator, geoDenominator := int64(1), int64(1) - if strings.EqualFold(model.InferenceGeo, "us") && price.geoSurcharge { - geoNumerator, geoDenominator = 11, 10 + regionNumerator, regionDenominator := int64(1), int64(1) + if strings.EqualFold(usage.ProviderRegion, "us") && price.usRegionSurcharge { + regionNumerator, regionDenominator = 11, 10 } scaled := func(tokens, rate, numerator, denominator int64) int64 { return tokens * rate * numerator / denominator } - cost := scaled(usage.InputTokens, inputRate, geoNumerator, geoDenominator) - cost += scaled(usage.OutputTokens, outputRate, geoNumerator, geoDenominator) - cost += scaled(usage.CacheCreation.Ephemeral5mInputTokens, inputRate*5, geoNumerator, geoDenominator*4) - cost += scaled(usage.CacheCreation.Ephemeral1hInputTokens, inputRate*2, geoNumerator, geoDenominator) - cost += scaled(usage.CacheReadInputTokens, inputRate, geoNumerator, geoDenominator*10) + cost := scaled(usage.InputTokens, inputRate, regionNumerator, regionDenominator) + cost += scaled(usage.OutputTokens, outputRate, regionNumerator, regionDenominator) + cost += scaled(usage.CacheCreation.Ephemeral5mInputTokens, inputRate*5, regionNumerator, regionDenominator*4) + cost += scaled(usage.CacheCreation.Ephemeral1hInputTokens, inputRate*2, regionNumerator, regionDenominator) + cost += scaled(usage.CacheReadInputTokens, inputRate, regionNumerator, regionDenominator*10) cost += usage.ServerToolUse.WebSearchRequests * webSearchRequestNanoUSD return cost, nil } @@ -192,19 +192,19 @@ func anthropicModelListPrice(id string) (modelListPrice, bool) { id = canonicalAnthropicModelID(id) switch id { case "claude-fable-5", "claude-mythos-5": - return modelListPrice{inputPerToken: 10_000, outputPerToken: 50_000, geoSurcharge: true}, true + return modelListPrice{inputPerToken: 10_000, outputPerToken: 50_000, usRegionSurcharge: true}, true case "claude-opus-5": - return modelListPrice{inputPerToken: 5_000, outputPerToken: 25_000, geoSurcharge: true, fastEligible: true}, true + return modelListPrice{inputPerToken: 5_000, outputPerToken: 25_000, usRegionSurcharge: true, fastEligible: true}, true case "claude-sonnet-5": - return modelListPrice{inputPerToken: 2_000, outputPerToken: 10_000, geoSurcharge: true}, true + return modelListPrice{inputPerToken: 2_000, outputPerToken: 10_000, usRegionSurcharge: true}, true case "claude-opus-4-8": - return modelListPrice{inputPerToken: 5_000, outputPerToken: 25_000, geoSurcharge: true, fastEligible: true}, true + return modelListPrice{inputPerToken: 5_000, outputPerToken: 25_000, usRegionSurcharge: true, fastEligible: true}, true case "claude-opus-4-7", "claude-opus-4-6", "claude-opus-4-5": - return modelListPrice{inputPerToken: 5_000, outputPerToken: 25_000, geoSurcharge: id == "claude-opus-4-7" || id == "claude-opus-4-6"}, true + return modelListPrice{inputPerToken: 5_000, outputPerToken: 25_000, usRegionSurcharge: id == "claude-opus-4-7" || id == "claude-opus-4-6"}, true case "claude-opus-4-1", "claude-opus-4": return modelListPrice{inputPerToken: 15_000, outputPerToken: 75_000}, true case "claude-sonnet-4-6": - return modelListPrice{inputPerToken: 3_000, outputPerToken: 15_000, geoSurcharge: true}, true + return modelListPrice{inputPerToken: 3_000, outputPerToken: 15_000, usRegionSurcharge: true}, true case "claude-sonnet-4-5", "claude-sonnet-4": return modelListPrice{inputPerToken: 3_000, outputPerToken: 15_000}, true case "claude-haiku-4-5": diff --git a/internal/domain/budget_test.go b/internal/domain/budget_test.go index 933973a..800d7cc 100644 --- a/internal/domain/budget_test.go +++ b/internal/domain/budget_test.go @@ -15,9 +15,10 @@ func TestModelUsageListCostNanoUSD(t *testing.T) { CacheReadInputTokens: 400, ServerToolUse: ServerToolUsage{WebSearchRequests: 2, WebFetchRequests: 7}, Speed: "standard", + ProviderRegion: "us", } cost, err := ModelUsageListCostNanoUSD( - Model{ID: "claude-opus-4-8-20260801", InferenceGeo: "us"}, + Model{ID: "claude-opus-4-8-20260801"}, usage, ) if err != nil { @@ -29,6 +30,13 @@ func TestModelUsageListCostNanoUSD(t *testing.T) { if amount := MonetaryAmountJSON(cost)["amount"]; amount != "3" { t.Fatalf("rounded monetary amount = %v, want 3 cents", amount) } + usage.ProviderRegion = "" + baseCost, err := ModelUsageListCostNanoUSD( + Model{ID: "claude-opus-4-8-20260801"}, usage, + ) + if err != nil || baseCost != 31_950_000 { + t.Fatalf("base list cost = %d nanoUSD, err=%v", baseCost, err) + } } func TestModelUsageListCostRejectsUnknownAndUnsupportedReportedFastMode(t *testing.T) { diff --git a/internal/domain/usage.go b/internal/domain/usage.go index 23b2904..4418f75 100644 --- a/internal/domain/usage.go +++ b/internal/domain/usage.go @@ -28,6 +28,10 @@ type TokenUsage struct { // intentionally not accumulated into Session usage; span events use it to // report the actual mode (which may differ from a requested fast fallback). Speed string + // ProviderRegion is the provider-reported region for one model request. It is + // retained only for internal list-cost accounting and is not accumulated or + // exposed through Mango's public Agent, Session, or usage resources. + ProviderRegion string } // ContextTokens returns the provider-measured context immediately after one diff --git a/internal/httpapi/agents_test.go b/internal/httpapi/agents_test.go index 3ed9c48..b9fa088 100644 --- a/internal/httpapi/agents_test.go +++ b/internal/httpapi/agents_test.go @@ -387,6 +387,29 @@ func TestAgents_ModelEffortAcceptsOfficialInputShapes(t *testing.T) { } } +func TestAgents_InferenceGeoIsNotPartOfModelConfiguration(t *testing.T) { + srv := newTestServer(t) + rec := do(srv, http.MethodPost, "/v1/agents", + `{"name":"Agent","model":{"id":"claude-opus-4-8","inference_geo":"us"}}`) + if rec.Code != http.StatusBadRequest { + t.Fatalf("status = %d, want 400: %s", rec.Code, rec.Body) + } + var body struct { + Type string `json:"type"` + Error struct { + Type string `json:"type"` + Message string `json:"message"` + } `json:"error"` + } + if err := json.Unmarshal(rec.Body.Bytes(), &body); err != nil { + t.Fatal(err) + } + if body.Type != "error" || body.Error.Type != "invalid_request_error" || + body.Error.Message != `unknown model field "inference_geo"` { + t.Fatalf("error envelope = %+v", body) + } +} + func TestAgents_MetadataValidationUsesResultingBag(t *testing.T) { srv := newTestServer(t) rec := do(srv, "POST", "/v1/agents", diff --git a/internal/httpapi/dto.go b/internal/httpapi/dto.go index 3f407a7..ce1045b 100644 --- a/internal/httpapi/dto.go +++ b/internal/httpapi/dto.go @@ -4,7 +4,6 @@ import ( "bytes" "encoding/json" "fmt" - "strings" "github.com/yanpgwang/mango/internal/domain" ) @@ -25,7 +24,7 @@ func parseModel(raw any) (domain.Model, error) { m := domain.Model{} for key := range v { switch key { - case "id", "effort", "speed", "inference_geo": + case "id", "effort", "speed": default: return domain.Model{}, domain.Validation(fmt.Sprintf("unknown model field %q", key)) } @@ -60,13 +59,6 @@ func parseModel(raw any) (domain.Model, error) { m.Speed = sp m.SpeedExplicit = true } - if rawGeo, present := v["inference_geo"]; present { - geo, ok := rawGeo.(string) - if !ok || strings.TrimSpace(geo) == "" { - return domain.Model{}, domain.Validation("model inference_geo must be a non-empty string") - } - m.InferenceGeo = geo - } if m.ID == "" { return domain.Model{}, domain.Validation("model id is required") } @@ -123,9 +115,6 @@ func agentToJSON(a domain.Agent) map[string]any { if a.Model.Speed != "" { model["speed"] = a.Model.Speed } - if a.Model.InferenceGeo != "" { - model["inference_geo"] = a.Model.InferenceGeo - } out := map[string]any{ "id": a.ID, "type": "agent", "version": a.Version, "name": a.Name, "model": model, "metadata": orEmptyMap(a.Metadata), "multiagent": a.Multiagent, diff --git a/internal/httpapi/openapi.yaml b/internal/httpapi/openapi.yaml index f1882c6..eee8b98 100644 --- a/internal/httpapi/openapi.yaml +++ b/internal/httpapi/openapi.yaml @@ -3515,10 +3515,6 @@ components: speed: type: string enum: [standard, fast] - inference_geo: - type: string - minLength: 1 - description: Geographic region pinned for every inference request made by this Agent. ModelEffortLevel: type: string enum: [low, medium, high, xhigh, max] @@ -3541,9 +3537,6 @@ components: speed: type: string enum: [standard, fast] - inference_geo: - type: string - minLength: 1 MonetaryAmount: type: object diff --git a/internal/httpapi/sdk_test.go b/internal/httpapi/sdk_test.go index fc3dace..670c20c 100644 --- a/internal/httpapi/sdk_test.go +++ b/internal/httpapi/sdk_test.go @@ -725,103 +725,6 @@ func TestSDK_SessionMultiagentRosterExpandsAndFreezesAgentSnapshots(t *testing.T } } -func TestSDK_AgentInferenceGeoRoundTripsAndClearsOnModelReplacement(t *testing.T) { - client, _ := sdkClientAndServer(t) - ctx := context.Background() - agent, err := client.Beta.Agents.New(ctx, anthropic.BetaAgentNewParams{ - Name: "Geo-pinned Agent", - Model: anthropic.BetaManagedAgentsModelConfigParams{ - ID: anthropic.BetaManagedAgentsModelClaudeOpus4_8, InferenceGeo: anthropic.String("us"), - }, - }) - if err != nil { - t.Fatalf("create geo-pinned Agent: %v", err) - } - if agent.Model.InferenceGeo != "us" || !agent.Model.JSON.InferenceGeo.Valid() { - t.Fatalf("created model = %s", agent.Model.RawJSON()) - } - - updated, err := client.Beta.Agents.Update(ctx, agent.ID, anthropic.BetaAgentUpdateParams{ - Model: anthropic.BetaManagedAgentsModelConfigParams{ - ID: anthropic.BetaManagedAgentsModelClaudeOpus4_8, - }, - }) - if err != nil { - t.Fatalf("clear inference_geo through whole-model replacement: %v", err) - } - if updated.Model.InferenceGeo != "" || updated.Model.JSON.InferenceGeo.Valid() { - t.Fatalf("updated model retained inference_geo: %s", updated.Model.RawJSON()) - } -} - -func TestSDK_SessionInferenceGeoOverridePreservesRosterInvariant(t *testing.T) { - client, server := sdkClientAndServer(t) - ctx := context.Background() - model := func(geo string) anthropic.BetaManagedAgentsModelConfigParams { - return anthropic.BetaManagedAgentsModelConfigParams{ - ID: anthropic.BetaManagedAgentsModelClaudeOpus4_8, InferenceGeo: anthropic.String(geo), - } - } - peer, err := client.Beta.Agents.New(ctx, anthropic.BetaAgentNewParams{ - Name: "Global peer", Model: model("global"), - }) - if err != nil { - t.Fatal(err) - } - coordinator, err := client.Beta.Agents.New(ctx, anthropic.BetaAgentNewParams{ - Name: "Global coordinator", Model: model("global"), - Multiagent: anthropic.BetaManagedAgentsMultiagentParams{ - Type: anthropic.BetaManagedAgentsMultiagentParamsTypeCoordinator, - Agents: []anthropic.BetaManagedAgentsMultiagentRosterEntryParamsUnion{{ - OfString: anthropic.String(peer.ID), - }}, - }, - }) - if err != nil { - t.Fatal(err) - } - environmentID := mustEnv(t, server.URL) - _, err = client.Beta.Sessions.New(ctx, anthropic.BetaSessionNewParams{ - Agent: anthropic.BetaSessionNewParamsAgentUnion{ - OfBetaManagedAgentsAgentWithOverridess: &anthropic.BetaManagedAgentsAgentWithOverridesParams{ - ID: coordinator.ID, Type: anthropic.BetaManagedAgentsAgentWithOverridesParamsTypeAgentWithOverrides, - Model: model("us"), - }, - }, - EnvironmentID: environmentID, - }) - assertAPIStatus(t, err, http.StatusBadRequest) - - selfEntry := anthropic.BetaManagedAgentsMultiagentRosterEntryParamsOfBetaManagedAgentsMultiagentSelfs( - anthropic.BetaManagedAgentsMultiagentSelfParamsTypeSelf, - ) - selfCoordinator, err := client.Beta.Agents.New(ctx, anthropic.BetaAgentNewParams{ - Name: "Self coordinator", Model: model("global"), - Multiagent: anthropic.BetaManagedAgentsMultiagentParams{ - Type: anthropic.BetaManagedAgentsMultiagentParamsTypeCoordinator, - Agents: []anthropic.BetaManagedAgentsMultiagentRosterEntryParamsUnion{selfEntry}, - }, - }) - if err != nil { - t.Fatal(err) - } - session, err := client.Beta.Sessions.New(ctx, anthropic.BetaSessionNewParams{ - Agent: anthropic.BetaSessionNewParamsAgentUnion{ - OfBetaManagedAgentsAgentWithOverridess: &anthropic.BetaManagedAgentsAgentWithOverridesParams{ - ID: selfCoordinator.ID, Type: anthropic.BetaManagedAgentsAgentWithOverridesParamsTypeAgentWithOverrides, - Model: model("us"), - }, - }, - EnvironmentID: environmentID, - }) - if err != nil { - t.Fatalf("self copies should inherit the session model override: %v", err) - } - if session.Agent.Model.InferenceGeo != "us" { - t.Fatalf("session model = %s", session.Agent.Model.RawJSON()) - } -} - func TestSDK_SkillReferencesPinAcrossAgentAndSessionSnapshots(t *testing.T) { client, server, sessions := sdkClientServerAndSessions(t) ctx := context.Background() diff --git a/internal/httpapi/testhelpers_test.go b/internal/httpapi/testhelpers_test.go index 2bcad69..7c827a7 100644 --- a/internal/httpapi/testhelpers_test.go +++ b/internal/httpapi/testhelpers_test.go @@ -418,12 +418,6 @@ func (s *testSessionService) Create( if input.Overrides != nil { snapshot = snapshot.WithOverrides(*input.Overrides) } - if snapshot.Model.InferenceGeo != agent.Model.InferenceGeo && - agent.Multiagent.HasExternalAgent(agent.ID) { - return domain.Session{}, domain.Validation( - "agent override model.inference_geo must match every independently referenced multiagent roster member", - ) - } snapshot.Skills, err = app.ResolveAgentSkillReferences(ctx, s.skillRef, snapshot.Skills) if err != nil { return domain.Session{}, err diff --git a/internal/model/anthropic.go b/internal/model/anthropic.go index e56bce6..cddc33a 100644 --- a/internal/model/anthropic.go +++ b/internal/model/anthropic.go @@ -107,7 +107,6 @@ type wireTool struct { type wireRequest struct { Model string `json:"model"` Speed string `json:"speed,omitempty"` - InferenceGeo string `json:"inference_geo,omitempty"` OutputConfig *wireOutputConfig `json:"output_config,omitempty"` System string `json:"system,omitempty"` MaxTokens int `json:"max_tokens"` @@ -129,6 +128,7 @@ type wireUsage struct { OutputTokens int64 `json:"output_tokens"` ServerToolUse wireServerToolUsage `json:"server_tool_use"` Speed string `json:"speed"` + InferenceGeo string `json:"inference_geo"` } type wireServerToolUsage struct { WebFetchRequests int64 `json:"web_fetch_requests"` @@ -149,6 +149,7 @@ type wireUsagePatch struct { OutputTokens *int64 `json:"output_tokens"` ServerToolUse *wireServerToolUsagePatch `json:"server_tool_use"` Speed *string `json:"speed"` + InferenceGeo *string `json:"inference_geo"` } type wireResponse struct { Content []json.RawMessage `json:"content"` @@ -169,7 +170,7 @@ func (a *Anthropic) buildWireRequest(req Request, stream bool) (wireRequest, err maxTokens = defaultMaxTokens } body := wireRequest{ - Model: model, InferenceGeo: req.InferenceGeo, System: req.System, + Model: model, System: req.System, MaxTokens: maxTokens, Stream: stream, } // high and standard are the Managed Agents and Messages defaults. Omitting @@ -607,6 +608,9 @@ func applyWireUsagePatch(usage *wireUsage, patch wireUsagePatch) { if patch.Speed != nil { usage.Speed = *patch.Speed } + if patch.InferenceGeo != nil { + usage.InferenceGeo = *patch.InferenceGeo + } } func usageFromWire(usage wireUsage) domain.TokenUsage { @@ -622,7 +626,8 @@ func usageFromWire(usage wireUsage) domain.TokenUsage { WebFetchRequests: usage.ServerToolUse.WebFetchRequests, WebSearchRequests: usage.ServerToolUse.WebSearchRequests, }, - Speed: usage.Speed, + Speed: usage.Speed, + ProviderRegion: usage.InferenceGeo, } } diff --git a/internal/model/anthropic_test.go b/internal/model/anthropic_test.go index e2f8004..48bed13 100644 --- a/internal/model/anthropic_test.go +++ b/internal/model/anthropic_test.go @@ -28,7 +28,7 @@ func TestAnthropic_SendsMessagesAndParsesResponse(t *testing.T) { b, _ := io.ReadAll(r.Body) _ = json.Unmarshal(b, &gotBody) w.Header().Set("Content-Type", "application/json") - _, _ = w.Write([]byte(`{"content":[{"type":"text","text":"hi back"}],"stop_reason":"end_turn","usage":{"cache_creation":{"ephemeral_1h_input_tokens":3,"ephemeral_5m_input_tokens":4},"cache_read_input_tokens":5,"input_tokens":11,"output_tokens":7,"server_tool_use":{"web_fetch_requests":2,"web_search_requests":1},"speed":"standard"}}`)) + _, _ = w.Write([]byte(`{"content":[{"type":"text","text":"hi back"}],"stop_reason":"end_turn","usage":{"cache_creation":{"ephemeral_1h_input_tokens":3,"ephemeral_5m_input_tokens":4},"cache_read_input_tokens":5,"input_tokens":11,"output_tokens":7,"server_tool_use":{"web_fetch_requests":2,"web_search_requests":1},"speed":"standard","inference_geo":"us"}}`)) })) defer srv.Close() @@ -39,12 +39,11 @@ func TestAnthropic_SendsMessagesAndParsesResponse(t *testing.T) { t.Fatal(err) } resp, err := c.CreateMessage(context.Background(), Request{ - Model: "claude-x", - Effort: "max", - Speed: "fast", - InferenceGeo: "us", - System: "sys", - Messages: []domain.Message{{Role: domain.RoleUser, Content: []domain.ContentBlock{{Type: "text", Text: "hi"}}}}, + Model: "claude-x", + Effort: "max", + Speed: "fast", + System: "sys", + Messages: []domain.Message{{Role: domain.RoleUser, Content: []domain.ContentBlock{{Type: "text", Text: "hi"}}}}, }) if err != nil { t.Fatal(err) @@ -64,8 +63,8 @@ func TestAnthropic_SendsMessagesAndParsesResponse(t *testing.T) { if gotBody["speed"] != "fast" { t.Errorf("body speed = %v, want fast", gotBody["speed"]) } - if gotBody["inference_geo"] != "us" { - t.Errorf("body inference_geo = %v, want us", gotBody["inference_geo"]) + if _, present := gotBody["inference_geo"]; present { + t.Errorf("provider-specific inference_geo leaked into request: %#v", gotBody) } outputConfig, ok := gotBody["output_config"].(map[string]any) if !ok || outputConfig["effort"] != "max" { @@ -80,7 +79,7 @@ func TestAnthropic_SendsMessagesAndParsesResponse(t *testing.T) { resp.Usage.CacheCreation.Ephemeral5mInputTokens != 4 || resp.Usage.ServerToolUse.WebFetchRequests != 2 || resp.Usage.ServerToolUse.WebSearchRequests != 1 || - resp.Usage.Speed != "standard" { + resp.Usage.Speed != "standard" || resp.Usage.ProviderRegion != "us" { t.Fatalf("usage = %#v", resp.Usage) } } @@ -172,13 +171,13 @@ func TestAnthropic_OmitsSemanticModelDefaults(t *testing.T) { t.Fatalf("default speed must be omitted for compatible endpoints: %s", encoded) } if _, present := wire["inference_geo"]; present { - t.Fatalf("unset inference_geo must be omitted: %s", encoded) + t.Fatalf("provider-specific inference_geo must not be configurable: %s", encoded) } } func TestDecodeMessageStream_AccumulatesUsage(t *testing.T) { stream := strings.NewReader( - "data: {\"type\":\"message_start\",\"message\":{\"usage\":{\"cache_creation\":{\"ephemeral_1h_input_tokens\":2,\"ephemeral_5m_input_tokens\":3},\"cache_read_input_tokens\":4,\"input_tokens\":10,\"output_tokens\":1}}}\n\n" + + "data: {\"type\":\"message_start\",\"message\":{\"usage\":{\"cache_creation\":{\"ephemeral_1h_input_tokens\":2,\"ephemeral_5m_input_tokens\":3},\"cache_read_input_tokens\":4,\"input_tokens\":10,\"output_tokens\":1,\"inference_geo\":\"us\"}}}\n\n" + "data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"text\",\"text\":\"\"}}\n\n" + "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"text_delta\",\"text\":\"ok\"}}\n\n" + "data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"end_turn\"},\"usage\":{\"input_tokens\":10,\"output_tokens\":6,\"server_tool_use\":{\"web_fetch_requests\":2,\"web_search_requests\":1}}}\n\n" + @@ -193,7 +192,8 @@ func TestDecodeMessageStream_AccumulatesUsage(t *testing.T) { resp.Usage.CacheCreation.Ephemeral1hInputTokens != 2 || resp.Usage.CacheCreation.Ephemeral5mInputTokens != 3 || resp.Usage.ServerToolUse.WebFetchRequests != 2 || - resp.Usage.ServerToolUse.WebSearchRequests != 1 { + resp.Usage.ServerToolUse.WebSearchRequests != 1 || + resp.Usage.ProviderRegion != "us" { t.Fatalf("stream usage = %#v", resp.Usage) } } diff --git a/internal/model/client.go b/internal/model/client.go index f6d2c7e..2b973f1 100644 --- a/internal/model/client.go +++ b/internal/model/client.go @@ -20,14 +20,13 @@ type ToolSchema struct { } type Request struct { - Model string - Effort string - Speed string - InferenceGeo string - System string - Messages []domain.Message - MaxTokens int - Tools []ToolSchema + Model string + Effort string + Speed string + System string + Messages []domain.Message + MaxTokens int + Tools []ToolSchema } type Response struct { diff --git a/internal/model/context.go b/internal/model/context.go index a537885..d101c49 100644 --- a/internal/model/context.go +++ b/internal/model/context.go @@ -180,16 +180,15 @@ func estimateFullRequestTokens(request Request) int { // can be estimated as a delta. func RequestContextFingerprint(request Request) string { value := struct { - Model string - Effort string - Speed string - InferenceGeo string - System string - MaxTokens int - Tools []ToolSchema + Model string + Effort string + Speed string + System string + MaxTokens int + Tools []ToolSchema }{ Model: request.Model, Effort: request.Effort, Speed: request.Speed, - InferenceGeo: request.InferenceGeo, System: request.System, + System: request.System, MaxTokens: RequestContextLimits(request).MaxOutputTokens, Tools: request.Tools, } diff --git a/internal/pg/advisor.go b/internal/pg/advisor.go index bfd5578..1fc307f 100644 --- a/internal/pg/advisor.go +++ b/internal/pg/advisor.go @@ -135,9 +135,8 @@ SELECT EXISTS( if usageModel == "" { usageModel = configured.Model } - inferenceGeo := session.AgentSnapshot.Model.InferenceGeo listCost, priceErr := domain.ModelResponseListCostNanoUSDAt( - domain.Model{ID: usageModel, InferenceGeo: inferenceGeo}, consultation.Usage, + domain.Model{ID: usageModel}, consultation.Usage, consultation.StopReason, recordedAt, ) listCostKnown := consultation.UsageKnown && priceErr == nil @@ -173,11 +172,11 @@ INSERT INTO session_threads ( } if _, err := tx.Exec(ctx, ` INSERT INTO model_request_usage ( - session_id, thread_id, request_event_id, model_id, inference_geo, + session_id, thread_id, request_event_id, model_id, stop_reason, usage, list_cost_nano_usd, created_at -) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9)`, +) VALUES ($1, $2, $3, $4, $5, $6, $7, $8)`, sessionID, thread.ID, consultation.UsageRequestID, usageModel, - inferenceGeo, consultation.StopReason, usageJSON, storedCost, recordedAt, + consultation.StopReason, usageJSON, storedCost, recordedAt, ); err != nil { return err } diff --git a/internal/pg/advisor_test.go b/internal/pg/advisor_test.go index 407feb0..7461354 100644 --- a/internal/pg/advisor_test.go +++ b/internal/pg/advisor_test.go @@ -17,9 +17,7 @@ func TestAdvisorConsultationsProjectThreadsEventsAndUsageIdempotently(t *testing session := newSession("sesn_advisor") session.AgentSnapshot = domain.Agent{ ID: session.AgentID, Version: session.AgentVersion, Name: "coordinator", - Model: domain.NormalizeModel(domain.Model{ - ID: "claude-sonnet-5", InferenceGeo: "us", - }), + Model: domain.NormalizeModel(domain.Model{ID: "claude-sonnet-5"}), Multiagent: &domain.Multiagent{Type: "coordinator", Agents: []domain.AgentReference{{ Type: "advisor", Model: "claude-opus-5", }}}, @@ -63,7 +61,7 @@ INSERT INTO agents ( t.Fatalf("primary Threads = %+v, err=%v", threads, err) } primary := threads[0] - executorUsage := domain.TokenUsage{InputTokens: 500, OutputTokens: 50} + executorUsage := domain.TokenUsage{InputTokens: 500, OutputTokens: 50, ProviderRegion: "us"} if err := store.AccountModelRequest( ctx, session.ID, @@ -79,13 +77,13 @@ INSERT INTO agents ( consultations := []domain.AdvisorConsultation{ advisorConsultationFixture( "sthr_advisor_plain", "sevt_advisor_usage_plain", "plain", - domain.TokenUsage{InputTokens: 1_000, OutputTokens: 100}, + domain.TokenUsage{InputTokens: 1_000, OutputTokens: 100, ProviderRegion: "us"}, []any{map[string]any{"type": "text", "text": "check the shutdown race"}}, true, ), advisorConsultationFixture( "sthr_advisor_second", "sevt_advisor_usage_second", "second", - domain.TokenUsage{InputTokens: 2_000, OutputTokens: 200}, + domain.TokenUsage{InputTokens: 2_000, OutputTokens: 200, ProviderRegion: "us"}, []any{map[string]any{"type": "text", "text": "challenge the locking assumption"}}, true, ), @@ -190,11 +188,11 @@ INSERT INTO agents ( if err != nil { t.Fatal(err) } - wantAdvisorUsage := domain.TokenUsage{InputTokens: 3_000, OutputTokens: 300} + wantAdvisorUsage := domain.TokenUsage{InputTokens: 3_000, OutputTokens: 300, ProviderRegion: "us"} wantUsage := executorUsage wantUsage.Add(wantAdvisorUsage) wantAdvisorCost, err := domain.ModelUsageListCostNanoUSD( - domain.Model{ID: "claude-opus-5", InferenceGeo: "us"}, wantAdvisorUsage, + domain.Model{ID: "claude-opus-5"}, wantAdvisorUsage, ) if err != nil { t.Fatal(err) @@ -281,25 +279,16 @@ INSERT INTO agents ( } var usageRows, terminationOutbox int - var inferenceGeo string if err := store.pool.QueryRow(ctx, ` SELECT count(*) FROM model_request_usage WHERE session_id = $1`, session.ID).Scan(&usageRows); err != nil { t.Fatal(err) } if err := store.pool.QueryRow(ctx, ` -SELECT inference_geo FROM model_request_usage -WHERE session_id = $1 AND thread_id = 'sthr_advisor_plain'`, session.ID).Scan(&inferenceGeo); err != nil { - t.Fatal(err) - } - if err := store.pool.QueryRow(ctx, ` SELECT count(*) FROM thread_orchestration_outbox WHERE session_id = $1`, session.ID).Scan(&terminationOutbox); err != nil { t.Fatal(err) } - if usageRows != 4 || terminationOutbox != 0 || inferenceGeo != "us" { - t.Fatalf( - "Advisor persistence rows usage=%d outbox=%d inference_geo=%q", - usageRows, terminationOutbox, inferenceGeo, - ) + if usageRows != 4 || terminationOutbox != 0 { + t.Fatalf("Advisor persistence rows usage=%d outbox=%d", usageRows, terminationOutbox) } if _, err := store.pool.Exec(ctx, ` INSERT INTO provider_transcript_turns ( diff --git a/internal/pg/budget.go b/internal/pg/budget.go index 0f5d529..a4821c5 100644 --- a/internal/pg/budget.go +++ b/internal/pg/budget.go @@ -166,12 +166,12 @@ func (s *Store) AccountModelRequest( } command, err := tx.Exec(ctx, ` INSERT INTO model_request_usage ( - session_id, thread_id, request_event_id, model_id, inference_geo, stop_reason, + session_id, thread_id, request_event_id, model_id, stop_reason, usage, list_cost_nano_usd, created_at -) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9) +) VALUES ($1, $2, $3, $4, $5, $6, $7, $8) ON CONFLICT (session_id, request_event_id) DO NOTHING`, - sessionID, threadID, requestEventID, model.ID, model.InferenceGeo, - stopReason, usageJSON, storedCost, pricedAt, + sessionID, threadID, requestEventID, model.ID, stopReason, + usageJSON, storedCost, pricedAt, ) if err != nil { return err diff --git a/internal/pg/migrations/00031_model_request_usage.sql b/internal/pg/migrations/00031_model_request_usage.sql index 183e8b5..950fdb9 100644 --- a/internal/pg/migrations/00031_model_request_usage.sql +++ b/internal/pg/migrations/00031_model_request_usage.sql @@ -6,11 +6,10 @@ -- idempotent while Session and Thread projections are updated atomically. CREATE TABLE model_request_usage ( session_id text NOT NULL REFERENCES sessions (id) ON DELETE CASCADE, - thread_id text NOT NULL, - request_event_id text NOT NULL, - model_id text NOT NULL, - inference_geo text NOT NULL DEFAULT '', - stop_reason text NOT NULL DEFAULT '', + thread_id text NOT NULL, + request_event_id text NOT NULL, + model_id text NOT NULL, + stop_reason text NOT NULL DEFAULT '', usage jsonb NOT NULL, list_cost_nano_usd bigint, created_at timestamptz NOT NULL, diff --git a/internal/pg/pgstore/models.go b/internal/pg/pgstore/models.go index c096b4a..a5ab374 100644 --- a/internal/pg/pgstore/models.go +++ b/internal/pg/pgstore/models.go @@ -183,7 +183,6 @@ type ModelRequestUsage struct { ThreadID string RequestEventID string ModelID string - InferenceGeo string StopReason string Usage []byte ListCostNanoUsd *int64 diff --git a/internal/temporal/activities.go b/internal/temporal/activities.go index 44769c3..bb985ce 100644 --- a/internal/temporal/activities.go +++ b/internal/temporal/activities.go @@ -901,10 +901,9 @@ func (a *Activities) PrepareTurn(ctx context.Context, in PrepareTurnInput) (Prep IsChild: executionThread != nil && executionThread.ParentThreadID != nil, SkillRuntimeRoot: runtimeSkills.Root, Request: model.Request{ - Model: executionAgent.Model.ID, - InferenceGeo: executionAgent.Model.InferenceGeo, - System: system, - Tools: toolSchemas, + Model: executionAgent.Model.ID, + System: system, + Tools: toolSchemas, }, } result.SessionOutputsEnabled = !selfHosted && !result.IsChild && @@ -1442,7 +1441,7 @@ func (a *Activities) EvaluateOutcome( return EvaluateOutcomeResult{FatalError: err.Error()}, nil } graderRequest := model.Request{ - Model: in.Model, Effort: in.Effort, Speed: in.Speed, InferenceGeo: in.InferenceGeo, + Model: in.Model, Effort: in.Effort, Speed: in.Speed, System: outcomeGraderSystem + " Return exactly one JSON object with " + `{"result":"satisfied|needs_revision|failed","explanation":"..."}.`, MaxTokens: 1024, diff --git a/internal/temporal/agent_workflow.go b/internal/temporal/agent_workflow.go index cb93072..c2a7eac 100644 --- a/internal/temporal/agent_workflow.go +++ b/internal/temporal/agent_workflow.go @@ -648,7 +648,6 @@ func runWorkflowTurnInternal( Model: prepared.Request.Model, Effort: prepared.Request.Effort, Speed: prepared.Request.Speed, - InferenceGeo: prepared.Request.InferenceGeo, Outcome: *prepared.Outcome, Candidate: candidate, Iteration: outcomeIteration, diff --git a/internal/temporal/agent_workflow_test.go b/internal/temporal/agent_workflow_test.go index 5587a21..5c72be5 100644 --- a/internal/temporal/agent_workflow_test.go +++ b/internal/temporal/agent_workflow_test.go @@ -147,9 +147,7 @@ func TestWorkflowTurn_AdmitsAndAccountsEveryModelRequestBeforeCompletion(t *test func(context.Context, PrepareTurnInput) (PrepareTurnResult, error) { return PrepareTurnResult{ ThreadID: "sthr_usage", - Request: model.Request{ - Model: "claude-opus-4-8", InferenceGeo: "us", - }, + Request: model.Request{Model: "claude-opus-4-8"}, }, nil }, activity.RegisterOptions{Name: ActivityPrepareTurn}, @@ -225,7 +223,6 @@ func TestWorkflowTurn_AdmitsAndAccountsEveryModelRequestBeforeCompletion(t *test require.Equal(t, "sthr_usage", accounted.ThreadID) require.NotEmpty(t, accounted.RequestEventID) require.Equal(t, "claude-opus-4-8", accounted.Model.ID) - require.Equal(t, "us", accounted.Model.InferenceGeo) require.Equal(t, "end_turn", accounted.StopReason) require.Equal(t, int64(7), accounted.Usage.InputTokens) require.Equal(t, int64(1), accounted.Usage.ServerToolUse.WebSearchRequests) @@ -247,7 +244,7 @@ func TestWorkflowTurn_AdvisorIsPrivatePortableToolWithIndependentRequest(t *test return PrepareTurnResult{ AttemptID: "ratm_advisor_workflow", ThreadID: "sthr_primary", Request: model.Request{ - Model: "executor-model", InferenceGeo: "us", + Model: "executor-model", System: "executor system", Tools: []model.ToolSchema{agentruntime.AdvisorToolSchema()}, Messages: []domain.Message{{ @@ -323,7 +320,6 @@ func TestWorkflowTurn_AdvisorIsPrivatePortableToolWithIndependentRequest(t *test require.Equal(t, 2, modelCalls) require.Equal(t, TurnToolAdvisor, advisorInput.ToolKind) require.Equal(t, "reviewer-model", advisorInput.AdvisorRequest.Model) - require.Equal(t, "us", advisorInput.AdvisorRequest.InferenceGeo) require.Empty(t, advisorInput.AdvisorRequest.Tools) require.Contains( t, diff --git a/internal/temporal/agent_workflow_tools.go b/internal/temporal/agent_workflow_tools.go index f5ca210..80c3742 100644 --- a/internal/temporal/agent_workflow_tools.go +++ b/internal/temporal/agent_workflow_tools.go @@ -279,7 +279,6 @@ func (t *workflowTurnState) accountOutcomeEvaluation( input.EndEventID, model.Request{ Model: input.Model, Effort: input.Effort, Speed: input.Speed, - InferenceGeo: input.InferenceGeo, }, evaluated.Usage, evaluated.StopReason, @@ -300,7 +299,6 @@ func (t *workflowTurnState) accountModelRequest( RequestEventID: requestEventID, Model: domain.Model{ ID: request.Model, Effort: request.Effort, Speed: request.Speed, - InferenceGeo: request.InferenceGeo, }, Usage: usage, StopReason: stopReason, }, diff --git a/internal/temporal/types.go b/internal/temporal/types.go index d718d56..a4e5c08 100644 --- a/internal/temporal/types.go +++ b/internal/temporal/types.go @@ -157,7 +157,6 @@ type EvaluateOutcomeInput struct { Model string `json:"model"` Effort string `json:"effort,omitempty"` Speed string `json:"speed,omitempty"` - InferenceGeo string `json:"inference_geo,omitempty"` Outcome domain.OutcomeSpec `json:"outcome"` Candidate []domain.Message `json:"candidate"` Iteration int `json:"iteration"` From 0fcdbf0def274f4b6128c63e7c690348e249a87c Mon Sep 17 00:00:00 2001 From: Yanpeng Wang Date: Thu, 27 Aug 2026 19:22:27 +0800 Subject: [PATCH 2/2] fix: isolate provider pricing metadata from context --- internal/domain/message.go | 8 ++++---- internal/domain/usage.go | 21 ++++++++++++++++++++- internal/model/context.go | 5 +++-- internal/model/context_test.go | 20 ++++++++++++++++++++ internal/pg/transcript_test.go | 6 +++--- 5 files changed, 50 insertions(+), 10 deletions(-) diff --git a/internal/domain/message.go b/internal/domain/message.go index f9da05c..d18e6eb 100644 --- a/internal/domain/message.go +++ b/internal/domain/message.go @@ -69,10 +69,10 @@ type Message struct { // provider usage had measured. ContentBlocks marks the response boundary when // adjacent assistant messages are merged to preserve role alternation. type ContextUsageAnchor struct { - Usage TokenUsage `json:"usage"` - RequestFingerprint string `json:"request_fingerprint"` - PrefixFingerprint string `json:"prefix_fingerprint"` - ContentBlocks int `json:"content_blocks"` + Usage ContextWindowUsage `json:"usage"` + RequestFingerprint string `json:"request_fingerprint"` + PrefixFingerprint string `json:"prefix_fingerprint"` + ContentBlocks int `json:"content_blocks"` } // ProviderToolUseMapping keeps provider-private tool ids separate from the diff --git a/internal/domain/usage.go b/internal/domain/usage.go index 4418f75..ccff136 100644 --- a/internal/domain/usage.go +++ b/internal/domain/usage.go @@ -34,11 +34,30 @@ type TokenUsage struct { ProviderRegion string } +// ContextWindowUsage is the provider usage subset needed to anchor context +// measurements. Pricing, routing, execution-mode, and server-tool facts do not +// belong in persisted conversation messages or subsequent model prompts. +type ContextWindowUsage struct { + CacheCreation CacheCreationUsage + CacheReadInputTokens int64 + InputTokens int64 + OutputTokens int64 +} + +func (u TokenUsage) ForContextWindow() ContextWindowUsage { + return ContextWindowUsage{ + CacheCreation: u.CacheCreation, + CacheReadInputTokens: u.CacheReadInputTokens, + InputTokens: u.InputTokens, + OutputTokens: u.OutputTokens, + } +} + // ContextTokens returns the provider-measured context immediately after one // response. Anthropic reports uncached input, cache creation, and cache reads // as disjoint input buckets; all occupy the request context. Output is included // because the assistant response becomes input to the next request. -func (u TokenUsage) ContextTokens() int64 { +func (u ContextWindowUsage) ContextTokens() int64 { return u.InputTokens + u.CacheCreation.Ephemeral1hInputTokens + u.CacheCreation.Ephemeral5mInputTokens + diff --git a/internal/model/context.go b/internal/model/context.go index d101c49..2242b63 100644 --- a/internal/model/context.go +++ b/internal/model/context.go @@ -200,13 +200,14 @@ func RequestContextFingerprint(request Request) string { // the precise merged-message boundary that the next provider request sees. func AnchoredAssistantMessage(request Request, response Response) domain.Message { message := domain.Message{Role: domain.RoleAssistant, Content: response.Content} - if response.Usage.ContextTokens() <= 0 { + contextUsage := response.Usage.ForContextWindow() + if contextUsage.ContextTokens() <= 0 { return message } anchored := appendAssistantForContext(request.Messages, message) last := len(anchored) - 1 message.ContextUsage = &domain.ContextUsageAnchor{ - Usage: response.Usage, + Usage: contextUsage, RequestFingerprint: RequestContextFingerprint(request), PrefixFingerprint: messagePrefixFingerprint( anchored, last, len(anchored[last].Content), diff --git a/internal/model/context_test.go b/internal/model/context_test.go index 6c9319e..10ef0c4 100644 --- a/internal/model/context_test.go +++ b/internal/model/context_test.go @@ -1,6 +1,7 @@ package model import ( + "encoding/json" "testing" "github.com/stretchr/testify/require" @@ -67,6 +68,25 @@ func TestMeasureRequestContextUsesLatestExactUsagePlusDelta(t *testing.T) { ) } +func TestAnchoredAssistantMessageExcludesPricingAndRoutingUsage(t *testing.T) { + anchored := AnchoredAssistantMessage(Request{Model: "claude-sonnet-5"}, Response{ + Content: []domain.ContentBlock{{Type: "text", Text: "answer"}}, + Usage: domain.TokenUsage{ + InputTokens: 10, OutputTokens: 2, Speed: "fast", ProviderRegion: "us", + ServerToolUse: domain.ServerToolUsage{WebSearchRequests: 1}, + }, + }) + require.NotNil(t, anchored.ContextUsage) + require.Equal(t, int64(12), anchored.ContextUsage.Usage.ContextTokens()) + + encoded, err := json.Marshal(anchored) + require.NoError(t, err) + require.NotContains(t, string(encoded), "ProviderRegion") + require.NotContains(t, string(encoded), "ServerToolUse") + require.NotContains(t, string(encoded), `"Speed"`) + require.NotContains(t, string(encoded), `"us"`) +} + func TestMeasureRequestContextRejectsStaleRequestAndPrefixAnchors(t *testing.T) { initial := Request{ Model: "claude-sonnet-4-6", System: "before", diff --git a/internal/pg/transcript_test.go b/internal/pg/transcript_test.go index 703d029..d85ff72 100644 --- a/internal/pg/transcript_test.go +++ b/internal/pg/transcript_test.go @@ -183,13 +183,13 @@ func TestProviderTranscriptContextIsIsolatedByThread(t *testing.T) { func TestAppendProviderMessagesPreservesLatestContextUsage(t *testing.T) { first := &domain.ContextUsageAnchor{ - Usage: domain.TokenUsage{InputTokens: 10}, + Usage: domain.ContextWindowUsage{InputTokens: 10}, RequestFingerprint: "request-1", PrefixFingerprint: "prefix-1", ContentBlocks: 1, } latest := &domain.ContextUsageAnchor{ - Usage: domain.TokenUsage{InputTokens: 20}, + Usage: domain.ContextWindowUsage{InputTokens: 20}, RequestFingerprint: "request-2", PrefixFingerprint: "prefix-2", ContentBlocks: 2, @@ -222,7 +222,7 @@ func TestAppendProviderMessagesPreservesLatestContextUsage(t *testing.T) { func TestCloseInterruptedProviderTranscript_PairsDanglingTools(t *testing.T) { anchor := &domain.ContextUsageAnchor{ - Usage: domain.TokenUsage{InputTokens: 25}, + Usage: domain.ContextWindowUsage{InputTokens: 25}, RequestFingerprint: "request", PrefixFingerprint: "prefix", ContentBlocks: 3,