Compare commits

..

12 Commits

Author SHA1 Message Date
Codex 530739baaf loop: refresh priorities after shipped capability
govulncheck / govulncheck (push) Waiting to run
Harness (E2E) / Harnesses (mock LLM) (push) Waiting to run
Harness (E2E) / Provider harnesses (live LLM conformance) (push) Waiting to run
Lint / golangci-lint (push) Waiting to run
Run Tests / Unit Tests (push) Waiting to run
Run Tests / Etcd Integration Tests (push) Waiting to run
2026-07-12 12:26:57 +00:00
Asim Aslam b5df7e0a71 Add Kubernetes CRD foundation (#4839)
Co-authored-by: Codex <codex@openai.com>
2026-07-12 13:01:09 +01:00
Asim Aslam 7e3d2d3b13 x402: harden spend cap — reject invalid amounts, require settler option (#4831)
Two spend-safety fixes from the gap audit (#4814):

- Client.Do refused a 402 only on the budget check, but parsed
  maxAmountRequired with a swallowed error, so a non-decimal, overflowing
  or negative amount became 0 and passed the cap trivially while Payer.Pay
  still signed against the string. Now reject any amount that is not a
  positive integer before signing.

- Require settled only when the facilitator implemented Settler; a
  verify-only facilitator served the resource while no funds moved. Add
  Config.RequireSettlement to fail closed in that case.

Tests cover invalid/negative/overflow amounts and the verify-only
fail-closed path.


Claude-Session: https://claude.ai/code/session_01CmdEY7pYmV5zzwCjNJ4ykL

Co-authored-by: Claude <noreply@anthropic.com>
2026-07-12 12:40:12 +01:00
Asim Aslam 3e4b13e2bd Fix grpcreflect JSON name lint (#4826)
* gateway/mcp: expose reflected gRPC services

* Fix grpcreflect JSON name lint

---------

Co-authored-by: Codex <codex@openai.com>
2026-07-12 11:28:53 +00:00
Asim Aslam 1b83cdff9c mcp: stdio/ws tool results are JSON + isError (fixes garbage to Claude Desktop) (#4825)
Closes #4813. The stdio transport is the path an external MCP host (Claude
Desktop) uses, and it emitted broken output:
- tool results were `fmt.Sprintf("%v", decodedJSON)` → Go map-syntax
  (`map[id:1 name:bob]`), not JSON. Now returned as JSON text.
- tool-execution failures were returned as JSON-RPC protocol errors; per the
  MCP spec they must be a result with `isError:true` so the agent can read the
  failure. Now they are (span/audit still record the error).

Both fixes are shared between stdio and websocket via a new `mcpToolResult`/
`mcpToolError` (dedupes the two transports). Added the missing stdio round-trip
tests (the package had zero) proving JSON output and the isError contract, using
an injected fake client; updated the websocket auth tests that asserted the old
protocol-error-on-tool-failure behavior.

Also fixes a pre-existing golangci-lint failure on master (unnecessary
`string(...)` conversion in grpcreflect.go from #4821) so the mcp package lints
clean — another one the required-checks gap let through.


Claude-Session: https://claude.ai/code/session_01CmdEY7pYmV5zzwCjNJ4ykL

Co-authored-by: Claude <noreply@anthropic.com>
2026-07-12 12:28:06 +01:00
Asim Aslam 3ef265f3c2 loop: refresh planner priorities after gRPC MCP (#4828)
Co-authored-by: Codex <codex@openai.com>
2026-07-12 12:27:04 +01:00
Asim Aslam e8977cf335 gateway/mcp: expose reflected gRPC services (#4821)
Co-authored-by: Codex <codex@openai.com>
2026-07-12 12:09:31 +01:00
Asim Aslam c6ab16f3bf loop: drop shipped x402 buyer priority (#4818)
Co-authored-by: Codex <codex@openai.com>
2026-07-12 11:44:34 +01:00
Asim Aslam f93f3c6045 examples: add x402 buyer agent (#4811)
Co-authored-by: Codex <codex@openai.com>
2026-07-12 11:20:25 +01:00
Asim Aslam 9b4b3ce827 loop: drop completed spend observability priority (#4808)
Co-authored-by: Codex <codex@openai.com>
2026-07-12 10:36:04 +01:00
Asim Aslam 4d6ebe1fd3 Observe agent x402 spend (#4806)
Co-authored-by: Codex <codex@openai.com>
2026-07-12 10:23:30 +01:00
Asim Aslam 26ab5a3bf0 loop: drop completed x402 buyer priority (#4804)
Co-authored-by: Codex <codex@openai.com>
2026-07-12 09:53:05 +01:00
32 changed files with 1452 additions and 66 deletions
+9 -5
View File
@@ -24,14 +24,18 @@ rewrites.
### Capability — the headline (roadmap: Now / Next)
1. **Agent spend observability** ([#4787](https://github.com/micro/go-micro/issues/4787)) — surface x402 spend in `RunInfo` and OpenTelemetry so payments are inspectable like every other agent action. (Follows #4786.)
2. **Example: an agent that pays for a paid tool** ([#4788](https://github.com/micro/go-micro/issues/4788)) — the runnable artifact that makes it real for a developer, against a mock facilitator (no live funds). (Follows #4786/#4787.)
3. **gRPC-reflection MCP** ([#4796](https://github.com/micro/go-micro/issues/4796)) — expose external reflected gRPC services as MCP tools, not only go-micro-native handlers. A large jump in what agents can operate without requiring teams to rewrite existing services.
4. **Kubernetes operator + CRDs foundation** ([#4797](https://github.com/micro/go-micro/issues/4797)) — add the first opt-in `Agent`, `Service`, and `Flow` resource foundation so the services → agents → workflows lifecycle has a native deployment path for Kubernetes users.
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, and provider-failure inspection. Further churn in those
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._
+3
View File
@@ -29,6 +29,9 @@ below is kept current between tags and rolled into the next version when it ship
### 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
+5
View File
@@ -451,6 +451,11 @@ func (a *agentImpl) spendWrap(next ai.ToolHandler) ai.ToolHandler {
amount, call.Name, a.spend, a.opts.MaxSpend))
}
a.spend += amount
if info, ok := ai.RunInfoFrom(ctx); ok {
info.Spent = a.spend
info.ToolSpend = amount
ctx = ai.WithRunInfo(ctx, info)
}
res := next(ctx, call)
if res.Refused != "" || toolErrorMessage(res) != "" {
a.spend -= amount
+11 -2
View File
@@ -211,7 +211,8 @@ func TestAgentPayerPaysX402ToolResultAndRetries(t *testing.T) {
defer srv.Close()
payer := &agentMockPayer{}
a := newTestAgent(Name("x402-payer"), Payer(payer), Budget(10), WithTool("paid.http", "paid http", nil, func(ctx context.Context, input map[string]any) (string, error) {
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
@@ -228,13 +229,21 @@ func TestAgentPayerPaysX402ToolResultAndRetries(t *testing.T) {
return string(body), nil
}))
res := a.toolHandler()(context.Background(), ai.ToolCall{ID: "pay-1", Name: "paid.http", Input: map[string]any{"url": srv.URL}})
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) {
+27 -2
View File
@@ -51,6 +51,8 @@ const (
AttrDispatch = "agent.dispatch"
AttrTrigger = "agent.trigger"
AttrRunEventKind = "agent.event.kind"
AttrSpend = "agent.spend"
AttrToolSpend = "agent.tool.spend"
)
type RunEvent struct {
@@ -73,6 +75,8 @@ type RunEvent struct {
Error string `json:"error,omitempty"`
ErrorKind string `json:"error_kind,omitempty"`
InputChars int `json:"input_chars,omitempty"`
Spent int64 `json:"spent,omitempty"`
ToolSpend int64 `json:"tool_spend,omitempty"`
}
type Usage = ai.Usage
@@ -111,6 +115,7 @@ type RunSummary struct {
LastKind string `json:"last_kind,omitempty"`
LastError string `json:"last_error,omitempty"`
LastErrorKind string `json:"last_error_kind,omitempty"`
Spent int64 `json:"spent,omitempty"`
}
func (a *agentImpl) tracer() trace.Tracer {
@@ -362,6 +367,7 @@ func (a *agentImpl) traceTool(next ai.ToolHandler) ai.ToolHandler {
return func(ctx context.Context, call ai.ToolCall) ai.ToolResult {
info, _ := ai.RunInfoFrom(ctx)
start := time.Now()
spentBefore := a.spend
if a.opts.TraceProvider == nil {
res := next(ctx, call)
@@ -371,7 +377,7 @@ func (a *agentImpl) traceTool(next ai.ToolHandler) ai.ToolHandler {
if toolAttempts <= 0 {
toolAttempts = 1
}
a.recordRunEvent(RunEvent{Time: time.Now(), RunID: info.RunID, ParentID: info.ParentID, Agent: info.Agent, Kind: "tool", Name: call.Name, Attempt: toolAttempts, MaxAttempts: a.opts.ToolMaxAttempts, LatencyMS: dur, Refused: res.Refused, Error: resErr, ErrorKind: classifyToolError(resErr)})
a.recordRunEvent(RunEvent{Time: time.Now(), RunID: info.RunID, ParentID: info.ParentID, Agent: info.Agent, Kind: "tool", Name: call.Name, Attempt: toolAttempts, MaxAttempts: a.opts.ToolMaxAttempts, LatencyMS: dur, Refused: res.Refused, Error: resErr, ErrorKind: classifyToolError(resErr), Spent: a.spend, ToolSpend: a.spend - spentBefore})
return res
}
@@ -385,6 +391,7 @@ func (a *agentImpl) traceTool(next ai.ToolHandler) ai.ToolHandler {
ctx, span := a.tracer().Start(ctx, spanNameToolCall, trace.WithAttributes(spanAttrs...))
res := next(ctx, call)
dur := time.Since(start).Milliseconds()
toolSpend := a.spend - spentBefore
attrs := []attribute.KeyValue{attribute.Int64(AttrLatencyMS, dur)}
toolAttempts := res.Attempts
if toolAttempts <= 0 {
@@ -394,6 +401,9 @@ func (a *agentImpl) traceTool(next ai.ToolHandler) ai.ToolHandler {
if a.opts.ToolMaxAttempts > 0 {
attrs = append(attrs, attribute.Int(AttrToolMaxAttempts, a.opts.ToolMaxAttempts))
}
if toolSpend > 0 {
attrs = append(attrs, attribute.Int64(AttrSpend, a.spend), attribute.Int64(AttrToolSpend, toolSpend))
}
if res.Refused != "" {
attrs = append(attrs, attribute.Bool(AttrGuardrailBlock, true), attribute.String(AttrRefusal, res.Refused))
}
@@ -409,7 +419,7 @@ func (a *agentImpl) traceTool(next ai.ToolHandler) ai.ToolHandler {
} else {
span.SetStatus(codes.Ok, "")
}
a.recordSpanEvent(span, RunEvent{Time: time.Now(), RunID: info.RunID, ParentID: info.ParentID, Agent: info.Agent, Kind: "tool", Name: call.Name, Attempt: toolAttempts, MaxAttempts: a.opts.ToolMaxAttempts, LatencyMS: dur, Refused: res.Refused, Error: resErr, ErrorKind: classifyToolError(resErr)})
a.recordSpanEvent(span, RunEvent{Time: time.Now(), RunID: info.RunID, ParentID: info.ParentID, Agent: info.Agent, Kind: "tool", Name: call.Name, Attempt: toolAttempts, MaxAttempts: a.opts.ToolMaxAttempts, LatencyMS: dur, Refused: res.Refused, Error: resErr, ErrorKind: classifyToolError(resErr), Spent: a.spend, ToolSpend: toolSpend})
span.End()
return res
}
@@ -496,6 +506,12 @@ func runEventAttributes(e RunEvent) []attribute.KeyValue {
if e.InputChars > 0 {
attrs = append(attrs, attribute.Int(AttrInputChars, e.InputChars))
}
if e.Spent > 0 {
attrs = append(attrs, attribute.Int64(AttrSpend, e.Spent))
}
if e.ToolSpend > 0 {
attrs = append(attrs, attribute.Int64(AttrToolSpend, e.ToolSpend))
}
attrs = appendUsage(attrs, e.Tokens)
if e.Refused != "" {
attrs = append(attrs, attribute.Bool(AttrGuardrailBlock, true), attribute.String(AttrRefusal, e.Refused))
@@ -530,6 +546,12 @@ func appendRunInfoAttributes(attrs []attribute.KeyValue, info ai.RunInfo) []attr
if info.Trigger != "" {
attrs = append(attrs, attribute.String(AttrTrigger, info.Trigger))
}
if info.Spent > 0 {
attrs = append(attrs, attribute.Int64(AttrSpend, info.Spent))
}
if info.ToolSpend > 0 {
attrs = append(attrs, attribute.Int64(AttrToolSpend, info.ToolSpend))
}
return attrs
}
@@ -616,6 +638,9 @@ func ListRunSummariesWithOptions(s store.Store, agentName string, opts RunListOp
if e.ErrorKind != "" {
summary.LastErrorKind = e.ErrorKind
}
if e.Spent > summary.Spent {
summary.Spent = e.Spent
}
}
if opts.Status != "" && summary.Status != opts.Status {
continue
+57
View File
@@ -421,6 +421,63 @@ func spanAttributes(attrs []attribute.KeyValue) map[string]string {
return out
}
func TestAgentOpenTelemetryToolSpanIncludesSpend(t *testing.T) {
exp := tracetest.NewInMemoryExporter()
tp := trace.NewTracerProvider(trace.WithSyncer(exp))
st := store.NewMemoryStore()
a := New(Name("spender"), Provider("oteltest"), Model("unit-model"), WithStore(st), TraceProvider(tp), MaxSpend(10), ToolSpend("probe", 7), WithTool("probe", "probe", nil, func(ctx context.Context, input map[string]any) (string, error) {
info, ok := ai.RunInfoFrom(ctx)
if !ok {
t.Fatal("RunInfo missing from paid tool context")
}
if info.Spent != 7 || info.ToolSpend != 7 {
t.Fatalf("RunInfo spend = (%d, %d), want (7, 7)", info.Spent, info.ToolSpend)
}
return "ok", nil
}))
if _, err := a.Ask(context.Background(), "hello"); err != nil {
t.Fatal(err)
}
var sawToolSpan bool
for _, s := range exp.GetSpans().Snapshots() {
if s.Name() != spanNameToolCall {
continue
}
sawToolSpan = true
attrs := spanAttributes(s.Attributes())
if attrs[AttrSpend] != "7" || attrs[AttrToolSpend] != "7" {
t.Fatalf("tool span missing spend attributes: %#v", attrs)
}
if !spanEventHasAttribute(s.Events(), "agent.tool", AttrToolSpend, "7") {
t.Fatalf("tool event missing spend attribute: %#v", s.Events())
}
}
if !sawToolSpan {
t.Fatal("tool span not emitted")
}
summaries, err := ListRunSummaries(st, "spender")
if err != nil {
t.Fatal(err)
}
if len(summaries) != 1 || summaries[0].Spent != 7 {
t.Fatalf("summary spend = %#v, want 7", summaries)
}
}
func spanEventHasAttribute(events []trace.Event, name, key, value string) bool {
for _, e := range events {
if e.Name != name {
continue
}
attrs := spanAttributes(e.Attributes)
if attrs[key] == value {
return true
}
}
return false
}
func TestAgentOpenTelemetrySpansDelegateLineage(t *testing.T) {
exp := tracetest.NewInMemoryExporter()
tp := trace.NewTracerProvider(trace.WithSyncer(exp))
+2
View File
@@ -135,6 +135,8 @@ type RunInfo struct {
VerificationFeedback string // feedback from the previous failed verifier attempt, when retrying a flow step
Dispatch string // how the run was dispatched (direct, broker, schedule, resume) when known
Trigger string // external trigger or schedule label that started the run, when known
Spent int64 // cumulative paid x402 spend in this run, in the asset's smallest unit
ToolSpend int64 // paid x402 spend attributed to the current tool call, in the asset's smallest unit
}
type runInfoKey struct{}
+3
View File
@@ -94,6 +94,9 @@ func writeAgentInspection(w io.Writer, name string, runs []goagent.RunSummary, a
if run.LastErrorKind != "" {
fmt.Fprintf(w, " error_kind=%s", run.LastErrorKind)
}
if run.Spent > 0 {
fmt.Fprintf(w, " spent=%d", run.Spent)
}
if run.LastError != "" {
fmt.Fprintf(w, " error=%q", run.LastError)
}
+2 -2
View File
@@ -11,13 +11,13 @@ import (
)
func TestWriteAgentInspectionIncludesActionableBreadcrumbs(t *testing.T) {
runs := []goagent.RunSummary{{RunID: "run-1", Status: "auth", Events: 4, LastKind: "model", LastError: "invalid API key", LastErrorKind: "auth", TraceID: "1234567890abcdef", Checkpoint: "failed", Stage: "ask"}}
runs := []goagent.RunSummary{{RunID: "run-1", Status: "auth", Events: 4, LastKind: "model", LastError: "invalid API key", LastErrorKind: "auth", TraceID: "1234567890abcdef", Checkpoint: "failed", Stage: "ask", Spent: 7}}
var out bytes.Buffer
if err := writeAgentInspection(&out, "support", runs, false); err != nil {
t.Fatal(err)
}
got := out.String()
for _, want := range []string{"Agent \"support\" runs", "run-1", "status=auth", "events=4", "last=model", "checkpoint=failed", "stage=ask", "error_kind=auth", `error="invalid API key"`, "trace=1234567890ab"} {
for _, want := range []string{"Agent \"support\" runs", "run-1", "status=auth", "events=4", "last=model", "checkpoint=failed", "stage=ask", "error_kind=auth", `error="invalid API key"`, "trace=1234567890ab", "spent=7"} {
if !strings.Contains(got, want) {
t.Fatalf("output missing %q:\n%s", want, got)
}
+31
View File
@@ -0,0 +1,31 @@
# 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.
+37
View File
@@ -0,0 +1,37 @@
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}
+37
View File
@@ -0,0 +1,37 @@
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}
+37
View File
@@ -0,0 +1,37 @@
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}
+8
View File
@@ -0,0 +1,8 @@
// 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
+87
View File
@@ -0,0 +1,87 @@
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")
}
}
+125
View File
@@ -0,0 +1,125 @@
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}
`
+138
View File
@@ -0,0 +1,138 @@
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
}
+4 -1
View File
@@ -12,6 +12,7 @@ provider-free unless the example README says otherwise.
| Prove the maintained 0→hero path | [`support`](./support/) | `go run ./examples/support` and `go test ./examples/support` | [`zero-to-hero` guide](../internal/website/docs/guides/zero-to-hero.md) |
| See planning and delegation | [`agent-plan-delegate`](./agent-plan-delegate/) | `go run ./examples/agent-plan-delegate` | [`plan-delegate` guide](../internal/website/docs/guides/plan-delegate.md) |
| Expose services through MCP | [`mcp/hello`](./mcp/hello/) | follow [`mcp`](./mcp/) setup | [`mcp/crud`](./mcp/crud/) and [`mcp/workflow`](./mcp/workflow/) |
| Try a paid tool with x402 | [`agent-x402-buyer`](./agent-x402-buyer/) | `go run ./examples/agent-x402-buyer` | [`Payments (x402)` guide](../internal/website/docs/guides/x402-payments.md) |
| Try A2A or gRPC interop next | [`agent-demo`](./agent-demo/) plus gateway docs | run the example, then use the gateway docs | [`grpc-interop`](./grpc-interop/) |
| Add workflow durability | [`flow-durable`](./flow-durable/) | `go run ./examples/flow-durable` | [`flow-loop`](./flow-loop/) |
@@ -28,7 +29,9 @@ provider-free unless the example README says otherwise.
4. **Interop next:** use [`mcp/hello`](./mcp/hello/), [`mcp/crud`](./mcp/crud/),
and [`mcp/workflow`](./mcp/workflow/) when you are ready to expose tools to
external AI clients.
5. **Workflow depth:** use [`flow-durable`](./flow-durable/) once the agent path
5. **Paid tools:** run [`agent-x402-buyer`](./agent-x402-buyer/) to see an
agent pay a local x402-protected tool with a mock facilitator and budget.
6. **Workflow depth:** use [`flow-durable`](./flow-durable/) once the agent path
needs checkpointed, resumable deterministic work.
## CLI wayfinding
+24
View File
@@ -0,0 +1,24 @@
# Agent x402 buyer
This example shows an agent paying for a paid HTTP tool with x402 without using
live funds or a live chain.
It starts a local paid endpoint guarded by `wrapper/x402` seller middleware and a
mock facilitator. A deterministic mock-model agent calls that endpoint as a tool,
receives the HTTP 402 challenge, pays with `AgentPayer`, stays inside
`AgentBudget`, retries the request, and prints the spend recorded for the run.
```bash
go run ./examples/agent-x402-buyer
```
Expected output includes:
- the paid tool response,
- one facilitator verify and settle call, and
- `run spend: 7 smallest units (budget 10)`.
The payment token and facilitator are intentionally local development fakes. To
settle real x402 payments, keep the same `AgentPayer` / `AgentBudget` shape but
replace the payer with a wallet-backed implementation and configure the seller
middleware with a hosted or self-run x402 facilitator.
+168
View File
@@ -0,0 +1,168 @@
// Agent x402 buyer — a provider-free example of an agent paying for a paid tool.
//
// Run:
//
// go run ./examples/agent-x402-buyer
//
// It starts a local HTTP tool protected by x402 middleware, then asks a
// deterministic mock-model agent to call that tool. The agent receives the 402
// challenge, uses AgentPayer and AgentBudget to pay within a local mock
// facilitator, retries the request, and prints the run spend.
package main
import (
"context"
"encoding/json"
"fmt"
"io"
"net/http"
"net/http/httptest"
"os"
"strings"
go_micro "go-micro.dev/v6"
"go-micro.dev/v6/agent"
"go-micro.dev/v6/ai"
"go-micro.dev/v6/store"
"go-micro.dev/v6/wrapper/x402"
)
const (
paidToolName = "paid.market_brief"
price = int64(7)
paymentToken = "dev-payment-token"
)
type devFacilitator struct {
verifyCount int
settleCount int
}
func (f *devFacilitator) Verify(ctx context.Context, payment string, req x402.Requirements) (x402.Result, error) {
f.verifyCount++
if payment != paymentToken {
return x402.Result{Valid: false, Reason: "unknown dev payment token"}, nil
}
return x402.Result{Valid: true, Payer: "dev-agent-wallet"}, nil
}
func (f *devFacilitator) Settle(ctx context.Context, payment string, req x402.Requirements) (x402.Result, error) {
f.settleCount++
return x402.Result{Valid: true, Settlement: "dev-settlement-001"}, nil
}
type devPayer struct{}
func (devPayer) Pay(ctx context.Context, req x402.Requirements) (string, error) {
return paymentToken, nil
}
type mockModel struct{ opts ai.Options }
func newMock(opts ...ai.Option) ai.Model {
m := &mockModel{}
_ = m.Init(opts...)
return m
}
func (m *mockModel) Init(opts ...ai.Option) error {
for _, o := range opts {
o(&m.opts)
}
return nil
}
func (m *mockModel) Options() ai.Options { return m.opts }
func (m *mockModel) String() string { return "agent-x402-buyer-mock" }
func (m *mockModel) Stream(context.Context, *ai.Request, ...ai.GenerateOption) (ai.Stream, error) {
return nil, fmt.Errorf("stream not supported by agent-x402-buyer mock")
}
func (m *mockModel) Generate(ctx context.Context, req *ai.Request, _ ...ai.GenerateOption) (*ai.Response, error) {
for _, tool := range req.Tools {
if tool.Name == paidToolName && m.opts.ToolHandler != nil {
out := m.opts.ToolHandler(ctx, ai.ToolCall{ID: "paid-brief", Name: tool.Name, Input: map[string]any{"url": req.Prompt}})
return &ai.Response{Answer: fmt.Sprintf("Paid tool returned: %s", out.Content)}, nil
}
}
return &ai.Response{Answer: "No paid tool was available."}, nil
}
func paidToolServer(fac *devFacilitator) *httptest.Server {
mux := http.NewServeMux()
paid := x402.Middleware(x402.Config{
PayTo: "0xMerchantDevWallet",
Network: "base-sepolia",
Amount: fmt.Sprint(price),
Description: "Local market brief for the x402 buyer example",
Facilitator: fac,
})
mux.Handle("/brief", paid(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
_ = json.NewEncoder(w).Encode(map[string]any{
"brief": "Mock demand is up 12% after the agent paid the local tool.",
"settlement": w.Header().Get(x402.PaymentResponseHeader),
})
})))
return httptest.NewServer(mux)
}
func run(w io.Writer) error {
ai.Register("agent-x402-buyer-mock", newMock)
fac := &devFacilitator{}
srv := paidToolServer(fac)
defer srv.Close()
st := store.NewMemoryStore()
buyer := agent.New(
agent.Name("x402-buyer"),
agent.Provider("agent-x402-buyer-mock"),
agent.Prompt("Call the paid market brief tool when given its URL."),
agent.WithStore(st),
go_micro.AgentPayer(devPayer{}),
go_micro.AgentBudget(10),
agent.WithTool(paidToolName, "Fetch a paid market brief over HTTP", map[string]any{
"url": map[string]any{"type": "string", "description": "Paid HTTP endpoint to call"},
}, func(ctx context.Context, input map[string]any) (string, error) {
url, _ := input["url"].(string)
resp, err := http.Get(url)
if err != nil {
return "", err
}
defer resp.Body.Close()
body, err := io.ReadAll(resp.Body)
if err != nil {
return "", err
}
return string(body), nil
}),
)
resp, err := buyer.Ask(context.Background(), srv.URL+"/brief")
if err != nil {
return err
}
events, err := agent.LoadRunEvents(st, "x402-buyer", resp.RunID)
if err != nil {
return err
}
var spent int64
for _, event := range events {
if event.Spent > spent {
spent = event.Spent
}
}
fmt.Fprintln(w, "Agent x402 buyer (provider: mock, funds: local dev token)")
fmt.Fprintln(w, strings.TrimSpace(resp.Reply))
fmt.Fprintf(w, "facilitator verify=%d settle=%d\n", fac.verifyCount, fac.settleCount)
fmt.Fprintf(w, "run spend: %d smallest units (budget 10)\n", spent)
return nil
}
func main() {
if err := run(os.Stdout); err != nil {
fmt.Println(err)
os.Exit(1)
}
}
+282
View File
@@ -0,0 +1,282 @@
package mcp
import (
"context"
"encoding/json"
"fmt"
"strings"
"time"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure"
reflectionpb "google.golang.org/grpc/reflection/grpc_reflection_v1alpha"
"google.golang.org/protobuf/encoding/protojson"
"google.golang.org/protobuf/proto"
"google.golang.org/protobuf/reflect/protodesc"
"google.golang.org/protobuf/reflect/protoreflect"
"google.golang.org/protobuf/reflect/protoregistry"
"google.golang.org/protobuf/types/descriptorpb"
"google.golang.org/protobuf/types/dynamicpb"
)
// ReflectedGRPCTarget describes an external gRPC server whose reflection
// catalog should be exposed as MCP tools. It is intentionally opt-in: teams can
// bridge existing reflected gRPC services without changing their servers or
// registering them in go-micro.
type ReflectedGRPCTarget struct {
// Name prefixes generated tools. When empty, Address is sanitized and used.
Name string
// Address is the host:port of the reflected gRPC server.
Address string
// DialOptions customize the connection. If none are supplied, an insecure
// transport is used for local/dev interoperability.
DialOptions []grpc.DialOption
// Timeout bounds reflection discovery and individual tool calls.
Timeout time.Duration
}
func (s *Server) discoverReflectedGRPC() error {
for _, target := range s.opts.ReflectedGRPCTargets {
if strings.TrimSpace(target.Address) == "" {
continue
}
tools, err := s.reflectedGRPCTools(target)
if err != nil {
return err
}
for _, tool := range tools {
s.tools[tool.Name] = tool
}
}
return nil
}
func (s *Server) reflectedGRPCTools(target ReflectedGRPCTarget) ([]*Tool, error) {
timeout := target.Timeout
if timeout == 0 {
timeout = 10 * time.Second
}
ctx, cancel := context.WithTimeout(s.opts.Context, timeout)
defer cancel()
dialOpts := target.DialOptions
if len(dialOpts) == 0 {
dialOpts = []grpc.DialOption{grpc.WithTransportCredentials(insecure.NewCredentials())}
}
conn, err := grpc.NewClient(target.Address, dialOpts...)
if err != nil {
return nil, fmt.Errorf("connect reflected grpc target %s: %w", target.Address, err)
}
defer conn.Close()
files, services, err := loadReflectedFiles(ctx, conn)
if err != nil {
return nil, fmt.Errorf("reflect grpc target %s: %w", target.Address, err)
}
prefix := target.Name
if prefix == "" {
prefix = sanitizeToolPart(target.Address)
}
var out []*Tool
for _, serviceName := range services {
desc, err := files.FindDescriptorByName(protoreflect.FullName(serviceName))
if err != nil {
continue
}
svc, ok := desc.(protoreflect.ServiceDescriptor)
if !ok {
continue
}
for i := 0; i < svc.Methods().Len(); i++ {
method := svc.Methods().Get(i)
if method.IsStreamingClient() || method.IsStreamingServer() {
continue
}
fullMethod := "/" + string(svc.FullName()) + "/" + string(method.Name())
toolName := prefix + "." + strings.ReplaceAll(string(svc.FullName()), ".", "_") + "." + string(method.Name())
input := method.Input()
out = append(out, &Tool{
Name: toolName,
Description: fmt.Sprintf("Call reflected gRPC method %s on %s", fullMethod, target.Address),
InputSchema: protoMessageSchema(input),
Handler: reflectedGRPCHandler(target, fullMethod, input, method.Output()),
})
}
}
return out, nil
}
func loadReflectedFiles(ctx context.Context, conn *grpc.ClientConn) (*protoregistryFiles, []string, error) {
client := reflectionpb.NewServerReflectionClient(conn)
stream, err := client.ServerReflectionInfo(ctx)
if err != nil {
return nil, nil, err
}
if err := stream.Send(&reflectionpb.ServerReflectionRequest{MessageRequest: &reflectionpb.ServerReflectionRequest_ListServices{ListServices: ""}}); err != nil {
return nil, nil, err
}
resp, err := stream.Recv()
if err != nil {
return nil, nil, err
}
list := resp.GetListServicesResponse()
if list == nil {
return nil, nil, fmt.Errorf("reflection list services returned %T", resp.MessageResponse)
}
set := &descriptorpb.FileDescriptorSet{}
seen := map[string]bool{}
var services []string
for _, svc := range list.Service {
name := svc.Name
if strings.HasPrefix(name, "grpc.reflection.") {
continue
}
services = append(services, name)
if err := requestFileContainingSymbol(ctx, client, name, set, seen); err != nil {
return nil, nil, err
}
}
files, err := newProtoregistryFiles(set)
if err != nil {
return nil, nil, err
}
return files, services, nil
}
func requestFileContainingSymbol(ctx context.Context, client reflectionpb.ServerReflectionClient, symbol string, set *descriptorpb.FileDescriptorSet, seen map[string]bool) error {
stream, err := client.ServerReflectionInfo(ctx)
if err != nil {
return err
}
if err := stream.Send(&reflectionpb.ServerReflectionRequest{MessageRequest: &reflectionpb.ServerReflectionRequest_FileContainingSymbol{FileContainingSymbol: symbol}}); err != nil {
return err
}
resp, err := stream.Recv()
if err != nil {
return err
}
fd := resp.GetFileDescriptorResponse()
if fd == nil {
return fmt.Errorf("reflection lookup for %s returned %T", symbol, resp.MessageResponse)
}
for _, raw := range fd.FileDescriptorProto {
var file descriptorpb.FileDescriptorProto
if err := proto.Unmarshal(raw, &file); err != nil {
return err
}
name := file.GetName()
if !seen[name] {
seen[name] = true
set.File = append(set.File, &file)
}
}
return nil
}
// protoregistryFiles is a narrow wrapper that keeps imports local to this file.
type protoregistryFiles struct{ files *protoregistry.Files }
func newProtoregistryFiles(set *descriptorpb.FileDescriptorSet) (*protoregistryFiles, error) {
files, err := protodesc.NewFiles(set)
if err != nil {
return nil, err
}
return &protoregistryFiles{files: files}, nil
}
func (p *protoregistryFiles) FindDescriptorByName(name protoreflect.FullName) (protoreflect.Descriptor, error) {
return p.files.FindDescriptorByName(name)
}
func reflectedGRPCHandler(target ReflectedGRPCTarget, fullMethod string, input, output protoreflect.MessageDescriptor) func(map[string]interface{}) (interface{}, error) {
return func(args map[string]interface{}) (interface{}, error) {
timeout := target.Timeout
if timeout == 0 {
timeout = 10 * time.Second
}
ctx, cancel := context.WithTimeout(context.Background(), timeout)
defer cancel()
dialOpts := target.DialOptions
if len(dialOpts) == 0 {
dialOpts = []grpc.DialOption{grpc.WithTransportCredentials(insecure.NewCredentials())}
}
conn, err := grpc.NewClient(target.Address, dialOpts...)
if err != nil {
return nil, err
}
defer conn.Close()
req := dynamicpb.NewMessage(input)
raw, err := json.Marshal(args)
if err != nil {
return nil, err
}
if err := protojson.Unmarshal(raw, req); err != nil {
return nil, err
}
rsp := dynamicpb.NewMessage(output)
if err := conn.Invoke(ctx, fullMethod, req, rsp); err != nil {
return nil, err
}
b, err := protojson.MarshalOptions{UseProtoNames: true, EmitUnpopulated: true}.Marshal(rsp)
if err != nil {
return nil, err
}
var out interface{}
if err := json.Unmarshal(b, &out); err != nil {
return nil, err
}
return out, nil
}
}
func protoMessageSchema(msg protoreflect.MessageDescriptor) map[string]interface{} {
schema := map[string]interface{}{"type": "object", "properties": map[string]interface{}{}}
props := schema["properties"].(map[string]interface{})
fields := msg.Fields()
for i := 0; i < fields.Len(); i++ {
field := fields.Get(i)
props[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)
}
+62
View File
@@ -0,0 +1,62 @@
package mcp
import (
"context"
"net"
"testing"
"time"
"google.golang.org/grpc"
helloworld "google.golang.org/grpc/examples/helloworld/helloworld"
"google.golang.org/grpc/reflection"
)
type reflectedGreeter struct {
helloworld.UnimplementedGreeterServer
}
func (reflectedGreeter) SayHello(_ context.Context, req *helloworld.HelloRequest) (*helloworld.HelloReply, error) {
return &helloworld.HelloReply{Message: "hello " + req.Name}, nil
}
func TestReflectedGRPCTargetDiscoversAndCallsUnaryTool(t *testing.T) {
lis, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
t.Fatal(err)
}
grpcServer := grpc.NewServer()
helloworld.RegisterGreeterServer(grpcServer, reflectedGreeter{})
reflection.Register(grpcServer)
go grpcServer.Serve(lis)
defer grpcServer.Stop()
s := newTestServer(Options{Context: context.Background()})
tools, err := s.reflectedGRPCTools(ReflectedGRPCTarget{
Name: "demo",
Address: lis.Addr().String(),
Timeout: 3 * time.Second,
})
if err != nil {
t.Fatalf("discover reflected tools: %v", err)
}
if len(tools) != 1 {
t.Fatalf("tools len = %d, want 1", len(tools))
}
tool := tools[0]
if tool.Name != "demo.helloworld_Greeter.SayHello" {
t.Fatalf("tool name = %q", tool.Name)
}
props := tool.InputSchema["properties"].(map[string]interface{})
if _, ok := props["name"]; !ok {
t.Fatalf("input schema missing name: %#v", tool.InputSchema)
}
out, err := tool.Handler(map[string]interface{}{"name": "Ada"})
if err != nil {
t.Fatalf("call reflected tool: %v", err)
}
got := out.(map[string]interface{})["message"]
if got != "hello Ada" {
t.Fatalf("message = %v, want hello Ada", got)
}
}
+9
View File
@@ -157,6 +157,11 @@ type Options struct {
// (the /mcp/call endpoint). Listing tools and health stay free.
// Opt-in: leave nil to disable payments.
Payment *x402.Config
// ReflectedGRPCTargets exposes unary methods from external gRPC servers
// that support server reflection as MCP tools. This bridges existing gRPC
// services into the agent tool catalog without requiring go-micro handlers.
ReflectedGRPCTargets []ReflectedGRPCTarget
}
// Server represents a running MCP gateway
@@ -286,6 +291,10 @@ func (s *Server) discoverServices() error {
s.toolsMu.Lock()
defer s.toolsMu.Unlock()
if err := s.discoverReflectedGRPC(); err != nil {
return err
}
for _, svc := range services {
// Get full service details
fullSvcs, err := s.opts.Registry.GetService(svc.Name)
+5 -19
View File
@@ -294,7 +294,9 @@ func (t *StdioTransport) handleToolsCall(req *JSONRPCRequest) {
AccountID: accountID, ScopesRequired: tool.Scopes,
Allowed: true, Duration: time.Since(start), Error: err.Error(),
})
t.sendError(req.ID, InternalError, "RPC call failed", 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()))
return
}
@@ -311,24 +313,8 @@ func (t *StdioTransport) handleToolsCall(req *JSONRPCRequest) {
Allowed: true, Duration: time.Since(start),
})
// 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,
})
// The downstream response is JSON — return it as JSON text, not %v.
t.sendResponse(req.ID, mcpToolResult(traceID, rsp.Data))
}
// sendResponse sends a JSON-RPC response
+120
View File
@@ -0,0 +1,120 @@
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")
}
}
+32
View File
@@ -0,0 +1,32 @@
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,
}
}
+4 -18
View File
@@ -275,7 +275,8 @@ func (wc *wsConn) handleToolsCall(req *JSONRPCRequest) {
AccountID: accountID, ScopesRequired: tool.Scopes,
Allowed: true, Duration: time.Since(start), Error: err.Error(),
})
wc.sendError(req.ID, InternalError, "RPC call failed", err.Error())
// Tool-execution failure → isError result (MCP spec), not a protocol error.
wc.sendResponse(req.ID, mcpToolError(traceID, "tool call failed: "+err.Error()))
return
}
@@ -291,23 +292,8 @@ func (wc *wsConn) handleToolsCall(req *JSONRPCRequest) {
Allowed: true, Duration: time.Since(start),
})
// 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,
})
// The downstream response is JSON — return it as JSON text, not %v.
wc.sendResponse(req.ID, mcpToolResult(traceID, rsp.Data))
}
// sendResponse sends a JSON-RPC success response.
+18 -15
View File
@@ -114,12 +114,13 @@ func TestWebSocket_ToolsCall_NoAuth(t *testing.T) {
"arguments": map[string]interface{}{"msg": "hi"},
})
// RPC will fail (no backend), but auth should pass (no auth configured)
if resp.Error == nil {
t.Fatal("expected RPC error (no backend)")
// 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)
}
if resp.Error.Code != InternalError {
t.Errorf("error code = %d, want %d", resp.Error.Code, InternalError)
if !isToolError(resp.Result) {
t.Fatalf("expected isError tool result, got %+v", resp.Result)
}
}
@@ -168,12 +169,13 @@ func TestWebSocket_ToolsCall_AuthRequired(t *testing.T) {
"arguments": map[string]interface{}{},
"_token": "valid-token",
})
// Auth passes, RPC fails (no backend)
if resp.Error == nil {
t.Fatal("expected RPC error")
// 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)
}
if resp.Error.Code != InternalError {
t.Errorf("error code = %d, want %d (RPC fail, not auth fail)", resp.Error.Code, InternalError)
if !isToolError(resp.Result) {
t.Fatalf("expected isError tool result, got %+v", resp.Result)
}
})
@@ -185,12 +187,13 @@ func TestWebSocket_ToolsCall_AuthRequired(t *testing.T) {
"name": "svc.Do",
"arguments": map[string]interface{}{},
})
// Auth passes via connection-level header, RPC fails (no backend)
if resp.Error == nil {
t.Fatal("expected RPC error")
// 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)
}
if resp.Error.Code != InternalError {
t.Errorf("error code = %d, want %d (RPC fail, not auth fail)", resp.Error.Code, InternalError)
if !isToolError(resp.Result) {
t.Fatalf("expected isError tool result, got %+v", resp.Result)
}
})
}
+11 -1
View File
@@ -8,6 +8,7 @@ import (
"io"
"net/http"
"strconv"
"strings"
"sync"
)
@@ -81,7 +82,16 @@ func (c *Client) Do(req *http.Request) (*http.Response, error) {
return resp, fmt.Errorf("x402: 402 response carried no requirements")
}
reqd := ch.Accepts[0]
amount, _ := strconv.ParseInt(reqd.MaxAmountRequired, 10, 64)
// 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)
}
// Spend cap: reserve before paying so concurrent calls cannot all pass
// the check and overspend the caller's allowance. Roll the reservation
+27
View File
@@ -172,6 +172,33 @@ 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) {
+13 -1
View File
@@ -126,6 +126,11 @@ 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 {
@@ -222,7 +227,14 @@ 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.
if s, ok := fac.(Settler); ok {
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 {
sres, err := s.Settle(r.Context(), payment, req)
if err != nil {
writeChallenge(w, req, "payment settlement failed: "+err.Error())
+54
View File
@@ -159,4 +159,58 @@ 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)