Compare commits

..

1 Commits

Author SHA1 Message Date
Codex 6d05ad0ecf loop: drop completed spend observability priority
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 09:35:40 +00:00
7 changed files with 4 additions and 557 deletions
+3 -8
View File
@@ -24,14 +24,9 @@ rewrites.
### Capability — the headline (roadmap: Now / Next)
1. **x402 buyer safety hard stop** ([#4814](https://github.com/micro/go-micro/issues/4814)) — close the budget-cap bypass and require a real settler for paid configs before autonomous paid-tool use becomes the next thing users copy into real agents.
2. **Kubernetes operator + CRDs foundation** ([#4797](https://github.com/micro/go-micro/issues/4797)) — add the first opt-in `Agent`, `Service`, and `Flow` resource foundation so the services → agents → workflows lifecycle has a native deployment path for Kubernetes users.
3. **A2A external-client conformance** ([#4815](https://github.com/micro/go-micro/issues/4815)) — make the gateway easier for non-go-micro agents to discover and stream from by serving the well-known agent card path and spec SSE events.
4. **MCP stdio/ws result conformance** ([#4813](https://github.com/micro/go-micro/issues/4813)) — return JSON tool results with explicit `isError` semantics across transports, backed by a stdio round-trip test.
### In flight — do not re-queue
- **gRPC-reflection MCP lint follow-up** ([#4824](https://github.com/micro/go-micro/issues/4824), [PR #4826](https://github.com/micro/go-micro/pull/4826)) — fixes the lint fallout from the merged reflected-gRPC MCP work.
1. **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.)
2. **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.
3. **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.
### Background — hardening & DX (roadmap: Ongoing; capped)
+1 -4
View File
@@ -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
-24
View File
@@ -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.
-168
View File
@@ -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)
}
}
-282
View File
@@ -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[string(field.JSONName())] = protoFieldSchema(field)
}
return schema
}
func protoFieldSchema(field protoreflect.FieldDescriptor) map[string]interface{} {
schema := map[string]interface{}{"type": protoJSONType(field)}
if field.IsList() {
schema["items"] = map[string]interface{}{"type": protoJSONType(field)}
}
if field.Kind() == protoreflect.MessageKind || field.Kind() == protoreflect.GroupKind {
schema = protoMessageSchema(field.Message())
}
return schema
}
func protoJSONType(field protoreflect.FieldDescriptor) string {
if field.IsList() {
return "array"
}
switch field.Kind() {
case protoreflect.BoolKind:
return "boolean"
case protoreflect.Int32Kind, protoreflect.Sint32Kind, protoreflect.Sfixed32Kind,
protoreflect.Uint32Kind, protoreflect.Fixed32Kind, protoreflect.Int64Kind,
protoreflect.Sint64Kind, protoreflect.Sfixed64Kind, protoreflect.Uint64Kind,
protoreflect.Fixed64Kind:
return "integer"
case protoreflect.FloatKind, protoreflect.DoubleKind:
return "number"
case protoreflect.MessageKind, protoreflect.GroupKind:
return "object"
default:
return "string"
}
}
func sanitizeToolPart(s string) string {
r := strings.NewReplacer(":", "_", "/", "_", ".", "_", "-", "_")
return r.Replace(s)
}
-62
View File
@@ -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)
}
}
-9
View File
@@ -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)