Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
11 changes: 10 additions & 1 deletion docs/adr/0029-llm-cost-and-usage-tracking-and-budgets.md
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
# 0029. LLM cost & token-usage tracking and budgets

- Status: Proposed
- Status: Accepted (Phase A — usage visibility — shipped; budget enforcement deferred)
- Date: 2026-06-02
- Deciders: FUSE maintainers

Expand Down Expand Up @@ -62,6 +62,15 @@ pricing) is a prerequisite and is left to the follow-up.

## More Information

- **Phase A (usage visibility) shipped**: `ai/chat` and `ai/agent` emit per-call token usage as
Prometheus counters `fuse_llm_tokens_total{function,provider,model,type}` (type ∈ prompt|
completion) and `fuse_llm_calls_total{function,provider,model,status}`
(`internal/metrics/registry.go`). A narrow `ai.UsageRecorder` port
(`internal/packages/functions/ai/usage.go`) keeps the ai package free of the prometheus
dependency; a metrics-backed adapter is injected via `packages.NewInternal`
(`internal/packages/usage_recorder.go`). Usage is still also returned in each node's `usage`
output. **Deferred**: Option B budget enforcement (needs a per-provider/per-model pricing table)
and Option C external metering.
- Current state: `pkg/llm/provider.go` (`Usage`), `internal/packages/functions/ai/agent.go`
(per-run aggregation), `internal/packages/functions/ai/chat.go`.
- Related: [ADR-0006](0006-llm-provider-abstraction-and-multi-provider-strategy.md),
Expand Down
17 changes: 9 additions & 8 deletions docs/adr/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -55,18 +55,19 @@ copying [`template.md`](template.md).
| 0026 | [Agent-as-orchestrator mode](0026-agent-as-orchestrator-mode.md) | Proposed | 2026-06-02 |
| 0027 | [Async tool invocation via a sub-execution channel](0027-async-tool-invocation-sub-execution-channel.md) | Proposed | 2026-06-02 |
| 0028 | [Agent prompt/context & conversation-memory model](0028-agent-prompt-context-and-memory-model.md) | Proposed | 2026-06-02 |
| 0029 | [LLM cost & token-usage tracking and budgets](0029-llm-cost-and-usage-tracking-and-budgets.md) | Proposed | 2026-06-02 |
| 0029 | [LLM cost & token-usage tracking and budgets](0029-llm-cost-and-usage-tracking-and-budgets.md) | Accepted | 2026-06-02 |
| 0030 | [Structured/JSON output enforcement for ai nodes](0030-structured-output-enforcement.md) | Proposed | 2026-06-02 |
| 0031 | [Settings, secrets & environments: a SecretStore seam](0031-settings-secrets-and-environments.md) | Accepted | 2026-06-02 |
| 0032 | [Sub-workflow composition: child workflows as first-class instances](0032-sub-workflow-composition.md) | Accepted | 2026-06-03 |
| 0033 | [Dependency injection & app composition with uber-go/fx](0033-dependency-injection-and-app-composition.md) | Accepted | 2026-06-03 |

### Proposed backlog (not yet implemented)

ADRs **0026–0030** are a cohesive **agent-capabilities** series exploring the next phase of the AI
agent ([ADR-0005](0005-ai-agents-as-workflow-nodes-phased-roadmap.md)) — orchestrator mode (0026),
async tool invocation (0027), prompt/context & memory (0028), cost/usage tracking & budgets (0029),
and structured-output enforcement (0030). They cross-reference one another and remain `Proposed`
(none implemented yet). **0025** (browser-automation package) is an independent, parallel stream,
not part of that series. These stay `Proposed` until scheduled; when one is implemented its status
moves to `Accepted` and its "More Information" records what shipped (as 0031 does).
ADRs **0026–0030** are a cohesive **agent-capabilities** series for the next phase of the AI agent
([ADR-0005](0005-ai-agents-as-workflow-nodes-phased-roadmap.md)) — orchestrator mode (0026), async
tool invocation (0027), prompt/context & memory (0028), cost/usage tracking & budgets (0029), and
structured-output enforcement (0030). The leaf capabilities are being implemented first: **0029**
Phase A (usage visibility) is `Accepted`/shipped; 0028 and 0030 follow; the larger orchestrator
(0026) + async-tools (0027) pair stays `Proposed` until reassessed. **0025** (browser-automation
package) is an independent, parallel stream, not part of that series. When an ADR is implemented its
status moves to `Accepted` and its "More Information" records what shipped (as 0031 does).
19 changes: 19 additions & 0 deletions internal/metrics/registry.go
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,12 @@ type FuseMetrics struct {
// Labels: function_id, status (success|error).
NodeExecDuration *prometheus.HistogramVec

// LLMTokens counts tokens consumed by ai/chat and ai/agent nodes (ADR-0029).
// Labels: function (ai/chat|ai/agent), provider, model, type (prompt|completion).
LLMTokens *prometheus.CounterVec
// LLMCalls counts LLM completion calls. Labels: function, provider, model, status (success|error).
LLMCalls *prometheus.CounterVec

registry *prometheus.Registry
}

Expand Down Expand Up @@ -57,6 +63,17 @@ func NewFuseMetrics() *FuseMetrics {
Help: "Duration of individual node (function) executions in seconds.",
Buckets: prometheus.DefBuckets,
}, []string{"function_id", "status"}),

LLMTokens: prometheus.NewCounterVec(prometheus.CounterOpts{
Namespace: "fuse",
Name: "llm_tokens_total",
Help: "Total LLM tokens consumed by ai nodes, by token type.",
}, []string{"function", "provider", "model", "type"}),
LLMCalls: prometheus.NewCounterVec(prometheus.CounterOpts{
Namespace: "fuse",
Name: "llm_calls_total",
Help: "Total LLM completion calls made by ai nodes.",
}, []string{"function", "provider", "model", "status"}),
}

reg.MustRegister(
Expand All @@ -65,6 +82,8 @@ func NewFuseMetrics() *FuseMetrics {
m.WorkflowsFailed,
m.WorkflowsCancelled,
m.NodeExecDuration,
m.LLMTokens,
m.LLMCalls,
)

return m
Expand Down
19 changes: 12 additions & 7 deletions internal/packages/functions/ai/agent.go
Original file line number Diff line number Diff line change
Expand Up @@ -60,8 +60,8 @@ func AgentFunctionMetadata() workflow.FunctionMetadata {
}

// makeAgentFunction builds the ai/agent function, closing over the provider
// registry and the tool registry.
func makeAgentFunction(providers llm.Registry, tools ToolRegistry) workflow.Function {
// registry, the tool registry, and the usage recorder (ADR-0029).
func makeAgentFunction(providers llm.Registry, tools ToolRegistry, usage UsageRecorder) workflow.Function {
return func(execInfo *workflow.ExecutionInfo) (workflow.FunctionResult, error) {
input := execInfo.Input

Expand All @@ -84,6 +84,7 @@ func makeAgentFunction(providers llm.Registry, tools ToolRegistry) workflow.Func
wfID: execInfo.WorkflowID,
execID: execInfo.ExecID,
environment: execInfo.Environment,
usage: usage,
}

messages := make([]llm.Message, 0, 4)
Expand Down Expand Up @@ -127,12 +128,13 @@ type agentExecutor struct {
wfID workflow.ID
execID workflow.ExecID
environment string
usage UsageRecorder
}

// run drives the reasoning loop until a final answer, an error, or the iteration
// limit. It returns exactly one FunctionOutput; the caller calls Finish once.
func (e *agentExecutor) run(ctx context.Context, messages []llm.Message) workflow.FunctionOutput {
usage := llm.Usage{}
totalUsage := llm.Usage{}
steps := make([]map[string]any, 0)

for i := 0; i < e.maxIters; i++ {
Expand All @@ -144,18 +146,21 @@ func (e *agentExecutor) run(ctx context.Context, messages []llm.Message) workflo
ToolChoice: "auto",
})
if err != nil {
e.usage.RecordCall(AgentFunctionID, e.provider.Name(), e.model, "error")
log.Error().Err(err).Str("provider", e.provider.Name()).Msg("ai/agent completion failed")
return errorOutput(fmt.Sprintf("ai/agent: completion failed: %v", err))
}
e.usage.RecordCall(AgentFunctionID, e.provider.Name(), e.model, "success")
e.usage.RecordUsage(AgentFunctionID, e.provider.Name(), e.model, resp.Usage)

usage.PromptTokens += resp.Usage.PromptTokens
usage.CompletionTokens += resp.Usage.CompletionTokens
usage.TotalTokens += resp.Usage.TotalTokens
totalUsage.PromptTokens += resp.Usage.PromptTokens
totalUsage.CompletionTokens += resp.Usage.CompletionTokens
totalUsage.TotalTokens += resp.Usage.TotalTokens

messages = append(messages, resp.Message)

if len(resp.Message.ToolCalls) == 0 {
return successOutput(resp.Message.Content, usage, steps)
return successOutput(resp.Message.Content, totalUsage, steps)
}

for _, tc := range resp.Message.ToolCalls {
Expand Down
2 changes: 1 addition & 1 deletion internal/packages/functions/ai/agent_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -108,7 +108,7 @@ func runAgent(t *testing.T, providers llm.Registry, tools ToolRegistry, input ma
execInfo := workflow.NewExecutionInfo("wf-1", workflow.NewExecID(1), "", fnInput)
execInfo.Finish = func(out workflow.FunctionOutput) { done <- out }

res, err := makeAgentFunction(providers, tools)(execInfo)
res, err := makeAgentFunction(providers, tools, NopUsageRecorder{})(execInfo)
require.NoError(t, err)
if !res.Async {
return res, workflow.FunctionOutput{}
Expand Down
8 changes: 6 additions & 2 deletions internal/packages/functions/ai/chat.go
Original file line number Diff line number Diff line change
Expand Up @@ -83,8 +83,9 @@ func ChatFunctionMetadata() workflow.FunctionMetadata {
}
}

// makeChatFunction builds the ai/chat function, closing over the provider registry.
func makeChatFunction(providers llm.Registry) workflow.Function {
// makeChatFunction builds the ai/chat function, closing over the provider registry and the
// usage recorder (ADR-0029).
func makeChatFunction(providers llm.Registry, usage UsageRecorder) workflow.Function {
return func(execInfo *workflow.ExecutionInfo) (workflow.FunctionResult, error) {
input := execInfo.Input

Expand Down Expand Up @@ -124,10 +125,13 @@ func makeChatFunction(providers llm.Registry) workflow.Function {

resp, err := provider.Chat(ctx, req)
if err != nil {
usage.RecordCall(ChatFunctionID, provider.Name(), req.Model, "error")
log.Error().Err(err).Str("provider", provider.Name()).Msg("ai/chat completion failed")
execInfo.Finish(workflow.NewFunctionOutput(workflow.FunctionError, map[string]any{"error": err.Error()}))
return
}
usage.RecordCall(ChatFunctionID, provider.Name(), req.Model, "success")
usage.RecordUsage(ChatFunctionID, provider.Name(), req.Model, resp.Usage)

execInfo.Finish(workflow.NewFunctionSuccessOutput(map[string]any{
"output": resp.Message.Content,
Expand Down
4 changes: 2 additions & 2 deletions internal/packages/functions/ai/chat_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@ func runChat(t *testing.T, reg llm.Registry, input map[string]any) (workflow.Fun
execInfo := workflow.NewExecutionInfo("wf-1", "exec-1", "", fnInput)
execInfo.Finish = func(out workflow.FunctionOutput) { done <- out }

res, err := makeChatFunction(reg)(execInfo)
res, err := makeChatFunction(reg, NopUsageRecorder{})(execInfo)
require.NoError(t, err)

if !res.Async {
Expand Down Expand Up @@ -119,7 +119,7 @@ func TestChat_ResolvesProviderForExecutionEnvironment(t *testing.T) {
execInfo := workflow.NewExecutionInfo("wf-1", "exec-1", "staging", fnInput)
execInfo.Finish = func(out workflow.FunctionOutput) { done <- out }

res, err := makeChatFunction(reg)(execInfo)
res, err := makeChatFunction(reg, NopUsageRecorder{})(execInfo)
require.NoError(t, err)
require.True(t, res.Async)
select {
Expand Down
12 changes: 8 additions & 4 deletions internal/packages/functions/ai/package.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,11 +13,15 @@ const PackageID = "fuse/pkg/ai"

// New creates a new ai Package. The LLM provider registry is closed over by the
// function implementations so they can resolve providers at execution time; the
// tool registry lets the agent expose existing functions as tools and invoke them.
func New(providers llm.Registry, tools ToolRegistry) *workflow.Package {
// tool registry lets the agent expose existing functions as tools and invoke them;
// the usage recorder surfaces token usage to observability (ADR-0029).
func New(providers llm.Registry, tools ToolRegistry, usage UsageRecorder) *workflow.Package {
if usage == nil {
usage = NopUsageRecorder{}
}
return workflow.NewPackage(
PackageID,
workflow.NewFunction(ChatFunctionID, ChatFunctionMetadata(), makeChatFunction(providers)),
workflow.NewFunction(AgentFunctionID, AgentFunctionMetadata(), makeAgentFunction(providers, tools)),
workflow.NewFunction(ChatFunctionID, ChatFunctionMetadata(), makeChatFunction(providers, usage)),
workflow.NewFunction(AgentFunctionID, AgentFunctionMetadata(), makeAgentFunction(providers, tools, usage)),
)
}
22 changes: 22 additions & 0 deletions internal/packages/functions/ai/usage.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,22 @@
package ai

import "github.com/open-source-cloud/fuse/pkg/llm"

// UsageRecorder records LLM token usage and call outcomes for observability (ADR-0029). It is a
// narrow port so this package stays free of the metrics/prometheus dependency; the engine injects
// a metrics-backed implementation, tests a no-op or a fake.
type UsageRecorder interface {
// RecordUsage records the tokens consumed by one completion for a given ai function.
RecordUsage(function, provider, model string, u llm.Usage)
// RecordCall records that a completion call was made and its outcome (success|error).
RecordCall(function, provider, model, status string)
}

// NopUsageRecorder is a UsageRecorder that does nothing (default for tests / when metrics are off).
type NopUsageRecorder struct{}

// RecordUsage does nothing.
func (NopUsageRecorder) RecordUsage(string, string, string, llm.Usage) {}

// RecordCall does nothing.
func (NopUsageRecorder) RecordCall(string, string, string, string) {}
75 changes: 75 additions & 0 deletions internal/packages/functions/ai/usage_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,75 @@
package ai

import (
"sync"
"testing"
"time"

"github.com/open-source-cloud/fuse/pkg/llm"
"github.com/open-source-cloud/fuse/pkg/workflow"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)

// fakeUsageRecorder captures RecordUsage / RecordCall invocations.
type fakeUsageRecorder struct {
mu sync.Mutex
usage []llm.Usage
status []string
}

func (f *fakeUsageRecorder) RecordUsage(_, _, _ string, u llm.Usage) {
f.mu.Lock()
defer f.mu.Unlock()
f.usage = append(f.usage, u)
}

func (f *fakeUsageRecorder) RecordCall(_, _, _, status string) {
f.mu.Lock()
defer f.mu.Unlock()
f.status = append(f.status, status)
}

func (f *fakeUsageRecorder) snapshot() ([]llm.Usage, []string) {
f.mu.Lock()
defer f.mu.Unlock()
return append([]llm.Usage(nil), f.usage...), append([]string(nil), f.status...)
}

func TestChat_RecordsUsageMetrics(t *testing.T) {
prov := &stubProvider{
name: "stub",
resp: llm.ChatResponse{
Message: llm.Message{Role: llm.RoleAssistant, Content: "hi"},
Usage: llm.Usage{PromptTokens: 3, CompletionTokens: 4, TotalTokens: 7},
},
}
reg := llm.NewStaticRegistry(map[string]llm.Provider{"stub": prov}, "stub")
rec := &fakeUsageRecorder{}

fnInput, err := workflow.NewFunctionInputWith(map[string]any{"input": "hello"})
require.NoError(t, err)
done := make(chan workflow.FunctionOutput, 1)
execInfo := workflow.NewExecutionInfo("wf-1", "exec-1", "", fnInput)
execInfo.Finish = func(out workflow.FunctionOutput) { done <- out }

res, err := makeChatFunction(reg, rec)(execInfo)
require.NoError(t, err)
require.True(t, res.Async)
select {
case <-done:
case <-time.After(2 * time.Second):
t.Fatal("timed out")
}

usage, status := rec.snapshot()
require.Len(t, usage, 1)
assert.Equal(t, 3, usage[0].PromptTokens)
assert.Equal(t, 4, usage[0].CompletionTokens)
assert.Equal(t, []string{"success"}, status)

// The no-op recorder must satisfy the interface and not panic.
var nop UsageRecorder = NopUsageRecorder{}
nop.RecordUsage(ChatFunctionID, "stub", "m", llm.Usage{PromptTokens: 1})
nop.RecordCall(ChatFunctionID, "stub", "m", "success")
}
12 changes: 8 additions & 4 deletions internal/packages/internal_packages.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package packages

import (
"github.com/open-source-cloud/fuse/internal/metrics"
"github.com/open-source-cloud/fuse/internal/packages/functions/ai"
"github.com/open-source-cloud/fuse/internal/packages/functions/debug"
"github.com/open-source-cloud/fuse/internal/packages/functions/http"
Expand All @@ -18,19 +19,22 @@ type (
)

// NewInternal creates new InternalPackages service. The LLM provider registry is
// injected so the ai package can expose chat/agent functions, and the package
// registry backs the agent's tool catalog (synchronous functions become tools).
func NewInternal(providers llm.Registry, registry Registry) InternalPackages {
// injected so the ai package can expose chat/agent functions, the package registry
// backs the agent's tool catalog (synchronous functions become tools), and the metrics
// recorder surfaces LLM token usage to observability (ADR-0029).
func NewInternal(providers llm.Registry, registry Registry, fuseMetrics *metrics.FuseMetrics) InternalPackages {
return &DefaultInternalPackages{
providers: providers,
tools: NewAgentToolRegistry(registry),
usage: newUsageRecorder(fuseMetrics),
}
}

// DefaultInternalPackages service for registering internal packages
type DefaultInternalPackages struct {
providers llm.Registry
tools ai.ToolRegistry
usage ai.UsageRecorder
}

// List returns the list of internal packages
Expand All @@ -40,6 +44,6 @@ func (p *DefaultInternalPackages) List() []*workflow.Package {
logic.New(),
http.New(),
system.New(),
ai.New(p.providers, p.tools),
ai.New(p.providers, p.tools, p.usage),
}
}
Loading
Loading