Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
30 commits
Select commit Hold shift + click to select a range
c35595f
test: define semantic routing behavior
dporkka Oct 4, 2026
f1c7126
test: route configured traffic through Bifrost auto policy
dporkka Oct 4, 2026
6360d29
feat: define semantic model routes
dporkka Oct 4, 2026
886d3af
feat: delegate semantic routes to Bifrost
dporkka Oct 4, 2026
55a860e
feat: propagate routing context through Bifrost
dporkka Oct 4, 2026
6b1a39a
feat: add Maxim Bifrost adaptive routing config
dporkka Oct 4, 2026
103435d
chore: add pinned Maxim Bifrost deployment
dporkka Oct 4, 2026
3cbc6fe
docs: document adaptive model routing deployment
dporkka Oct 4, 2026
e5a0881
test: target Maxim Bifrost OpenAI endpoint
dporkka Oct 4, 2026
40d2698
docs: configure Maxim Bifrost defaults
dporkka Oct 4, 2026
55443e4
test: require routing decision outcome schema
dporkka Oct 4, 2026
7e7807a
test: define external spend authority behavior
dporkka Oct 4, 2026
f6f9534
test: define routing telemetry and outcome feedback
dporkka Oct 4, 2026
08aed57
feat: persist routing decisions and verifier outcomes
dporkka Oct 4, 2026
ba59197
feat: annotate routing policy and spend authority
dporkka Oct 4, 2026
bcf176a
feat: expose routing provenance and spend authority
dporkka Oct 4, 2026
d5f6ffd
feat: make gateway spend authority explicit
dporkka Oct 4, 2026
907c097
feat: persist routing decisions and task outcomes
dporkka Oct 4, 2026
42f336a
feat: delegate gateway spend and record routing decisions
dporkka Oct 4, 2026
c054d91
feat: feed verifier outcomes back into routing telemetry
dporkka Oct 4, 2026
aa35fb7
feat: propagate failed and paused outcomes to routing telemetry
dporkka Oct 4, 2026
e12e108
fix: close routing-aware runner summary block
dporkka Oct 4, 2026
ebf91c3
fix: keep gateway spend out of local cost ledger
dporkka Oct 4, 2026
dc70eed
feat: add routing decisions to canonical schema
dporkka Oct 4, 2026
529db39
test: verify routing provenance and spend authority
dporkka Oct 4, 2026
4199fc1
docs: document routing telemetry and spend authority
dporkka Oct 4, 2026
bc68bd9
style: gofmt cost authority tests
dporkka Oct 4, 2026
572d7b2
style: gofmt routing telemetry tests
dporkka Oct 4, 2026
a3ecacb
fix: remove unused routing telemetry helper
dporkka Oct 4, 2026
602be80
docs: define routing telemetry ownership
dporkka Oct 4, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 11 additions & 4 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -93,11 +93,17 @@ WORKER_HEALTH_PORT=8081
# AI PROVIDERS (via Bifrost gateway or direct)
# =============================================================================

# Bifrost AI Gateway (recommended -- unified interface)
BIFROST_URL=http://localhost:8083
# Maxim Bifrost AI Gateway (recommended -- unified interface).
# BIFROST_API_KEY should be a Bifrost sk-bf-* virtual key in production.
BIFROST_URL=http://localhost:8083/v1
BIFROST_API_KEY=
BIFROST_ENCRYPTION_KEY=

# Direct provider keys (fallback)
# Bifrost routes hosted inference through OpenRouter by default.
OPENROUTER_API_KEY=

# Direct provider keys (fallback). The OpenAI key is also used by the sample
# Bifrost config for its low-cost semantic complexity embedding classifier.
OPENAI_API_KEY=sk-...
OPENAI_BASE_URL=https://api.openai.com/v1

Expand All @@ -113,7 +119,8 @@ GROQ_BASE_URL=https://api.groq.com/openai/v1
FIREWORKS_API_KEY=
FIREWORKS_BASE_URL=https://api.fireworks.ai/inference/v1

# Default model for agent runs
# Default model for agent runs. When Bifrost is configured, semantic routing
# is preferred and this remains the legacy direct-provider fallback.
DEFAULT_MODEL=gpt-4o
DEFAULT_PROVIDER=openai

Expand Down
162 changes: 162 additions & 0 deletions apps/api/internal/agentrunner/routing_telemetry.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,162 @@
package agentrunner

import (
"context"
"encoding/json"
"fmt"
"strings"
"time"

"github.com/google/uuid"

"github.com/ai-dev-control-plane/api/internal/modelrouter"
"github.com/ai-dev-control-plane/models"
)

func (r *Runner) recordRoutingDecision(ctx context.Context, run *models.AgentRun, task *models.Task, stepNumber int, result *modelrouter.CallResult) error {
if r == nil || r.db == nil || run == nil || task == nil || result == nil {
return nil
}
if stepNumber <= 0 {
if err := r.db.QueryRowContext(ctx, `
SELECT COALESCE(MAX(step_number), 0) + 1
FROM agent_steps
WHERE agent_run_id = $1
`, run.ID).Scan(&stepNumber); err != nil {
return fmt.Errorf("resolve routing step number: %w", err)
}
}

route := strings.TrimSpace(result.Route)
if route == "" {
route = modelrouter.RouteAuto
}
routeSource := strings.TrimSpace(result.RouteSource)
if routeSource == "" {
routeSource = modelrouter.RouteSourceLegacyLocal
}
policyVersion := strings.TrimSpace(result.PolicyVersion)
if policyVersion == "" {
policyVersion = modelrouter.SemanticRoutingPolicyVersion
}
spendAuthority := strings.TrimSpace(result.SpendAuthority)
if spendAuthority == "" {
spendAuthority = modelrouter.SpendAuthorityDevPlane
}
model := strings.TrimSpace(result.Model)
if model == "" && run.Model != nil {
model = strings.TrimSpace(*run.Model)
}
if model == "" {
model = "unknown"
}
provider := strings.TrimSpace(result.Provider)
if provider == "" && run.Provider != nil {
provider = strings.TrimSpace(*run.Provider)
}
if provider == "" {
provider = "unknown"
}

_, err := r.db.ExecContext(ctx, `
INSERT INTO routing_decisions (
id, agent_run_id, task_id, step_number, route, route_source,
policy_version, task_type, difficulty, agent_role, risk_level,
model, provider, spend_authority, prompt_tokens, completion_tokens,
total_tokens, estimated_cost, latency_ms, provider_call_succeeded,
outcome_status, created_at
) VALUES (
$1, $2, $3, $4, $5, $6,
$7, $8, $9, $10, $11,
$12, $13, $14, $15, $16,
$17, $18, $19, $20,
'pending', $21
)
`, uuid.New().String(), run.ID, task.ID, stepNumber, route, routeSource,
policyVersion, taskTypeForRole(run.AgentRole), difficultyForRisk(string(task.RiskLevel)), run.AgentRole, string(task.RiskLevel),
model, provider, spendAuthority, result.PromptTokens, result.CompletionTokens,
result.TotalTokens, result.Cost, result.LatencyMs, true, time.Now().UTC())
if err != nil {
return fmt.Errorf("record routing decision: %w", err)
}
return nil
}

func (r *Runner) recordRoutingOutcome(ctx context.Context, runID, status string, verifierResults map[string]any, outcomeErr string) error {
if r == nil || r.db == nil || strings.TrimSpace(runID) == "" {
return nil
}

var verifierPassed any
encodedVerifier := "{}"
if verifierResults != nil {
verifierPassed = finalChecksPassed(verifierResults)
encoded, err := json.Marshal(verifierResults)
if err != nil {
return fmt.Errorf("marshal routing verifier results: %w", err)
}
encodedVerifier = string(encoded)
}

now := time.Now().UTC()
_, err := r.db.ExecContext(ctx, `
UPDATE routing_decisions
SET outcome_status = $1,
verifier_passed = COALESCE($2, verifier_passed),
verifier_results = CASE WHEN $2 IS NULL THEN verifier_results ELSE $3 END,
error = CASE WHEN $4 = '' THEN error ELSE $4 END,
outcome_recorded_at = $5
WHERE agent_run_id = $6
AND outcome_status IN ('pending', 'paused')
`, status, verifierPassed, encodedVerifier, outcomeErr, now, runID)
if err != nil {
return fmt.Errorf("record routing outcome: %w", err)
}
return nil
}

func (r *Runner) markRoutingHumanIntervention(ctx context.Context, runID string) error {
if r == nil || r.db == nil || strings.TrimSpace(runID) == "" {
return nil
}
_, err := r.db.ExecContext(ctx, `
UPDATE routing_decisions
SET outcome_status = 'paused',
human_intervention_required = true,
outcome_recorded_at = $1
WHERE agent_run_id = $2
AND outcome_status = 'pending'
`, time.Now().UTC(), runID)
if err != nil {
return fmt.Errorf("mark routing human intervention: %w", err)
}
return nil
}

func finalChecksPassed(results map[string]any) bool {
if results == nil {
return false
}
if passed, ok := results["passed"].(bool); ok {
return passed
}

tests, ok := results["tests"].(map[string]any)
if !ok {
return false
}
if passed, ok := tests["passed"].(bool); ok {
return passed
}
output, _ := tests["output"].(string)
if strings.TrimSpace(output) == "" {
return false
}
var payload struct {
Passed bool `json:"passed"`
}
if err := json.Unmarshal([]byte(output), &payload); err != nil {
return false
}
return payload.Passed
}
183 changes: 183 additions & 0 deletions apps/api/internal/agentrunner/routing_telemetry_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,183 @@
package agentrunner

import (
"context"
"database/sql"
"encoding/json"
"log/slog"
"testing"

"github.com/ai-dev-control-plane/api/internal/modelrouter"
"github.com/ai-dev-control-plane/api/internal/tools"
"github.com/ai-dev-control-plane/models"
)

func TestNextModelActionSendsStableAutoRouteContext(t *testing.T) {
provider := &fakeModelProvider{responses: []string{
`{"action":"final_response","content":"done"}`,
}}
runner := NewRunner(nil, tools.NewWorkspaceTools(slog.Default()), allowAllPolicies(), nil, nil, slog.Default()).
WithModelRouter(modelrouter.NewRouter(testRouterConfig(), provider))

run := &models.AgentRun{ID: "run-1", TaskID: "task-1", AgentRole: models.AgentRoleSecurity}
task := &models.Task{ID: "task-1", RiskLevel: models.RiskLevelCritical}

if _, _, err := runner.nextModelAction(context.Background(), run, task, "system", nil, nil); err != nil {
t.Fatalf("nextModelAction() error: %v", err)
}
if len(provider.calls) != 1 {
t.Fatalf("provider calls = %d, want 1", len(provider.calls))
}
req := provider.calls[0]
if req.Route != modelrouter.RouteAuto {
t.Fatalf("route = %q, want %q", req.Route, modelrouter.RouteAuto)
}
if req.RoutingMetadata["agent-role"] != models.AgentRoleSecurity {
t.Fatalf("agent-role metadata = %q", req.RoutingMetadata["agent-role"])
}
if req.RoutingMetadata["risk"] != string(models.RiskLevelCritical) {
t.Fatalf("risk metadata = %q", req.RoutingMetadata["risk"])
}
}

func TestRoutingDecisionTelemetryCapturesVerifierOutcome(t *testing.T) {
db := setupRunnerOrchestrationDB(t)
defer db.Close()
createRoutingDecisionTestTable(t, db)

runner := NewRunner(db, tools.NewWorkspaceTools(slog.Default()), allowAllPolicies(), nil, nil, slog.Default())
run := &models.AgentRun{ID: "run-1", TaskID: "task-1", AgentRole: models.AgentRoleImplementer}
task := &models.Task{ID: "task-1", RiskLevel: models.RiskLevelHigh}
result := &modelrouter.CallResult{
Route: modelrouter.RouteCodingDeep,
RouteSource: modelrouter.RouteSourceExplicit,
PolicyVersion: modelrouter.SemanticRoutingPolicyVersion,
Model: "deepseek/deepseek-v4.1-flash",
Provider: "bifrost",
SpendAuthority: modelrouter.SpendAuthorityGateway,
PromptTokens: 100,
CompletionTokens: 25,
TotalTokens: 125,
Cost: 0,
LatencyMs: 420,
}

if err := runner.recordRoutingDecision(context.Background(), run, task, 7, result); err != nil {
t.Fatalf("recordRoutingDecision() error: %v", err)
}

checks := map[string]any{
"tests": map[string]any{"passed": true, "exit_code": float64(0)},
"passed": true,
}
if err := runner.recordRoutingOutcome(context.Background(), run.ID, models.AgentRunStatusCompleted, checks, ""); err != nil {
t.Fatalf("recordRoutingOutcome() error: %v", err)
}

var route, source, policyVersion, taskType, difficulty, role, risk, model, provider, spendAuthority, outcome string
var step, prompt, completion, total, latency int
var estimatedCost float64
var providerSucceeded, verifierPassed, humanIntervention bool
var verifierJSON string
if err := db.QueryRow(`
SELECT step_number, route, route_source, policy_version, task_type, difficulty,
agent_role, risk_level, model, provider, spend_authority,
prompt_tokens, completion_tokens, total_tokens, estimated_cost, latency_ms,
provider_call_succeeded, outcome_status, verifier_passed,
verifier_results, human_intervention_required
FROM routing_decisions WHERE agent_run_id = 'run-1'
`).Scan(
&step, &route, &source, &policyVersion, &taskType, &difficulty,
&role, &risk, &model, &provider, &spendAuthority,
&prompt, &completion, &total, &estimatedCost, &latency,
&providerSucceeded, &outcome, &verifierPassed, &verifierJSON, &humanIntervention,
); err != nil {
t.Fatalf("query routing decision: %v", err)
}

if step != 7 || route != modelrouter.RouteCodingDeep || source != modelrouter.RouteSourceExplicit {
t.Fatalf("routing identity = step %d route %q source %q", step, route, source)
}
if policyVersion != modelrouter.SemanticRoutingPolicyVersion {
t.Fatalf("policy version = %q", policyVersion)
}
if taskType != modelrouter.TaskTypeCode || difficulty != modelrouter.DifficultyHard {
t.Fatalf("task routing context = %q/%q", taskType, difficulty)
}
if role != models.AgentRoleImplementer || risk != string(models.RiskLevelHigh) {
t.Fatalf("agent context = %q/%q", role, risk)
}
if model != result.Model || provider != result.Provider || spendAuthority != modelrouter.SpendAuthorityGateway {
t.Fatalf("resolved target = %q/%q authority=%q", provider, model, spendAuthority)
}
if prompt != 100 || completion != 25 || total != 125 || estimatedCost != 0 || latency != 420 {
t.Fatalf("usage telemetry mismatch")
}
if !providerSucceeded || outcome != models.AgentRunStatusCompleted || !verifierPassed || humanIntervention {
t.Fatalf("outcome telemetry = provider=%v outcome=%q verifier=%v intervention=%v", providerSucceeded, outcome, verifierPassed, humanIntervention)
}
var decoded map[string]any
if err := json.Unmarshal([]byte(verifierJSON), &decoded); err != nil {
t.Fatalf("decode verifier results: %v", err)
}
if passed, _ := decoded["passed"].(bool); !passed {
t.Fatalf("verifier results did not preserve passed=true: %s", verifierJSON)
}
}

func TestFinalChecksPassedUsesStructuredToolResult(t *testing.T) {
if finalChecksPassed(map[string]any{
"tests": map[string]any{
"output": `{"passed":false,"exit_code":1}`,
"error": "<nil>",
},
}) {
t.Fatal("expected failed test payload to produce verifier failure even when tool call returned nil error")
}
if !finalChecksPassed(map[string]any{
"tests": map[string]any{
"output": `{"passed":true,"exit_code":0}`,
"error": "<nil>",
},
}) {
t.Fatal("expected passed test payload to produce verifier success")
}
}

func createRoutingDecisionTestTable(t *testing.T, db *sql.DB) {
t.Helper()
_, err := db.Exec(`
CREATE TABLE routing_decisions (
id TEXT PRIMARY KEY,
agent_run_id TEXT NOT NULL,
task_id TEXT NOT NULL,
step_number INTEGER NOT NULL,
route TEXT NOT NULL,
route_source TEXT NOT NULL,
policy_version TEXT NOT NULL,
task_type TEXT NOT NULL,
difficulty TEXT NOT NULL,
agent_role TEXT NOT NULL,
risk_level TEXT NOT NULL,
model TEXT NOT NULL,
provider TEXT NOT NULL,
spend_authority TEXT NOT NULL,
prompt_tokens INTEGER NOT NULL DEFAULT 0,
completion_tokens INTEGER NOT NULL DEFAULT 0,
total_tokens INTEGER NOT NULL DEFAULT 0,
estimated_cost REAL NOT NULL DEFAULT 0,
latency_ms INTEGER NOT NULL DEFAULT 0,
provider_call_succeeded BOOLEAN NOT NULL DEFAULT true,
outcome_status TEXT NOT NULL DEFAULT 'pending',
verifier_passed BOOLEAN,
verifier_results TEXT DEFAULT '{}',
human_intervention_required BOOLEAN NOT NULL DEFAULT false,
error TEXT,
created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
outcome_recorded_at DATETIME
)
`)
if err != nil {
t.Fatalf("create routing_decisions test table: %v", err)
}
}
Loading
Loading