Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| ae04469743 |
@@ -21,8 +21,12 @@ changes, architectural rewrites. Those go to the human.
|
||||
|
||||
## Work queue (ranked)
|
||||
|
||||
1. **Broaden provider-backed agent streaming end to end** ([#3857](https://github.com/micro/go-micro/issues/3857)) — #3861 closed the durable agent resume discoverability gap, and there are no open `codex` PRs in flight. Streaming is now the highest-value remaining Next-phase increment because the README, website, and roadmap all sell an interactive services → agents → workflows lifecycle; responsive `micro chat`, provider-backed `ai.Stream`, and A2A long-task UX are part of the adoption on-ramp, not just internal depth. Keep this scoped to additive provider-backed coverage, mock/live-provider tests, and discoverable support status so developers can trust the chat/agent/A2A path without a public-API rewrite.
|
||||
2. **Add an AP2 mandate layer over A2A and x402** ([#3552](https://github.com/micro/go-micro/issues/3552)) — this remains a forward interop investment, not a Now/Next blocker: Go Micro already has A2A agents and x402 paid tools, so a small signed-mandate foundation can keep agent payments aligned with the open-protocol story without pulling the queue ahead of adoption, resilience, streaming, or durable agent operation. Keep it additive and opt-in while the AP2/FIDO work settles.
|
||||
1. **Fix atlascloud plan-delegate missing delegated notify side effect** ([#3805](https://github.com/micro/go-micro/issues/3805)) — the universe concierge notification gap closed in #3811, so the remaining top live-harness blocker is the delegated-plan variant: agents can create the expected work but still miss the observable notify side effect. Keep this first because it protects the adoption-critical 0→hero promise with a real side-effect assertion rather than a weakened text-only pass.
|
||||
2. **Isolate file-store tests from shared default directory** ([#3751](https://github.com/micro/go-micro/issues/3751)) — repeated `go test -race -cover ./...` failures still point at file-store reliability under the shared default directory/table setup. This remains the highest-value store/CI stability issue because a flaky evaluator erodes the loop's ability to ship adoption work safely.
|
||||
3. **Stabilize file-store suffix expiry test timing** ([#3780](https://github.com/micro/go-micro/issues/3780)) — the newer store failure is a narrower timing-sensitive suffix-expiry assertion under `-race -cover`. Keep it adjacent to #3751 but separate because it may need a focused TTL/assertion fix even if directory isolation improves the broader file-store tests.
|
||||
4. **Propagate agent run cancellation and deadlines through model and tool calls** ([#3544](https://github.com/micro/go-micro/issues/3544)) — once the live harness and red CI blockers are cleared, the highest-value remaining resilience gap is predictable failure semantics across agent runs, model calls, tool calls, plan/delegate, and flow handoffs. Tool retries, live-provider deadline tuning, delegated-plan completion, and side-effect enforcement are in place; the lifecycle still needs cancellation/deadline propagation so work fails safely instead of becoming opaque loops.
|
||||
5. **Emit OpenTelemetry spans for agent run timelines** ([#3525](https://github.com/micro/go-micro/issues/3525)) — recent work made runs inspectable, correlated trace metadata through scheduled dispatch, verified restart resume, added opt-in tool retries, hardened provider conformance, and fixed provider-emitted text tool calls. The next Next-phase step is to turn that RunInfo foundation into standard OTel spans for agent runs, model calls, tool calls, checkpoint/resume, cancellation/deadlines, and failures.
|
||||
6. **Add an AP2 mandate layer over A2A and x402** ([#3552](https://github.com/micro/go-micro/issues/3552)) — this is a forward interop investment, not a Now-phase blocker: Go Micro already has A2A agents and x402 paid tools, so a small signed-mandate foundation can keep agent payments aligned with the open-protocol story without pulling the queue away from adoption, resilience, or observability. Keep it additive and opt-in while the AP2/FIDO work settles.
|
||||
|
||||
_Seeded by Claude Code from the roadmap + open issues; thereafter maintained by the
|
||||
architecture-review pass._
|
||||
|
||||
@@ -443,7 +443,6 @@ resp, _ := m.Generate(ctx, &ai.Request{Prompt: "hello"})
|
||||
- [multi-service](examples/multi-service/) — Multiple services in one binary
|
||||
- [mcp](examples/mcp/) — MCP integration with AI agents
|
||||
- [agent-plan-delegate](examples/agent-plan-delegate/) — Agent planning and multi-agent delegation
|
||||
- [agent-durable](examples/agent-durable/) — Checkpoint and resume an agent run without replaying completed tool side effects
|
||||
- [grpc-interop](examples/grpc-interop/) — Call go-micro from any gRPC client
|
||||
|
||||
See [all examples](examples/README.md).
|
||||
|
||||
@@ -45,7 +45,6 @@ const (
|
||||
AttrFlowStep = "agent.flow.step"
|
||||
AttrDispatch = "agent.dispatch"
|
||||
AttrTrigger = "agent.trigger"
|
||||
AttrRunEventKind = "agent.event.kind"
|
||||
)
|
||||
|
||||
type RunEvent struct {
|
||||
@@ -324,7 +323,6 @@ func runEventAttributes(e RunEvent) []attribute.KeyValue {
|
||||
attrs := []attribute.KeyValue{
|
||||
attribute.String(AttrRunID, e.RunID),
|
||||
attribute.String(AttrAgentName, e.Agent),
|
||||
attribute.String(AttrRunEventKind, e.Kind),
|
||||
}
|
||||
if e.ParentID != "" {
|
||||
attrs = append(attrs, attribute.String(AttrParentRunID, e.ParentID))
|
||||
|
||||
+1
-2
@@ -279,8 +279,7 @@ func spanEventHasRunInfo(events []trace.Event, name, runID, agentName string) bo
|
||||
continue
|
||||
}
|
||||
attrs := spanAttributes(event.Attributes)
|
||||
wantKind := strings.TrimPrefix(name, "agent.")
|
||||
if attrs[AttrRunID] == runID && attrs[AttrAgentName] == agentName && attrs[AttrRunEventKind] == wantKind {
|
||||
if attrs[AttrRunID] == runID && attrs[AttrAgentName] == agentName {
|
||||
return true
|
||||
}
|
||||
}
|
||||
|
||||
@@ -132,34 +132,6 @@ func TestToolCallTimeoutPropagatesDeadlineToCustomTool(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestAskCancellationDuringToolCallFailsRun(t *testing.T) {
|
||||
fakeGen = func(ctx context.Context, opts ai.Options, req *ai.Request) (*ai.Response, error) {
|
||||
if opts.ToolHandler == nil {
|
||||
t.Fatal("missing tool handler")
|
||||
}
|
||||
res := opts.ToolHandler(ctx, ai.ToolCall{ID: "call-1", Name: "cancel-self"})
|
||||
if !strings.Contains(res.Content, context.Canceled.Error()) {
|
||||
t.Fatalf("tool result = %q, want cancellation error", res.Content)
|
||||
}
|
||||
return &ai.Response{Reply: "should not succeed"}, nil
|
||||
}
|
||||
defer func() { fakeGen = nil }()
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
a := newTestAgent(
|
||||
Name("cancel-during-tool"),
|
||||
WithTool("cancel-self", "cancel the run context", nil, func(context.Context, map[string]any) (string, error) {
|
||||
cancel()
|
||||
return "", context.Canceled
|
||||
}),
|
||||
)
|
||||
|
||||
_, err := a.Ask(ctx, "cancel during tool")
|
||||
if !errors.Is(err, context.Canceled) {
|
||||
t.Fatalf("Ask error = %v, want context canceled", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAskCheckpointRecordsTerminalOperationalFailureStatus(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
|
||||
+4
-8
@@ -99,19 +99,15 @@ func GenerateWithRetry(ctx context.Context, m Model, req *Request, policy Genera
|
||||
}
|
||||
resp, err := m.Generate(callCtx, req, opts...)
|
||||
cancel()
|
||||
|
||||
// Caller cancellation/deadline always wins and is not retried, even if
|
||||
// a provider or tool loop swallowed the canceled tool result and returned
|
||||
// a final response. This keeps agent runs from appearing successful after
|
||||
// their controlling context was abandoned.
|
||||
if ctxErr := ctx.Err(); ctxErr != nil {
|
||||
return nil, ctxErr
|
||||
}
|
||||
if err == nil {
|
||||
return resp, nil
|
||||
}
|
||||
last = err
|
||||
|
||||
// Caller cancellation/deadline always wins and is not retried.
|
||||
if ctx.Err() != nil {
|
||||
return nil, ctx.Err()
|
||||
}
|
||||
transient := IsTransientError(err)
|
||||
if attempt == policy.MaxAttempts || !transient {
|
||||
if attempt > 1 || transient {
|
||||
|
||||
@@ -16,7 +16,6 @@ type preflightCheck struct {
|
||||
OK bool
|
||||
Detail string
|
||||
Fix string
|
||||
Next string
|
||||
}
|
||||
|
||||
type preflightDeps struct {
|
||||
@@ -53,9 +52,6 @@ func runAgentPreflight(w io.Writer, deps preflightDeps) error {
|
||||
if !check.OK && check.Fix != "" {
|
||||
fmt.Fprintf(w, " Fix: %s\n", check.Fix)
|
||||
}
|
||||
if !check.OK && check.Next != "" {
|
||||
fmt.Fprintf(w, " Next: %s\n", check.Next)
|
||||
}
|
||||
}
|
||||
if failures > 0 {
|
||||
return fmt.Errorf("first-agent preflight failed: %d check(s) need attention", failures)
|
||||
@@ -91,23 +87,19 @@ func agentPreflightChecks(deps preflightDeps) []preflightCheck {
|
||||
func checkGoToolchain(deps preflightDeps) preflightCheck {
|
||||
path, err := deps.lookPath("go")
|
||||
if err != nil {
|
||||
return preflightCheck{Name: "Go toolchain", Detail: "go was not found on PATH", Fix: "Install Go 1.24 or newer from https://go.dev/doc/install and ensure go is on PATH.", Next: "After installing Go, rerun micro agent preflight, then continue with docs/guides/your-first-agent.html."}
|
||||
return preflightCheck{Name: "Go toolchain", Fix: "Install Go 1.24 or newer and ensure go is on PATH."}
|
||||
}
|
||||
out, err := deps.commandOutput("go", "version")
|
||||
if err != nil {
|
||||
return preflightCheck{Name: "Go toolchain", Detail: strings.TrimSpace(string(out)), Fix: "Ensure the go command runs successfully (try `go version`) before starting the agent walkthrough.", Next: "Use docs/guides/debugging-agents.html after the toolchain check passes if an agent run still fails."}
|
||||
return preflightCheck{Name: "Go toolchain", Detail: strings.TrimSpace(string(out)), Fix: "Ensure the go command runs successfully."}
|
||||
}
|
||||
version := firstLine(out)
|
||||
if !goVersionAtLeast(version, 1, 24) {
|
||||
return preflightCheck{Name: "Go toolchain", Detail: fmt.Sprintf("%s (%s)", version, path), Fix: "Upgrade to Go 1.24 or newer before running generated services.", Next: "Rerun micro agent preflight, then continue with docs/guides/your-first-agent.html."}
|
||||
}
|
||||
return preflightCheck{Name: "Go toolchain", OK: true, Detail: fmt.Sprintf("%s (%s)", version, path)}
|
||||
return preflightCheck{Name: "Go toolchain", OK: true, Detail: fmt.Sprintf("%s (%s)", firstLine(out), path)}
|
||||
}
|
||||
|
||||
func checkMicroBinary(deps preflightDeps) preflightCheck {
|
||||
exe, err := deps.executable()
|
||||
if err != nil || exe == "" {
|
||||
return preflightCheck{Name: "micro binary", Detail: "micro executable path is unavailable", Fix: "Install the micro CLI or run this check through `go run ./cmd/micro agent preflight` from the repository.", Next: "Then follow docs/getting-started.html for the scaffold -> run path."}
|
||||
return preflightCheck{Name: "micro binary", Fix: "Install the micro CLI or run this check through go run ./cmd/micro agent preflight."}
|
||||
}
|
||||
version := deps.version()
|
||||
if version == "" {
|
||||
@@ -125,7 +117,7 @@ func checkProviderKey(deps preflightDeps) preflightCheck {
|
||||
}
|
||||
}
|
||||
if len(found) == 0 {
|
||||
return preflightCheck{Name: "provider API key", Detail: "no supported provider key found", Fix: "Export MICRO_AI_API_KEY or a provider key such as ANTHROPIC_API_KEY before running provider-backed agents.", Next: "For a no-secret path, run the mock-model walkthrough in docs/guides/no-secret-first-agent.html; for real providers, see docs/guides/debugging-agents.html#provider-failures."}
|
||||
return preflightCheck{Name: "provider API key", Detail: "no supported provider key found", Fix: "Export MICRO_AI_API_KEY or a provider key such as ANTHROPIC_API_KEY before running provider-backed agents."}
|
||||
}
|
||||
return preflightCheck{Name: "provider API key", OK: true, Detail: "found " + strings.Join(found, ", ")}
|
||||
}
|
||||
@@ -133,7 +125,7 @@ func checkProviderKey(deps preflightDeps) preflightCheck {
|
||||
func checkPortAvailable(deps preflightDeps, addr, use string) preflightCheck {
|
||||
ln, err := deps.listen("tcp", addr)
|
||||
if err != nil {
|
||||
return preflightCheck{Name: "local port " + addr, Detail: "busy or unavailable for " + use, Fix: "Stop the process using " + addr + " (for example, `lsof -i :8080`) or run `micro run --address` with a free port.", Next: "Once the gateway starts, open http://localhost:8080/agent or continue with docs/guides/your-first-agent.html#chat-with-your-agent."}
|
||||
return preflightCheck{Name: "local port " + addr, Detail: "busy or unavailable for " + use, Fix: "Stop the process using " + addr + " or run micro run --address with a free port."}
|
||||
}
|
||||
_ = ln.Close()
|
||||
return preflightCheck{Name: "local port " + addr, OK: true, Detail: "available for " + use}
|
||||
@@ -146,18 +138,3 @@ func firstLine(b []byte) string {
|
||||
}
|
||||
return s
|
||||
}
|
||||
|
||||
func goVersionAtLeast(line string, wantMajor, wantMinor int) bool {
|
||||
idx := strings.Index(line, "go1.")
|
||||
if idx < 0 {
|
||||
return false
|
||||
}
|
||||
var major, minor int
|
||||
if _, err := fmt.Sscanf(line[idx:], "go%d.%d", &major, &minor); err != nil {
|
||||
return false
|
||||
}
|
||||
if major != wantMajor {
|
||||
return major > wantMajor
|
||||
}
|
||||
return minor >= wantMinor
|
||||
}
|
||||
|
||||
@@ -61,59 +61,13 @@ func TestRunAgentPreflightReportsActionableFailures(t *testing.T) {
|
||||
t.Fatal("runAgentPreflight() error = nil")
|
||||
}
|
||||
got := out.String()
|
||||
for _, want := range []string{"✗ Go toolchain", "go was not found on PATH", "https://go.dev/doc/install", "docs/guides/your-first-agent.html", "✗ micro binary", "go run ./cmd/micro agent preflight", "✗ provider API key", "docs/guides/no-secret-first-agent.html", "docs/guides/debugging-agents.html#provider-failures", "✗ local port :8080", "lsof -i :8080", "micro run --address"} {
|
||||
for _, want := range []string{"✗ Go toolchain", "Install Go 1.24", "✗ micro binary", "✗ provider API key", "ANTHROPIC_API_KEY", "✗ local port :8080", "micro run --address"} {
|
||||
if !strings.Contains(got, want) {
|
||||
t.Fatalf("output missing %q:\n%s", want, got)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestRunAgentPreflightReportsOldGoVersion(t *testing.T) {
|
||||
deps := preflightDeps{
|
||||
lookPath: func(name string) (string, error) { return "/usr/bin/" + name, nil },
|
||||
commandOutput: func(name string, args ...string) ([]byte, error) {
|
||||
return []byte("go version go1.23.9 linux/amd64\n"), nil
|
||||
},
|
||||
executable: func() (string, error) { return "/usr/local/bin/micro", nil },
|
||||
getenv: func(key string) string {
|
||||
if key == "ANTHROPIC_API_KEY" {
|
||||
return "set"
|
||||
}
|
||||
return ""
|
||||
},
|
||||
listen: func(network, address string) (net.Listener, error) { return stubListener{}, nil },
|
||||
}
|
||||
|
||||
var out bytes.Buffer
|
||||
err := runAgentPreflight(&out, deps)
|
||||
if err == nil {
|
||||
t.Fatal("runAgentPreflight() error = nil")
|
||||
}
|
||||
got := out.String()
|
||||
for _, want := range []string{"✗ Go toolchain", "go1.23.9", "Upgrade to Go 1.24 or newer", "Rerun micro agent preflight"} {
|
||||
if !strings.Contains(got, want) {
|
||||
t.Fatalf("output missing %q:\n%s", want, got)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestGoVersionAtLeast(t *testing.T) {
|
||||
tests := []struct {
|
||||
line string
|
||||
want bool
|
||||
}{
|
||||
{line: "go version go1.24.0 linux/amd64", want: true},
|
||||
{line: "go version go1.25.1 linux/amd64", want: true},
|
||||
{line: "go version go1.23.9 linux/amd64", want: false},
|
||||
{line: "unexpected", want: false},
|
||||
}
|
||||
for _, tt := range tests {
|
||||
if got := goVersionAtLeast(tt.line, 1, 24); got != tt.want {
|
||||
t.Fatalf("goVersionAtLeast(%q) = %v, want %v", tt.line, got, tt.want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestFirstLine(t *testing.T) {
|
||||
if got := firstLine([]byte("one\ntwo")); got != "one" {
|
||||
t.Fatalf("firstLine() = %q", got)
|
||||
|
||||
@@ -148,17 +148,13 @@ func main() {
|
||||
fmt.Fprintf(os.Stderr, "content-type = %q, want text/event-stream\n", ct)
|
||||
os.Exit(1)
|
||||
}
|
||||
summary, err := readSSESummary(res.Body)
|
||||
payload, err := readSSEData(res.Body)
|
||||
if err != nil {
|
||||
fmt.Fprintln(os.Stderr, err)
|
||||
os.Exit(1)
|
||||
}
|
||||
if summary.State != "completed" {
|
||||
fmt.Fprintf(os.Stderr, "stream final state = %q, want completed; payload: %s\n", summary.State, summary.Payload)
|
||||
os.Exit(1)
|
||||
}
|
||||
if !summary.HasArtifactText {
|
||||
fmt.Fprintf(os.Stderr, "stream completed without artifact text: %s\n", summary.Payload)
|
||||
if !strings.Contains(payload, "a2a-fallback-ok") {
|
||||
fmt.Fprintf(os.Stderr, "stream payload missing marker: %s\n", payload)
|
||||
os.Exit(1)
|
||||
}
|
||||
if !sawTool || !sawRunInfo {
|
||||
@@ -168,16 +164,10 @@ func main() {
|
||||
fmt.Println("\n\033[32m✓ A2A message/stream fell back to Ask and preserved tool/run metadata\033[0m")
|
||||
}
|
||||
|
||||
type streamSummary struct {
|
||||
Payload string
|
||||
State string
|
||||
HasArtifactText bool
|
||||
}
|
||||
|
||||
func readSSESummary(r io.Reader) (streamSummary, error) {
|
||||
func readSSEData(r io.Reader) (string, error) {
|
||||
scanner := bufio.NewScanner(r)
|
||||
var event strings.Builder
|
||||
var summary streamSummary
|
||||
var payload strings.Builder
|
||||
seen := false
|
||||
flush := func() error {
|
||||
data := strings.TrimSpace(event.String())
|
||||
@@ -185,40 +175,19 @@ func readSSESummary(r io.Reader) (streamSummary, error) {
|
||||
if data == "" {
|
||||
return nil
|
||||
}
|
||||
var envelope struct {
|
||||
Result struct {
|
||||
Status struct {
|
||||
State string `json:"state"`
|
||||
} `json:"status"`
|
||||
Artifacts []struct {
|
||||
Parts []struct {
|
||||
Text string `json:"text"`
|
||||
} `json:"parts"`
|
||||
} `json:"artifacts"`
|
||||
} `json:"result"`
|
||||
}
|
||||
if err := json.Unmarshal([]byte(data), &envelope); err != nil {
|
||||
if !json.Valid([]byte(data)) {
|
||||
return fmt.Errorf("SSE data event is not JSON: %s", data)
|
||||
}
|
||||
seen = true
|
||||
summary.Payload += data + "\n"
|
||||
if envelope.Result.Status.State != "" {
|
||||
summary.State = envelope.Result.Status.State
|
||||
}
|
||||
for _, artifact := range envelope.Result.Artifacts {
|
||||
for _, part := range artifact.Parts {
|
||||
if strings.TrimSpace(part.Text) != "" {
|
||||
summary.HasArtifactText = true
|
||||
}
|
||||
}
|
||||
}
|
||||
payload.WriteString(data)
|
||||
payload.WriteByte('\n')
|
||||
return nil
|
||||
}
|
||||
for scanner.Scan() {
|
||||
line := scanner.Text()
|
||||
if strings.TrimSpace(line) == "" {
|
||||
if err := flush(); err != nil {
|
||||
return streamSummary{}, err
|
||||
return "", err
|
||||
}
|
||||
continue
|
||||
}
|
||||
@@ -232,13 +201,13 @@ func readSSESummary(r io.Reader) (streamSummary, error) {
|
||||
event.WriteString(strings.TrimSpace(data))
|
||||
}
|
||||
if err := scanner.Err(); err != nil {
|
||||
return streamSummary{}, err
|
||||
return "", err
|
||||
}
|
||||
if err := flush(); err != nil {
|
||||
return streamSummary{}, err
|
||||
return "", err
|
||||
}
|
||||
if !seen {
|
||||
return streamSummary{}, errors.New("no SSE data received")
|
||||
return "", errors.New("no SSE data received")
|
||||
}
|
||||
return summary, nil
|
||||
return payload.String(), nil
|
||||
}
|
||||
|
||||
@@ -5,26 +5,19 @@ import (
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestReadSSESummaryUsesCompletedTaskInvariants(t *testing.T) {
|
||||
summary, err := readSSESummary(strings.NewReader("data: {\"jsonrpc\":\"2.0\",\"result\":{\"status\":{\"state\":\"working\"}}}\n\n" +
|
||||
"data: {\"jsonrpc\":\"2.0\",\"result\":{\"status\":{\"state\":\"completed\"},\"artifacts\":[{\"parts\":[{\"kind\":\"text\",\"text\":\"provider-specific answer\"}]}]}}\n\n"))
|
||||
func TestReadSSEDataAcceptsMultipleJSONEvents(t *testing.T) {
|
||||
payload, err := readSSEData(strings.NewReader("event: status\ndata: {\"phase\":\"started\"}\n\ndata:{\"result\":\"a2a-fallback-ok\"}\n\n"))
|
||||
if err != nil {
|
||||
t.Fatalf("readSSESummary() error = %v", err)
|
||||
t.Fatalf("readSSEData returned error: %v", err)
|
||||
}
|
||||
if summary.State != "completed" {
|
||||
t.Fatalf("State = %q, want completed", summary.State)
|
||||
}
|
||||
if !summary.HasArtifactText {
|
||||
t.Fatal("HasArtifactText = false, want true")
|
||||
}
|
||||
if strings.Contains(summary.Payload, "a2a-fallback-ok") {
|
||||
t.Fatalf("test fixture should not rely on marker text: %s", summary.Payload)
|
||||
if !strings.Contains(payload, "started") || !strings.Contains(payload, "a2a-fallback-ok") {
|
||||
t.Fatalf("payload = %q, want both event payloads", payload)
|
||||
}
|
||||
}
|
||||
|
||||
func TestReadSSESummaryRejectsNonJSONData(t *testing.T) {
|
||||
_, err := readSSESummary(strings.NewReader("data: not-json\n\n"))
|
||||
func TestReadSSEDataRejectsInvalidJSONEvent(t *testing.T) {
|
||||
_, err := readSSEData(strings.NewReader("data: {\"ok\":true}\n\ndata: {bad json}\n\n"))
|
||||
if err == nil {
|
||||
t.Fatal("readSSESummary() error = nil, want non-JSON error")
|
||||
t.Fatal("readSSEData returned nil error for invalid JSON event")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -127,10 +127,6 @@ func (s *TaskService) count() int {
|
||||
return len(s.tasks)
|
||||
}
|
||||
|
||||
const delegatedNotifyTask = "Use the notify Send tool exactly once to tell owner@acme.com: The launch plan is ready. Do not answer until the notify tool call has succeeded."
|
||||
|
||||
const commsPrompt = "You handle outbound notifications. When asked to notify someone, you must call the notify Send tool exactly once before replying. Never claim a notification was sent unless the notify tool returned success."
|
||||
|
||||
type SendRequest struct {
|
||||
To string `json:"to" description:"Recipient address"`
|
||||
Message string `json:"message" description:"Message body"`
|
||||
@@ -279,12 +275,12 @@ func (m *mockModel) Generate(ctx context.Context, req *ai.Request, _ ...ai.Gener
|
||||
if m.unknownDelegateOnce && !m.emittedUnknownDelegate {
|
||||
m.emittedUnknownDelegate = true
|
||||
m.call("conductor", "atlascloud_delegate", map[string]any{
|
||||
"task": delegatedNotifyTask,
|
||||
"task": "Notify owner@acme.com that the launch plan is ready",
|
||||
"to": "comms",
|
||||
})
|
||||
} else {
|
||||
m.call("conductor", del, map[string]any{
|
||||
"task": delegatedNotifyTask,
|
||||
"task": "Notify owner@acme.com that the launch plan is ready",
|
||||
"to": "comms",
|
||||
})
|
||||
}
|
||||
@@ -358,7 +354,7 @@ func runPlanDelegate(provider string) error {
|
||||
agent.Name("comms"),
|
||||
agent.Address("127.0.0.1:0"),
|
||||
agent.Services("notify"),
|
||||
agent.Prompt(commsPrompt),
|
||||
agent.Prompt("You handle outbound notifications. Use the notify service."),
|
||||
agent.Provider(provider), agent.APIKey(apiKey),
|
||||
agent.WithRegistry(reg), agent.WithClient(cl), agent.WithStore(mem),
|
||||
agent.WithCheckpoint(commsCheckpoint),
|
||||
@@ -418,7 +414,7 @@ func runPlanDelegate(provider string) error {
|
||||
}()
|
||||
|
||||
if err := waitForPlanDelegateExecution(executeDone, taskSvc, notifySvc, func(ctx context.Context) error {
|
||||
_, err := conductor.Ask(ctx, "The Design, Build, and Ship tasks already exist, but the owner notification is still missing. Delegate exactly one notification to the \"comms\" agent now with this exact subtask: "+delegatedNotifyTask+" Do not create more tasks and do not answer until comms has handled the notification.")
|
||||
_, err := conductor.Ask(ctx, "The Design, Build, and Ship tasks already exist, but the owner notification is still missing. Delegate exactly one notification to the \"comms\" agent now: ask comms to notify owner@acme.com that the launch plan is ready. Do not create more tasks and do not answer until comms has handled the notification.")
|
||||
return err
|
||||
}); err != nil {
|
||||
return err
|
||||
@@ -451,12 +447,9 @@ func waitForPlanDelegateExecution(done <-chan error, taskSvc *TaskService, notif
|
||||
tasks := taskSvc.count()
|
||||
notify := notifySvc.count()
|
||||
if err != nil {
|
||||
if isClientTimeout(err) {
|
||||
if tasks == 3 && notify == 1 {
|
||||
fmt.Printf("\n\033[33mwarning:\033[0m flow execute returned after completed side effects: %v\n", err)
|
||||
return nil
|
||||
}
|
||||
return classifiedPlanDelegateTimeout(tasks, notify, err)
|
||||
if isClientTimeout(err) && tasks == 3 && notify == 1 {
|
||||
fmt.Printf("\n\033[33mwarning:\033[0m flow execute returned after completed side effects: %v\n", err)
|
||||
return nil
|
||||
}
|
||||
return fmt.Errorf("flow execute after side effects tasks=%d notify=%d: %w", tasks, notify, err)
|
||||
}
|
||||
@@ -484,10 +477,6 @@ func waitForPlanDelegateExecution(done <-chan error, taskSvc *TaskService, notif
|
||||
}
|
||||
}
|
||||
|
||||
func classifiedPlanDelegateTimeout(tasks, notify int, err error) error {
|
||||
return fmt.Errorf("provider latency/outage during plan-delegate before required side effects completed (tasks=%d/3 notify=%d/1); retry live provider or inspect provider logs if this recurs: %w", tasks, notify, err)
|
||||
}
|
||||
|
||||
func isClientTimeout(err error) bool {
|
||||
msg := strings.ToLower(err.Error())
|
||||
return strings.Contains(msg, "request timeout") || strings.Contains(msg, "code=408") || strings.Contains(msg, "code\":408")
|
||||
|
||||
@@ -67,7 +67,7 @@ func TestPlanDelegateEndToEnd(t *testing.T) {
|
||||
agent.Name("comms"),
|
||||
agent.Address("127.0.0.1:0"),
|
||||
agent.Services("notify"),
|
||||
agent.Prompt(commsPrompt),
|
||||
agent.Prompt("You handle outbound notifications."),
|
||||
agent.Provider("mock"),
|
||||
agent.WithRegistry(reg),
|
||||
agent.WithClient(cl),
|
||||
@@ -149,7 +149,7 @@ func TestFlowDispatchesToAgentEndToEnd(t *testing.T) {
|
||||
agent.Name("comms"),
|
||||
agent.Address("127.0.0.1:0"),
|
||||
agent.Services("notify"),
|
||||
agent.Prompt(commsPrompt),
|
||||
agent.Prompt("You handle outbound notifications."),
|
||||
agent.Provider("mock"),
|
||||
agent.WithRegistry(reg),
|
||||
agent.WithClient(cl),
|
||||
@@ -338,7 +338,7 @@ func TestPlanDelegateExecutionAcceptsClientTimeoutAfterSideEffects(t *testing.T)
|
||||
}
|
||||
}
|
||||
|
||||
func TestPlanDelegateExecutionClassifiesClientTimeoutBeforeSideEffects(t *testing.T) {
|
||||
func TestPlanDelegateExecutionRejectsClientTimeoutBeforeSideEffects(t *testing.T) {
|
||||
done := make(chan error, 1)
|
||||
done <- errors.New(`{"id":"go.micro.client","code":408,"detail":"<nil>","status":"Request Timeout"}`)
|
||||
|
||||
@@ -346,35 +346,8 @@ func TestPlanDelegateExecutionClassifiesClientTimeoutBeforeSideEffects(t *testin
|
||||
if err == nil {
|
||||
t.Fatal("waitForPlanDelegateExecution returned nil, want timeout before side effects to fail")
|
||||
}
|
||||
for _, want := range []string{
|
||||
"provider latency/outage during plan-delegate",
|
||||
"tasks=0/3 notify=0/1",
|
||||
"retry live provider or inspect provider logs",
|
||||
"Request Timeout",
|
||||
} {
|
||||
if got := err.Error(); !strings.Contains(got, want) {
|
||||
t.Fatalf("error = %q, want %q", got, want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestPlanDelegateExecutionClassifiesPartialClientTimeout(t *testing.T) {
|
||||
taskSvc := new(TaskService)
|
||||
for _, title := range []string{"Design", "Build", "Ship"} {
|
||||
var rsp AddResponse
|
||||
if err := taskSvc.Add(context.Background(), &AddRequest{Title: title}, &rsp); err != nil {
|
||||
t.Fatalf("Add(%q): %v", title, err)
|
||||
}
|
||||
}
|
||||
done := make(chan error, 1)
|
||||
done <- errors.New(`{"id":"go.micro.client","code":408,"detail":"<nil>","status":"Request Timeout"}`)
|
||||
|
||||
err := waitForPlanDelegateExecution(done, taskSvc, new(NotifyService), nil)
|
||||
if err == nil {
|
||||
t.Fatal("waitForPlanDelegateExecution returned nil, want timeout before notify to fail")
|
||||
}
|
||||
if got := err.Error(); !strings.Contains(got, "tasks=3/3 notify=0/1") {
|
||||
t.Fatalf("error = %q, want partial side-effect counts", got)
|
||||
if got := err.Error(); !strings.Contains(got, "tasks=0 notify=0") {
|
||||
t.Fatalf("error = %q, want side-effect counts", got)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -16,7 +16,6 @@ agents → workflows lifecycle.
|
||||
| First service-backed agent | [`examples/agent-demo`](https://github.com/micro/go-micro/tree/master/examples/agent-demo) | Multi-service project/task/team app with agent playground integration. |
|
||||
| 0→hero lifecycle | [`examples/support`](https://github.com/micro/go-micro/tree/master/examples/support) | No-secret support-desk story: typed services, an agent, an event-driven flow, and a guardrail. |
|
||||
| Planning and delegation | [`examples/agent-plan-delegate`](https://github.com/micro/go-micro/tree/master/examples/agent-plan-delegate) | Two agents collaborate through `plan` and `delegate` over normal Go Micro RPC. |
|
||||
| Durable agent runs | [`examples/agent-durable`](https://github.com/micro/go-micro/tree/master/examples/agent-durable) | Checkpoint and resume a model-directed run without replaying completed tool side effects. |
|
||||
| Durable workflows | [`examples/flow-durable`](https://github.com/micro/go-micro/tree/master/examples/flow-durable) | Ordered, checkpointed flow steps resume without duplicating completed side effects. |
|
||||
| AI-callable services | [`examples/mcp`](https://github.com/micro/go-micro/tree/master/examples/mcp) | MCP examples that expose service endpoints as model tools. |
|
||||
|
||||
@@ -36,11 +35,7 @@ agents → workflows lifecycle.
|
||||
[`examples/agent-plan-delegate`](https://github.com/micro/go-micro/tree/master/examples/agent-plan-delegate).
|
||||
- [Agents and Workflows](../guides/agents-and-workflows.html) → run
|
||||
[`examples/flow-durable`](https://github.com/micro/go-micro/tree/master/examples/flow-durable)
|
||||
for deterministic checkpointed steps,
|
||||
[`examples/agent-durable`](https://github.com/micro/go-micro/tree/master/examples/agent-durable)
|
||||
for model-directed checkpointed runs, and
|
||||
[`examples/support`](https://github.com/micro/go-micro/tree/master/examples/support)
|
||||
for the full services → agents → workflows lifecycle.
|
||||
and [`examples/support`](https://github.com/micro/go-micro/tree/master/examples/support).
|
||||
|
||||
## Repository examples
|
||||
|
||||
|
||||
@@ -25,7 +25,7 @@ Before your first provider-backed agent run, check the local path with:
|
||||
micro agent preflight
|
||||
```
|
||||
|
||||
The preflight is read-only: it verifies Go 1.24+, the `micro` binary, provider-key setup, and whether the default `micro run` gateway port is free, without calling an LLM provider. When a check fails it prints the exact fix plus the next guide to open, so the scaffold → run → chat path stays walkable.
|
||||
The preflight is read-only: it verifies Go, the `micro` binary, provider-key setup, and whether the default `micro run` gateway port is free, without calling an LLM provider.
|
||||
|
||||
## Install
|
||||
|
||||
|
||||
@@ -17,9 +17,7 @@ micro inspect ... # read the recorded run or workflow history
|
||||
|
||||
Debug the lifecycle in the same order Go Micro runs it: first prove the service is
|
||||
registered and callable, then inspect the agent run that chose tools, then inspect
|
||||
any workflow that handed off to the agent. If the first local run fails before a
|
||||
chat turn, run `micro agent preflight`; failed checks include `Fix:` and `Next:`
|
||||
lines for Go, CLI installation, provider-key setup, and the local gateway port.
|
||||
any workflow that handed off to the agent.
|
||||
|
||||
## 1. Reproduce one small turn
|
||||
|
||||
|
||||
@@ -89,6 +89,6 @@ CI keeps those CLI boundaries present with:
|
||||
go test ./cmd/micro -run TestFirstAgentWalkthroughCLIBoundaries -count=1
|
||||
```
|
||||
|
||||
If `micro agent preflight` reports a missing provider key, you can still use this no-secret path because it runs against the mock model; the command now prints this guide as the next step for that failure. If chat behaves unexpectedly, continue to
|
||||
If chat behaves unexpectedly, continue to
|
||||
[Debugging your agent](debugging-agents.html) for provider checks, run history,
|
||||
memory, and tool-call inspection.
|
||||
|
||||
@@ -53,7 +53,7 @@ Run the read-only first-agent preflight before starting the walkthrough. The sam
|
||||
micro agent preflight
|
||||
```
|
||||
|
||||
It checks Go 1.24+, the `micro` binary, provider-key setup, and the default local gateway port without contacting a provider. Failed checks include a `Fix:` line and a `Next:` line that points back to this guide, the no-secret walkthrough, or the debugging guide.
|
||||
It checks Go, the `micro` binary, provider-key setup, and the default local gateway port without contacting a provider.
|
||||
|
||||
## 1. Create a workspace
|
||||
|
||||
|
||||
+47
-32
@@ -3,6 +3,7 @@ package store
|
||||
import (
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
@@ -11,20 +12,15 @@ import (
|
||||
"github.com/kr/pretty"
|
||||
)
|
||||
|
||||
func newTestFileStore(t *testing.T, opts ...Option) Store {
|
||||
t.Helper()
|
||||
opts = append(opts, DirOption(t.TempDir()))
|
||||
s := NewStore(opts...)
|
||||
t.Cleanup(func() {
|
||||
if err := s.Close(); err != nil {
|
||||
t.Errorf("failed to close file store: %v", err)
|
||||
}
|
||||
})
|
||||
return s
|
||||
func cleanup(db string, s Store) {
|
||||
s.Close()
|
||||
dir := filepath.Join(DefaultDir, db+"/")
|
||||
os.RemoveAll(dir)
|
||||
}
|
||||
|
||||
func TestFileStoreReInit(t *testing.T) {
|
||||
s := newTestFileStore(t, Table("aaa"))
|
||||
s := NewStore(Table("aaa"))
|
||||
defer cleanup(DefaultDatabase, s)
|
||||
s.Init(Table("bbb"))
|
||||
if s.Options().Table != "bbb" {
|
||||
t.Error("Init didn't reinitialise the store")
|
||||
@@ -32,22 +28,26 @@ func TestFileStoreReInit(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestFileStoreBasic(t *testing.T) {
|
||||
s := newTestFileStore(t)
|
||||
s := NewStore()
|
||||
defer cleanup(DefaultDatabase, s)
|
||||
fileTest(s, t)
|
||||
}
|
||||
|
||||
func TestFileStoreTable(t *testing.T) {
|
||||
s := newTestFileStore(t, Table("testTable"))
|
||||
s := NewStore(Table("testTable"))
|
||||
defer cleanup(DefaultDatabase, s)
|
||||
fileTest(s, t)
|
||||
}
|
||||
|
||||
func TestFileStoreDatabase(t *testing.T) {
|
||||
s := newTestFileStore(t, Database("testdb"))
|
||||
s := NewStore(Database("testdb"))
|
||||
defer cleanup("testdb", s)
|
||||
fileTest(s, t)
|
||||
}
|
||||
|
||||
func TestFileStoreDatabaseTable(t *testing.T) {
|
||||
s := newTestFileStore(t, Table("testTable"), Database("testdb"))
|
||||
s := NewStore(Table("testTable"), Database("testdb"))
|
||||
defer cleanup("testdb", s)
|
||||
fileTest(s, t)
|
||||
}
|
||||
|
||||
@@ -135,22 +135,22 @@ func fileTest(s Store, t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// Write records with suffix matches and an already-expired record. Avoid
|
||||
// wall-clock boundary sleeps here: under -race/-cover, sleeping exactly the
|
||||
// TTL made this assertion flaky on slower CI runners.
|
||||
// Write 3 records with various expiry and get with Suffix
|
||||
records = []*Record{
|
||||
{
|
||||
Key: "foo",
|
||||
Value: []byte("foofoo"),
|
||||
},
|
||||
{
|
||||
Key: "barfoo",
|
||||
Value: []byte("barfoobarfoo"),
|
||||
Expiry: -time.Second,
|
||||
Key: "barfoo",
|
||||
Value: []byte("barfoobarfoo"),
|
||||
|
||||
Expiry: time.Millisecond * 100,
|
||||
},
|
||||
{
|
||||
Key: "bazbarfoo",
|
||||
Value: []byte("bazbarfoobazbarfoo"),
|
||||
Key: "bazbarfoo",
|
||||
Value: []byte("bazbarfoobazbarfoo"),
|
||||
Expiry: 2 * time.Millisecond * 100,
|
||||
},
|
||||
}
|
||||
for _, r := range records {
|
||||
@@ -160,24 +160,39 @@ func fileTest(s Store, t *testing.T) {
|
||||
}
|
||||
if results, err := s.Read("foo", ReadSuffix()); err != nil {
|
||||
t.Errorf("Couldn't read all \"foo\" keys, got %# v (%s)", spew.Sdump(results), err)
|
||||
} else if len(results) != 2 {
|
||||
t.Errorf("Expected 2 unexpired suffix items, got %d (%# v)", len(results), spew.Sdump(results))
|
||||
}
|
||||
if err := s.Delete("bazbarfoo"); err != nil {
|
||||
t.Errorf("Delete failed (%v)", err)
|
||||
} else {
|
||||
if len(results) != 3 {
|
||||
t.Errorf("Expected 3 items, got %d", len(results))
|
||||
// t.Logf("Table test: %v\n", spew.Sdump(results))
|
||||
}
|
||||
}
|
||||
time.Sleep(time.Millisecond * 100)
|
||||
if results, err := s.Read("foo", ReadSuffix()); err != nil {
|
||||
t.Errorf("Couldn't read all \"foo\" keys, got %# v (%s)", spew.Sdump(results), err)
|
||||
} else if len(results) != 1 {
|
||||
t.Errorf("Expected 1 unexpired suffix item, got %d (%# v)", len(results), spew.Sdump(results))
|
||||
} else {
|
||||
if len(results) != 2 {
|
||||
t.Errorf("Expected 2 items, got %d", len(results))
|
||||
// t.Logf("Table test: %v\n", spew.Sdump(results))
|
||||
}
|
||||
}
|
||||
time.Sleep(time.Millisecond * 100)
|
||||
if results, err := s.Read("foo", ReadSuffix()); err != nil {
|
||||
t.Errorf("Couldn't read all \"foo\" keys, got %# v (%s)", spew.Sdump(results), err)
|
||||
} else {
|
||||
if len(results) != 1 {
|
||||
t.Errorf("Expected 1 item, got %d", len(results))
|
||||
// t.Logf("Table test: %# v\n", spew.Sdump(results))
|
||||
}
|
||||
}
|
||||
if err := s.Delete("foo"); err != nil {
|
||||
t.Errorf("Delete failed (%v)", err)
|
||||
}
|
||||
if results, err := s.Read("foo", ReadSuffix()); err != nil {
|
||||
t.Errorf("Couldn't read all \"foo\" keys, got %# v (%s)", spew.Sdump(results), err)
|
||||
} else if len(results) != 0 {
|
||||
t.Errorf("Expected 0 items, got %d (%# v)", len(results), spew.Sdump(results))
|
||||
} else {
|
||||
if len(results) != 0 {
|
||||
t.Errorf("Expected 0 items, got %d (%# v)", len(results), spew.Sdump(results))
|
||||
}
|
||||
}
|
||||
|
||||
// Test Table, Suffix and WriteOptions
|
||||
|
||||
Reference in New Issue
Block a user