项目文件夹

文件
T
Asim Aslam 1bc886fa82 Introduce Agent abstraction and integrate with chat router (#2939)
* docs: Agent interface design sketch

Proposes Agent as a top-level abstraction alongside Service in the
micro package. Agent manages services — scoped tools, system prompt,
conversation memory, registry-discoverable.

Design only, no implementation.

https://claude.ai/code/session_01QTp4SshuVmLAvvEGJe4TJd

* feat: Agent as a first-class abstraction

Introduce micro.NewAgent() alongside micro.New() — Agent is to
intelligence what Service is to capability.

Agent interface:
- Chat(ctx, message) (*Response, error) — core interaction method
- Run() — registers in registry, subscribes to broker, blocks
- Stop() — graceful shutdown
- Scoped tools — only sees endpoints of its assigned services
- Persistent memory — conversation history stored in store
- Agent-to-agent — communication via broker topics

Top-level API:
  agent := micro.NewAgent("task-mgr",
      micro.AgentServices("task"),
      micro.AgentPrompt("You manage tasks."),
      micro.AgentProvider("anthropic"),
  )
  agent.Run()

https://claude.ai/code/session_01QTp4SshuVmLAvvEGJe4TJd

* feat: wire agents into chat router, add micro agent CLI, expose Flow

Three top-level abstractions:
  micro.New("task")              — Service (capability)
  micro.NewAgent("task-mgr")     — Agent (intelligence)
  micro.NewFlow("onboard-user")  — Flow (event-driven orchestration)

micro chat as router:
- Discovers agents from registry on startup
- Single agent: routes directly
- Multiple agents: LLM classifies intent, dispatches to right agent
  via route_to_agent tool
- No agents: falls back to current direct-service behaviour
- Banner shows discovered agents

micro agent CLI:
- micro agent list — shows registered agents and their services
- micro agent describe <name> — shows agent details from registry

https://claude.ai/code/session_01QTp4SshuVmLAvvEGJe4TJd

* feat: move flow to top level, update docs for three abstractions

Package structure now consistent:
  service/   — Service (capability)
  agent/     — Agent (intelligence)
  flow/      — Flow (event-driven orchestration)

ai/flow/ kept as backward-compatible re-export.

Updated across all surfaces:
- CLAUDE.md: added agent/ and flow/ to project structure
- README: added "Building Agents" section with NewAgent() examples,
  updated features table (Agents, Flows, Chat router), CLI table
  (agent list, agent describe), docs links
- Website: features grid shows Services, Agents, Flows as the three
  pillars alongside generation, MCP, and pluggable architecture
- micro.go: Flow imported from top-level flow/ package

https://claude.ai/code/session_01QTp4SshuVmLAvvEGJe4TJd

* docs: rewrite getting-started, fix ai-integration import paths

Getting started now covers all three abstractions:
- Service (write handlers, micro run, templates)
- Agent (micro.NewAgent, scoped tools, memory, CLI)
- Flow (event-driven LLM orchestration)
Leads with prompt-based generation, then manual service creation.

ai-integration.md: fixed flow import path from go-micro.dev/v5/ai/flow
to go-micro.dev/v5/flow, updated stack diagram to show agent/flow/chat.

https://claude.ai/code/session_01QTp4SshuVmLAvvEGJe4TJd

* blog: Introducing micro.NewAgent()

Post 16 — announces Agent as a first-class abstraction. Shows the
API (NewAgent, AgentServices, AgentPrompt, AgentProvider), scoped
tools, persistent memory, multi-service agents, multi-agent systems,
and the three-abstraction comparison table (Service/Agent/Flow).

https://claude.ai/code/session_01QTp4SshuVmLAvvEGJe4TJd

* feat: fix agent registration, blog post 16

Agent registration:
- Add node with address so mDNS can discover agents
- Store type and services in node metadata (mDNS requirement)
- Connect broker before subscribing, non-fatal if broker unavailable
- Print registration confirmation on Run()

Agent/chat discovery:
- Check both service-level and node-level metadata for type=agent
  (mDNS stores metadata on nodes, not services)

Blog post 16: "Introducing micro.NewAgent()" — announces the Agent
abstraction with code examples, comparison table, multi-agent patterns.

Tested end-to-end: micro run → micro agent list discovers the agent →
micro chat routes to it → agent calls service endpoints.

https://claude.ai/code/session_01QTp4SshuVmLAvvEGJe4TJd

* feat: agents are proper services with RPC Chat endpoint

Refactored agent to use server.Server instead of fake registry entries.
An agent now:
- Creates a real RPC server with server.Name(agentName)
- Registers an Agent.Chat handler callable via standard RPC
- Sets server metadata type=agent, services=x,y for discovery
- No more fake addresses or broker hacks

micro chat calls agents via RPC (client.Call) instead of creating
local agent instances. The registry stays clean — agents are real
services with real endpoints.

Removed broker dependency from agent options. Agent-to-agent
communication is just RPC like everything else.

https://claude.ai/code/session_01QTp4SshuVmLAvvEGJe4TJd

* feat: agent uses proto-defined RPC interface

Added agent/proto/agent.proto with Agent service definition:
  rpc Chat(ChatRequest) returns (ChatResponse)

Agent now implements the generated AgentHandler interface and
registers via pb.RegisterAgentHandler. The Chat endpoint is a
standard proto-based RPC callable by any go-micro client.

Renamed the programmatic API from Chat() to Ask() to avoid
collision with the proto handler method name.

micro chat calls agents via standard RPC with JSON-encoded
request/response — no special types needed on the caller side.

https://claude.ai/code/session_01QTp4SshuVmLAvvEGJe4TJd

* feat: generate agent alongside services, update all docs

micro run --prompt now generates an agent binary that manages all
the generated services. The agent reads MICRO_AI_PROVIDER and
MICRO_AI_API_KEY from the environment. micro run propagates these
when started with --prompt.

Run banner shows services and agents separately.

Updated README, getting-started guide, and landing page to show
the complete flow: generate → services + agent start → micro chat
routes to agent → agent orchestrates services.

https://claude.ai/code/session_01QTp4SshuVmLAvvEGJe4TJd

---------

Co-authored-by: Claude <noreply@anthropic.com>
2026-06-05 10:25:14 +01:00

214 行
5.4 KiB
Go

// Package flow provides event-driven LLM orchestration for go-micro
// services. A Flow subscribes to a broker topic, feeds each event
// into an LLM with all registered services as tools, and lets the
// model decide which RPCs to call.
//
// Usage:
//
// f := flow.New("onboard-user",
// flow.Trigger("events.user.created"),
// flow.Prompt("New user created: {{.Data}}. Send welcome email and create workspace."),
// flow.Provider("anthropic"),
// flow.APIKey(key),
// )
// f.Register(service)
// service.Run()
package flow
import (
"bytes"
"context"
"encoding/json"
"fmt"
"sync"
"text/template"
"time"
"go-micro.dev/v5/ai"
"go-micro.dev/v5/broker"
"go-micro.dev/v5/client"
"go-micro.dev/v5/logger"
"go-micro.dev/v5/registry"
// Register default providers.
_ "go-micro.dev/v5/ai/anthropic"
_ "go-micro.dev/v5/ai/atlascloud"
_ "go-micro.dev/v5/ai/gemini"
_ "go-micro.dev/v5/ai/groq"
_ "go-micro.dev/v5/ai/mistral"
_ "go-micro.dev/v5/ai/openai"
_ "go-micro.dev/v5/ai/together"
)
// Flow is an event-driven LLM orchestration unit. It subscribes to
// a broker topic, discovers services as tools, and feeds each event
// into an LLM that decides which RPCs to call.
type Flow struct {
name string
opts Options
model ai.Model
toolSet *ai.Tools
tmpl *template.Template
log logger.Logger
mu sync.Mutex
results []Result
}
// Result records one flow execution.
type Result struct {
FlowName string `json:"flow"`
Trigger string `json:"trigger"`
Prompt string `json:"prompt"`
Reply string `json:"reply,omitempty"`
Answer string `json:"answer,omitempty"`
ToolCalls []string `json:"tool_calls,omitempty"`
Error string `json:"error,omitempty"`
Timestamp time.Time `json:"timestamp"`
Duration float64 `json:"duration_seconds"`
}
// New creates a Flow with the given name and options.
func New(name string, opts ...Option) *Flow {
o := Options{
Provider: "openai",
SystemPrompt: "You are a service orchestrator. Use the available tools to fulfill the request. Explain what you do.",
HistoryLimit: 20,
}
for _, opt := range opts {
opt(&o)
}
var tmpl *template.Template
if o.Prompt != "" {
var err error
tmpl, err = template.New(name).Parse(o.Prompt)
if err != nil {
tmpl = template.Must(template.New(name).Parse("{{.Data}}"))
}
}
return &Flow{
name: name,
opts: o,
tmpl: tmpl,
log: logger.DefaultLogger,
}
}
// Register wires the flow into a running service. It sets up the
// model, discovers tools from the registry, and subscribes to the
// trigger topic on the broker. Call this before service.Run().
func (f *Flow) Register(reg registry.Registry, br broker.Broker, cl client.Client) error {
f.toolSet = ai.NewTools(reg, ai.ToolClient(cl))
var modelOpts []ai.Option
if f.opts.APIKey != "" {
modelOpts = append(modelOpts, ai.WithAPIKey(f.opts.APIKey))
}
if f.opts.Model != "" {
modelOpts = append(modelOpts, ai.WithModel(f.opts.Model))
}
if f.opts.BaseURL != "" {
modelOpts = append(modelOpts, ai.WithBaseURL(f.opts.BaseURL))
}
modelOpts = append(modelOpts, ai.WithTools(f.toolSet))
f.model = ai.New(f.opts.Provider, modelOpts...)
if f.model == nil {
return fmt.Errorf("unknown provider: %s", f.opts.Provider)
}
if f.opts.TriggerTopic != "" {
_, err := br.Subscribe(f.opts.TriggerTopic, func(p broker.Event) error {
data := string(p.Message().Body)
if err := f.Execute(context.Background(), data); err != nil {
f.log.Logf(logger.ErrorLevel, "Flow %s failed: %v", f.name, err)
}
return nil
})
if err != nil {
return fmt.Errorf("subscribe to %s: %w", f.opts.TriggerTopic, err)
}
f.log.Logf(logger.InfoLevel, "Flow %s subscribed to %s", f.name, f.opts.TriggerTopic)
}
return nil
}
// Execute runs the flow once with the given input data. This is
// called automatically on each broker event, but can also be
// invoked directly for testing or one-shot use.
func (f *Flow) Execute(ctx context.Context, data string) error {
start := time.Now()
discovered, err := f.toolSet.Discover()
if err != nil {
return fmt.Errorf("discover tools: %w", err)
}
prompt := data
if f.tmpl != nil {
var buf bytes.Buffer
f.tmpl.Execute(&buf, map[string]string{"Data": data})
prompt = buf.String()
}
resp, err := f.model.Generate(ctx, &ai.Request{
Prompt: prompt,
SystemPrompt: f.opts.SystemPrompt,
Tools: discovered,
})
result := Result{
FlowName: f.name,
Trigger: f.opts.TriggerTopic,
Prompt: prompt,
Timestamp: start,
Duration: time.Since(start).Seconds(),
}
if err != nil {
result.Error = err.Error()
f.record(result)
return err
}
result.Reply = resp.Reply
result.Answer = resp.Answer
for _, tc := range resp.ToolCalls {
args, _ := json.Marshal(tc.Input)
result.ToolCalls = append(result.ToolCalls, fmt.Sprintf("%s(%s)", tc.Name, args))
}
f.record(result)
f.log.Logf(logger.InfoLevel, "Flow %s completed in %.1fs: %d tool calls",
f.name, result.Duration, len(result.ToolCalls))
return nil
}
// Results returns a copy of all recorded execution results.
func (f *Flow) Results() []Result {
f.mu.Lock()
defer f.mu.Unlock()
out := make([]Result, len(f.results))
copy(out, f.results)
return out
}
// Name returns the flow name.
func (f *Flow) Name() string {
return f.name
}
func (f *Flow) record(r Result) {
f.mu.Lock()
f.results = append(f.results, r)
f.mu.Unlock()
if f.opts.OnResult != nil {
f.opts.OnResult(r)
}
}