项目文件夹

文件
Asim Aslam c7657f73f4
goreleaser / goreleaser (push) Has been cancelled
Refactor agent plan storage, update docs, and release v6 (#2977)
* test(harness): read agent plan from the scoped store

The store-scoping change moved an agent's plan from the default table
key agent/{name}/plan to its own table (database "agent", table {name},
key "plan"). The plan-delegate harness tests still read the old key and
failed with 'not found'; read through store.Scope(mem, "agent", name)
like the agent does.

* docs: orient agents-first across README, landing, and docs overview

Lead with agents (then services and flows), surface MCP + A2A as the
interop story, and frame agents as services. Landing hero and feature
grid reordered agents-first with an A2A gateway card.

* v6: module path go-micro.dev/v6, TLS secure by default, NewService

Cut v6. Three breaking changes, bundled so the major bump is paid once:

- Module path go-micro.dev/v5 -> go-micro.dev/v6 across all imports + go.mod.
- TLS verification on by default (was off). MICRO_TLS_SECURE removed;
  MICRO_TLS_INSECURE=true opts out for self-signed/dev.
- micro.NewService(name, opts...) is the canonical service constructor,
  symmetric with NewAgent/NewFlow; micro.New kept as a deprecated alias;
  the old name-less NewService(opts...) removed. Generators emit NewService.

Also ports the JWT auth token provider in-module (go-micro.dev/v6/auth/jwt/token
on golang-jwt/jwt/v5), dropping the v5-pinned github.com/micro/plugins/v5/auth/jwt
and the deprecated dgrijalva/jwt-go.

Docs/README/landing updated to v6 and @latest; v5->v6 migration guide added;
CHANGELOG cut as [6.0.0]. Blog posts left at their historical versions.

---------

Co-authored-by: Claude <noreply@anthropic.com>
2026-06-18 11:55:35 +01:00

205 行
6.3 KiB
Go

package main
import (
"context"
"testing"
"time"
"go-micro.dev/v6/agent"
"go-micro.dev/v6/ai"
"go-micro.dev/v6/broker"
"go-micro.dev/v6/client"
"go-micro.dev/v6/flow"
"go-micro.dev/v6/registry"
"go-micro.dev/v6/selector"
"go-micro.dev/v6/service"
"go-micro.dev/v6/store"
)
// waitForService polls the registry until name is registered, instead of
// sleeping. Keeps the test deterministic.
func waitForService(t *testing.T, reg registry.Registry, name string) {
t.Helper()
deadline := time.Now().Add(5 * time.Second)
for time.Now().Before(deadline) {
if svcs, err := reg.GetService(name); err == nil && len(svcs) > 0 && len(svcs[0].Nodes) > 0 {
return
}
time.Sleep(20 * time.Millisecond)
}
t.Fatalf("service %q never registered", name)
}
// TestPlanDelegateEndToEnd runs the whole feature against the real stack
// — real services, a shared in-memory registry, real RPC, the real agent
// loop, real store — with only the LLM mocked. No mDNS, no sleeps.
func TestPlanDelegateEndToEnd(t *testing.T) {
ai.Register("mock", newMock)
// Shared infrastructure: one in-memory registry, a client bound to
// it, and an in-memory store. Everything resolves through these.
reg := registry.NewMemoryRegistry()
cl := client.NewClient(
client.Registry(reg),
client.Selector(selector.NewSelector(selector.Registry(reg))),
)
mem := store.NewMemoryStore()
// Real services on the shared registry/client.
taskSvc := new(TaskService)
task := service.New(service.Name("task"), service.Registry(reg), service.Client(cl))
if err := task.Handle(taskSvc); err != nil {
t.Fatalf("handle task: %v", err)
}
go task.Run()
notifySvc := new(NotifyService)
notify := service.New(service.Name("notify"), service.Registry(reg), service.Client(cl))
if err := notify.Handle(notifySvc); err != nil {
t.Fatalf("handle notify: %v", err)
}
go notify.Run()
// Real comms agent (owns notify), registered so delegate reaches it over RPC.
comms := agent.New(
agent.Name("comms"),
agent.Services("notify"),
agent.Prompt("You handle outbound notifications."),
agent.Provider("mock"),
agent.WithRegistry(reg),
agent.WithClient(cl),
agent.WithStore(mem),
)
go comms.Run()
defer comms.Stop()
waitForService(t, reg, "task")
waitForService(t, reg, "notify")
waitForService(t, reg, "comms")
// Real conductor agent (owns task), driven programmatically.
conductor := agent.New(
agent.Name("conductor"),
agent.Services("task"),
agent.Prompt("Plan first, create tasks, delegate notifications to the comms agent."),
agent.Provider("mock"),
agent.WithRegistry(reg),
agent.WithClient(cl),
agent.WithStore(mem),
)
resp, err := conductor.Ask(context.Background(),
"Create three launch tasks: Design, Build, and Ship. Then notify owner@acme.com that the plan is ready.")
if err != nil {
t.Fatalf("Ask: %v", err)
}
if resp.Reply == "" {
t.Error("conductor returned an empty reply")
}
// Tasks were created via real RPC into the task service.
if n := taskSvc.count(); n != 3 {
t.Errorf("task service has %d tasks, want 3", n)
}
// The plan was persisted to the real store, in the agent's scoped table.
if recs, err := store.Scope(mem, "agent", "conductor").Read("plan"); err != nil || len(recs) == 0 {
t.Errorf("plan not persisted to store: err=%v recs=%d", err, len(recs))
}
// Delegation reached the comms agent over RPC, which called notify.
if n := notifySvc.count(); n != 1 {
t.Errorf("notify service called %d times, want 1 (delegation did not reach comms)", n)
}
}
// TestFlowDispatchesToAgentEndToEnd proves "Flow triggers, Agent reasons":
// a workflow event hands off to the registered conductor agent, which then
// plans, creates tasks, and delegates to comms — all over real RPC. Only
// the LLM is mocked.
func TestFlowDispatchesToAgentEndToEnd(t *testing.T) {
ai.Register("mock", newMock)
reg := registry.NewMemoryRegistry()
cl := client.NewClient(
client.Registry(reg),
client.Selector(selector.NewSelector(selector.Registry(reg))),
)
mem := store.NewMemoryStore()
taskSvc := new(TaskService)
task := service.New(service.Name("task"), service.Registry(reg), service.Client(cl))
if err := task.Handle(taskSvc); err != nil {
t.Fatalf("handle task: %v", err)
}
go task.Run()
notifySvc := new(NotifyService)
notify := service.New(service.Name("notify"), service.Registry(reg), service.Client(cl))
if err := notify.Handle(notifySvc); err != nil {
t.Fatalf("handle notify: %v", err)
}
go notify.Run()
comms := agent.New(
agent.Name("comms"),
agent.Services("notify"),
agent.Prompt("You handle outbound notifications."),
agent.Provider("mock"),
agent.WithRegistry(reg),
agent.WithClient(cl),
agent.WithStore(mem),
)
go comms.Run()
defer comms.Stop()
// Unlike the previous test, the conductor must be registered (running)
// so the flow can reach it over RPC.
conductor := agent.New(
agent.Name("conductor"),
agent.Services("task"),
agent.Prompt("Plan first, create tasks, delegate notifications to the comms agent."),
agent.Provider("mock"),
agent.WithRegistry(reg),
agent.WithClient(cl),
agent.WithStore(mem),
)
go conductor.Run()
defer conductor.Stop()
waitForService(t, reg, "task")
waitForService(t, reg, "notify")
waitForService(t, reg, "comms")
waitForService(t, reg, "conductor")
// A workflow that hands each event to the conductor agent.
f := flow.New("onboard",
flow.Agent("conductor"),
flow.Prompt("Get the launch ready: {{.Data}}"),
)
if err := f.Register(reg, broker.DefaultBroker, cl); err != nil {
t.Fatalf("flow register: %v", err)
}
// Fire the workflow (as a broker event would).
if err := f.Execute(context.Background(), "three tasks then notify owner@acme.com"); err != nil {
t.Fatalf("flow execute: %v", err)
}
// The flow recorded the agent's reply.
if rs := f.Results(); len(rs) != 1 || rs[0].Reply == "" {
t.Errorf("flow result = %+v, want one result with a reply", rs)
}
// The agent ran end to end: tasks created, plan stored, comms notified.
if n := taskSvc.count(); n != 3 {
t.Errorf("task service has %d tasks, want 3", n)
}
if recs, err := store.Scope(mem, "agent", "conductor").Read("plan"); err != nil || len(recs) == 0 {
t.Errorf("plan not persisted: err=%v recs=%d", err, len(recs))
}
if n := notifySvc.count(); n != 1 {
t.Errorf("notify called %d times, want 1 (flow->agent->delegate->comms chain broken)", n)
}
}