Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| bb5164c3a8 |
+19
-33
@@ -1,41 +1,27 @@
|
||||
# Priorities
|
||||
|
||||
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.
|
||||
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.
|
||||
|
||||
**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.
|
||||
**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.
|
||||
|
||||
**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.
|
||||
**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.
|
||||
|
||||
## Work queue (ranked)
|
||||
|
||||
### Capability — the headline (roadmap: Now / Next)
|
||||
1. **Normalize provider error status in agent inspection** ([#4777](https://github.com/micro/go-micro/issues/4777)) — #4775 closed the retry/timeout controls gap, so the next highest-value user-facing gap is operability when a provider still fails: classify timeout, cancellation, rate-limit, auth/configuration, and transient provider errors into stable RunInfo/inspect/tracing details so developers know whether to retry, fix credentials, raise deadlines, or wait out an outage without changing provider public APIs.
|
||||
|
||||
1. **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.
|
||||
2. **AP2 mandate foundation for agent payments** ([#4841](https://github.com/micro/go-micro/issues/4841)) — add opt-in checkout/payment mandate signing and verification so A2A-carried payment authority can settle over x402 without changing defaults.
|
||||
3. **Kubernetes CRD reconciler foundation** ([#4842](https://github.com/micro/go-micro/issues/4842)) — turn the shipped alpha `Agent`, `Service`, and `Flow` CRDs into a minimally runnable native deployment path with workload reconciliation and status conditions.
|
||||
|
||||
### In flight — do not re-queue
|
||||
|
||||
_None right now._
|
||||
|
||||
### 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, provider-failure inspection, x402 buyer safety, gRPC-reflection MCP,
|
||||
MCP result conformance, and the alpha Kubernetes CRD surface. Further churn in those
|
||||
areas should be marked `needs-human` unless it unlocks a clear user-visible capability._
|
||||
_Seeded by Claude Code from the roadmap + open issues; thereafter maintained by the
|
||||
architecture-review pass._
|
||||
|
||||
@@ -18,49 +18,15 @@ 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/`)
|
||||
|
||||
### Security
|
||||
- **x402 spend-cap hardening** — the paying `Client` now refuses a 402 whose `maxAmountRequired` is not a positive integer (a swallowed parse error or negative amount previously bypassed the budget cap), and a new `Config.RequireSettlement` fails closed when a paid request is served by a verify-only facilitator that never captures funds. (`wrapper/x402/`)
|
||||
|
||||
---
|
||||
|
||||
## [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
|
||||
|
||||
+30
-41
@@ -32,55 +32,44 @@ default.
|
||||
and history, end to end.
|
||||
5. Battle-tested: works across every provider, fails safely, observable.
|
||||
|
||||
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.
|
||||
## Now — hardening
|
||||
|
||||
## Now — capability
|
||||
- **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.
|
||||
|
||||
- **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.
|
||||
## Shipped agent depth
|
||||
|
||||
## Next — reach & deployment
|
||||
- **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.
|
||||
|
||||
- **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.
|
||||
## Next — agentic depth
|
||||
|
||||
## Later — exploratory
|
||||
- **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.
|
||||
|
||||
- **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).
|
||||
## Later
|
||||
|
||||
## Ongoing — hardening & DX (background, not the headline)
|
||||
- Memory management (summarization, retrieval/RAG); human-in-the-loop pause/resume;
|
||||
richer A2A live-stream reconnection (`tasks/resubscribe`) and `input-required`
|
||||
handoffs.
|
||||
|
||||
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.
|
||||
## 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.
|
||||
|
||||
## How it's sustained
|
||||
|
||||
|
||||
@@ -5,8 +5,6 @@ import (
|
||||
"crypto/sha256"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
@@ -14,7 +12,6 @@ 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
|
||||
@@ -131,7 +128,6 @@ 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)
|
||||
@@ -177,64 +173,6 @@ 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.
|
||||
@@ -451,11 +389,6 @@ 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,17 +2,12 @@ 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
|
||||
@@ -186,132 +181,3 @@ 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,7 +10,6 @@ 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"
|
||||
)
|
||||
|
||||
@@ -106,10 +105,6 @@ 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".
|
||||
@@ -255,18 +250,6 @@ 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.
|
||||
|
||||
+3
-37
@@ -51,8 +51,6 @@ const (
|
||||
AttrDispatch = "agent.dispatch"
|
||||
AttrTrigger = "agent.trigger"
|
||||
AttrRunEventKind = "agent.event.kind"
|
||||
AttrSpend = "agent.spend"
|
||||
AttrToolSpend = "agent.tool.spend"
|
||||
)
|
||||
|
||||
type RunEvent struct {
|
||||
@@ -75,8 +73,6 @@ 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
|
||||
@@ -86,8 +82,7 @@ type Usage = ai.Usage
|
||||
type RunListOptions struct {
|
||||
// Status, when set, keeps only runs with the matching status
|
||||
// (for example "running", "done", "canceled", "timeout",
|
||||
// "rate_limited", "auth", "configuration", "unavailable",
|
||||
// "provider_error", "error", or "refused").
|
||||
// "rate_limited", "error", or "refused").
|
||||
Status string
|
||||
// TraceID, when set, keeps only runs correlated with this trace id.
|
||||
// A prefix is accepted so operators can paste the shortened trace id
|
||||
@@ -115,7 +110,6 @@ 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 {
|
||||
@@ -367,7 +361,6 @@ 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)
|
||||
@@ -377,7 +370,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), Spent: a.spend, ToolSpend: a.spend - spentBefore})
|
||||
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)})
|
||||
return res
|
||||
}
|
||||
|
||||
@@ -391,7 +384,6 @@ 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 {
|
||||
@@ -401,9 +393,6 @@ 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))
|
||||
}
|
||||
@@ -419,7 +408,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), Spent: a.spend, ToolSpend: toolSpend})
|
||||
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)})
|
||||
span.End()
|
||||
return res
|
||||
}
|
||||
@@ -506,12 +495,6 @@ 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))
|
||||
@@ -546,12 +529,6 @@ 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
|
||||
}
|
||||
|
||||
@@ -638,9 +615,6 @@ 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
|
||||
@@ -688,14 +662,6 @@ func runErrorStatus(kind string) string {
|
||||
return "timeout"
|
||||
case ai.ErrorKindRateLimited:
|
||||
return "rate_limited"
|
||||
case ai.ErrorKindAuth:
|
||||
return "auth"
|
||||
case ai.ErrorKindConfiguration:
|
||||
return "configuration"
|
||||
case ai.ErrorKindUnavailable:
|
||||
return "unavailable"
|
||||
case ai.ErrorKindProvider:
|
||||
return "provider_error"
|
||||
default:
|
||||
return "error"
|
||||
}
|
||||
|
||||
+1
-61
@@ -421,63 +421,6 @@ 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))
|
||||
@@ -720,10 +663,7 @@ func TestRunStatusClassifiesOperationalErrorKinds(t *testing.T) {
|
||||
{name: "canceled", kind: ai.ErrorKindCanceled, want: "canceled"},
|
||||
{name: "timeout", kind: ai.ErrorKindTimeout, want: "timeout"},
|
||||
{name: "rate limited", kind: ai.ErrorKindRateLimited, want: "rate_limited"},
|
||||
{name: "auth", kind: ai.ErrorKindAuth, want: "auth"},
|
||||
{name: "configuration", kind: ai.ErrorKindConfiguration, want: "configuration"},
|
||||
{name: "unavailable", kind: ai.ErrorKindUnavailable, want: "unavailable"},
|
||||
{name: "provider", kind: ai.ErrorKindProvider, want: "provider_error"},
|
||||
{name: "provider", kind: ai.ErrorKindProvider, want: "error"},
|
||||
}
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
|
||||
@@ -35,7 +35,7 @@ func TestRegisteredProviders(t *testing.T) {
|
||||
}
|
||||
|
||||
got = ai.RegisteredProviders("stream")
|
||||
want = []string{"anthropic", "atlascloud", "gemini", "groq", "minimax", "mistral", "openai", "together"}
|
||||
want = []string{"anthropic", "atlascloud", "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, Stream: true}},
|
||||
{Provider: "gemini", Capabilities: ai.Capabilities{Model: 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", "gemini", "groq", "minimax", "mistral", "openai", "test-stream", "together"}
|
||||
want := []string{"anthropic", "atlascloud", "groq", "minimax", "mistral", "openai", "test-stream", "together"}
|
||||
if !reflect.DeepEqual(got, want) {
|
||||
t.Fatalf("RegisteredProviders(stream) = %#v, want %#v", got, want)
|
||||
}
|
||||
|
||||
+4
-136
@@ -10,7 +10,6 @@
|
||||
package gemini
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
@@ -26,7 +25,6 @@ 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.
|
||||
@@ -71,7 +69,9 @@ func (p *Provider) Generate(ctx context.Context, req *ai.Request, opts ...ai.Gen
|
||||
})
|
||||
}
|
||||
|
||||
contents := geminiContents(req)
|
||||
contents := []map[string]any{
|
||||
{"role": "user", "parts": []map[string]any{{"text": req.Prompt}}},
|
||||
}
|
||||
|
||||
apiReq := map[string]any{
|
||||
"contents": contents,
|
||||
@@ -135,121 +135,7 @@ 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) {
|
||||
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()
|
||||
return nil, fmt.Errorf("%w: gemini provider", ai.ErrStreamingUnsupported)
|
||||
}
|
||||
|
||||
func (p *Provider) callAPI(ctx context.Context, req map[string]any) (*ai.Response, []map[string]any, error) {
|
||||
@@ -338,21 +224,3 @@ 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
|
||||
}
|
||||
|
||||
+6
-103
@@ -2,12 +2,7 @@ package gemini
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"io"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"go-micro.dev/v6/ai"
|
||||
@@ -86,108 +81,16 @@ func TestProvider_Generate_NoAPIKey(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
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)
|
||||
}
|
||||
func TestProvider_Stream_NotImplemented(t *testing.T) {
|
||||
p := NewProvider()
|
||||
|
||||
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"},
|
||||
},
|
||||
req := &ai.Request{
|
||||
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")
|
||||
}
|
||||
|
||||
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")
|
||||
_, err := p.Stream(context.Background(), req)
|
||||
if !errors.Is(err, ai.ErrStreamingUnsupported) {
|
||||
t.Fatalf("Stream error = %v, want ErrStreamingUnsupported", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -135,8 +135,6 @@ 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{}
|
||||
|
||||
+6
-16
@@ -85,14 +85,12 @@ func parseRetryAfter(value string, now time.Time) time.Duration {
|
||||
type ErrorKind string
|
||||
|
||||
const (
|
||||
ErrorKindUnknown ErrorKind = "unknown"
|
||||
ErrorKindCanceled ErrorKind = "canceled"
|
||||
ErrorKindTimeout ErrorKind = "timeout"
|
||||
ErrorKindRateLimited ErrorKind = "rate_limited"
|
||||
ErrorKindUnavailable ErrorKind = "unavailable"
|
||||
ErrorKindAuth ErrorKind = "auth"
|
||||
ErrorKindConfiguration ErrorKind = "configuration"
|
||||
ErrorKindProvider ErrorKind = "provider"
|
||||
ErrorKindUnknown ErrorKind = "unknown"
|
||||
ErrorKindCanceled ErrorKind = "canceled"
|
||||
ErrorKindTimeout ErrorKind = "timeout"
|
||||
ErrorKindRateLimited ErrorKind = "rate_limited"
|
||||
ErrorKindUnavailable ErrorKind = "unavailable"
|
||||
ErrorKindProvider ErrorKind = "provider"
|
||||
)
|
||||
|
||||
// ClassifiedError is implemented by errors that expose a stable ErrorKind.
|
||||
@@ -278,10 +276,6 @@ func ClassifyError(err error) ErrorKind {
|
||||
switch {
|
||||
case code == 429:
|
||||
return ErrorKindRateLimited
|
||||
case code == 401 || code == 403:
|
||||
return ErrorKindAuth
|
||||
case code == 400 || code == 404:
|
||||
return ErrorKindConfiguration
|
||||
case code >= 500:
|
||||
return ErrorKindUnavailable
|
||||
case code > 0:
|
||||
@@ -294,10 +288,6 @@ func ClassifyError(err error) ErrorKind {
|
||||
return ErrorKindRateLimited
|
||||
case strings.Contains(msg, "timeout") || strings.Contains(msg, "deadline"):
|
||||
return ErrorKindTimeout
|
||||
case strings.Contains(msg, "unauthorized") || strings.Contains(msg, "forbidden") || strings.Contains(msg, "invalid api key") || strings.Contains(msg, "api key") || strings.Contains(msg, "credential"):
|
||||
return ErrorKindAuth
|
||||
case strings.Contains(msg, "missing") || strings.Contains(msg, "not configured") || strings.Contains(msg, "configuration") || strings.Contains(msg, "unsupported model") || strings.Contains(msg, "model not found"):
|
||||
return ErrorKindConfiguration
|
||||
case strings.Contains(msg, "temporar") || strings.Contains(msg, "unavailable"):
|
||||
return ErrorKindUnavailable
|
||||
default:
|
||||
|
||||
+1
-5
@@ -230,13 +230,9 @@ func TestClassifyErrorDistinguishesOperationalOutcomes(t *testing.T) {
|
||||
{name: "canceled", err: context.Canceled, want: ErrorKindCanceled},
|
||||
{name: "timeout", err: context.DeadlineExceeded, want: ErrorKindTimeout},
|
||||
{name: "rate limit status", err: statusErr(429), want: ErrorKindRateLimited},
|
||||
{name: "auth status", err: statusErr(401), want: ErrorKindAuth},
|
||||
{name: "configuration status", err: statusErr(400), want: ErrorKindConfiguration},
|
||||
{name: "unavailable status", err: statusErr(503), want: ErrorKindUnavailable},
|
||||
{name: "provider status", err: statusErr(409), want: ErrorKindProvider},
|
||||
{name: "provider status", err: statusErr(400), want: ErrorKindProvider},
|
||||
{name: "rate limit text", err: errors.New("rate limit exceeded"), want: ErrorKindRateLimited},
|
||||
{name: "auth text", err: errors.New("invalid API key"), want: ErrorKindAuth},
|
||||
{name: "configuration text", err: errors.New("model not found"), want: ErrorKindConfiguration},
|
||||
}
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
|
||||
@@ -216,7 +216,6 @@ 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) {
|
||||
@@ -257,6 +256,24 @@ 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")
|
||||
|
||||
@@ -43,7 +43,7 @@ It reads durable local run history, so it works after the agent or flow has stop
|
||||
func inspectAgentFlags() []cli.Flag {
|
||||
return []cli.Flag{
|
||||
&cli.BoolFlag{Name: "json", Usage: "Print run summaries as JSON for automation"},
|
||||
&cli.StringFlag{Name: "status", Usage: "Only show runs with this status (running, done, canceled, timeout, rate_limited, auth, configuration, unavailable, provider_error, error, refused)"},
|
||||
&cli.StringFlag{Name: "status", Usage: "Only show runs with this status (running, done, error, refused)"},
|
||||
&cli.StringFlag{Name: "trace", Usage: "Only show runs whose trace id matches this full id or prefix"},
|
||||
&cli.IntFlag{Name: "limit", Usage: "Show the most recently updated N runs"},
|
||||
}
|
||||
@@ -91,12 +91,6 @@ func writeAgentInspection(w io.Writer, name string, runs []goagent.RunSummary, a
|
||||
if run.Stage != "" {
|
||||
fmt.Fprintf(w, " stage=%s", run.Stage)
|
||||
}
|
||||
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", Spent: 7}}
|
||||
runs := []goagent.RunSummary{{RunID: "run-1", Status: "error", Events: 4, LastKind: "tool", LastError: "boom", TraceID: "1234567890abcdef", Checkpoint: "failed", Stage: "ask"}}
|
||||
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", "spent=7"} {
|
||||
for _, want := range []string{"Agent \"support\" runs", "run-1", "status=error", "events=4", "last=tool", "checkpoint=failed", "stage=ask", `error="boom"`, "trace=1234567890ab", `micro agent history support run-1`, `micro.AgentResume(ctx, agent, "run-1")`, `micro.ResumeStreamAsk(ctx, agent, "run-1")`} {
|
||||
if !strings.Contains(got, want) {
|
||||
t.Fatalf("output missing %q:\n%s", want, got)
|
||||
}
|
||||
|
||||
@@ -1,31 +0,0 @@
|
||||
# Kubernetes deployment foundation (alpha)
|
||||
|
||||
This package is the first opt-in Kubernetes foundation for the Go Micro lifecycle:
|
||||
`Service`, `Agent`, and `Flow` resources. It is intentionally experimental and
|
||||
additive. Nothing in the Go Micro runtime installs these resources or changes
|
||||
production defaults.
|
||||
|
||||
## What is included
|
||||
|
||||
- Alpha CRD manifests in `config/crd/` for `agents.micro.dev`,
|
||||
`services.micro.dev`, and `flows.micro.dev`.
|
||||
- A small dependency-free mapper that turns a desired Go Micro resource into the
|
||||
Kubernetes `Deployment` shape an operator reconciliation loop will own.
|
||||
- Unit tests that validate the structural CRD fragments and dry-run the
|
||||
Agent-to-Deployment mapping.
|
||||
|
||||
## Local validation
|
||||
|
||||
```sh
|
||||
go test ./deploy/kubernetes
|
||||
```
|
||||
|
||||
If you have a Kubernetes cluster and `kubectl` available, you can also perform a
|
||||
server-side dry run of the CRDs:
|
||||
|
||||
```sh
|
||||
kubectl apply --dry-run=server -f deploy/kubernetes/config/crd/
|
||||
```
|
||||
|
||||
The manifests are `v1alpha1`; expect the API shape to evolve before this becomes
|
||||
a production operator.
|
||||
@@ -1,37 +0,0 @@
|
||||
apiVersion: apiextensions.k8s.io/v1
|
||||
kind: CustomResourceDefinition
|
||||
metadata:
|
||||
name: agents.micro.dev
|
||||
spec:
|
||||
group: micro.dev
|
||||
scope: Namespaced
|
||||
names:
|
||||
plural: agents
|
||||
singular: agent
|
||||
kind: Agent
|
||||
shortNames: [magent]
|
||||
versions:
|
||||
- name: v1alpha1
|
||||
served: true
|
||||
storage: true
|
||||
schema:
|
||||
openAPIV3Schema:
|
||||
type: object
|
||||
required: [spec]
|
||||
properties:
|
||||
spec:
|
||||
type: object
|
||||
required: [image]
|
||||
properties:
|
||||
image: {type: string, minLength: 1}
|
||||
command:
|
||||
type: array
|
||||
items: {type: string}
|
||||
args:
|
||||
type: array
|
||||
items: {type: string}
|
||||
replicas: {type: integer, minimum: 0}
|
||||
registry: {type: string}
|
||||
env:
|
||||
type: object
|
||||
additionalProperties: {type: string}
|
||||
@@ -1,37 +0,0 @@
|
||||
apiVersion: apiextensions.k8s.io/v1
|
||||
kind: CustomResourceDefinition
|
||||
metadata:
|
||||
name: flows.micro.dev
|
||||
spec:
|
||||
group: micro.dev
|
||||
scope: Namespaced
|
||||
names:
|
||||
plural: flows
|
||||
singular: flow
|
||||
kind: Flow
|
||||
shortNames: [mflow]
|
||||
versions:
|
||||
- name: v1alpha1
|
||||
served: true
|
||||
storage: true
|
||||
schema:
|
||||
openAPIV3Schema:
|
||||
type: object
|
||||
required: [spec]
|
||||
properties:
|
||||
spec:
|
||||
type: object
|
||||
required: [image]
|
||||
properties:
|
||||
image: {type: string, minLength: 1}
|
||||
command:
|
||||
type: array
|
||||
items: {type: string}
|
||||
args:
|
||||
type: array
|
||||
items: {type: string}
|
||||
replicas: {type: integer, minimum: 0}
|
||||
registry: {type: string}
|
||||
env:
|
||||
type: object
|
||||
additionalProperties: {type: string}
|
||||
@@ -1,37 +0,0 @@
|
||||
apiVersion: apiextensions.k8s.io/v1
|
||||
kind: CustomResourceDefinition
|
||||
metadata:
|
||||
name: services.micro.dev
|
||||
spec:
|
||||
group: micro.dev
|
||||
scope: Namespaced
|
||||
names:
|
||||
plural: services
|
||||
singular: service
|
||||
kind: Service
|
||||
shortNames: [mservice]
|
||||
versions:
|
||||
- name: v1alpha1
|
||||
served: true
|
||||
storage: true
|
||||
schema:
|
||||
openAPIV3Schema:
|
||||
type: object
|
||||
required: [spec]
|
||||
properties:
|
||||
spec:
|
||||
type: object
|
||||
required: [image]
|
||||
properties:
|
||||
image: {type: string, minLength: 1}
|
||||
command:
|
||||
type: array
|
||||
items: {type: string}
|
||||
args:
|
||||
type: array
|
||||
items: {type: string}
|
||||
replicas: {type: integer, minimum: 0}
|
||||
registry: {type: string}
|
||||
env:
|
||||
type: object
|
||||
additionalProperties: {type: string}
|
||||
@@ -1,8 +0,0 @@
|
||||
// Package kubernetes contains the experimental Kubernetes deployment foundation
|
||||
// for Go Micro services, agents, and flows.
|
||||
//
|
||||
// The package is intentionally small and additive: it exposes alpha custom
|
||||
// resource manifests and a dry-run mapper that turns a resource spec into the
|
||||
// Deployment shape an operator would reconcile. It does not install an operator
|
||||
// or change any runtime defaults.
|
||||
package kubernetes
|
||||
@@ -1,87 +0,0 @@
|
||||
package kubernetes
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestCRDManifestsAreStructural(t *testing.T) {
|
||||
for _, kind := range []Kind{KindAgent, KindService, KindFlow} {
|
||||
manifest := CRDManifests[kind]
|
||||
if manifest == "" {
|
||||
t.Fatalf("missing manifest for %s", kind)
|
||||
}
|
||||
checks := []string{
|
||||
"apiVersion: apiextensions.k8s.io/v1",
|
||||
"kind: CustomResourceDefinition",
|
||||
"group: micro.dev",
|
||||
"kind: " + string(kind),
|
||||
"name: v1alpha1",
|
||||
"served: true",
|
||||
"storage: true",
|
||||
"openAPIV3Schema:",
|
||||
"type: object",
|
||||
"required: [image]",
|
||||
}
|
||||
for _, check := range checks {
|
||||
if !strings.Contains(manifest, check) {
|
||||
t.Fatalf("%s manifest missing %q:\n%s", kind, check, manifest)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestMapDeploymentForAgent(t *testing.T) {
|
||||
deployment, err := MapDeployment(Resource{
|
||||
Kind: KindAgent,
|
||||
Name: "support-agent",
|
||||
Namespace: "agents",
|
||||
Spec: WorkloadSpec{
|
||||
Image: "ghcr.io/acme/support-agent:v1",
|
||||
Replicas: 2,
|
||||
Registry: "kubernetes",
|
||||
Environment: map[string]string{
|
||||
"MODEL": "gpt-5.5",
|
||||
},
|
||||
},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("MapDeployment returned error: %v", err)
|
||||
}
|
||||
if deployment.Name != "support-agent" || deployment.Namespace != "agents" {
|
||||
t.Fatalf("unexpected identity: %+v", deployment)
|
||||
}
|
||||
if deployment.Replicas != 2 {
|
||||
t.Fatalf("replicas = %d, want 2", deployment.Replicas)
|
||||
}
|
||||
if got := deployment.Labels["micro.dev/kind"]; got != "agent" {
|
||||
t.Fatalf("micro.dev/kind label = %q, want agent", got)
|
||||
}
|
||||
container := deployment.Pod.Container
|
||||
if container.Image != "ghcr.io/acme/support-agent:v1" {
|
||||
t.Fatalf("image = %q", container.Image)
|
||||
}
|
||||
if got := container.Environment["MICRO_REGISTRY"]; got != "kubernetes" {
|
||||
t.Fatalf("MICRO_REGISTRY = %q, want kubernetes", got)
|
||||
}
|
||||
if got := container.Environment["MODEL"]; got != "gpt-5.5" {
|
||||
t.Fatalf("MODEL = %q, want gpt-5.5", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestMapDeploymentDefaultsAndValidation(t *testing.T) {
|
||||
deployment, err := MapDeployment(Resource{Kind: KindService, Name: "api", Spec: WorkloadSpec{Image: "api:latest"}})
|
||||
if err != nil {
|
||||
t.Fatalf("MapDeployment returned error: %v", err)
|
||||
}
|
||||
if deployment.Namespace != "default" || deployment.Replicas != 1 {
|
||||
t.Fatalf("defaults = namespace %q replicas %d", deployment.Namespace, deployment.Replicas)
|
||||
}
|
||||
|
||||
if _, err := MapDeployment(Resource{Kind: KindFlow, Name: "ingest"}); err == nil {
|
||||
t.Fatal("MapDeployment without image succeeded")
|
||||
}
|
||||
if _, err := MapDeployment(Resource{Kind: "Job", Name: "job", Spec: WorkloadSpec{Image: "job:latest"}}); err == nil {
|
||||
t.Fatal("MapDeployment with unsupported kind succeeded")
|
||||
}
|
||||
}
|
||||
@@ -1,125 +0,0 @@
|
||||
package kubernetes
|
||||
|
||||
// CRDManifests contains the alpha CRDs for Go Micro lifecycle resources.
|
||||
var CRDManifests = map[Kind]string{
|
||||
KindAgent: agentCRD,
|
||||
KindService: serviceCRD,
|
||||
KindFlow: flowCRD,
|
||||
}
|
||||
|
||||
const agentCRD = `apiVersion: apiextensions.k8s.io/v1
|
||||
kind: CustomResourceDefinition
|
||||
metadata:
|
||||
name: agents.micro.dev
|
||||
spec:
|
||||
group: micro.dev
|
||||
scope: Namespaced
|
||||
names:
|
||||
plural: agents
|
||||
singular: agent
|
||||
kind: Agent
|
||||
shortNames: [magent]
|
||||
versions:
|
||||
- name: v1alpha1
|
||||
served: true
|
||||
storage: true
|
||||
schema:
|
||||
openAPIV3Schema:
|
||||
type: object
|
||||
required: [spec]
|
||||
properties:
|
||||
spec:
|
||||
type: object
|
||||
required: [image]
|
||||
properties:
|
||||
image: {type: string, minLength: 1}
|
||||
command:
|
||||
type: array
|
||||
items: {type: string}
|
||||
args:
|
||||
type: array
|
||||
items: {type: string}
|
||||
replicas: {type: integer, minimum: 0}
|
||||
registry: {type: string}
|
||||
env:
|
||||
type: object
|
||||
additionalProperties: {type: string}
|
||||
`
|
||||
|
||||
const serviceCRD = `apiVersion: apiextensions.k8s.io/v1
|
||||
kind: CustomResourceDefinition
|
||||
metadata:
|
||||
name: services.micro.dev
|
||||
spec:
|
||||
group: micro.dev
|
||||
scope: Namespaced
|
||||
names:
|
||||
plural: services
|
||||
singular: service
|
||||
kind: Service
|
||||
shortNames: [mservice]
|
||||
versions:
|
||||
- name: v1alpha1
|
||||
served: true
|
||||
storage: true
|
||||
schema:
|
||||
openAPIV3Schema:
|
||||
type: object
|
||||
required: [spec]
|
||||
properties:
|
||||
spec:
|
||||
type: object
|
||||
required: [image]
|
||||
properties:
|
||||
image: {type: string, minLength: 1}
|
||||
command:
|
||||
type: array
|
||||
items: {type: string}
|
||||
args:
|
||||
type: array
|
||||
items: {type: string}
|
||||
replicas: {type: integer, minimum: 0}
|
||||
registry: {type: string}
|
||||
env:
|
||||
type: object
|
||||
additionalProperties: {type: string}
|
||||
`
|
||||
|
||||
const flowCRD = `apiVersion: apiextensions.k8s.io/v1
|
||||
kind: CustomResourceDefinition
|
||||
metadata:
|
||||
name: flows.micro.dev
|
||||
spec:
|
||||
group: micro.dev
|
||||
scope: Namespaced
|
||||
names:
|
||||
plural: flows
|
||||
singular: flow
|
||||
kind: Flow
|
||||
shortNames: [mflow]
|
||||
versions:
|
||||
- name: v1alpha1
|
||||
served: true
|
||||
storage: true
|
||||
schema:
|
||||
openAPIV3Schema:
|
||||
type: object
|
||||
required: [spec]
|
||||
properties:
|
||||
spec:
|
||||
type: object
|
||||
required: [image]
|
||||
properties:
|
||||
image: {type: string, minLength: 1}
|
||||
command:
|
||||
type: array
|
||||
items: {type: string}
|
||||
args:
|
||||
type: array
|
||||
items: {type: string}
|
||||
replicas: {type: integer, minimum: 0}
|
||||
registry: {type: string}
|
||||
env:
|
||||
type: object
|
||||
additionalProperties: {type: string}
|
||||
`
|
||||
@@ -1,138 +0,0 @@
|
||||
package kubernetes
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"sort"
|
||||
"strings"
|
||||
)
|
||||
|
||||
const (
|
||||
// Group is the API group for the alpha Go Micro Kubernetes resources.
|
||||
Group = "micro.dev"
|
||||
// Version is the current alpha API version for the CRDs in this package.
|
||||
Version = "v1alpha1"
|
||||
)
|
||||
|
||||
// Kind identifies a Go Micro lifecycle resource that can be reconciled toward a
|
||||
// Kubernetes Deployment.
|
||||
type Kind string
|
||||
|
||||
const (
|
||||
KindAgent Kind = "Agent"
|
||||
KindService Kind = "Service"
|
||||
KindFlow Kind = "Flow"
|
||||
)
|
||||
|
||||
// WorkloadSpec is the common alpha spec shared by Agent, Service, and Flow CRDs.
|
||||
type WorkloadSpec struct {
|
||||
Image string `json:"image"`
|
||||
Command []string `json:"command,omitempty"`
|
||||
Args []string `json:"args,omitempty"`
|
||||
Replicas int32 `json:"replicas,omitempty"`
|
||||
Registry string `json:"registry,omitempty"`
|
||||
Environment map[string]string `json:"env,omitempty"`
|
||||
}
|
||||
|
||||
// Resource is the minimal desired state for a Go Micro lifecycle resource.
|
||||
type Resource struct {
|
||||
Kind Kind
|
||||
Name string
|
||||
Namespace string
|
||||
Spec WorkloadSpec
|
||||
}
|
||||
|
||||
// Deployment is a small, dependency-free representation of the Kubernetes
|
||||
// Deployment fields the alpha reconciler skeleton owns.
|
||||
type Deployment struct {
|
||||
Name string
|
||||
Namespace string
|
||||
Labels map[string]string
|
||||
Replicas int32
|
||||
Pod PodTemplate
|
||||
}
|
||||
|
||||
// PodTemplate describes the pod fields emitted by MapDeployment.
|
||||
type PodTemplate struct {
|
||||
Labels map[string]string
|
||||
Container Container
|
||||
}
|
||||
|
||||
// Container describes the single Go Micro workload container.
|
||||
type Container struct {
|
||||
Name string
|
||||
Image string
|
||||
Command []string
|
||||
Args []string
|
||||
Environment map[string]string
|
||||
}
|
||||
|
||||
// MapDeployment maps a Go Micro alpha resource to the Deployment shape an
|
||||
// operator reconciliation loop would apply.
|
||||
func MapDeployment(resource Resource) (Deployment, error) {
|
||||
if resource.Kind != KindAgent && resource.Kind != KindService && resource.Kind != KindFlow {
|
||||
return Deployment{}, fmt.Errorf("unsupported kind %q", resource.Kind)
|
||||
}
|
||||
name := strings.TrimSpace(resource.Name)
|
||||
if name == "" {
|
||||
return Deployment{}, fmt.Errorf("name is required")
|
||||
}
|
||||
image := strings.TrimSpace(resource.Spec.Image)
|
||||
if image == "" {
|
||||
return Deployment{}, fmt.Errorf("spec.image is required")
|
||||
}
|
||||
|
||||
namespace := strings.TrimSpace(resource.Namespace)
|
||||
if namespace == "" {
|
||||
namespace = "default"
|
||||
}
|
||||
replicas := resource.Spec.Replicas
|
||||
if replicas == 0 {
|
||||
replicas = 1
|
||||
}
|
||||
|
||||
labels := map[string]string{
|
||||
"app.kubernetes.io/name": name,
|
||||
"app.kubernetes.io/managed-by": "go-micro",
|
||||
"micro.dev/kind": strings.ToLower(string(resource.Kind)),
|
||||
}
|
||||
env := copyMap(resource.Spec.Environment)
|
||||
if resource.Spec.Registry != "" {
|
||||
env["MICRO_REGISTRY"] = resource.Spec.Registry
|
||||
}
|
||||
|
||||
return Deployment{
|
||||
Name: name,
|
||||
Namespace: namespace,
|
||||
Labels: copyMap(labels),
|
||||
Replicas: replicas,
|
||||
Pod: PodTemplate{
|
||||
Labels: copyMap(labels),
|
||||
Container: Container{
|
||||
Name: name,
|
||||
Image: image,
|
||||
Command: append([]string(nil), resource.Spec.Command...),
|
||||
Args: append([]string(nil), resource.Spec.Args...),
|
||||
Environment: env,
|
||||
},
|
||||
},
|
||||
}, nil
|
||||
}
|
||||
|
||||
// EnvironmentKeys returns stable environment variable keys from a mapped
|
||||
// container. It is useful for deterministic validation and rendering.
|
||||
func (c Container) EnvironmentKeys() []string {
|
||||
keys := make([]string, 0, len(c.Environment))
|
||||
for key := range c.Environment {
|
||||
keys = append(keys, key)
|
||||
}
|
||||
sort.Strings(keys)
|
||||
return keys
|
||||
}
|
||||
|
||||
func copyMap(in map[string]string) map[string]string {
|
||||
out := make(map[string]string, len(in))
|
||||
for k, v := range in {
|
||||
out[k] = v
|
||||
}
|
||||
return out
|
||||
}
|
||||
+1
-4
@@ -12,7 +12,6 @@ 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/) |
|
||||
|
||||
@@ -29,9 +28,7 @@ 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. **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
|
||||
5. **Workflow depth:** use [`flow-durable`](./flow-durable/) once the agent path
|
||||
needs checkpointed, resumable deterministic work.
|
||||
|
||||
## CLI wayfinding
|
||||
|
||||
@@ -1,24 +0,0 @@
|
||||
# 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.
|
||||
@@ -1,168 +0,0 @@
|
||||
// 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)
|
||||
}
|
||||
}
|
||||
@@ -1,282 +0,0 @@
|
||||
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[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)
|
||||
}
|
||||
@@ -1,62 +0,0 @@
|
||||
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,11 +157,6 @@ 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
|
||||
@@ -291,10 +286,6 @@ 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)
|
||||
|
||||
+19
-5
@@ -294,9 +294,7 @@ func (t *StdioTransport) handleToolsCall(req *JSONRPCRequest) {
|
||||
AccountID: accountID, ScopesRequired: tool.Scopes,
|
||||
Allowed: true, Duration: time.Since(start), Error: err.Error(),
|
||||
})
|
||||
// A tool-execution failure is reported as an isError result, not a
|
||||
// JSON-RPC protocol error (per the MCP spec), so the agent can read it.
|
||||
t.sendResponse(req.ID, mcpToolError(traceID, "tool call failed: "+err.Error()))
|
||||
t.sendError(req.ID, InternalError, "RPC call failed", err.Error())
|
||||
return
|
||||
}
|
||||
|
||||
@@ -313,8 +311,24 @@ func (t *StdioTransport) handleToolsCall(req *JSONRPCRequest) {
|
||||
Allowed: true, Duration: time.Since(start),
|
||||
})
|
||||
|
||||
// The downstream response is JSON — return it as JSON text, not %v.
|
||||
t.sendResponse(req.ID, mcpToolResult(traceID, rsp.Data))
|
||||
// Parse response
|
||||
var result interface{}
|
||||
if err := json.Unmarshal(rsp.Data, &result); err != nil {
|
||||
// If unmarshal fails, return raw data
|
||||
result = map[string]interface{}{
|
||||
"data": string(rsp.Data),
|
||||
}
|
||||
}
|
||||
|
||||
t.sendResponse(req.ID, map[string]interface{}{
|
||||
"content": []interface{}{
|
||||
map[string]interface{}{
|
||||
"type": "text",
|
||||
"text": fmt.Sprintf("%v", result),
|
||||
},
|
||||
},
|
||||
"trace_id": traceID,
|
||||
})
|
||||
}
|
||||
|
||||
// sendResponse sends a JSON-RPC response
|
||||
|
||||
@@ -1,120 +0,0 @@
|
||||
package mcp
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"testing"
|
||||
|
||||
"go-micro.dev/v6/client"
|
||||
)
|
||||
|
||||
// fakeCallClient overrides Call to return canned data or an error; NewRequest
|
||||
// and the rest are promoted from the embedded real client.
|
||||
type fakeCallClient struct {
|
||||
client.Client
|
||||
data []byte
|
||||
err error
|
||||
}
|
||||
|
||||
func (f *fakeCallClient) Call(ctx context.Context, req client.Request, rsp interface{}, opts ...client.CallOption) error {
|
||||
if f.err != nil {
|
||||
return f.err
|
||||
}
|
||||
if r, ok := rsp.(*struct{ Data []byte }); ok {
|
||||
r.Data = f.data
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// isToolError reports whether an MCP tools/call result carries isError:true.
|
||||
func isToolError(result interface{}) bool {
|
||||
m, ok := result.(map[string]interface{})
|
||||
if !ok {
|
||||
return false
|
||||
}
|
||||
b, _ := m["isError"].(bool)
|
||||
return b
|
||||
}
|
||||
|
||||
// toolResultText extracts the first text content of an MCP tools/call result.
|
||||
func toolResultText(t *testing.T, result interface{}) string {
|
||||
t.Helper()
|
||||
m, ok := result.(map[string]interface{})
|
||||
if !ok {
|
||||
t.Fatalf("result is not a map: %#v", result)
|
||||
}
|
||||
content, ok := m["content"].([]interface{})
|
||||
if !ok || len(content) == 0 {
|
||||
t.Fatalf("result has no content: %#v", result)
|
||||
}
|
||||
first, _ := content[0].(map[string]interface{})
|
||||
text, _ := first["text"].(string)
|
||||
return text
|
||||
}
|
||||
|
||||
// driveStdio sends one JSON-RPC request through a StdioTransport and returns the
|
||||
// decoded response, capturing the transport's stdout into a buffer.
|
||||
func driveStdio(t *testing.T, s *Server, method string, id interface{}, params interface{}) JSONRPCResponse {
|
||||
t.Helper()
|
||||
tr := NewStdioTransport(s)
|
||||
var out bytes.Buffer
|
||||
tr.writer = bufio.NewWriter(&out)
|
||||
raw, _ := json.Marshal(params)
|
||||
tr.handleRequest(&JSONRPCRequest{JSONRPC: "2.0", ID: id, Method: method, Params: raw})
|
||||
var resp JSONRPCResponse
|
||||
if err := json.Unmarshal(bytes.TrimSpace(out.Bytes()), &resp); err != nil {
|
||||
t.Fatalf("decode stdio response: %v (raw=%q)", err, out.String())
|
||||
}
|
||||
return resp
|
||||
}
|
||||
|
||||
// The stdio transport is the path an external MCP host (Claude Desktop) uses.
|
||||
// It must return tool output as JSON text, not fmt.Sprintf("%v", ...) which
|
||||
// yields Go map-syntax and is unparseable by a real client.
|
||||
func TestStdio_ToolsCall_ReturnsJSONNotGoSyntax(t *testing.T) {
|
||||
s := newTestServer(Options{})
|
||||
s.opts.Client = &fakeCallClient{Client: client.DefaultClient, data: []byte(`{"id":1,"name":"bob"}`)}
|
||||
s.tools["svc.Echo"] = &Tool{Name: "svc.Echo", Service: "svc", Endpoint: "Echo"}
|
||||
|
||||
resp := driveStdio(t, s, "tools/call", 1, map[string]interface{}{
|
||||
"name": "svc.Echo",
|
||||
"arguments": map[string]interface{}{"msg": "hi"},
|
||||
})
|
||||
if resp.Error != nil {
|
||||
t.Fatalf("unexpected protocol error: %+v", resp.Error)
|
||||
}
|
||||
text := toolResultText(t, resp.Result)
|
||||
// The bug returned Go map-syntax ("map[id:1 name:bob]"), which fails to parse.
|
||||
var got map[string]interface{}
|
||||
if err := json.Unmarshal([]byte(text), &got); err != nil {
|
||||
t.Fatalf("tool result text is not JSON (the %%v bug): %q", text)
|
||||
}
|
||||
if got["name"] != "bob" {
|
||||
t.Errorf("result = %v, want name=bob", got)
|
||||
}
|
||||
}
|
||||
|
||||
// A tool-execution failure must be an MCP isError result, not a JSON-RPC
|
||||
// protocol error, so the agent can read the failure.
|
||||
func TestStdio_ToolsCall_FailureIsIsErrorResult(t *testing.T) {
|
||||
s := newTestServer(Options{})
|
||||
s.opts.Client = &fakeCallClient{Client: client.DefaultClient, err: errors.New("backend down")}
|
||||
s.tools["svc.Echo"] = &Tool{Name: "svc.Echo", Service: "svc", Endpoint: "Echo"}
|
||||
|
||||
resp := driveStdio(t, s, "tools/call", 1, map[string]interface{}{
|
||||
"name": "svc.Echo",
|
||||
"arguments": map[string]interface{}{},
|
||||
})
|
||||
if resp.Error != nil {
|
||||
t.Fatalf("tool failure returned a protocol error, want isError result: %+v", resp.Error)
|
||||
}
|
||||
if !isToolError(resp.Result) {
|
||||
t.Fatalf("expected isError result, got %+v", resp.Result)
|
||||
}
|
||||
if text := toolResultText(t, resp.Result); text == "" {
|
||||
t.Error("isError result should carry the error text")
|
||||
}
|
||||
}
|
||||
@@ -1,32 +0,0 @@
|
||||
package mcp
|
||||
|
||||
// MCP tools/call result shaping, shared by the stdio and websocket JSON-RPC
|
||||
// transports. Kept in one place so both transports produce spec-shaped results.
|
||||
|
||||
// mcpToolResult builds a successful MCP tools/call result. The downstream RPC
|
||||
// response body (data) is JSON, so it is returned as JSON text — NOT
|
||||
// fmt.Sprintf("%v", ...) of a decoded value, which produces Go map-syntax
|
||||
// (map[id:1 name:bob]) instead of JSON and is what an external MCP client
|
||||
// (e.g. Claude Desktop over stdio) would otherwise receive.
|
||||
func mcpToolResult(traceID string, data []byte) map[string]interface{} {
|
||||
return map[string]interface{}{
|
||||
"content": []interface{}{
|
||||
map[string]interface{}{"type": "text", "text": string(data)},
|
||||
},
|
||||
"trace_id": traceID,
|
||||
}
|
||||
}
|
||||
|
||||
// mcpToolError builds an MCP tools/call result for a tool-EXECUTION failure.
|
||||
// Per the MCP spec a tool that fails returns a normal result with isError:true
|
||||
// (the error as text content), NOT a JSON-RPC protocol error — that way the
|
||||
// agent can read the failure instead of seeing a transport-level error.
|
||||
func mcpToolError(traceID, msg string) map[string]interface{} {
|
||||
return map[string]interface{}{
|
||||
"content": []interface{}{
|
||||
map[string]interface{}{"type": "text", "text": msg},
|
||||
},
|
||||
"isError": true,
|
||||
"trace_id": traceID,
|
||||
}
|
||||
}
|
||||
@@ -275,8 +275,7 @@ func (wc *wsConn) handleToolsCall(req *JSONRPCRequest) {
|
||||
AccountID: accountID, ScopesRequired: tool.Scopes,
|
||||
Allowed: true, Duration: time.Since(start), Error: err.Error(),
|
||||
})
|
||||
// Tool-execution failure → isError result (MCP spec), not a protocol error.
|
||||
wc.sendResponse(req.ID, mcpToolError(traceID, "tool call failed: "+err.Error()))
|
||||
wc.sendError(req.ID, InternalError, "RPC call failed", err.Error())
|
||||
return
|
||||
}
|
||||
|
||||
@@ -292,8 +291,23 @@ func (wc *wsConn) handleToolsCall(req *JSONRPCRequest) {
|
||||
Allowed: true, Duration: time.Since(start),
|
||||
})
|
||||
|
||||
// The downstream response is JSON — return it as JSON text, not %v.
|
||||
wc.sendResponse(req.ID, mcpToolResult(traceID, rsp.Data))
|
||||
// Parse response
|
||||
var result interface{}
|
||||
if err := json.Unmarshal(rsp.Data, &result); err != nil {
|
||||
result = map[string]interface{}{
|
||||
"data": string(rsp.Data),
|
||||
}
|
||||
}
|
||||
|
||||
wc.sendResponse(req.ID, map[string]interface{}{
|
||||
"content": []interface{}{
|
||||
map[string]interface{}{
|
||||
"type": "text",
|
||||
"text": fmt.Sprintf("%v", result),
|
||||
},
|
||||
},
|
||||
"trace_id": traceID,
|
||||
})
|
||||
}
|
||||
|
||||
// sendResponse sends a JSON-RPC success response.
|
||||
|
||||
@@ -114,13 +114,12 @@ func TestWebSocket_ToolsCall_NoAuth(t *testing.T) {
|
||||
"arguments": map[string]interface{}{"msg": "hi"},
|
||||
})
|
||||
|
||||
// No auth required → the tool runs; the RPC fails (no backend), which the
|
||||
// MCP spec surfaces as an isError result, not a JSON-RPC protocol error.
|
||||
if resp.Error != nil {
|
||||
t.Fatalf("expected no protocol error, got %+v", resp.Error)
|
||||
// RPC will fail (no backend), but auth should pass (no auth configured)
|
||||
if resp.Error == nil {
|
||||
t.Fatal("expected RPC error (no backend)")
|
||||
}
|
||||
if !isToolError(resp.Result) {
|
||||
t.Fatalf("expected isError tool result, got %+v", resp.Result)
|
||||
if resp.Error.Code != InternalError {
|
||||
t.Errorf("error code = %d, want %d", resp.Error.Code, InternalError)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -169,13 +168,12 @@ func TestWebSocket_ToolsCall_AuthRequired(t *testing.T) {
|
||||
"arguments": map[string]interface{}{},
|
||||
"_token": "valid-token",
|
||||
})
|
||||
// Auth passes → the tool runs; RPC fails (no backend) → isError result,
|
||||
// not a JSON-RPC protocol error (which would mean auth failed).
|
||||
if resp.Error != nil {
|
||||
t.Fatalf("expected no protocol error (auth passed), got %+v", resp.Error)
|
||||
// Auth passes, RPC fails (no backend)
|
||||
if resp.Error == nil {
|
||||
t.Fatal("expected RPC error")
|
||||
}
|
||||
if !isToolError(resp.Result) {
|
||||
t.Fatalf("expected isError tool result, got %+v", resp.Result)
|
||||
if resp.Error.Code != InternalError {
|
||||
t.Errorf("error code = %d, want %d (RPC fail, not auth fail)", resp.Error.Code, InternalError)
|
||||
}
|
||||
})
|
||||
|
||||
@@ -187,13 +185,12 @@ func TestWebSocket_ToolsCall_AuthRequired(t *testing.T) {
|
||||
"name": "svc.Do",
|
||||
"arguments": map[string]interface{}{},
|
||||
})
|
||||
// Auth passes via connection-level header → tool runs; RPC fails (no
|
||||
// backend) → isError result, not a JSON-RPC protocol error.
|
||||
if resp.Error != nil {
|
||||
t.Fatalf("expected no protocol error (auth passed), got %+v", resp.Error)
|
||||
// Auth passes via connection-level header, RPC fails (no backend)
|
||||
if resp.Error == nil {
|
||||
t.Fatal("expected RPC error")
|
||||
}
|
||||
if !isToolError(resp.Result) {
|
||||
t.Fatalf("expected isError tool result, got %+v", resp.Result)
|
||||
if resp.Error.Code != InternalError {
|
||||
t.Errorf("error code = %d, want %d (RPC fail, not auth fail)", resp.Error.Code, InternalError)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
@@ -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 | Yes | No |
|
||||
| `gemini` | Yes | No | No | No | 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 | ✅ Verified when configured | ⚠️ Unverified |
|
||||
| `gemini` | ✅ Verified when configured | — Unsupported | — Unsupported | ⚠️ Unverified | ⚠️ 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,7 +12,6 @@ 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"
|
||||
)
|
||||
|
||||
@@ -131,13 +130,6 @@ 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) }
|
||||
|
||||
|
||||
+1
-11
@@ -8,7 +8,6 @@ import (
|
||||
"io"
|
||||
"net/http"
|
||||
"strconv"
|
||||
"strings"
|
||||
"sync"
|
||||
)
|
||||
|
||||
@@ -82,16 +81,7 @@ func (c *Client) Do(req *http.Request) (*http.Response, error) {
|
||||
return resp, fmt.Errorf("x402: 402 response carried no requirements")
|
||||
}
|
||||
reqd := ch.Accepts[0]
|
||||
// The amount governs the whole spend cap, so it must be a real positive
|
||||
// integer. A swallowed parse error (non-decimal, overflow, empty) would
|
||||
// yield 0 and pass the budget check trivially, and a negative amount would
|
||||
// inflate the remaining allowance — either way the cap is defeated. Refuse
|
||||
// before signing anything.
|
||||
amount, err := strconv.ParseInt(strings.TrimSpace(reqd.MaxAmountRequired), 10, 64)
|
||||
if err != nil || amount <= 0 {
|
||||
return resp, fmt.Errorf("x402: refusing to pay %s: invalid maxAmountRequired %q",
|
||||
reqd.Resource, reqd.MaxAmountRequired)
|
||||
}
|
||||
amount, _ := strconv.ParseInt(reqd.MaxAmountRequired, 10, 64)
|
||||
|
||||
// Spend cap: reserve before paying so concurrent calls cannot all pass
|
||||
// the check and overspend the caller's allowance. Roll the reservation
|
||||
|
||||
@@ -172,33 +172,6 @@ func TestClientBudgetReservationRollsBackOnPayError(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// A 402 whose maxAmountRequired is not a positive integer must be refused
|
||||
// before any payment — otherwise a swallowed parse error (0) or a negative
|
||||
// amount defeats the spend cap. The payer is never called and nothing is spent.
|
||||
func TestClientRefusesInvalidAmount(t *testing.T) {
|
||||
for _, amount := range []string{"abc", "-100", "99999999999999999999999999", "0x10", "1.5"} {
|
||||
srv := paidServer(amount)
|
||||
|
||||
payer := &mockPayer{}
|
||||
c := &Client{Payer: payer, Budget: 1_000_000}
|
||||
req, _ := http.NewRequest(http.MethodGet, srv.URL, nil)
|
||||
resp, err := c.Do(req)
|
||||
if resp != nil {
|
||||
resp.Body.Close()
|
||||
}
|
||||
if err == nil {
|
||||
t.Errorf("amount %q: expected refusal, got nil error", amount)
|
||||
}
|
||||
if payer.calls != 0 {
|
||||
t.Errorf("amount %q: payer called %d times, want 0", amount, payer.calls)
|
||||
}
|
||||
if c.Spent() != 0 {
|
||||
t.Errorf("amount %q: spent %d, want 0", amount, c.Spent())
|
||||
}
|
||||
srv.Close()
|
||||
}
|
||||
}
|
||||
|
||||
type payerFunc func(context.Context, Requirements) (string, error)
|
||||
|
||||
func (f payerFunc) Pay(ctx context.Context, req Requirements) (string, error) {
|
||||
|
||||
+1
-13
@@ -126,11 +126,6 @@ type Config struct {
|
||||
// FacilitatorURL is the verify/settle endpoint used when Facilitator
|
||||
// is nil (e.g. Coinbase CDP or Alchemy).
|
||||
FacilitatorURL string `json:"facilitator,omitempty"`
|
||||
// RequireSettlement fails closed when a paid request cannot be settled:
|
||||
// if the facilitator only verifies (does not implement Settler), Require
|
||||
// refuses to serve rather than releasing the resource while no funds move.
|
||||
// Leave false only for verify-only flows where authorization is enough.
|
||||
RequireSettlement bool `json:"requireSettlement,omitempty"`
|
||||
}
|
||||
|
||||
func (c Config) network() string {
|
||||
@@ -227,14 +222,7 @@ func (c Config) Require(w http.ResponseWriter, r *http.Request, amount, resource
|
||||
}
|
||||
// Capture the funds when the facilitator can settle. Verify alone only
|
||||
// authorizes the "exact" transfer; settlement broadcasts it.
|
||||
s, canSettle := fac.(Settler)
|
||||
if c.RequireSettlement && !canSettle {
|
||||
// Fail closed: a paid config must not serve the resource on a
|
||||
// verify-only facilitator, or it gives the tool away for free.
|
||||
writeChallenge(w, req, "payment settlement unavailable")
|
||||
return false
|
||||
}
|
||||
if canSettle {
|
||||
if s, ok := fac.(Settler); ok {
|
||||
sres, err := s.Settle(r.Context(), payment, req)
|
||||
if err != nil {
|
||||
writeChallenge(w, req, "payment settlement failed: "+err.Error())
|
||||
|
||||
@@ -159,58 +159,4 @@ func TestCDPAuthorizeAttachesBearer(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// TestRequireSettlementFailsClosed checks that a paid config with
|
||||
// RequireSettlement refuses to serve when the facilitator only verifies (does
|
||||
// not settle) — otherwise the resource is released while no funds move.
|
||||
func TestRequireSettlementFailsClosed(t *testing.T) {
|
||||
// mockFacilitator implements Verify but not Settler.
|
||||
cfg := Config{PayTo: "0xpay", Facilitator: mockFacilitator{valid: true}, RequireSettlement: true}
|
||||
r := httptest.NewRequest(http.MethodGet, "/tool", nil)
|
||||
r.Header.Set(PaymentHeader, "eyJ4IjoxfQ==")
|
||||
rec := httptest.NewRecorder()
|
||||
|
||||
if cfg.Require(rec, r, "10000", "chat") {
|
||||
t.Fatal("Require should fail closed when settlement is required but unavailable")
|
||||
}
|
||||
if rec.Code != http.StatusPaymentRequired {
|
||||
t.Errorf("status = %d, want 402", rec.Code)
|
||||
}
|
||||
|
||||
// Without RequireSettlement the verify-only facilitator still serves.
|
||||
cfg.RequireSettlement = false
|
||||
rec = httptest.NewRecorder()
|
||||
r = httptest.NewRequest(http.MethodGet, "/tool", nil)
|
||||
r.Header.Set(PaymentHeader, "eyJ4IjoxfQ==")
|
||||
if !cfg.Require(rec, r, "10000", "chat") {
|
||||
t.Fatalf("verify-only should serve when settlement is not required; body=%s", rec.Body.String())
|
||||
}
|
||||
}
|
||||
|
||||
// TestRequireSettlementServesWithSettler checks that a paid config with
|
||||
// RequireSettlement serves when the facilitator can settle.
|
||||
func TestRequireSettlementServesWithSettler(t *testing.T) {
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
switch r.URL.Path {
|
||||
case "/verify":
|
||||
_ = json.NewEncoder(w).Encode(map[string]any{"isValid": true})
|
||||
case "/settle":
|
||||
_ = json.NewEncoder(w).Encode(map[string]any{"success": true, "transaction": "0xabc"})
|
||||
}
|
||||
}))
|
||||
defer srv.Close()
|
||||
|
||||
// HTTPFacilitator implements Settler.
|
||||
cfg := Config{PayTo: "0xpay", FacilitatorURL: srv.URL, RequireSettlement: true}
|
||||
r := httptest.NewRequest(http.MethodGet, "/tool", nil)
|
||||
r.Header.Set(PaymentHeader, "eyJ4IjoxfQ==")
|
||||
rec := httptest.NewRecorder()
|
||||
|
||||
if !cfg.Require(rec, r, "10000", "chat") {
|
||||
t.Fatalf("Require should serve with a settling facilitator; body=%s", rec.Body.String())
|
||||
}
|
||||
if got := rec.Header().Get(PaymentResponseHeader); got != "0xabc" {
|
||||
t.Errorf("settlement header = %q, want 0xabc", got)
|
||||
}
|
||||
}
|
||||
|
||||
var _ Settler = (*HTTPFacilitator)(nil)
|
||||
|
||||
Reference in New Issue
Block a user