Compare commits
14 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 0209532a7a | |||
| e8977cf335 | |||
| c6ab16f3bf | |||
| f93f3c6045 | |||
| 9b4b3ce827 | |||
| 4d6ebe1fd3 | |||
| 26ab5a3bf0 | |||
| d584d372cd | |||
| 521aff145f | |||
| eafa186894 | |||
| 91c57663cc | |||
| cb49c1c2d8 | |||
| db31341b30 | |||
| a583d5741d |
+33
-19
@@ -1,27 +1,41 @@
|
||||
# Priorities
|
||||
|
||||
The ranked work queue for the autonomous improvement loop. The
|
||||
**architecture-review** pass (the *architect*) owns this file: each run it turns
|
||||
the [roadmap](../../ROADMAP.md) plus an internal scan (gaps in the
|
||||
services → agents → workflows lifecycle, API coherence, drift, tech debt, test and
|
||||
DX friction) into a single ordered list — highest-value first — and links each
|
||||
item to a tracking issue. The hourly **continuous-improvement** pass works the
|
||||
**top item whose issue is still open**. So the architect decides *what*, and the
|
||||
increment loop *builds* it.
|
||||
The ranked work queue for the autonomous improvement loop. The **planner** owns
|
||||
this file: each run it turns the [roadmap](../../ROADMAP.md) plus an internal scan
|
||||
into a single ordered list — highest-value first — each item linked to a tracking
|
||||
issue. The **builder** works the top item whose issue is still open. So the
|
||||
planner decides *what*, the builder *builds* it.
|
||||
|
||||
**Reading / editing.** An item is done when its linked issue closes (the increment
|
||||
that builds it adds `Closes #<issue>`). Roadmap phase (Now → Next → Later) is the
|
||||
primary ordering; internal findings are interleaved by value, not kept in a
|
||||
separate list. The human can reorder this list — or the issues — at any time to
|
||||
redirect the loop; direction always wins.
|
||||
**Bias to capability, not busy-work.** The top of this queue is net-new capability
|
||||
from the roadmap's *Now/Next* items. Hardening/conformance/DX polish is background
|
||||
work (roadmap *Ongoing*) — kept low here and capped, never allowed to crowd out
|
||||
capability. If an area has had several increments with no user-visible gain, it is done
|
||||
for now; rank real-headroom capability instead.
|
||||
|
||||
**Off-limits to the loop** (the architect proposes these as notes, never as queue
|
||||
items the loop can auto-merge): brand/positioning copy, breaking public-API
|
||||
changes, architectural rewrites. Those go to the human.
|
||||
**Reading / editing.** An item is done when its linked issue closes (the PR that
|
||||
builds it adds `Closes #<issue>`). The human can reorder this list or the issues at
|
||||
any time — direction always wins.
|
||||
|
||||
**Off-limits to the loop** (planner proposes as notes, never auto-merged queue
|
||||
items): brand/positioning copy, breaking public-API changes, architectural
|
||||
rewrites.
|
||||
|
||||
## Work queue (ranked)
|
||||
|
||||
1. **Add Gemini provider streaming support** ([#4784](https://github.com/micro/go-micro/issues/4784)) — #4782 closed the provider failure inspection gap, so the next highest-value user-facing gap is the last plain chat streaming hole in provider coverage: implement usable Gemini `ai.Stream` support so `micro chat`, agent streaming, and A2A streaming behave consistently across supported providers, with focused parser/error coverage and no broad public API changes.
|
||||
### Capability — the headline (roadmap: Now / Next)
|
||||
|
||||
_Seeded by Claude Code from the roadmap + open issues; thereafter maintained by the
|
||||
architecture-review pass._
|
||||
1. **x402 buyer safety hard stop** ([#4814](https://github.com/micro/go-micro/issues/4814)) — close the budget-cap bypass and require a real settler for paid configs before autonomous paid-tool use becomes the next thing users copy into real agents.
|
||||
2. **Kubernetes operator + CRDs foundation** ([#4797](https://github.com/micro/go-micro/issues/4797)) — add the first opt-in `Agent`, `Service`, and `Flow` resource foundation so the services → agents → workflows lifecycle has a native deployment path for Kubernetes users.
|
||||
3. **A2A external-client conformance** ([#4815](https://github.com/micro/go-micro/issues/4815)) — make the gateway easier for non-go-micro agents to discover and stream from by serving the well-known agent card path and spec SSE events.
|
||||
4. **MCP stdio/ws result conformance** ([#4813](https://github.com/micro/go-micro/issues/4813)) — return JSON tool results with explicit `isError` semantics across transports, backed by a stdio round-trip test.
|
||||
|
||||
### In flight — do not re-queue
|
||||
|
||||
- **gRPC-reflection MCP lint follow-up** ([#4824](https://github.com/micro/go-micro/issues/4824), [PR #4826](https://github.com/micro/go-micro/pull/4826)) — fixes the lint fallout from the merged reflected-gRPC MCP work.
|
||||
|
||||
### Background — hardening & DX (roadmap: Ongoing; capped)
|
||||
|
||||
_Background hardening is intentionally empty right now. Recent work covered first-agent
|
||||
wayfinding, plan/delegate recovery, provider fallback repair, streaming, memory
|
||||
compaction, retry controls, and provider-failure inspection. Further churn in those
|
||||
areas should be marked `needs-human` unless it unlocks a clear user-visible capability._
|
||||
|
||||
@@ -18,15 +18,46 @@ below is kept current between tags and rolled into the next version when it ship
|
||||
## [Unreleased]
|
||||
|
||||
### Added
|
||||
- **Gemini streaming support** — the Gemini provider now supports streaming model responses. (`ai/gemini/`)
|
||||
- **Model retry jitter controls** — model retry behavior can now use jitter controls to reduce synchronized retry bursts. (`ai/`, `agent/`)
|
||||
- **Compacted memory summaries** — agent memory now exposes compacted run summaries for easier inspection and recovery. (`agent/`)
|
||||
- **CLI input resume for agent runs** — the CLI can resume agent runs that require additional user input. (`cmd/micro/`, `agent/`)
|
||||
|
||||
### Changed
|
||||
- **Remote agent chat streaming** — `micro chat` now streams replies from remote agents instead of waiting for the full response. (`cmd/micro/`, `agent/`)
|
||||
|
||||
### Fixed
|
||||
- **Provider failure inspection metadata** — provider failures recorded during agent runs now retain classification metadata for inspection. (`agent/`, `ai/`)
|
||||
|
||||
---
|
||||
|
||||
## [6.7.0] - July 2026
|
||||
|
||||
### Added
|
||||
- **A2A streaming conformance harness** — A2A streaming behavior is now covered by focused conformance checks. (`gateway/a2a/`, `internal/harness/`)
|
||||
- **Agent x402 spend budget guardrail** — agents now have spend budget guardrails for x402-paid tool calls. (`agent/`, `gateway/`)
|
||||
- **First-agent chat/inspect fixture** — the maintained first-agent CLI fixture now covers chat and inspect boundaries together. (`internal/harness/`, `cmd/micro/`)
|
||||
- **Zero-to-hero inspect transcript check** — the 0→hero harness now verifies the inspect transcript path stays visible in the lifecycle walkthrough. (`internal/harness/zero-to-hero-ci/`, `internal/website/docs/`)
|
||||
|
||||
### Changed
|
||||
- **Agent stream run context propagation** — agent streams now preserve run context through streaming paths for more complete tracing and inspection. (`agent/`)
|
||||
- **Postgres store pgx v5 migration** — the Postgres store now uses pgx v5. (`store/postgres/`, `go.mod`)
|
||||
- **Plan-delegate plan persistence** — plan/delegate runs now persist plan state more defensively across harness scenarios. (`agent/`, `internal/harness/`)
|
||||
|
||||
### Fixed
|
||||
- **Nested tool-call markup rejection** — agent argument parsing now rejects nested tool-call markup instead of accepting ambiguous tool input. (`agent/`)
|
||||
- **Retry cancellation during backoff** — retry backoff now respects cancellation more reliably. (`agent/`, `ai/`)
|
||||
- **Plan-delegate mock recovery regression gate** — the harness now catches plan/delegate mock recovery regressions before they ship. (`internal/harness/`, `agent/`)
|
||||
- **First-agent fixture registration wait** — first-agent fixture registration is less race-prone during harness runs. (`internal/harness/`)
|
||||
- **Memory stream Nack ordering** — memory stream Nack handling now preserves ordering more reliably. (`broker/memory/`)
|
||||
- **Zero-to-hero fixture output race** — 0→hero fixture output is less race-prone during harness runs. (`internal/harness/zero-to-hero-ci/`)
|
||||
|
||||
### Documentation
|
||||
- **First-agent quickcheck wayfinding** — public docs now keep the quickcheck path discoverable from the first-agent route. (`README.md`, `internal/website/docs/`)
|
||||
- **Ordered 0→hero transcript** — docs and harness checks now keep the 0→hero transcript order explicit. (`internal/website/docs/`, `internal/harness/`)
|
||||
- **First-agent debug breadcrumbs** — docs now surface the first-agent debug smoke path more clearly. (`internal/website/docs/`)
|
||||
- **README badge cleanup** — the README no longer shows the Go Report Card badge. (`README.md`)
|
||||
|
||||
---
|
||||
|
||||
## [6.6.0] - July 2026
|
||||
|
||||
+41
-30
@@ -32,44 +32,55 @@ default.
|
||||
and history, end to end.
|
||||
5. Battle-tested: works across every provider, fails safely, observable.
|
||||
|
||||
## Now — hardening
|
||||
The forward work is **net-new capability**, not more hardening. Maintenance
|
||||
(conformance, resilience, DX polish) continues in the background (see *Ongoing*
|
||||
below) — but it is not the roadmap. This capability work is.
|
||||
|
||||
- **Cross-provider conformance** — the same agent scenario across all seven
|
||||
providers, gated on keys, on a schedule.
|
||||
- **Failure & resilience** — timeouts, rate limits, cancellation, deadline/context
|
||||
propagation, retry/backoff.
|
||||
- **Getting-started contract** — define and CI-verify the 0→1 and 0→hero flows.
|
||||
## Now — capability
|
||||
|
||||
## Shipped agent depth
|
||||
- **Agents that pay (x402 buyer in the runtime).** The seller side ships (paid
|
||||
tools via the `wrapper/x402` middleware) and the buyer `x402.Client` (a
|
||||
budget-capped `Payer` that turns a `402` into pay-and-retry) exists — but an
|
||||
agent can't yet *autonomously* pay for a paid tool. Wire the buyer into the
|
||||
agent tool loop: a budget-capped `AgentPayer` so an agent that hits a
|
||||
payment-required tool settles it within budget and retries, with the spend
|
||||
gated (like `ApproveTool`) and observable in `RunInfo`/traces. This makes
|
||||
go-micro a runtime for **autonomous agent commerce**. *(flagship — decomposed
|
||||
into issues in the loop queue)*
|
||||
- **AP2 mandate foundation** ([#3552](https://github.com/micro/go-micro/issues/3552))
|
||||
— verifiable payment **mandates** (a Checkout Mandate and a Payment Mandate),
|
||||
signed and attached over A2A, with the Payment Mandate naming an x402 rail. The
|
||||
authorization/audit layer above A2A + x402 that positions go-micro early in the
|
||||
emerging agent-payments standard (Google's AP2, standardized via FIDO).
|
||||
Additive and opt-in.
|
||||
|
||||
- **Durable agent loop** — opt-in `Checkpoint` support lets agent `Ask` and
|
||||
streaming runs persist, list pending work, and resume without replaying completed
|
||||
tool calls. Human-input pauses resume through explicit input helpers.
|
||||
- **Agent observability** — agent `RunInfo` now feeds OpenTelemetry spans/events
|
||||
across runs, model turns, tool calls, retries, delegation lineage, and resume
|
||||
checkpoints.
|
||||
## Next — reach & deployment
|
||||
|
||||
## Next — agentic depth
|
||||
- **gRPC-reflection MCP** — derive MCP tools from *any* gRPC service via server
|
||||
reflection, not just go-micro-native handlers. Point the gateway at an external
|
||||
gRPC service and its methods become agent tools — a large jump in what an agent
|
||||
can operate.
|
||||
- **Kubernetes operator + CRDs** — `Agent`, `Service`, and `Flow` as first-class
|
||||
Kubernetes resources; an operator reconciles them into Deployments wired to the
|
||||
registry. The production deployment story for teams already on K8s.
|
||||
|
||||
- **Streaming** — broaden provider-backed `ai.Stream` coverage and keep chat/A2A streaming end to end.
|
||||
- **Resume operations polish** — keep improving CLI/docs breadcrumbs for finding
|
||||
pending agent runs and deciding whether to call resume, resume-input, or stream
|
||||
resume in production.
|
||||
- **Observability hardening** — keep span attributes and run inspection coherent
|
||||
across agents, flows, and gateways as more providers and workflow paths are
|
||||
exercised.
|
||||
## Later — exploratory
|
||||
|
||||
## Later
|
||||
- **Runtime-fitness loop** — a persistently-running dogfood app (Mu) plus an
|
||||
operator/canary loop role, so the autonomous loop evolves go-micro against
|
||||
**real runtime signal** (latency, errors, cost) with canary + rollback — not
|
||||
just green CI. The demand signal the loop is missing today.
|
||||
- **HTTP/3 transport**; richer A2A live-stream reconnection (`tasks/resubscribe`,
|
||||
`input-required` handoffs); memory management (summarization, retrieval/RAG).
|
||||
|
||||
- Memory management (summarization, retrieval/RAG); human-in-the-loop pause/resume;
|
||||
richer A2A live-stream reconnection (`tasks/resubscribe`) and `input-required`
|
||||
handoffs.
|
||||
## Ongoing — hardening & DX (background, not the headline)
|
||||
|
||||
## Developer experience (ongoing)
|
||||
|
||||
- A seamless CLI inner loop (scaffold → run → chat → inspect → deploy); UI
|
||||
discipline (trim what isn't great); a maintained real-world example that doubles
|
||||
as the 0→hero reference; docs kept in lockstep with the code.
|
||||
Continuous but **capped** so it never crowds out capability: cross-provider
|
||||
conformance, failure/resilience (timeouts, cancellation, retry/backoff), the
|
||||
0→1 and 0→hero getting-started contract, streaming/observability coherence, and a
|
||||
seamless CLI inner loop (scaffold → run → chat → inspect → deploy). Real, but
|
||||
maintenance — the loop should spend the majority of its cycles on the capability above,
|
||||
not here.
|
||||
|
||||
## How it's sustained
|
||||
|
||||
|
||||
@@ -5,6 +5,8 @@ import (
|
||||
"crypto/sha256"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
@@ -12,6 +14,7 @@ import (
|
||||
codecBytes "go-micro.dev/v6/codec/bytes"
|
||||
"go-micro.dev/v6/gateway/a2a"
|
||||
"go-micro.dev/v6/store"
|
||||
"go-micro.dev/v6/wrapper/x402"
|
||||
)
|
||||
|
||||
// Built-in agent tools. These are not service endpoints — they are
|
||||
@@ -128,6 +131,7 @@ func (a *agentImpl) toolHandler() ai.ToolHandler {
|
||||
// so the result runs plan → step → loop → approve → checkpoint → base.
|
||||
h := a.baseHandler()
|
||||
h = a.toolTimeoutWrap(h)
|
||||
h = a.x402PayWrap(h)
|
||||
h = a.toolRetryWrap(h)
|
||||
h = a.checkpointToolWrap(h)
|
||||
h = a.approveWrap(h)
|
||||
@@ -173,6 +177,64 @@ func (a *agentImpl) toolTimeoutWrap(next ai.ToolHandler) ai.ToolHandler {
|
||||
}
|
||||
}
|
||||
|
||||
// x402PayWrap pays an x402 Payment Required tool result and retries the
|
||||
// underlying HTTP tool once. Tools that proxy HTTP paid resources can return the
|
||||
// raw x402 402 challenge body and include a "url" input; the agent then uses
|
||||
// wrapper/x402.Client so payer and budget semantics stay in one place.
|
||||
func (a *agentImpl) x402PayWrap(next ai.ToolHandler) ai.ToolHandler {
|
||||
return func(ctx context.Context, call ai.ToolCall) ai.ToolResult {
|
||||
res := next(ctx, call)
|
||||
if res.Refused != "" || !isX402Challenge(res.Content) {
|
||||
return res
|
||||
}
|
||||
url, _ := call.Input["url"].(string)
|
||||
if url == "" {
|
||||
return errResult(call.ID, "x402: payment required but tool result did not include a retryable url input")
|
||||
}
|
||||
budget := a.opts.Budget
|
||||
if budget > 0 {
|
||||
remaining := budget - a.spend
|
||||
if remaining <= 0 {
|
||||
return refused(call.ID, ai.RefusedSpendBudget, fmt.Sprintf(
|
||||
"x402 spend budget exceeded: no budget remaining for %s (spent %d of %d)",
|
||||
call.Name, a.spend, budget))
|
||||
}
|
||||
budget = remaining
|
||||
}
|
||||
client := &x402.Client{Payer: a.opts.Payer, Budget: budget}
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil)
|
||||
if err != nil {
|
||||
return errResult(call.ID, err.Error())
|
||||
}
|
||||
resp, err := client.Do(req)
|
||||
if err != nil {
|
||||
if strings.Contains(err.Error(), "would exceed budget") {
|
||||
return refused(call.ID, ai.RefusedSpendBudget, err.Error())
|
||||
}
|
||||
return errResult(call.ID, err.Error())
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
body, err := io.ReadAll(resp.Body)
|
||||
if err != nil {
|
||||
return errResult(call.ID, err.Error())
|
||||
}
|
||||
a.spend += client.Spent()
|
||||
var value any
|
||||
if err := json.Unmarshal(body, &value); err != nil {
|
||||
value = string(body)
|
||||
}
|
||||
return ai.ToolResult{ID: call.ID, Value: value, Content: string(body), Attempts: 2}
|
||||
}
|
||||
}
|
||||
|
||||
func isX402Challenge(content string) bool {
|
||||
var ch struct {
|
||||
X402Version int `json:"x402Version"`
|
||||
Accepts []x402.Requirements `json:"accepts"`
|
||||
}
|
||||
return json.Unmarshal([]byte(content), &ch) == nil && ch.X402Version > 0 && len(ch.Accepts) > 0
|
||||
}
|
||||
|
||||
// toolRetryWrap retries transient tool failures with bounded backoff. It is
|
||||
// opt-in because tools can have side effects; guardrail refusals and caller
|
||||
// cancellation are never retried.
|
||||
@@ -389,6 +451,11 @@ func (a *agentImpl) spendWrap(next ai.ToolHandler) ai.ToolHandler {
|
||||
amount, call.Name, a.spend, a.opts.MaxSpend))
|
||||
}
|
||||
a.spend += amount
|
||||
if info, ok := ai.RunInfoFrom(ctx); ok {
|
||||
info.Spent = a.spend
|
||||
info.ToolSpend = amount
|
||||
ctx = ai.WithRunInfo(ctx, info)
|
||||
}
|
||||
res := next(ctx, call)
|
||||
if res.Refused != "" || toolErrorMessage(res) != "" {
|
||||
a.spend -= amount
|
||||
|
||||
@@ -2,12 +2,17 @@ package agent
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"io"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"go-micro.dev/v6/ai"
|
||||
"go-micro.dev/v6/registry"
|
||||
"go-micro.dev/v6/store"
|
||||
"go-micro.dev/v6/wrapper/x402"
|
||||
)
|
||||
|
||||
// toolContent runs a tool call through a handler and returns the content
|
||||
@@ -181,3 +186,132 @@ func TestNestedTextToolCallArgumentsAreRefused(t *testing.T) {
|
||||
t.Fatalf("content = %q, want nested tool-call refusal", content)
|
||||
}
|
||||
}
|
||||
|
||||
type agentMockPayer struct{ calls int }
|
||||
|
||||
func (p *agentMockPayer) Pay(ctx context.Context, req x402.Requirements) (string, error) {
|
||||
p.calls++
|
||||
return "paid", nil
|
||||
}
|
||||
|
||||
func TestAgentPayerPaysX402ToolResultAndRetries(t *testing.T) {
|
||||
paid := false
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
if r.Header.Get(x402.PaymentHeader) == "paid" {
|
||||
paid = true
|
||||
_, _ = w.Write([]byte(`{"ok":true}`))
|
||||
return
|
||||
}
|
||||
w.WriteHeader(http.StatusPaymentRequired)
|
||||
json.NewEncoder(w).Encode(map[string]any{
|
||||
"x402Version": x402.Version,
|
||||
"accepts": []x402.Requirements{{Scheme: "exact", Network: "base", MaxAmountRequired: "7", Resource: r.URL.String(), PayTo: "0xmerchant"}},
|
||||
})
|
||||
}))
|
||||
defer srv.Close()
|
||||
|
||||
payer := &agentMockPayer{}
|
||||
st := store.NewMemoryStore()
|
||||
a := newTestAgent(Name("x402-payer"), WithStore(st), Payer(payer), Budget(10), WithTool("paid.http", "paid http", nil, func(ctx context.Context, input map[string]any) (string, error) {
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodGet, srv.URL, nil)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
resp, err := http.DefaultClient.Do(req)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
body, err := io.ReadAll(resp.Body)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
return string(body), nil
|
||||
}))
|
||||
|
||||
ctx := ai.WithRunInfo(context.Background(), ai.RunInfo{RunID: "run-paid", Agent: "x402-payer"})
|
||||
res := a.toolHandler()(ctx, ai.ToolCall{ID: "pay-1", Name: "paid.http", Input: map[string]any{"url": srv.URL}})
|
||||
if !paid || payer.calls != 1 {
|
||||
t.Fatalf("payment not made: paid=%v payer.calls=%d", paid, payer.calls)
|
||||
}
|
||||
if res.Content != `{"ok":true}` || res.Attempts != 2 {
|
||||
t.Fatalf("result = %+v, want paid response with retry attempt", res)
|
||||
}
|
||||
events, err := LoadRunEvents(st, "x402-payer", "run-paid")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(events) != 1 || events[0].Spent != 7 || events[0].ToolSpend != 7 {
|
||||
t.Fatalf("spend events = %#v, want one tool event with spent/tool_spend 7", events)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAgentPayerRefusesX402OverBudget(t *testing.T) {
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
w.WriteHeader(http.StatusPaymentRequired)
|
||||
json.NewEncoder(w).Encode(map[string]any{
|
||||
"x402Version": x402.Version,
|
||||
"accepts": []x402.Requirements{{Scheme: "exact", Network: "base", MaxAmountRequired: "70", Resource: r.URL.String(), PayTo: "0xmerchant"}},
|
||||
})
|
||||
}))
|
||||
defer srv.Close()
|
||||
|
||||
payer := &agentMockPayer{}
|
||||
a := newTestAgent(Name("x402-over-budget"), Payer(payer), Budget(10), WithTool("paid.http", "paid http", nil, func(ctx context.Context, input map[string]any) (string, error) {
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodGet, srv.URL, nil)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
resp, err := http.DefaultClient.Do(req)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
body, err := io.ReadAll(resp.Body)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
return string(body), nil
|
||||
}))
|
||||
|
||||
res := a.toolHandler()(context.Background(), ai.ToolCall{ID: "pay-1", Name: "paid.http", Input: map[string]any{"url": srv.URL}})
|
||||
if payer.calls != 0 {
|
||||
t.Fatalf("payer called despite over-budget refusal")
|
||||
}
|
||||
if res.Refused != ai.RefusedSpendBudget || !strings.Contains(res.Content, "would exceed budget") {
|
||||
t.Fatalf("result = %+v, want budget refusal", res)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAgentPayerRequiredWithoutPayerReturnsClearError(t *testing.T) {
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
w.WriteHeader(http.StatusPaymentRequired)
|
||||
json.NewEncoder(w).Encode(map[string]any{
|
||||
"x402Version": x402.Version,
|
||||
"accepts": []x402.Requirements{{Scheme: "exact", Network: "base", MaxAmountRequired: "7", Resource: r.URL.String(), PayTo: "0xmerchant"}},
|
||||
})
|
||||
}))
|
||||
defer srv.Close()
|
||||
|
||||
a := newTestAgent(Name("x402-no-payer"), Budget(10), WithTool("paid.http", "paid http", nil, func(ctx context.Context, input map[string]any) (string, error) {
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodGet, srv.URL, nil)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
resp, err := http.DefaultClient.Do(req)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
body, err := io.ReadAll(resp.Body)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
return string(body), nil
|
||||
}))
|
||||
|
||||
res := a.toolHandler()(context.Background(), ai.ToolCall{ID: "pay-1", Name: "paid.http", Input: map[string]any{"url": srv.URL}})
|
||||
if !strings.Contains(res.Content, "no Payer configured") {
|
||||
t.Fatalf("content = %q, want no payer error", res.Content)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -10,6 +10,7 @@ import (
|
||||
"go-micro.dev/v6/flow"
|
||||
"go-micro.dev/v6/registry"
|
||||
"go-micro.dev/v6/store"
|
||||
"go-micro.dev/v6/wrapper/x402"
|
||||
"go.opentelemetry.io/otel/trace"
|
||||
)
|
||||
|
||||
@@ -105,6 +106,10 @@ type Options struct {
|
||||
// unit (0 = disabled). ToolSpend lists known paid tools and their prices.
|
||||
MaxSpend int64
|
||||
ToolSpend map[string]int64
|
||||
// Payer lets the agent settle x402 Payment Required challenges from tools.
|
||||
// Budget bounds autonomous x402 payments per Ask (0 = unlimited).
|
||||
Payer x402.Payer
|
||||
Budget int64
|
||||
|
||||
// A2AAddress, if set, makes Run serve this agent over the A2A protocol
|
||||
// on that address directly (no separate gateway), e.g. ":4000".
|
||||
@@ -250,6 +255,18 @@ func ToolSpend(tool string, amount int64) Option {
|
||||
}
|
||||
}
|
||||
|
||||
// Payer configures the wallet/signing hook used to settle x402-paid tools.
|
||||
// Without a payer, payment-required tool results are returned as clear errors.
|
||||
func Payer(p x402.Payer) Option {
|
||||
return func(o *Options) { o.Payer = p }
|
||||
}
|
||||
|
||||
// Budget bounds autonomous x402 payments per Ask, in the asset's smallest
|
||||
// unit (0 = unlimited). The budget is enforced by wrapper/x402.Client.
|
||||
func Budget(amount int64) Option {
|
||||
return func(o *Options) { o.Budget = amount }
|
||||
}
|
||||
|
||||
// LoopLimit sets how many times the agent may repeat the same tool call
|
||||
// (same name and arguments) in one Ask before it is refused as a
|
||||
// no-progress loop. 0 disables loop detection.
|
||||
|
||||
+27
-2
@@ -51,6 +51,8 @@ const (
|
||||
AttrDispatch = "agent.dispatch"
|
||||
AttrTrigger = "agent.trigger"
|
||||
AttrRunEventKind = "agent.event.kind"
|
||||
AttrSpend = "agent.spend"
|
||||
AttrToolSpend = "agent.tool.spend"
|
||||
)
|
||||
|
||||
type RunEvent struct {
|
||||
@@ -73,6 +75,8 @@ type RunEvent struct {
|
||||
Error string `json:"error,omitempty"`
|
||||
ErrorKind string `json:"error_kind,omitempty"`
|
||||
InputChars int `json:"input_chars,omitempty"`
|
||||
Spent int64 `json:"spent,omitempty"`
|
||||
ToolSpend int64 `json:"tool_spend,omitempty"`
|
||||
}
|
||||
|
||||
type Usage = ai.Usage
|
||||
@@ -111,6 +115,7 @@ type RunSummary struct {
|
||||
LastKind string `json:"last_kind,omitempty"`
|
||||
LastError string `json:"last_error,omitempty"`
|
||||
LastErrorKind string `json:"last_error_kind,omitempty"`
|
||||
Spent int64 `json:"spent,omitempty"`
|
||||
}
|
||||
|
||||
func (a *agentImpl) tracer() trace.Tracer {
|
||||
@@ -362,6 +367,7 @@ func (a *agentImpl) traceTool(next ai.ToolHandler) ai.ToolHandler {
|
||||
return func(ctx context.Context, call ai.ToolCall) ai.ToolResult {
|
||||
info, _ := ai.RunInfoFrom(ctx)
|
||||
start := time.Now()
|
||||
spentBefore := a.spend
|
||||
|
||||
if a.opts.TraceProvider == nil {
|
||||
res := next(ctx, call)
|
||||
@@ -371,7 +377,7 @@ func (a *agentImpl) traceTool(next ai.ToolHandler) ai.ToolHandler {
|
||||
if toolAttempts <= 0 {
|
||||
toolAttempts = 1
|
||||
}
|
||||
a.recordRunEvent(RunEvent{Time: time.Now(), RunID: info.RunID, ParentID: info.ParentID, Agent: info.Agent, Kind: "tool", Name: call.Name, Attempt: toolAttempts, MaxAttempts: a.opts.ToolMaxAttempts, LatencyMS: dur, Refused: res.Refused, Error: resErr, ErrorKind: classifyToolError(resErr)})
|
||||
a.recordRunEvent(RunEvent{Time: time.Now(), RunID: info.RunID, ParentID: info.ParentID, Agent: info.Agent, Kind: "tool", Name: call.Name, Attempt: toolAttempts, MaxAttempts: a.opts.ToolMaxAttempts, LatencyMS: dur, Refused: res.Refused, Error: resErr, ErrorKind: classifyToolError(resErr), Spent: a.spend, ToolSpend: a.spend - spentBefore})
|
||||
return res
|
||||
}
|
||||
|
||||
@@ -385,6 +391,7 @@ func (a *agentImpl) traceTool(next ai.ToolHandler) ai.ToolHandler {
|
||||
ctx, span := a.tracer().Start(ctx, spanNameToolCall, trace.WithAttributes(spanAttrs...))
|
||||
res := next(ctx, call)
|
||||
dur := time.Since(start).Milliseconds()
|
||||
toolSpend := a.spend - spentBefore
|
||||
attrs := []attribute.KeyValue{attribute.Int64(AttrLatencyMS, dur)}
|
||||
toolAttempts := res.Attempts
|
||||
if toolAttempts <= 0 {
|
||||
@@ -394,6 +401,9 @@ func (a *agentImpl) traceTool(next ai.ToolHandler) ai.ToolHandler {
|
||||
if a.opts.ToolMaxAttempts > 0 {
|
||||
attrs = append(attrs, attribute.Int(AttrToolMaxAttempts, a.opts.ToolMaxAttempts))
|
||||
}
|
||||
if toolSpend > 0 {
|
||||
attrs = append(attrs, attribute.Int64(AttrSpend, a.spend), attribute.Int64(AttrToolSpend, toolSpend))
|
||||
}
|
||||
if res.Refused != "" {
|
||||
attrs = append(attrs, attribute.Bool(AttrGuardrailBlock, true), attribute.String(AttrRefusal, res.Refused))
|
||||
}
|
||||
@@ -409,7 +419,7 @@ func (a *agentImpl) traceTool(next ai.ToolHandler) ai.ToolHandler {
|
||||
} else {
|
||||
span.SetStatus(codes.Ok, "")
|
||||
}
|
||||
a.recordSpanEvent(span, RunEvent{Time: time.Now(), RunID: info.RunID, ParentID: info.ParentID, Agent: info.Agent, Kind: "tool", Name: call.Name, Attempt: toolAttempts, MaxAttempts: a.opts.ToolMaxAttempts, LatencyMS: dur, Refused: res.Refused, Error: resErr, ErrorKind: classifyToolError(resErr)})
|
||||
a.recordSpanEvent(span, RunEvent{Time: time.Now(), RunID: info.RunID, ParentID: info.ParentID, Agent: info.Agent, Kind: "tool", Name: call.Name, Attempt: toolAttempts, MaxAttempts: a.opts.ToolMaxAttempts, LatencyMS: dur, Refused: res.Refused, Error: resErr, ErrorKind: classifyToolError(resErr), Spent: a.spend, ToolSpend: toolSpend})
|
||||
span.End()
|
||||
return res
|
||||
}
|
||||
@@ -496,6 +506,12 @@ func runEventAttributes(e RunEvent) []attribute.KeyValue {
|
||||
if e.InputChars > 0 {
|
||||
attrs = append(attrs, attribute.Int(AttrInputChars, e.InputChars))
|
||||
}
|
||||
if e.Spent > 0 {
|
||||
attrs = append(attrs, attribute.Int64(AttrSpend, e.Spent))
|
||||
}
|
||||
if e.ToolSpend > 0 {
|
||||
attrs = append(attrs, attribute.Int64(AttrToolSpend, e.ToolSpend))
|
||||
}
|
||||
attrs = appendUsage(attrs, e.Tokens)
|
||||
if e.Refused != "" {
|
||||
attrs = append(attrs, attribute.Bool(AttrGuardrailBlock, true), attribute.String(AttrRefusal, e.Refused))
|
||||
@@ -530,6 +546,12 @@ func appendRunInfoAttributes(attrs []attribute.KeyValue, info ai.RunInfo) []attr
|
||||
if info.Trigger != "" {
|
||||
attrs = append(attrs, attribute.String(AttrTrigger, info.Trigger))
|
||||
}
|
||||
if info.Spent > 0 {
|
||||
attrs = append(attrs, attribute.Int64(AttrSpend, info.Spent))
|
||||
}
|
||||
if info.ToolSpend > 0 {
|
||||
attrs = append(attrs, attribute.Int64(AttrToolSpend, info.ToolSpend))
|
||||
}
|
||||
return attrs
|
||||
}
|
||||
|
||||
@@ -616,6 +638,9 @@ func ListRunSummariesWithOptions(s store.Store, agentName string, opts RunListOp
|
||||
if e.ErrorKind != "" {
|
||||
summary.LastErrorKind = e.ErrorKind
|
||||
}
|
||||
if e.Spent > summary.Spent {
|
||||
summary.Spent = e.Spent
|
||||
}
|
||||
}
|
||||
if opts.Status != "" && summary.Status != opts.Status {
|
||||
continue
|
||||
|
||||
@@ -421,6 +421,63 @@ func spanAttributes(attrs []attribute.KeyValue) map[string]string {
|
||||
return out
|
||||
}
|
||||
|
||||
func TestAgentOpenTelemetryToolSpanIncludesSpend(t *testing.T) {
|
||||
exp := tracetest.NewInMemoryExporter()
|
||||
tp := trace.NewTracerProvider(trace.WithSyncer(exp))
|
||||
st := store.NewMemoryStore()
|
||||
a := New(Name("spender"), Provider("oteltest"), Model("unit-model"), WithStore(st), TraceProvider(tp), MaxSpend(10), ToolSpend("probe", 7), WithTool("probe", "probe", nil, func(ctx context.Context, input map[string]any) (string, error) {
|
||||
info, ok := ai.RunInfoFrom(ctx)
|
||||
if !ok {
|
||||
t.Fatal("RunInfo missing from paid tool context")
|
||||
}
|
||||
if info.Spent != 7 || info.ToolSpend != 7 {
|
||||
t.Fatalf("RunInfo spend = (%d, %d), want (7, 7)", info.Spent, info.ToolSpend)
|
||||
}
|
||||
return "ok", nil
|
||||
}))
|
||||
if _, err := a.Ask(context.Background(), "hello"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
var sawToolSpan bool
|
||||
for _, s := range exp.GetSpans().Snapshots() {
|
||||
if s.Name() != spanNameToolCall {
|
||||
continue
|
||||
}
|
||||
sawToolSpan = true
|
||||
attrs := spanAttributes(s.Attributes())
|
||||
if attrs[AttrSpend] != "7" || attrs[AttrToolSpend] != "7" {
|
||||
t.Fatalf("tool span missing spend attributes: %#v", attrs)
|
||||
}
|
||||
if !spanEventHasAttribute(s.Events(), "agent.tool", AttrToolSpend, "7") {
|
||||
t.Fatalf("tool event missing spend attribute: %#v", s.Events())
|
||||
}
|
||||
}
|
||||
if !sawToolSpan {
|
||||
t.Fatal("tool span not emitted")
|
||||
}
|
||||
summaries, err := ListRunSummaries(st, "spender")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(summaries) != 1 || summaries[0].Spent != 7 {
|
||||
t.Fatalf("summary spend = %#v, want 7", summaries)
|
||||
}
|
||||
}
|
||||
|
||||
func spanEventHasAttribute(events []trace.Event, name, key, value string) bool {
|
||||
for _, e := range events {
|
||||
if e.Name != name {
|
||||
continue
|
||||
}
|
||||
attrs := spanAttributes(e.Attributes)
|
||||
if attrs[key] == value {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
func TestAgentOpenTelemetrySpansDelegateLineage(t *testing.T) {
|
||||
exp := tracetest.NewInMemoryExporter()
|
||||
tp := trace.NewTracerProvider(trace.WithSyncer(exp))
|
||||
|
||||
@@ -35,7 +35,7 @@ func TestRegisteredProviders(t *testing.T) {
|
||||
}
|
||||
|
||||
got = ai.RegisteredProviders("stream")
|
||||
want = []string{"anthropic", "atlascloud", "groq", "minimax", "mistral", "openai", "together"}
|
||||
want = []string{"anthropic", "atlascloud", "gemini", "groq", "minimax", "mistral", "openai", "together"}
|
||||
if !reflect.DeepEqual(got, want) {
|
||||
t.Fatalf("RegisteredProviders(stream) = %#v, want %#v", got, want)
|
||||
}
|
||||
@@ -46,7 +46,7 @@ func TestCapabilityRows(t *testing.T) {
|
||||
want := []ai.CapabilityRow{
|
||||
{Provider: "anthropic", Capabilities: ai.Capabilities{Model: true, Stream: true, ToolStream: true}},
|
||||
{Provider: "atlascloud", Capabilities: ai.Capabilities{Model: true, Image: true, Video: true, Stream: true}},
|
||||
{Provider: "gemini", Capabilities: ai.Capabilities{Model: true}},
|
||||
{Provider: "gemini", Capabilities: ai.Capabilities{Model: true, Stream: true}},
|
||||
{Provider: "groq", Capabilities: ai.Capabilities{Model: true, Stream: true, ToolStream: true}},
|
||||
{Provider: "minimax", Capabilities: ai.Capabilities{Model: true, Stream: true, ToolStream: true}},
|
||||
{Provider: "mistral", Capabilities: ai.Capabilities{Model: true, Stream: true, ToolStream: true}},
|
||||
@@ -90,7 +90,7 @@ func TestRegisterStream(t *testing.T) {
|
||||
}
|
||||
|
||||
got := ai.RegisteredProviders("stream")
|
||||
want := []string{"anthropic", "atlascloud", "groq", "minimax", "mistral", "openai", "test-stream", "together"}
|
||||
want := []string{"anthropic", "atlascloud", "gemini", "groq", "minimax", "mistral", "openai", "test-stream", "together"}
|
||||
if !reflect.DeepEqual(got, want) {
|
||||
t.Fatalf("RegisteredProviders(stream) = %#v, want %#v", got, want)
|
||||
}
|
||||
|
||||
+136
-4
@@ -10,6 +10,7 @@
|
||||
package gemini
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
@@ -25,6 +26,7 @@ func init() {
|
||||
ai.Register("gemini", func(opts ...ai.Option) ai.Model {
|
||||
return NewProvider(opts...)
|
||||
})
|
||||
ai.RegisterStream("gemini")
|
||||
}
|
||||
|
||||
// Provider implements the ai.Model interface for Google Gemini.
|
||||
@@ -69,9 +71,7 @@ func (p *Provider) Generate(ctx context.Context, req *ai.Request, opts ...ai.Gen
|
||||
})
|
||||
}
|
||||
|
||||
contents := []map[string]any{
|
||||
{"role": "user", "parts": []map[string]any{{"text": req.Prompt}}},
|
||||
}
|
||||
contents := geminiContents(req)
|
||||
|
||||
apiReq := map[string]any{
|
||||
"contents": contents,
|
||||
@@ -135,7 +135,121 @@ func (p *Provider) Generate(ctx context.Context, req *ai.Request, opts ...ai.Gen
|
||||
}
|
||||
|
||||
func (p *Provider) Stream(ctx context.Context, req *ai.Request, opts ...ai.GenerateOption) (ai.Stream, error) {
|
||||
return nil, fmt.Errorf("%w: gemini provider", ai.ErrStreamingUnsupported)
|
||||
apiReq := map[string]any{
|
||||
"contents": geminiContents(req),
|
||||
}
|
||||
if req.SystemPrompt != "" {
|
||||
apiReq["system_instruction"] = map[string]any{
|
||||
"parts": []map[string]any{{"text": req.SystemPrompt}},
|
||||
}
|
||||
}
|
||||
if p.opts.MaxTokens > 0 {
|
||||
apiReq["generationConfig"] = map[string]any{"maxOutputTokens": p.opts.MaxTokens}
|
||||
}
|
||||
|
||||
reqBody, err := json.Marshal(apiReq)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to marshal stream request: %w", err)
|
||||
}
|
||||
|
||||
apiURL := strings.TrimRight(p.opts.BaseURL, "/") +
|
||||
"/v1beta/models/" + p.opts.Model + ":streamGenerateContent?alt=sse"
|
||||
httpReq, err := http.NewRequestWithContext(ctx, http.MethodPost, apiURL, bytes.NewReader(reqBody))
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to create stream request: %w", err)
|
||||
}
|
||||
httpReq.Header.Set("Content-Type", "application/json")
|
||||
httpReq.Header.Set("Accept", "text/event-stream")
|
||||
httpReq.Header.Set("x-goog-api-key", p.opts.APIKey)
|
||||
|
||||
httpResp, err := http.DefaultClient.Do(httpReq)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("stream API request failed: %w", err)
|
||||
}
|
||||
if httpResp.StatusCode != http.StatusOK {
|
||||
defer httpResp.Body.Close()
|
||||
respBody, _ := io.ReadAll(httpResp.Body)
|
||||
return nil, ai.NewHTTPError(httpResp, respBody)
|
||||
}
|
||||
return &streamReader{body: httpResp.Body, scanner: bufio.NewScanner(httpResp.Body)}, nil
|
||||
}
|
||||
|
||||
type streamReader struct {
|
||||
body io.ReadCloser
|
||||
scanner *bufio.Scanner
|
||||
closed bool
|
||||
}
|
||||
|
||||
func (s *streamReader) Recv() (*ai.Response, error) {
|
||||
for s.scanner.Scan() {
|
||||
line := strings.TrimSpace(s.scanner.Text())
|
||||
if line == "" || strings.HasPrefix(line, ":") || strings.HasPrefix(line, "event:") {
|
||||
continue
|
||||
}
|
||||
if !strings.HasPrefix(line, "data:") {
|
||||
continue
|
||||
}
|
||||
data := strings.TrimSpace(strings.TrimPrefix(line, "data:"))
|
||||
if data == "[DONE]" {
|
||||
return nil, io.EOF
|
||||
}
|
||||
|
||||
var chunk struct {
|
||||
Error *struct {
|
||||
Code int `json:"code"`
|
||||
Message string `json:"message"`
|
||||
Status string `json:"status"`
|
||||
} `json:"error"`
|
||||
Candidates []struct {
|
||||
Content struct {
|
||||
Parts []struct {
|
||||
Text string `json:"text"`
|
||||
} `json:"parts"`
|
||||
} `json:"content"`
|
||||
} `json:"candidates"`
|
||||
UsageMetadata *struct {
|
||||
PromptTokenCount int `json:"promptTokenCount"`
|
||||
CandidatesTokenCount int `json:"candidatesTokenCount"`
|
||||
TotalTokenCount int `json:"totalTokenCount"`
|
||||
} `json:"usageMetadata"`
|
||||
}
|
||||
if err := json.Unmarshal([]byte(data), &chunk); err != nil {
|
||||
return nil, fmt.Errorf("failed to parse stream chunk: %w", err)
|
||||
}
|
||||
if chunk.Error != nil {
|
||||
return nil, fmt.Errorf("gemini stream error (%s): %s", chunk.Error.Status, chunk.Error.Message)
|
||||
}
|
||||
for _, candidate := range chunk.Candidates {
|
||||
var parts []string
|
||||
for _, part := range candidate.Content.Parts {
|
||||
if part.Text != "" {
|
||||
parts = append(parts, part.Text)
|
||||
}
|
||||
}
|
||||
if len(parts) > 0 {
|
||||
return &ai.Response{Reply: strings.Join(parts, "")}, nil
|
||||
}
|
||||
}
|
||||
if chunk.UsageMetadata != nil {
|
||||
return &ai.Response{Usage: ai.Usage{
|
||||
InputTokens: chunk.UsageMetadata.PromptTokenCount,
|
||||
OutputTokens: chunk.UsageMetadata.CandidatesTokenCount,
|
||||
TotalTokens: chunk.UsageMetadata.TotalTokenCount,
|
||||
}}, nil
|
||||
}
|
||||
}
|
||||
if err := s.scanner.Err(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return nil, io.EOF
|
||||
}
|
||||
|
||||
func (s *streamReader) Close() error {
|
||||
if s.closed {
|
||||
return nil
|
||||
}
|
||||
s.closed = true
|
||||
return s.body.Close()
|
||||
}
|
||||
|
||||
func (p *Provider) callAPI(ctx context.Context, req map[string]any) (*ai.Response, []map[string]any, error) {
|
||||
@@ -224,3 +338,21 @@ type functionCallPB struct {
|
||||
Name string `json:"name"`
|
||||
Args map[string]any `json:"args"`
|
||||
}
|
||||
|
||||
func geminiContents(req *ai.Request) []map[string]any {
|
||||
contents := make([]map[string]any, 0, len(req.Messages)+1)
|
||||
for _, m := range req.Messages {
|
||||
role := m.Role
|
||||
if role == "assistant" {
|
||||
role = "model"
|
||||
}
|
||||
if role == "system" || role == "" {
|
||||
continue
|
||||
}
|
||||
contents = append(contents, map[string]any{"role": role, "parts": []map[string]any{{"text": fmt.Sprint(m.Content)}}})
|
||||
}
|
||||
if req.Prompt != "" {
|
||||
contents = append(contents, map[string]any{"role": "user", "parts": []map[string]any{{"text": req.Prompt}}})
|
||||
}
|
||||
return contents
|
||||
}
|
||||
|
||||
+103
-6
@@ -2,7 +2,12 @@ package gemini
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"io"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"go-micro.dev/v6/ai"
|
||||
@@ -81,16 +86,108 @@ func TestProvider_Generate_NoAPIKey(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestProvider_Stream_NotImplemented(t *testing.T) {
|
||||
p := NewProvider()
|
||||
func TestProvider_Stream(t *testing.T) {
|
||||
var sawRequest bool
|
||||
ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
sawRequest = true
|
||||
if r.URL.Path != "/v1beta/models/gemini-2.5-flash:streamGenerateContent" {
|
||||
t.Fatalf("path = %s, want streamGenerateContent", r.URL.Path)
|
||||
}
|
||||
if r.URL.Query().Get("alt") != "sse" {
|
||||
t.Fatalf("alt = %q, want sse", r.URL.Query().Get("alt"))
|
||||
}
|
||||
if got := r.Header.Get("Accept"); got != "text/event-stream" {
|
||||
t.Fatalf("Accept = %q, want text/event-stream", got)
|
||||
}
|
||||
if got := r.Header.Get("x-goog-api-key"); got != "test-key" {
|
||||
t.Fatalf("x-goog-api-key = %q, want test-key", got)
|
||||
}
|
||||
|
||||
req := &ai.Request{
|
||||
var body map[string]any
|
||||
if err := json.NewDecoder(r.Body).Decode(&body); err != nil {
|
||||
t.Fatalf("decode request: %v", err)
|
||||
}
|
||||
contents, ok := body["contents"].([]any)
|
||||
if !ok || len(contents) != 3 {
|
||||
t.Fatalf("contents = %#v, want history + prompt", body["contents"])
|
||||
}
|
||||
second := contents[1].(map[string]any)
|
||||
if second["role"] != "model" {
|
||||
t.Fatalf("assistant history role = %#v, want model", second["role"])
|
||||
}
|
||||
|
||||
w.Header().Set("Content-Type", "text/event-stream")
|
||||
_, _ = w.Write([]byte("data: {\"candidates\":[{\"content\":{\"parts\":[{\"text\":\"hel\"}]}}]}\n\n"))
|
||||
_, _ = w.Write([]byte("data: {\"candidates\":[{\"content\":{\"parts\":[{\"text\":\"lo\"}]}}],\"usageMetadata\":{\"promptTokenCount\":3,\"candidatesTokenCount\":2,\"totalTokenCount\":5}}\n\n"))
|
||||
_, _ = w.Write([]byte("data: [DONE]\n\n"))
|
||||
}))
|
||||
defer ts.Close()
|
||||
|
||||
p := NewProvider(ai.WithAPIKey("test-key"), ai.WithBaseURL(ts.URL))
|
||||
stream, err := p.Stream(context.Background(), &ai.Request{
|
||||
Messages: []ai.Message{
|
||||
{Role: "user", Content: "previous question"},
|
||||
{Role: "assistant", Content: "previous answer"},
|
||||
},
|
||||
Prompt: "Hello",
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("Stream returned error: %v", err)
|
||||
}
|
||||
defer stream.Close()
|
||||
if !sawRequest {
|
||||
t.Fatal("server did not receive stream request")
|
||||
}
|
||||
|
||||
_, err := p.Stream(context.Background(), req)
|
||||
if !errors.Is(err, ai.ErrStreamingUnsupported) {
|
||||
t.Fatalf("Stream error = %v, want ErrStreamingUnsupported", err)
|
||||
first, err := stream.Recv()
|
||||
if err != nil || first.Reply != "hel" {
|
||||
t.Fatalf("first chunk = %#v, %v; want hel", first, err)
|
||||
}
|
||||
second, err := stream.Recv()
|
||||
if err != nil || second.Reply != "lo" {
|
||||
t.Fatalf("second chunk = %#v, %v; want lo", second, err)
|
||||
}
|
||||
if _, err := stream.Recv(); !errors.Is(err, io.EOF) {
|
||||
t.Fatalf("final error = %v, want EOF", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestProvider_StreamPropagatesMalformedChunk(t *testing.T) {
|
||||
ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
w.Header().Set("Content-Type", "text/event-stream")
|
||||
_, _ = w.Write([]byte("data: {bad json}\n\n"))
|
||||
}))
|
||||
defer ts.Close()
|
||||
|
||||
p := NewProvider(ai.WithAPIKey("test-key"), ai.WithBaseURL(ts.URL))
|
||||
stream, err := p.Stream(context.Background(), &ai.Request{Prompt: "Hello"})
|
||||
if err != nil {
|
||||
t.Fatalf("Stream returned error: %v", err)
|
||||
}
|
||||
defer stream.Close()
|
||||
|
||||
if _, err := stream.Recv(); err == nil {
|
||||
t.Fatal("Recv returned nil error for malformed chunk")
|
||||
}
|
||||
}
|
||||
|
||||
func TestProvider_StreamPropagatesProviderError(t *testing.T) {
|
||||
ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
http.Error(w, "quota exhausted", http.StatusTooManyRequests)
|
||||
}))
|
||||
defer ts.Close()
|
||||
|
||||
p := NewProvider(ai.WithAPIKey("test-key"), ai.WithBaseURL(ts.URL))
|
||||
stream, err := p.Stream(context.Background(), &ai.Request{Prompt: "Hello"})
|
||||
if err == nil {
|
||||
_ = stream.Close()
|
||||
t.Fatal("Stream returned nil error for provider failure")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "429") || !strings.Contains(err.Error(), "quota exhausted") {
|
||||
t.Fatalf("Stream error = %v, want provider status and body", err)
|
||||
}
|
||||
if strings.Contains(err.Error(), "test-key") {
|
||||
t.Fatal("stream error leaked API key")
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -135,6 +135,8 @@ type RunInfo struct {
|
||||
VerificationFeedback string // feedback from the previous failed verifier attempt, when retrying a flow step
|
||||
Dispatch string // how the run was dispatched (direct, broker, schedule, resume) when known
|
||||
Trigger string // external trigger or schedule label that started the run, when known
|
||||
Spent int64 // cumulative paid x402 spend in this run, in the asset's smallest unit
|
||||
ToolSpend int64 // paid x402 spend attributed to the current tool call, in the asset's smallest unit
|
||||
}
|
||||
|
||||
type runInfoKey struct{}
|
||||
|
||||
@@ -216,6 +216,7 @@ func TestConfiguredProviderStreamsSkipWithoutCredentials(t *testing.T) {
|
||||
{provider: "together", keyEnv: "TOGETHER_API_KEY", modelEnv: "TOGETHER_MODEL"},
|
||||
{provider: "atlascloud", keyEnv: "ATLASCLOUD_API_KEY", modelEnv: "ATLASCLOUD_MODEL"},
|
||||
{provider: "anthropic", keyEnv: "ANTHROPIC_API_KEY", modelEnv: "ANTHROPIC_MODEL"},
|
||||
{provider: "gemini", keyEnv: "GEMINI_API_KEY", modelEnv: "GEMINI_MODEL"},
|
||||
} {
|
||||
tc := tc
|
||||
t.Run(tc.provider, func(t *testing.T) {
|
||||
@@ -256,24 +257,6 @@ func TestConfiguredProviderStreamsSkipWithoutCredentials(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestUnsupportedProvidersReturnStreamingUnsupportedAndStayUnregistered(t *testing.T) {
|
||||
for _, provider := range []string{"gemini"} {
|
||||
provider := provider
|
||||
t.Run(provider, func(t *testing.T) {
|
||||
if caps := ai.ProviderCapabilities(provider); caps.Stream {
|
||||
t.Fatalf("ProviderCapabilities(%q).Stream = true, want false", provider)
|
||||
}
|
||||
_, err := ai.New(provider, ai.WithAPIKey("test-key")).Stream(context.Background(), &ai.Request{Prompt: "Hello"})
|
||||
if !errors.Is(err, ai.ErrStreamingUnsupported) {
|
||||
t.Fatalf("Stream error = %v, want ErrStreamingUnsupported", err)
|
||||
}
|
||||
if err != nil && strings.Contains(err.Error(), "test-key") {
|
||||
t.Fatal("streaming unsupported error leaked API key")
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func conformingStreamProviders(t *testing.T) []string {
|
||||
t.Helper()
|
||||
providers := ai.RegisteredProviders("stream")
|
||||
|
||||
@@ -94,6 +94,9 @@ func writeAgentInspection(w io.Writer, name string, runs []goagent.RunSummary, a
|
||||
if run.LastErrorKind != "" {
|
||||
fmt.Fprintf(w, " error_kind=%s", run.LastErrorKind)
|
||||
}
|
||||
if run.Spent > 0 {
|
||||
fmt.Fprintf(w, " spent=%d", run.Spent)
|
||||
}
|
||||
if run.LastError != "" {
|
||||
fmt.Fprintf(w, " error=%q", run.LastError)
|
||||
}
|
||||
|
||||
@@ -11,13 +11,13 @@ import (
|
||||
)
|
||||
|
||||
func TestWriteAgentInspectionIncludesActionableBreadcrumbs(t *testing.T) {
|
||||
runs := []goagent.RunSummary{{RunID: "run-1", Status: "auth", Events: 4, LastKind: "model", LastError: "invalid API key", LastErrorKind: "auth", TraceID: "1234567890abcdef", Checkpoint: "failed", Stage: "ask"}}
|
||||
runs := []goagent.RunSummary{{RunID: "run-1", Status: "auth", Events: 4, LastKind: "model", LastError: "invalid API key", LastErrorKind: "auth", TraceID: "1234567890abcdef", Checkpoint: "failed", Stage: "ask", Spent: 7}}
|
||||
var out bytes.Buffer
|
||||
if err := writeAgentInspection(&out, "support", runs, false); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
got := out.String()
|
||||
for _, want := range []string{"Agent \"support\" runs", "run-1", "status=auth", "events=4", "last=model", "checkpoint=failed", "stage=ask", "error_kind=auth", `error="invalid API key"`, "trace=1234567890ab"} {
|
||||
for _, want := range []string{"Agent \"support\" runs", "run-1", "status=auth", "events=4", "last=model", "checkpoint=failed", "stage=ask", "error_kind=auth", `error="invalid API key"`, "trace=1234567890ab", "spent=7"} {
|
||||
if !strings.Contains(got, want) {
|
||||
t.Fatalf("output missing %q:\n%s", want, got)
|
||||
}
|
||||
|
||||
+4
-1
@@ -12,6 +12,7 @@ provider-free unless the example README says otherwise.
|
||||
| Prove the maintained 0→hero path | [`support`](./support/) | `go run ./examples/support` and `go test ./examples/support` | [`zero-to-hero` guide](../internal/website/docs/guides/zero-to-hero.md) |
|
||||
| See planning and delegation | [`agent-plan-delegate`](./agent-plan-delegate/) | `go run ./examples/agent-plan-delegate` | [`plan-delegate` guide](../internal/website/docs/guides/plan-delegate.md) |
|
||||
| Expose services through MCP | [`mcp/hello`](./mcp/hello/) | follow [`mcp`](./mcp/) setup | [`mcp/crud`](./mcp/crud/) and [`mcp/workflow`](./mcp/workflow/) |
|
||||
| Try a paid tool with x402 | [`agent-x402-buyer`](./agent-x402-buyer/) | `go run ./examples/agent-x402-buyer` | [`Payments (x402)` guide](../internal/website/docs/guides/x402-payments.md) |
|
||||
| Try A2A or gRPC interop next | [`agent-demo`](./agent-demo/) plus gateway docs | run the example, then use the gateway docs | [`grpc-interop`](./grpc-interop/) |
|
||||
| Add workflow durability | [`flow-durable`](./flow-durable/) | `go run ./examples/flow-durable` | [`flow-loop`](./flow-loop/) |
|
||||
|
||||
@@ -28,7 +29,9 @@ provider-free unless the example README says otherwise.
|
||||
4. **Interop next:** use [`mcp/hello`](./mcp/hello/), [`mcp/crud`](./mcp/crud/),
|
||||
and [`mcp/workflow`](./mcp/workflow/) when you are ready to expose tools to
|
||||
external AI clients.
|
||||
5. **Workflow depth:** use [`flow-durable`](./flow-durable/) once the agent path
|
||||
5. **Paid tools:** run [`agent-x402-buyer`](./agent-x402-buyer/) to see an
|
||||
agent pay a local x402-protected tool with a mock facilitator and budget.
|
||||
6. **Workflow depth:** use [`flow-durable`](./flow-durable/) once the agent path
|
||||
needs checkpointed, resumable deterministic work.
|
||||
|
||||
## CLI wayfinding
|
||||
|
||||
@@ -0,0 +1,24 @@
|
||||
# Agent x402 buyer
|
||||
|
||||
This example shows an agent paying for a paid HTTP tool with x402 without using
|
||||
live funds or a live chain.
|
||||
|
||||
It starts a local paid endpoint guarded by `wrapper/x402` seller middleware and a
|
||||
mock facilitator. A deterministic mock-model agent calls that endpoint as a tool,
|
||||
receives the HTTP 402 challenge, pays with `AgentPayer`, stays inside
|
||||
`AgentBudget`, retries the request, and prints the spend recorded for the run.
|
||||
|
||||
```bash
|
||||
go run ./examples/agent-x402-buyer
|
||||
```
|
||||
|
||||
Expected output includes:
|
||||
|
||||
- the paid tool response,
|
||||
- one facilitator verify and settle call, and
|
||||
- `run spend: 7 smallest units (budget 10)`.
|
||||
|
||||
The payment token and facilitator are intentionally local development fakes. To
|
||||
settle real x402 payments, keep the same `AgentPayer` / `AgentBudget` shape but
|
||||
replace the payer with a wallet-backed implementation and configure the seller
|
||||
middleware with a hosted or self-run x402 facilitator.
|
||||
@@ -0,0 +1,168 @@
|
||||
// Agent x402 buyer — a provider-free example of an agent paying for a paid tool.
|
||||
//
|
||||
// Run:
|
||||
//
|
||||
// go run ./examples/agent-x402-buyer
|
||||
//
|
||||
// It starts a local HTTP tool protected by x402 middleware, then asks a
|
||||
// deterministic mock-model agent to call that tool. The agent receives the 402
|
||||
// challenge, uses AgentPayer and AgentBudget to pay within a local mock
|
||||
// facilitator, retries the request, and prints the run spend.
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"os"
|
||||
"strings"
|
||||
|
||||
go_micro "go-micro.dev/v6"
|
||||
"go-micro.dev/v6/agent"
|
||||
"go-micro.dev/v6/ai"
|
||||
"go-micro.dev/v6/store"
|
||||
"go-micro.dev/v6/wrapper/x402"
|
||||
)
|
||||
|
||||
const (
|
||||
paidToolName = "paid.market_brief"
|
||||
price = int64(7)
|
||||
paymentToken = "dev-payment-token"
|
||||
)
|
||||
|
||||
type devFacilitator struct {
|
||||
verifyCount int
|
||||
settleCount int
|
||||
}
|
||||
|
||||
func (f *devFacilitator) Verify(ctx context.Context, payment string, req x402.Requirements) (x402.Result, error) {
|
||||
f.verifyCount++
|
||||
if payment != paymentToken {
|
||||
return x402.Result{Valid: false, Reason: "unknown dev payment token"}, nil
|
||||
}
|
||||
return x402.Result{Valid: true, Payer: "dev-agent-wallet"}, nil
|
||||
}
|
||||
|
||||
func (f *devFacilitator) Settle(ctx context.Context, payment string, req x402.Requirements) (x402.Result, error) {
|
||||
f.settleCount++
|
||||
return x402.Result{Valid: true, Settlement: "dev-settlement-001"}, nil
|
||||
}
|
||||
|
||||
type devPayer struct{}
|
||||
|
||||
func (devPayer) Pay(ctx context.Context, req x402.Requirements) (string, error) {
|
||||
return paymentToken, nil
|
||||
}
|
||||
|
||||
type mockModel struct{ opts ai.Options }
|
||||
|
||||
func newMock(opts ...ai.Option) ai.Model {
|
||||
m := &mockModel{}
|
||||
_ = m.Init(opts...)
|
||||
return m
|
||||
}
|
||||
|
||||
func (m *mockModel) Init(opts ...ai.Option) error {
|
||||
for _, o := range opts {
|
||||
o(&m.opts)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
func (m *mockModel) Options() ai.Options { return m.opts }
|
||||
func (m *mockModel) String() string { return "agent-x402-buyer-mock" }
|
||||
func (m *mockModel) Stream(context.Context, *ai.Request, ...ai.GenerateOption) (ai.Stream, error) {
|
||||
return nil, fmt.Errorf("stream not supported by agent-x402-buyer mock")
|
||||
}
|
||||
|
||||
func (m *mockModel) Generate(ctx context.Context, req *ai.Request, _ ...ai.GenerateOption) (*ai.Response, error) {
|
||||
for _, tool := range req.Tools {
|
||||
if tool.Name == paidToolName && m.opts.ToolHandler != nil {
|
||||
out := m.opts.ToolHandler(ctx, ai.ToolCall{ID: "paid-brief", Name: tool.Name, Input: map[string]any{"url": req.Prompt}})
|
||||
return &ai.Response{Answer: fmt.Sprintf("Paid tool returned: %s", out.Content)}, nil
|
||||
}
|
||||
}
|
||||
return &ai.Response{Answer: "No paid tool was available."}, nil
|
||||
}
|
||||
|
||||
func paidToolServer(fac *devFacilitator) *httptest.Server {
|
||||
mux := http.NewServeMux()
|
||||
paid := x402.Middleware(x402.Config{
|
||||
PayTo: "0xMerchantDevWallet",
|
||||
Network: "base-sepolia",
|
||||
Amount: fmt.Sprint(price),
|
||||
Description: "Local market brief for the x402 buyer example",
|
||||
Facilitator: fac,
|
||||
})
|
||||
mux.Handle("/brief", paid(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
_ = json.NewEncoder(w).Encode(map[string]any{
|
||||
"brief": "Mock demand is up 12% after the agent paid the local tool.",
|
||||
"settlement": w.Header().Get(x402.PaymentResponseHeader),
|
||||
})
|
||||
})))
|
||||
return httptest.NewServer(mux)
|
||||
}
|
||||
|
||||
func run(w io.Writer) error {
|
||||
ai.Register("agent-x402-buyer-mock", newMock)
|
||||
|
||||
fac := &devFacilitator{}
|
||||
srv := paidToolServer(fac)
|
||||
defer srv.Close()
|
||||
|
||||
st := store.NewMemoryStore()
|
||||
buyer := agent.New(
|
||||
agent.Name("x402-buyer"),
|
||||
agent.Provider("agent-x402-buyer-mock"),
|
||||
agent.Prompt("Call the paid market brief tool when given its URL."),
|
||||
agent.WithStore(st),
|
||||
go_micro.AgentPayer(devPayer{}),
|
||||
go_micro.AgentBudget(10),
|
||||
agent.WithTool(paidToolName, "Fetch a paid market brief over HTTP", map[string]any{
|
||||
"url": map[string]any{"type": "string", "description": "Paid HTTP endpoint to call"},
|
||||
}, func(ctx context.Context, input map[string]any) (string, error) {
|
||||
url, _ := input["url"].(string)
|
||||
resp, err := http.Get(url)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
body, err := io.ReadAll(resp.Body)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
return string(body), nil
|
||||
}),
|
||||
)
|
||||
|
||||
resp, err := buyer.Ask(context.Background(), srv.URL+"/brief")
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
events, err := agent.LoadRunEvents(st, "x402-buyer", resp.RunID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
var spent int64
|
||||
for _, event := range events {
|
||||
if event.Spent > spent {
|
||||
spent = event.Spent
|
||||
}
|
||||
}
|
||||
|
||||
fmt.Fprintln(w, "Agent x402 buyer (provider: mock, funds: local dev token)")
|
||||
fmt.Fprintln(w, strings.TrimSpace(resp.Reply))
|
||||
fmt.Fprintf(w, "facilitator verify=%d settle=%d\n", fac.verifyCount, fac.settleCount)
|
||||
fmt.Fprintf(w, "run spend: %d smallest units (budget 10)\n", spent)
|
||||
return nil
|
||||
}
|
||||
|
||||
func main() {
|
||||
if err := run(os.Stdout); err != nil {
|
||||
fmt.Println(err)
|
||||
os.Exit(1)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,282 @@
|
||||
package mcp
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"google.golang.org/grpc"
|
||||
"google.golang.org/grpc/credentials/insecure"
|
||||
reflectionpb "google.golang.org/grpc/reflection/grpc_reflection_v1alpha"
|
||||
"google.golang.org/protobuf/encoding/protojson"
|
||||
"google.golang.org/protobuf/proto"
|
||||
"google.golang.org/protobuf/reflect/protodesc"
|
||||
"google.golang.org/protobuf/reflect/protoreflect"
|
||||
"google.golang.org/protobuf/reflect/protoregistry"
|
||||
"google.golang.org/protobuf/types/descriptorpb"
|
||||
"google.golang.org/protobuf/types/dynamicpb"
|
||||
)
|
||||
|
||||
// ReflectedGRPCTarget describes an external gRPC server whose reflection
|
||||
// catalog should be exposed as MCP tools. It is intentionally opt-in: teams can
|
||||
// bridge existing reflected gRPC services without changing their servers or
|
||||
// registering them in go-micro.
|
||||
type ReflectedGRPCTarget struct {
|
||||
// Name prefixes generated tools. When empty, Address is sanitized and used.
|
||||
Name string
|
||||
// Address is the host:port of the reflected gRPC server.
|
||||
Address string
|
||||
// DialOptions customize the connection. If none are supplied, an insecure
|
||||
// transport is used for local/dev interoperability.
|
||||
DialOptions []grpc.DialOption
|
||||
// Timeout bounds reflection discovery and individual tool calls.
|
||||
Timeout time.Duration
|
||||
}
|
||||
|
||||
func (s *Server) discoverReflectedGRPC() error {
|
||||
for _, target := range s.opts.ReflectedGRPCTargets {
|
||||
if strings.TrimSpace(target.Address) == "" {
|
||||
continue
|
||||
}
|
||||
tools, err := s.reflectedGRPCTools(target)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
for _, tool := range tools {
|
||||
s.tools[tool.Name] = tool
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *Server) reflectedGRPCTools(target ReflectedGRPCTarget) ([]*Tool, error) {
|
||||
timeout := target.Timeout
|
||||
if timeout == 0 {
|
||||
timeout = 10 * time.Second
|
||||
}
|
||||
ctx, cancel := context.WithTimeout(s.opts.Context, timeout)
|
||||
defer cancel()
|
||||
|
||||
dialOpts := target.DialOptions
|
||||
if len(dialOpts) == 0 {
|
||||
dialOpts = []grpc.DialOption{grpc.WithTransportCredentials(insecure.NewCredentials())}
|
||||
}
|
||||
conn, err := grpc.NewClient(target.Address, dialOpts...)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("connect reflected grpc target %s: %w", target.Address, err)
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
files, services, err := loadReflectedFiles(ctx, conn)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("reflect grpc target %s: %w", target.Address, err)
|
||||
}
|
||||
|
||||
prefix := target.Name
|
||||
if prefix == "" {
|
||||
prefix = sanitizeToolPart(target.Address)
|
||||
}
|
||||
|
||||
var out []*Tool
|
||||
for _, serviceName := range services {
|
||||
desc, err := files.FindDescriptorByName(protoreflect.FullName(serviceName))
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
svc, ok := desc.(protoreflect.ServiceDescriptor)
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
for i := 0; i < svc.Methods().Len(); i++ {
|
||||
method := svc.Methods().Get(i)
|
||||
if method.IsStreamingClient() || method.IsStreamingServer() {
|
||||
continue
|
||||
}
|
||||
fullMethod := "/" + string(svc.FullName()) + "/" + string(method.Name())
|
||||
toolName := prefix + "." + strings.ReplaceAll(string(svc.FullName()), ".", "_") + "." + string(method.Name())
|
||||
input := method.Input()
|
||||
out = append(out, &Tool{
|
||||
Name: toolName,
|
||||
Description: fmt.Sprintf("Call reflected gRPC method %s on %s", fullMethod, target.Address),
|
||||
InputSchema: protoMessageSchema(input),
|
||||
Handler: reflectedGRPCHandler(target, fullMethod, input, method.Output()),
|
||||
})
|
||||
}
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
func loadReflectedFiles(ctx context.Context, conn *grpc.ClientConn) (*protoregistryFiles, []string, error) {
|
||||
client := reflectionpb.NewServerReflectionClient(conn)
|
||||
stream, err := client.ServerReflectionInfo(ctx)
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
if err := stream.Send(&reflectionpb.ServerReflectionRequest{MessageRequest: &reflectionpb.ServerReflectionRequest_ListServices{ListServices: ""}}); err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
resp, err := stream.Recv()
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
list := resp.GetListServicesResponse()
|
||||
if list == nil {
|
||||
return nil, nil, fmt.Errorf("reflection list services returned %T", resp.MessageResponse)
|
||||
}
|
||||
|
||||
set := &descriptorpb.FileDescriptorSet{}
|
||||
seen := map[string]bool{}
|
||||
var services []string
|
||||
for _, svc := range list.Service {
|
||||
name := svc.Name
|
||||
if strings.HasPrefix(name, "grpc.reflection.") {
|
||||
continue
|
||||
}
|
||||
services = append(services, name)
|
||||
if err := requestFileContainingSymbol(ctx, client, name, set, seen); err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
}
|
||||
files, err := newProtoregistryFiles(set)
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
return files, services, nil
|
||||
}
|
||||
|
||||
func requestFileContainingSymbol(ctx context.Context, client reflectionpb.ServerReflectionClient, symbol string, set *descriptorpb.FileDescriptorSet, seen map[string]bool) error {
|
||||
stream, err := client.ServerReflectionInfo(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if err := stream.Send(&reflectionpb.ServerReflectionRequest{MessageRequest: &reflectionpb.ServerReflectionRequest_FileContainingSymbol{FileContainingSymbol: symbol}}); err != nil {
|
||||
return err
|
||||
}
|
||||
resp, err := stream.Recv()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
fd := resp.GetFileDescriptorResponse()
|
||||
if fd == nil {
|
||||
return fmt.Errorf("reflection lookup for %s returned %T", symbol, resp.MessageResponse)
|
||||
}
|
||||
for _, raw := range fd.FileDescriptorProto {
|
||||
var file descriptorpb.FileDescriptorProto
|
||||
if err := proto.Unmarshal(raw, &file); err != nil {
|
||||
return err
|
||||
}
|
||||
name := file.GetName()
|
||||
if !seen[name] {
|
||||
seen[name] = true
|
||||
set.File = append(set.File, &file)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// protoregistryFiles is a narrow wrapper that keeps imports local to this file.
|
||||
type protoregistryFiles struct{ files *protoregistry.Files }
|
||||
|
||||
func newProtoregistryFiles(set *descriptorpb.FileDescriptorSet) (*protoregistryFiles, error) {
|
||||
files, err := protodesc.NewFiles(set)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &protoregistryFiles{files: files}, nil
|
||||
}
|
||||
|
||||
func (p *protoregistryFiles) FindDescriptorByName(name protoreflect.FullName) (protoreflect.Descriptor, error) {
|
||||
return p.files.FindDescriptorByName(name)
|
||||
}
|
||||
|
||||
func reflectedGRPCHandler(target ReflectedGRPCTarget, fullMethod string, input, output protoreflect.MessageDescriptor) func(map[string]interface{}) (interface{}, error) {
|
||||
return func(args map[string]interface{}) (interface{}, error) {
|
||||
timeout := target.Timeout
|
||||
if timeout == 0 {
|
||||
timeout = 10 * time.Second
|
||||
}
|
||||
ctx, cancel := context.WithTimeout(context.Background(), timeout)
|
||||
defer cancel()
|
||||
dialOpts := target.DialOptions
|
||||
if len(dialOpts) == 0 {
|
||||
dialOpts = []grpc.DialOption{grpc.WithTransportCredentials(insecure.NewCredentials())}
|
||||
}
|
||||
conn, err := grpc.NewClient(target.Address, dialOpts...)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
req := dynamicpb.NewMessage(input)
|
||||
raw, err := json.Marshal(args)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if err := protojson.Unmarshal(raw, req); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
rsp := dynamicpb.NewMessage(output)
|
||||
if err := conn.Invoke(ctx, fullMethod, req, rsp); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
b, err := protojson.MarshalOptions{UseProtoNames: true, EmitUnpopulated: true}.Marshal(rsp)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
var out interface{}
|
||||
if err := json.Unmarshal(b, &out); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
}
|
||||
|
||||
func protoMessageSchema(msg protoreflect.MessageDescriptor) map[string]interface{} {
|
||||
schema := map[string]interface{}{"type": "object", "properties": map[string]interface{}{}}
|
||||
props := schema["properties"].(map[string]interface{})
|
||||
fields := msg.Fields()
|
||||
for i := 0; i < fields.Len(); i++ {
|
||||
field := fields.Get(i)
|
||||
props[string(field.JSONName())] = protoFieldSchema(field)
|
||||
}
|
||||
return schema
|
||||
}
|
||||
|
||||
func protoFieldSchema(field protoreflect.FieldDescriptor) map[string]interface{} {
|
||||
schema := map[string]interface{}{"type": protoJSONType(field)}
|
||||
if field.IsList() {
|
||||
schema["items"] = map[string]interface{}{"type": protoJSONType(field)}
|
||||
}
|
||||
if field.Kind() == protoreflect.MessageKind || field.Kind() == protoreflect.GroupKind {
|
||||
schema = protoMessageSchema(field.Message())
|
||||
}
|
||||
return schema
|
||||
}
|
||||
|
||||
func protoJSONType(field protoreflect.FieldDescriptor) string {
|
||||
if field.IsList() {
|
||||
return "array"
|
||||
}
|
||||
switch field.Kind() {
|
||||
case protoreflect.BoolKind:
|
||||
return "boolean"
|
||||
case protoreflect.Int32Kind, protoreflect.Sint32Kind, protoreflect.Sfixed32Kind,
|
||||
protoreflect.Uint32Kind, protoreflect.Fixed32Kind, protoreflect.Int64Kind,
|
||||
protoreflect.Sint64Kind, protoreflect.Sfixed64Kind, protoreflect.Uint64Kind,
|
||||
protoreflect.Fixed64Kind:
|
||||
return "integer"
|
||||
case protoreflect.FloatKind, protoreflect.DoubleKind:
|
||||
return "number"
|
||||
case protoreflect.MessageKind, protoreflect.GroupKind:
|
||||
return "object"
|
||||
default:
|
||||
return "string"
|
||||
}
|
||||
}
|
||||
|
||||
func sanitizeToolPart(s string) string {
|
||||
r := strings.NewReplacer(":", "_", "/", "_", ".", "_", "-", "_")
|
||||
return r.Replace(s)
|
||||
}
|
||||
@@ -0,0 +1,62 @@
|
||||
package mcp
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"google.golang.org/grpc"
|
||||
helloworld "google.golang.org/grpc/examples/helloworld/helloworld"
|
||||
"google.golang.org/grpc/reflection"
|
||||
)
|
||||
|
||||
type reflectedGreeter struct {
|
||||
helloworld.UnimplementedGreeterServer
|
||||
}
|
||||
|
||||
func (reflectedGreeter) SayHello(_ context.Context, req *helloworld.HelloRequest) (*helloworld.HelloReply, error) {
|
||||
return &helloworld.HelloReply{Message: "hello " + req.Name}, nil
|
||||
}
|
||||
|
||||
func TestReflectedGRPCTargetDiscoversAndCallsUnaryTool(t *testing.T) {
|
||||
lis, err := net.Listen("tcp", "127.0.0.1:0")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
grpcServer := grpc.NewServer()
|
||||
helloworld.RegisterGreeterServer(grpcServer, reflectedGreeter{})
|
||||
reflection.Register(grpcServer)
|
||||
go grpcServer.Serve(lis)
|
||||
defer grpcServer.Stop()
|
||||
|
||||
s := newTestServer(Options{Context: context.Background()})
|
||||
tools, err := s.reflectedGRPCTools(ReflectedGRPCTarget{
|
||||
Name: "demo",
|
||||
Address: lis.Addr().String(),
|
||||
Timeout: 3 * time.Second,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("discover reflected tools: %v", err)
|
||||
}
|
||||
if len(tools) != 1 {
|
||||
t.Fatalf("tools len = %d, want 1", len(tools))
|
||||
}
|
||||
tool := tools[0]
|
||||
if tool.Name != "demo.helloworld_Greeter.SayHello" {
|
||||
t.Fatalf("tool name = %q", tool.Name)
|
||||
}
|
||||
props := tool.InputSchema["properties"].(map[string]interface{})
|
||||
if _, ok := props["name"]; !ok {
|
||||
t.Fatalf("input schema missing name: %#v", tool.InputSchema)
|
||||
}
|
||||
|
||||
out, err := tool.Handler(map[string]interface{}{"name": "Ada"})
|
||||
if err != nil {
|
||||
t.Fatalf("call reflected tool: %v", err)
|
||||
}
|
||||
got := out.(map[string]interface{})["message"]
|
||||
if got != "hello Ada" {
|
||||
t.Fatalf("message = %v, want hello Ada", got)
|
||||
}
|
||||
}
|
||||
@@ -157,6 +157,11 @@ type Options struct {
|
||||
// (the /mcp/call endpoint). Listing tools and health stay free.
|
||||
// Opt-in: leave nil to disable payments.
|
||||
Payment *x402.Config
|
||||
|
||||
// ReflectedGRPCTargets exposes unary methods from external gRPC servers
|
||||
// that support server reflection as MCP tools. This bridges existing gRPC
|
||||
// services into the agent tool catalog without requiring go-micro handlers.
|
||||
ReflectedGRPCTargets []ReflectedGRPCTarget
|
||||
}
|
||||
|
||||
// Server represents a running MCP gateway
|
||||
@@ -286,6 +291,10 @@ func (s *Server) discoverServices() error {
|
||||
s.toolsMu.Lock()
|
||||
defer s.toolsMu.Unlock()
|
||||
|
||||
if err := s.discoverReflectedGRPC(); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
for _, svc := range services {
|
||||
// Get full service details
|
||||
fullSvcs, err := s.opts.Registry.GetService(svc.Name)
|
||||
|
||||
@@ -40,7 +40,7 @@ The built-in providers currently register these capability interfaces:
|
||||
| --- | --- | --- | --- | --- | --- |
|
||||
| `anthropic` | Yes | No | No | Yes | Yes |
|
||||
| `atlascloud` | Yes | Yes | Yes | Yes | No |
|
||||
| `gemini` | Yes | No | No | No | No |
|
||||
| `gemini` | Yes | No | No | Yes | No |
|
||||
| `groq` | Yes | No | No | Yes | Yes |
|
||||
| `minimax` | Yes | No | No | Yes | Yes |
|
||||
| `mistral` | Yes | No | No | Yes | Yes |
|
||||
|
||||
@@ -59,7 +59,7 @@ previous section.
|
||||
| --- | --- | --- | --- | --- | --- |
|
||||
| `anthropic` | ✅ Verified when configured | — Unsupported | — Unsupported | ✅ Verified when configured | ⚠️ Unverified |
|
||||
| `openai` | ✅ Verified when configured | ✅ Registered | — Unsupported | ⚠️ Unverified | ⚠️ Unverified |
|
||||
| `gemini` | ✅ Verified when configured | — Unsupported | — Unsupported | ⚠️ Unverified | ⚠️ Unverified |
|
||||
| `gemini` | ✅ Verified when configured | — Unsupported | — Unsupported | ✅ Verified when configured | ⚠️ Unverified |
|
||||
| `groq` | ✅ Verified when configured | — Unsupported | — Unsupported | ⚠️ Unverified | ⚠️ Unverified |
|
||||
| `mistral` | ✅ Verified when configured | — Unsupported | — Unsupported | ⚠️ Unverified | ⚠️ Unverified |
|
||||
| `together` | ✅ Verified when configured | — Unsupported | — Unsupported | ⚠️ Unverified | ⚠️ Unverified |
|
||||
|
||||
@@ -12,6 +12,7 @@ import (
|
||||
"go-micro.dev/v6/server"
|
||||
"go-micro.dev/v6/service"
|
||||
"go-micro.dev/v6/store"
|
||||
"go-micro.dev/v6/wrapper/x402"
|
||||
"go.opentelemetry.io/otel/trace"
|
||||
)
|
||||
|
||||
@@ -130,6 +131,13 @@ func AgentToolSpend(tool string, amount int64) AgentOption {
|
||||
return agent.ToolSpend(tool, amount)
|
||||
}
|
||||
|
||||
// AgentPayer configures the wallet/signing hook used to settle x402-paid tools.
|
||||
func AgentPayer(p x402.Payer) AgentOption { return agent.Payer(p) }
|
||||
|
||||
// AgentBudget bounds autonomous x402 payments per Ask, in the asset's smallest
|
||||
// unit (0 = unlimited).
|
||||
func AgentBudget(amount int64) AgentOption { return agent.Budget(amount) }
|
||||
|
||||
// AgentModelCallTimeout sets the timeout for each provider Generate call.
|
||||
func AgentModelCallTimeout(d time.Duration) AgentOption { return agent.ModelCallTimeout(d) }
|
||||
|
||||
|
||||
Reference in New Issue
Block a user