micro--go-micro
1bc886fa82
* 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>
214 行
5.4 KiB
Go
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)
|
|
}
|
|
}
|