Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| add3e0ca81 |
@@ -21,8 +21,8 @@ changes, architectural rewrites. Those go to the human.
|
||||
|
||||
## Work queue (ranked)
|
||||
|
||||
1. **Emit OpenTelemetry spans from agent run history** ([#4218](https://github.com/micro/go-micro/issues/4218)) — streaming coverage shipped via #4226, so close the next agentic-depth operability seam called out by the roadmap and blog: `RunInfo`/history and production traces should tell the same story for model steps, tool calls, retries, delegation, and failures, making `micro inspect` and deployed traces feel like one debugging path.
|
||||
2. **Add examples wayfinding index for first-agent adoption** ([#4223](https://github.com/micro/go-micro/issues/4223)) — keep developer adoption weighted with internal hardening: the README and getting-started guide now have a stronger first-agent path, but examples remain spread across docs, CLI output, and directories. A single CI-guarded examples map should make the smallest no-secret agent, the 0→hero support app, and next interop examples discoverable from one place.
|
||||
1. **Verify no-secret agent debugging walkthrough** ([#4142](https://github.com/micro/go-micro/issues/4142)) — with the plan-delegate completion regression fixed in #4146, the highest-value remaining Now-phase gap is the developer adoption path immediately after the first chat. Extend the maintained provider-free 0→1 route into `micro inspect agent`, run history, memory, and provider checks so renamed commands or stale debugging docs fail in CI where new developers most need confidence.
|
||||
2. **Add durable agent checkpoint resume smoke coverage** ([#4148](https://github.com/micro/go-micro/issues/4148)) — once the first-agent debugging seam is protected, move to the top Next-phase harness gap: prove an interrupted agent run can resume from persisted state with enough run/step history for inspect/debugging. This keeps the services → agents → workflows lifecycle cohesive by giving agents the same durability story flows already have, without taking on a breaking API redesign.
|
||||
|
||||
_Seeded by Claude Code from the roadmap + open issues; thereafter maintained by the
|
||||
architecture-review pass._
|
||||
|
||||
@@ -16,35 +16,17 @@ next version when it ships.
|
||||
|
||||
## [Unreleased]
|
||||
|
||||
### Added
|
||||
- **StreamAsk close cancellation** — agent streaming calls now cancel promptly when their runner closes, avoiding orphaned stream work. (`agent/`)
|
||||
- **Agent resume pending helper** — agent durability now has a focused helper for resuming pending checkpointed runs. (`agent/`)
|
||||
|
||||
### Fixed
|
||||
- **Plan/delegate retry idempotency** — agent retries now preserve side-effect and notification dedupe across conformance retry paths, including completion and owner-notification edge cases. (`agent/`, `internal/harness/`)
|
||||
|
||||
---
|
||||
|
||||
## [6.3.17] - July 2026
|
||||
|
||||
### Added
|
||||
- **First-agent examples CLI wayfinding** — `micro examples` now prints the maintained provider-free first-agent examples in copy/paste order. (`cmd/micro/`)
|
||||
- **0→hero CLI entrypoint** — `micro zero-to-hero` now points developers at the maintained no-secret services → agents → workflows harness and runnable examples. (`cmd/micro/`)
|
||||
- **First-agent tutorial smoke harness** — the first-agent tutorial path now has smoke coverage to keep the no-secret on-ramp runnable. (`internal/harness/`)
|
||||
- **No-secret agent debugging smoke** — the no-secret agent debugging path now has smoke coverage for the first-agent troubleshooting flow. (`internal/harness/`)
|
||||
- **Durable checkpoint resume smoke coverage** — durable agent resume after checkpointing now has focused smoke coverage. (`agent/`, `internal/harness/`)
|
||||
|
||||
### Fixed
|
||||
- **Plan/delegate notify replays** — duplicate and replayed plan-delegate notifications are now idempotent, so resumed runs do not duplicate completed notifications. (`agent/`, `internal/harness/`)
|
||||
- **Provider conformance scheduling** — provider conformance workflow dispatches now guard their scheduling path more reliably. (`.github/workflows/`)
|
||||
- **Plan/delegate notification completion** — delegated notifications now preserve plan completion state more reliably, including duplicate, paraphrased, and delegated-owner notification paths. (`agent/`, `internal/harness/`)
|
||||
- **AtlasCloud tool fallback** — AtlasCloud built-in tool schemas and follow-up tool fallback handling now recover conformance delegate retries more reliably. (`ai/atlascloud/`, `agent/`)
|
||||
- **Agent conformance retry completion** — conformance retry prompts and completion handling are more deterministic for delegated agent runs. (`agent/`, `internal/harness/`)
|
||||
|
||||
### Documentation
|
||||
- **First-agent quickstart numbering** — the first-agent on-ramp numbering is consistent across the README and website docs. (`README.md`, `internal/website/docs/`)
|
||||
- **First-agent inspect command** — docs now use the maintained `micro inspect agent <name>` form. (`README.md`, `internal/website/docs/`)
|
||||
- **`micro loop` quickstart wayfinding** — docs now surface the loop quickstart from the public docs index and README wayfinding. (`README.md`, `internal/website/docs/`)
|
||||
|
||||
---
|
||||
|
||||
|
||||
@@ -31,7 +31,6 @@ Running Go Micro in production, or building on it and want help? Paid **support,
|
||||
- [Building Agents](#building-agents) — [Plan & Delegate](#plan--delegate), [Pluggable](#batteries-included-pluggable), [Paid tools (x402)](#paid-tools-x402), [A2A](#reachable-by-other-agents-a2a)
|
||||
- [Features](#features)
|
||||
- [CLI](#cli)
|
||||
- [Autonomous improvement loop](#autonomous-improvement-loop)
|
||||
- [Multi-Service Projects](#multi-service-projects)
|
||||
- [Data Model](#data-model)
|
||||
- [AI Providers](#ai-providers)
|
||||
@@ -106,25 +105,6 @@ walkable agent path in this order:
|
||||
services → agents → workflows loop with scaffold, run, chat, inspect, flow
|
||||
history, and deploy dry-run commands that match the maintained harness.
|
||||
|
||||
### Autonomous improvement loop
|
||||
|
||||
Want the same services → agents → workflows lifecycle applied to your
|
||||
repository? `micro loop` scaffolds the autonomous improvement loop used by Go
|
||||
Micro itself: a North Star, ranked issue queue, role prompts, GitHub Actions
|
||||
workflows, and verification for CI-gated PRs.
|
||||
|
||||
```bash
|
||||
micro loop init --roles all
|
||||
micro loop verify
|
||||
```
|
||||
|
||||
Before turning on the schedule, configure a dispatch token such as
|
||||
`CODEX_TRIGGER_TOKEN`, protect the default branch with required CI checks
|
||||
(`go build ./...`, `go test ./...`, and `golangci-lint run ./...` for this
|
||||
repository), and seed `.github/loop/PRIORITIES.md` with one scoped issue per
|
||||
increment. See the [`micro loop` quickstart](internal/website/docs/guides/micro-loop.md)
|
||||
for the setup checklist and operating model.
|
||||
|
||||
### Generate from a prompt — with an LLM key
|
||||
|
||||
Set a provider key, describe what you want, and the AI designs services, writes handlers, compiles, and starts them:
|
||||
|
||||
@@ -244,31 +244,6 @@ func Pending(ctx context.Context, ag Agent) ([]flow.Run, error) {
|
||||
return a.pending(ctx)
|
||||
}
|
||||
|
||||
// ResumePending resumes every checkpointed agent run that has not completed
|
||||
// yet, in the same oldest-first order returned by Pending.
|
||||
//
|
||||
// It is a convenience for service startup and recovery loops: after recreating
|
||||
// an agent with the same checkpoint store, call ResumePending to drain the
|
||||
// durable backlog without listing and resuming each run manually. If any run
|
||||
// fails again, ResumePending stops and returns that run id with the error so
|
||||
// callers can log, alert, or retry later without hiding the failing run.
|
||||
func ResumePending(ctx context.Context, ag Agent) (string, error) {
|
||||
a, ok := ag.(*agentImpl)
|
||||
if !ok {
|
||||
return "", fmt.Errorf("agent resume pending: unsupported agent implementation %T", ag)
|
||||
}
|
||||
runs, err := a.pending(ctx)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
for _, run := range runs {
|
||||
if _, err := a.resume(ctx, run.ID); err != nil {
|
||||
return run.ID, err
|
||||
}
|
||||
}
|
||||
return "", nil
|
||||
}
|
||||
|
||||
func (a *agentImpl) ask(ctx context.Context, message, parentRunID string) (*Response, error) {
|
||||
a.mu.Lock()
|
||||
defer a.mu.Unlock()
|
||||
|
||||
+8
-59
@@ -2,7 +2,6 @@ package agent
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"strings"
|
||||
@@ -592,9 +591,6 @@ func (a *agentImpl) handleDelegate(ctx context.Context, call ai.ToolCall) ai.Too
|
||||
return errResult(call.ID, "task is required")
|
||||
}
|
||||
to, _ := input["to"].(string)
|
||||
if cached, ok := a.cachedDelegateResult(call.ID, to, task); ok {
|
||||
return cached
|
||||
}
|
||||
|
||||
// An external agent on another framework, addressed by A2A URL.
|
||||
if strings.HasPrefix(to, "http://") || strings.HasPrefix(to, "https://") {
|
||||
@@ -602,7 +598,9 @@ func (a *agentImpl) handleDelegate(ctx context.Context, call ai.ToolCall) ai.Too
|
||||
if err != nil {
|
||||
return errResult(call.ID, "delegate to A2A agent "+to+": "+err.Error())
|
||||
}
|
||||
return a.storeDelegateResult(call.ID, to, task, map[string]any{"agent": to, "reply": reply})
|
||||
out := map[string]any{"agent": to, "reply": reply}
|
||||
b, _ := json.Marshal(out)
|
||||
return ai.ToolResult{ID: call.ID, Value: out, Content: string(b)}
|
||||
}
|
||||
|
||||
// Delegate-first: an existing agent that owns the domain handles it.
|
||||
@@ -611,7 +609,9 @@ func (a *agentImpl) handleDelegate(ctx context.Context, call ai.ToolCall) ai.Too
|
||||
if err != nil {
|
||||
return errResult(call.ID, "delegate to agent "+to+": "+err.Error())
|
||||
}
|
||||
return a.storeDelegateResult(call.ID, to, task, map[string]any{"agent": to, "reply": reply})
|
||||
out := map[string]any{"agent": to, "reply": reply}
|
||||
b, _ := json.Marshal(out)
|
||||
return ai.ToolResult{ID: call.ID, Value: out, Content: string(b)}
|
||||
}
|
||||
|
||||
// Otherwise create a focused, ephemeral sub-agent. Fresh context:
|
||||
@@ -644,60 +644,9 @@ func (a *agentImpl) handleDelegate(ctx context.Context, call ai.ToolCall) ai.Too
|
||||
if err != nil {
|
||||
return errResult(call.ID, "sub-agent: "+err.Error())
|
||||
}
|
||||
return a.storeDelegateResult(call.ID, to, task, map[string]any{"reply": resp.Reply})
|
||||
}
|
||||
|
||||
func (a *agentImpl) cachedDelegateResult(id, to, task string) (ai.ToolResult, bool) {
|
||||
recs, err := a.stateStore().Read(delegateResultKey(to, task))
|
||||
if err != nil || len(recs) == 0 {
|
||||
return ai.ToolResult{}, false
|
||||
}
|
||||
var out map[string]any
|
||||
if err := json.Unmarshal(recs[0].Value, &out); err != nil {
|
||||
return ai.ToolResult{}, false
|
||||
}
|
||||
out := map[string]any{"reply": resp.Reply}
|
||||
b, _ := json.Marshal(out)
|
||||
return ai.ToolResult{ID: id, Value: out, Content: string(b)}, true
|
||||
}
|
||||
|
||||
func (a *agentImpl) storeDelegateResult(id, to, task string, out map[string]any) ai.ToolResult {
|
||||
b, _ := json.Marshal(out)
|
||||
_ = a.stateStore().Write(&store.Record{Key: delegateResultKey(to, task), Value: b})
|
||||
return ai.ToolResult{ID: id, Value: out, Content: string(b)}
|
||||
}
|
||||
|
||||
func delegateResultKey(to, task string) string {
|
||||
fp := normalizeDelegateTarget(to) + "\x00" + normalizeDelegateTask(task)
|
||||
sum := sha256.Sum256([]byte(fp))
|
||||
return fmt.Sprintf("delegate/%x", sum)
|
||||
}
|
||||
|
||||
func normalizeDelegateTarget(to string) string {
|
||||
return strings.Join(strings.Fields(strings.ToLower(strings.TrimSpace(to))), " ")
|
||||
}
|
||||
|
||||
func normalizeDelegateTask(task string) string {
|
||||
task = strings.ToLower(strings.TrimSpace(task))
|
||||
task = strings.Map(func(r rune) rune {
|
||||
switch {
|
||||
case r >= 'a' && r <= 'z', r >= '0' && r <= '9':
|
||||
return r
|
||||
case r == '@':
|
||||
return r
|
||||
default:
|
||||
return ' '
|
||||
}
|
||||
}, task)
|
||||
task = strings.Join(strings.Fields(task), " ")
|
||||
if strings.Contains(task, "notify") &&
|
||||
strings.Contains(task, "owner") &&
|
||||
strings.Contains(task, "acme") &&
|
||||
strings.Contains(task, "launch") &&
|
||||
strings.Contains(task, "plan") &&
|
||||
(strings.Contains(task, "ready") || strings.Contains(task, "readiness") || strings.Contains(task, "prepared") || strings.Contains(task, "complete")) {
|
||||
return "notify owner@acme.com launch-plan-ready"
|
||||
}
|
||||
return task
|
||||
return ai.ToolResult{ID: call.ID, Value: out, Content: string(b)}
|
||||
}
|
||||
|
||||
// isAgent reports whether name resolves to a registered agent (a
|
||||
|
||||
@@ -182,31 +182,6 @@ func TestBuiltinsAccessor(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestDelegateResultCacheReusesLaunchReadinessParaphrases(t *testing.T) {
|
||||
mem := store.NewMemoryStore()
|
||||
a := New(Name("planner"), WithStore(mem)).(*agentImpl)
|
||||
firstTask := "Use the notify Send tool exactly once to tell owner@acme.com: The launch plan is ready."
|
||||
first := a.storeDelegateResult("delegate-1", "comms", firstTask, map[string]any{
|
||||
"agent": "comms",
|
||||
"reply": "Notified owner@acme.com.",
|
||||
})
|
||||
if first.Content == "" {
|
||||
t.Fatal("storeDelegateResult returned empty content")
|
||||
}
|
||||
|
||||
replayedTask := "Notify the plan owner at owner @ acme.com that launch readiness is prepared and complete."
|
||||
cached, ok := a.cachedDelegateResult("delegate-2", " COMMS ", replayedTask)
|
||||
if !ok {
|
||||
t.Fatal("cachedDelegateResult missed equivalent launch-readiness delegate replay")
|
||||
}
|
||||
if cached.ID != "delegate-2" {
|
||||
t.Fatalf("cached result ID = %q, want replay call ID", cached.ID)
|
||||
}
|
||||
if !containsStr(cached.Content, "Notified owner@acme.com") {
|
||||
t.Fatalf("cached result content = %q, want original delegate reply", cached.Content)
|
||||
}
|
||||
}
|
||||
|
||||
func TestIsAgent(t *testing.T) {
|
||||
reg := registry.NewMemoryRegistry()
|
||||
|
||||
|
||||
+1
-5
@@ -50,13 +50,9 @@ func (a *agentImpl) saveRun(ctx context.Context, run flow.Run) error {
|
||||
return fmt.Errorf("agent %s checkpoint save: %w", a.opts.Name, err)
|
||||
}
|
||||
if info, ok := ai.RunInfoFrom(ctx); ok {
|
||||
stage := run.State.Stage
|
||||
if stage == "" && len(run.Steps) > 0 {
|
||||
stage = run.Steps[0].Name
|
||||
}
|
||||
a.recordTimelineEvent(ctx, RunEvent{
|
||||
Time: time.Now(), RunID: info.RunID, ParentID: info.ParentID, Agent: info.Agent,
|
||||
Kind: "checkpoint", Name: stage, Status: run.Status,
|
||||
Kind: "checkpoint", Name: run.State.Stage, Status: run.Status,
|
||||
})
|
||||
}
|
||||
return nil
|
||||
|
||||
@@ -5,7 +5,6 @@ import (
|
||||
"errors"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"go-micro.dev/v6/ai"
|
||||
"go-micro.dev/v6/client"
|
||||
@@ -288,8 +287,7 @@ func TestCheckpointContinuesRunThroughSeveralSingleStepTurns(t *testing.T) {
|
||||
|
||||
func TestResumeFailedCheckpointAfterFreshAgentRestart(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
st := store.NewMemoryStore()
|
||||
cp := flow.StoreCheckpoint(st, "restart-resume-agent")
|
||||
cp := flow.StoreCheckpoint(store.NewMemoryStore(), "restart-resume-agent")
|
||||
toolRuns := 0
|
||||
modelCalls := 0
|
||||
failFirst := true
|
||||
@@ -310,7 +308,7 @@ func TestResumeFailedCheckpointAfterFreshAgentRestart(t *testing.T) {
|
||||
defer func() { fakeGen = nil }()
|
||||
|
||||
newAgent := func() *agentImpl {
|
||||
return newTestAgent(Name("restart-resume-agent"), WithStore(st), WithCheckpoint(cp),
|
||||
return newTestAgent(Name("restart-resume-agent"), WithCheckpoint(cp),
|
||||
WithTool("external.provision", "provision service once", nil, func(context.Context, map[string]any) (string, error) {
|
||||
toolRuns++
|
||||
return "provisioned", nil
|
||||
@@ -332,19 +330,6 @@ func TestResumeFailedCheckpointAfterFreshAgentRestart(t *testing.T) {
|
||||
if len(runs) != 1 {
|
||||
t.Fatalf("Pending before restart returned %d runs, want 1", len(runs))
|
||||
}
|
||||
summaries, err := ListRunSummaries(st, "restart-resume-agent")
|
||||
if err != nil {
|
||||
t.Fatalf("ListRunSummaries before restart: %v", err)
|
||||
}
|
||||
if len(summaries) != 1 {
|
||||
t.Fatalf("run summaries before restart = %d, want 1", len(summaries))
|
||||
}
|
||||
if summaries[0].RunID != runs[0].ID || summaries[0].Status != "error" || summaries[0].Checkpoint != "failed" || summaries[0].Stage != agentAskStep {
|
||||
t.Fatalf("summary before restart = %#v, want failed ask checkpoint for %s", summaries[0], runs[0].ID)
|
||||
}
|
||||
if summaries[0].Events < 4 || summaries[0].LastError == "" {
|
||||
t.Fatalf("summary before restart lacks debug history/error: %#v", summaries[0])
|
||||
}
|
||||
|
||||
restarted := newAgent()
|
||||
resp, err := Resume(ctx, restarted, runs[0].ID)
|
||||
@@ -367,34 +352,6 @@ func TestResumeFailedCheckpointAfterFreshAgentRestart(t *testing.T) {
|
||||
if loaded.Status != "done" || loaded.ParentID != runs[0].ParentID {
|
||||
t.Fatalf("loaded run status/parent = %s/%s, want done/%s", loaded.Status, loaded.ParentID, runs[0].ParentID)
|
||||
}
|
||||
summaries, err = ListRunSummaries(st, "restart-resume-agent")
|
||||
if err != nil {
|
||||
t.Fatalf("ListRunSummaries after restart: %v", err)
|
||||
}
|
||||
if len(summaries) != 1 {
|
||||
t.Fatalf("run summaries after restart = %d, want 1", len(summaries))
|
||||
}
|
||||
if summaries[0].RunID != runs[0].ID || summaries[0].Status != "done" || summaries[0].Checkpoint != "done" || summaries[0].Stage != agentAskStep {
|
||||
t.Fatalf("summary after restart = %#v, want done ask checkpoint for %s", summaries[0], runs[0].ID)
|
||||
}
|
||||
if summaries[0].Events < 7 {
|
||||
t.Fatalf("summary after restart recorded %d events, want durable failure/resume/done history", summaries[0].Events)
|
||||
}
|
||||
events, err := LoadRunEvents(st, "restart-resume-agent", runs[0].ID)
|
||||
if err != nil {
|
||||
t.Fatalf("LoadRunEvents after restart: %v", err)
|
||||
}
|
||||
seen := map[string]bool{"run": false, "tool": false, "checkpoint": false, "error": false, "resume": false, "done": false}
|
||||
for _, e := range events {
|
||||
if _, ok := seen[e.Kind]; ok {
|
||||
seen[e.Kind] = true
|
||||
}
|
||||
}
|
||||
for kind, ok := range seen {
|
||||
if !ok {
|
||||
t.Fatalf("events after restart missing %s: %#v", kind, events)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestResumeFailedCheckpointDoesNotDuplicateCompactedMemory(t *testing.T) {
|
||||
@@ -463,51 +420,6 @@ func countMemoryContent(messages []ai.Message, needle string) int {
|
||||
return count
|
||||
}
|
||||
|
||||
func TestResumePendingResumesOldestAgentRunsUntilFailure(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
cp := flow.StoreCheckpoint(store.NewMemoryStore(), "resume-pending-agent")
|
||||
base := time.Date(2026, 7, 7, 12, 0, 0, 0, time.UTC)
|
||||
for _, run := range []flow.Run{
|
||||
{ID: "run-ok", Flow: "resume-pending-agent", Status: "failed", State: flow.State{Stage: agentAskStep, Data: []byte("ok")}, Started: base},
|
||||
{ID: "run-blocked", Flow: "resume-pending-agent", Status: "failed", State: flow.State{Stage: agentAskStep, Data: []byte("block")}, Started: base.Add(time.Minute)},
|
||||
{ID: "run-later", Flow: "resume-pending-agent", Status: "failed", State: flow.State{Stage: agentAskStep, Data: []byte("later")}, Started: base.Add(2 * time.Minute)},
|
||||
} {
|
||||
if err := cp.Save(ctx, run); err != nil {
|
||||
t.Fatalf("Save(%s): %v", run.ID, err)
|
||||
}
|
||||
}
|
||||
|
||||
var prompts []string
|
||||
fakeGen = func(ctx context.Context, opts ai.Options, req *ai.Request) (*ai.Response, error) {
|
||||
prompts = append(prompts, req.Prompt)
|
||||
if req.Prompt == "block" {
|
||||
return nil, errors.New("still blocked")
|
||||
}
|
||||
return &ai.Response{Reply: req.Prompt + " resumed"}, nil
|
||||
}
|
||||
defer func() { fakeGen = nil }()
|
||||
|
||||
a := newTestAgent(Name("resume-pending-agent"), WithCheckpoint(cp))
|
||||
failedRun, err := ResumePending(ctx, a)
|
||||
if err == nil {
|
||||
t.Fatal("ResumePending succeeded, want blocked run error")
|
||||
}
|
||||
if failedRun != "run-blocked" {
|
||||
t.Fatalf("failed run = %q, want run-blocked", failedRun)
|
||||
}
|
||||
if got, want := strings.Join(prompts, ","), "ok,block"; got != want {
|
||||
t.Fatalf("prompts = %q, want %q", got, want)
|
||||
}
|
||||
loaded, ok, err := cp.Load(ctx, "run-ok")
|
||||
if err != nil || !ok || loaded.Status != "done" {
|
||||
t.Fatalf("run-ok loaded=%v err=%v status=%q, want done", ok, err, loaded.Status)
|
||||
}
|
||||
loaded, ok, err = cp.Load(ctx, "run-later")
|
||||
if err != nil || !ok || loaded.Status != "failed" {
|
||||
t.Fatalf("run-later loaded=%v err=%v status=%q, want still failed", ok, err, loaded.Status)
|
||||
}
|
||||
}
|
||||
|
||||
func TestPendingReturnsUnfinishedAgentRuns(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
cp := flow.StoreCheckpoint(store.NewMemoryStore(), "pending-agent")
|
||||
|
||||
+7
-119
@@ -180,7 +180,7 @@ func runAgentConformanceScenario(t *testing.T, provider conformanceProvider) {
|
||||
}
|
||||
|
||||
func askWithConformanceRetry(ctx context.Context, a Agent, initialPrompt string, sawTool, sawBlockedDelegate *bool) (*Response, error) {
|
||||
const maxAttempts = 4
|
||||
const maxAttempts = 3
|
||||
prompt := initialPrompt
|
||||
var resp *Response
|
||||
for attempt := 1; attempt <= maxAttempts; attempt++ {
|
||||
@@ -198,11 +198,7 @@ func askWithConformanceRetry(ctx context.Context, a Agent, initialPrompt string,
|
||||
if attempt == maxAttempts {
|
||||
break
|
||||
}
|
||||
prompt = nextConformanceRetryPrompt(sawRequiredTool, sawRequiredDelegate, hasMarker, attempt+1)
|
||||
}
|
||||
missing := missingConformanceRequirements(sawTool, sawBlockedDelegate, responseHasConformanceMarker(resp))
|
||||
if len(missing) > 0 {
|
||||
return resp, fmt.Errorf("provider conformance incomplete after %d attempts: missing %s", maxAttempts, strings.Join(missing, ", "))
|
||||
prompt = nextConformanceRetryPrompt(sawRequiredTool, sawRequiredDelegate, hasMarker)
|
||||
}
|
||||
return resp, nil
|
||||
}
|
||||
@@ -211,30 +207,10 @@ func askWithConformanceToolRetry(ctx context.Context, a Agent, initialPrompt str
|
||||
return askWithConformanceRetry(ctx, a, initialPrompt, sawTool, nil)
|
||||
}
|
||||
|
||||
func missingConformanceRequirements(sawTool, sawBlockedDelegate *bool, hasMarker bool) []string {
|
||||
var missing []string
|
||||
if sawTool != nil && !*sawTool {
|
||||
missing = append(missing, "conformance_echo")
|
||||
}
|
||||
if sawBlockedDelegate != nil && !*sawBlockedDelegate {
|
||||
missing = append(missing, "guarded delegate")
|
||||
}
|
||||
if !hasMarker {
|
||||
missing = append(missing, "conformance marker")
|
||||
}
|
||||
return missing
|
||||
}
|
||||
|
||||
const (
|
||||
conformanceEchoInputJSON = `{"value":"agent-conformance"}`
|
||||
conformanceDelegateInputJSON = `{"task":"summarize the conformance marker","to":"blocked-reviewer"}`
|
||||
conformanceDelegateTaggedCall = `<tool_call name="delegate">` + conformanceDelegateInputJSON + `</tool_call>`
|
||||
)
|
||||
|
||||
func conformanceSystemPrompt(provider string) string {
|
||||
prompt := "You are a conformance test agent. Create a short plan, use conformance_echo exactly once with input " + conformanceEchoInputJSON + ", then attempt to delegate a summary to blocked-reviewer with input " + conformanceDelegateInputJSON + ". You must complete both tool calls before any final answer; a final answer that only mentions the steps without calling both tools is invalid. If the delegate is refused, explain the refusal and answer with the echo result."
|
||||
prompt := "You are a conformance test agent. Create a short plan, use conformance_echo exactly once with input {\"value\":\"agent-conformance\"}, then attempt to delegate a summary to blocked-reviewer with input {\"task\":\"summarize the conformance marker\",\"to\":\"blocked-reviewer\"}. If the delegate is refused, explain the refusal and answer with the echo result."
|
||||
if provider == "atlascloud" {
|
||||
prompt += " AtlasCloud/minimax conformance note: the delegate attempt is mandatory after conformance_echo. If native tool_calls are unavailable, emit the delegate as " + conformanceDelegateTaggedCall + " rather than answering in prose."
|
||||
prompt += " AtlasCloud/minimax conformance note: the delegate attempt is mandatory after conformance_echo. If native tool_calls are unavailable, emit the delegate as <tool_call name=\"delegate\">{\"task\":\"summarize the conformance marker\",\"to\":\"blocked-reviewer\"}</tool_call> rather than answering in prose."
|
||||
}
|
||||
return prompt
|
||||
}
|
||||
@@ -243,7 +219,6 @@ func TestAgentProviderConformanceAtlasCloudPromptRequiresTaggedDelegateFallback(
|
||||
prompt := conformanceSystemPrompt("atlascloud")
|
||||
for _, want := range []string{
|
||||
"delegate attempt is mandatory",
|
||||
"You must complete both tool calls before any final answer",
|
||||
"<tool_call name=\"delegate\">",
|
||||
`{"task":"summarize the conformance marker","to":"blocked-reviewer"}`,
|
||||
} {
|
||||
@@ -257,45 +232,12 @@ func TestAgentProviderConformanceAtlasCloudPromptRequiresTaggedDelegateFallback(
|
||||
}
|
||||
}
|
||||
|
||||
func TestAgentProviderConformanceRetryPromptsRequireBothTools(t *testing.T) {
|
||||
for name, prompt := range map[string]string{
|
||||
"missing tool": nextConformanceRetryPrompt(false, false, false, 2),
|
||||
"missing delegate": nextConformanceRetryPrompt(true, false, true, 2),
|
||||
} {
|
||||
for _, want := range []string{
|
||||
"delegate exactly once",
|
||||
conformanceDelegateTaggedCall,
|
||||
"do not",
|
||||
} {
|
||||
if !strings.Contains(prompt, want) {
|
||||
t.Fatalf("%s retry prompt %q missing %q", name, prompt, want)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestAgentProviderConformanceFinalDelegateRetryUsesTaggedCall(t *testing.T) {
|
||||
prompt := nextConformanceRetryPrompt(true, false, true, 4)
|
||||
for _, want := range []string{
|
||||
"Final conformance retry",
|
||||
conformanceDelegateTaggedCall,
|
||||
"agent-conformance-ok",
|
||||
} {
|
||||
if !strings.Contains(prompt, want) {
|
||||
t.Fatalf("final delegate retry prompt %q missing %q", prompt, want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func nextConformanceRetryPrompt(sawTool, sawBlockedDelegate, hasMarker bool, attempt int) string {
|
||||
if attempt >= 4 && sawTool && !sawBlockedDelegate {
|
||||
return "Final conformance retry: emit exactly this tagged tool call so the harness can execute the guarded delegate refusal, then include agent-conformance-ok and the refusal in the final answer: " + conformanceDelegateTaggedCall
|
||||
}
|
||||
func nextConformanceRetryPrompt(sawTool, sawBlockedDelegate, hasMarker bool) string {
|
||||
switch {
|
||||
case !sawTool:
|
||||
return "The previous response did not call the required conformance_echo tool. Retry the same conformance check now: first call conformance_echo exactly once with input " + conformanceEchoInputJSON + ", then call delegate exactly once with input " + conformanceDelegateInputJSON + "; do not provide a final answer until both tool calls have been attempted. If native delegate tool_calls are unavailable after conformance_echo, emit exactly " + conformanceDelegateTaggedCall + ". The delegate is expected to be refused by policy; include that refusal and the agent-conformance marker in the final answer."
|
||||
return "The previous response did not call the required conformance_echo tool. Retry the same conformance check now: you must call conformance_echo exactly once with input {\"value\":\"agent-conformance\"} before any final answer, then include the tool result marker in the final answer."
|
||||
case !sawBlockedDelegate:
|
||||
return "The previous response called conformance_echo but did not attempt the required guarded delegation. Continue the same conformance check now: call delegate exactly once with input " + conformanceDelegateInputJSON + "; do not answer in prose until that delegate call has been attempted. If native tool_calls are unavailable, emit exactly " + conformanceDelegateTaggedCall + ". The delegate is expected to be refused by policy; include that refusal and the agent-conformance marker in the final answer."
|
||||
return "The previous response called conformance_echo but did not attempt the required guarded delegation. Continue the same conformance check now: call delegate exactly once with input {\"task\":\"summarize the conformance marker\",\"to\":\"blocked-reviewer\"}; do not answer in prose until that delegate call has been attempted. If native tool_calls are unavailable, emit exactly <tool_call name=\"delegate\">{\"task\":\"summarize the conformance marker\",\"to\":\"blocked-reviewer\"}</tool_call>. The delegate is expected to be refused by policy; include that refusal and the agent-conformance marker in the final answer."
|
||||
case !hasMarker:
|
||||
return "The previous response completed the required tool calls but omitted the conformance marker. Continue the same conformance check now: do not call more tools; answer with the prior echo result marker agent-conformance-ok and mention the guarded delegate refusal."
|
||||
default:
|
||||
@@ -516,60 +458,6 @@ func TestAgentProviderConformanceRetriesMissingDelegate(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestAgentProviderConformanceFailsWhenDelegateStillMissing(t *testing.T) {
|
||||
var attempts int
|
||||
fakeGen = func(ctx context.Context, opts ai.Options, req *ai.Request) (*ai.Response, error) {
|
||||
attempts++
|
||||
if err := validateConformanceRequest(req, opts); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
echo := opts.ToolHandler(ctx, ai.ToolCall{
|
||||
ID: fmt.Sprintf("fake-call-%d", attempts),
|
||||
Name: "conformance_echo",
|
||||
Input: map[string]any{"value": "agent-conformance"},
|
||||
})
|
||||
return &ai.Response{
|
||||
Reply: "called conformance_echo with agent-conformance-ok but skipped delegate",
|
||||
Answer: echo.Content,
|
||||
ToolCalls: []ai.ToolCall{
|
||||
{ID: fmt.Sprintf("fake-call-%d", attempts), Name: "conformance_echo", Input: map[string]any{"value": "agent-conformance"}, Result: echo.Content},
|
||||
},
|
||||
}, nil
|
||||
}
|
||||
defer func() { fakeGen = nil }()
|
||||
|
||||
var sawTool bool
|
||||
var sawBlockedDelegate bool
|
||||
a := New(
|
||||
Name("conformance-retry-delegate-exhausted"),
|
||||
Provider("fake"),
|
||||
WithRegistry(registry.NewMemoryRegistry()),
|
||||
WithStore(store.NewMemoryStore()),
|
||||
WithMemory(NewInMemory(4)),
|
||||
ApproveTool(func(tool string, input map[string]any) (bool, string) {
|
||||
if tool == "delegate" {
|
||||
sawBlockedDelegate = true
|
||||
return false, "cross-provider conformance blocks delegate side effects"
|
||||
}
|
||||
return true, ""
|
||||
}),
|
||||
WithTool("conformance_echo", "Echo a conformance value.", map[string]any{
|
||||
"value": map[string]any{"type": "string"},
|
||||
}, func(ctx context.Context, input map[string]any) (string, error) {
|
||||
sawTool = true
|
||||
return `{"marker":"agent-conformance-ok"}`, nil
|
||||
}),
|
||||
)
|
||||
|
||||
_, err := askWithConformanceRetry(context.Background(), a, "Run the provider conformance check.", &sawTool, &sawBlockedDelegate)
|
||||
if err == nil || !strings.Contains(err.Error(), "guarded delegate") {
|
||||
t.Fatalf("Ask error = %v, want missing guarded delegate", err)
|
||||
}
|
||||
if attempts != 4 {
|
||||
t.Fatalf("attempts = %d, want retries through max attempts", attempts)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAgentExecutesProviderTextToolCallFallback(t *testing.T) {
|
||||
fakeGen = func(ctx context.Context, opts ai.Options, req *ai.Request) (*ai.Response, error) {
|
||||
if opts.ToolHandler == nil {
|
||||
|
||||
+5
-27
@@ -36,8 +36,6 @@ const (
|
||||
AttrTotalTokens = "agent.tokens.total"
|
||||
AttrAttempt = "agent.model.attempt"
|
||||
AttrMaxAttempts = "agent.model.max_attempts"
|
||||
AttrToolAttempt = "agent.tool.attempt"
|
||||
AttrToolMaxAttempts = "agent.tool.max_attempts"
|
||||
AttrToolName = "agent.tool.name"
|
||||
AttrDelegate = "agent.delegate"
|
||||
AttrGuardrailBlock = "agent.guardrail.block"
|
||||
@@ -366,11 +364,7 @@ func (a *agentImpl) traceTool(next ai.ToolHandler) ai.ToolHandler {
|
||||
res := next(ctx, call)
|
||||
dur := time.Since(start).Milliseconds()
|
||||
resErr := resultError(res)
|
||||
toolAttempts := res.Attempts
|
||||
if toolAttempts <= 0 {
|
||||
toolAttempts = 1
|
||||
}
|
||||
a.recordRunEvent(RunEvent{Time: time.Now(), RunID: info.RunID, ParentID: info.ParentID, Agent: info.Agent, Kind: "tool", Name: call.Name, Attempt: toolAttempts, MaxAttempts: a.opts.ToolMaxAttempts, LatencyMS: dur, Refused: res.Refused, Error: resErr, ErrorKind: classifyToolError(resErr)})
|
||||
a.recordRunEvent(RunEvent{Time: time.Now(), RunID: info.RunID, ParentID: info.ParentID, Agent: info.Agent, Kind: "tool", Name: call.Name, LatencyMS: dur, Refused: res.Refused, Error: resErr, ErrorKind: classifyToolError(resErr)})
|
||||
return res
|
||||
}
|
||||
|
||||
@@ -384,14 +378,6 @@ func (a *agentImpl) traceTool(next ai.ToolHandler) ai.ToolHandler {
|
||||
res := next(ctx, call)
|
||||
dur := time.Since(start).Milliseconds()
|
||||
attrs := []attribute.KeyValue{attribute.Int64(AttrLatencyMS, dur)}
|
||||
toolAttempts := res.Attempts
|
||||
if toolAttempts <= 0 {
|
||||
toolAttempts = 1
|
||||
}
|
||||
attrs = append(attrs, attribute.Int(AttrToolAttempt, toolAttempts))
|
||||
if a.opts.ToolMaxAttempts > 0 {
|
||||
attrs = append(attrs, attribute.Int(AttrToolMaxAttempts, a.opts.ToolMaxAttempts))
|
||||
}
|
||||
if res.Refused != "" {
|
||||
attrs = append(attrs, attribute.Bool(AttrGuardrailBlock, true), attribute.String(AttrRefusal, res.Refused))
|
||||
}
|
||||
@@ -407,7 +393,7 @@ func (a *agentImpl) traceTool(next ai.ToolHandler) ai.ToolHandler {
|
||||
} else {
|
||||
span.SetStatus(codes.Ok, "")
|
||||
}
|
||||
a.recordSpanEvent(span, RunEvent{Time: time.Now(), RunID: info.RunID, ParentID: info.ParentID, Agent: info.Agent, Kind: "tool", Name: call.Name, Attempt: toolAttempts, MaxAttempts: a.opts.ToolMaxAttempts, LatencyMS: dur, Refused: res.Refused, Error: resErr, ErrorKind: classifyToolError(resErr)})
|
||||
a.recordSpanEvent(span, RunEvent{Time: time.Now(), RunID: info.RunID, ParentID: info.ParentID, Agent: info.Agent, Kind: "tool", Name: call.Name, LatencyMS: dur, Refused: res.Refused, Error: resErr, ErrorKind: classifyToolError(resErr)})
|
||||
span.End()
|
||||
return res
|
||||
}
|
||||
@@ -475,18 +461,10 @@ func runEventAttributes(e RunEvent) []attribute.KeyValue {
|
||||
attrs = append(attrs, attribute.String(AttrModel, e.Model))
|
||||
}
|
||||
if e.Attempt > 0 {
|
||||
if e.Kind == "tool" {
|
||||
attrs = append(attrs, attribute.Int(AttrToolAttempt, e.Attempt))
|
||||
} else {
|
||||
attrs = append(attrs, attribute.Int(AttrAttempt, e.Attempt))
|
||||
}
|
||||
attrs = append(attrs, attribute.Int(AttrAttempt, e.Attempt))
|
||||
}
|
||||
if e.MaxAttempts > 0 {
|
||||
if e.Kind == "tool" {
|
||||
attrs = append(attrs, attribute.Int(AttrToolMaxAttempts, e.MaxAttempts))
|
||||
} else {
|
||||
attrs = append(attrs, attribute.Int(AttrMaxAttempts, e.MaxAttempts))
|
||||
}
|
||||
attrs = append(attrs, attribute.Int(AttrMaxAttempts, e.MaxAttempts))
|
||||
}
|
||||
if e.LatencyMS > 0 {
|
||||
attrs = append(attrs, attribute.Int64(AttrLatencyMS, e.LatencyMS))
|
||||
@@ -646,7 +624,7 @@ func runStatus(events []RunEvent) string {
|
||||
if e.Error != "" || e.Kind == "error" {
|
||||
status = runErrorStatus(e.ErrorKind)
|
||||
}
|
||||
if e.Kind == "done" {
|
||||
if e.Kind == "done" && status == "running" {
|
||||
status = "done"
|
||||
}
|
||||
}
|
||||
|
||||
@@ -142,82 +142,6 @@ func TestAgentOpenTelemetrySpans(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestAgentOpenTelemetryToolRetryAttempts(t *testing.T) {
|
||||
exp := tracetest.NewInMemoryExporter()
|
||||
tp := trace.NewTracerProvider(trace.WithSyncer(exp))
|
||||
st := store.NewMemoryStore()
|
||||
calls := 0
|
||||
a := New(
|
||||
Name("tool-retry-otel"),
|
||||
Provider("oteltest"),
|
||||
WithStore(st),
|
||||
TraceProvider(tp),
|
||||
ToolRetry(3, time.Millisecond),
|
||||
WithTool("probe", "probe", nil, func(context.Context, map[string]any) (string, error) {
|
||||
calls++
|
||||
if calls == 1 {
|
||||
return "", errors.New("rate limit exceeded")
|
||||
}
|
||||
return "ok", nil
|
||||
}),
|
||||
)
|
||||
if _, err := a.Ask(context.Background(), "hello"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if calls != 2 {
|
||||
t.Fatalf("tool calls = %d, want retry success after 2 attempts", calls)
|
||||
}
|
||||
|
||||
var sawToolSpan bool
|
||||
for _, span := range exp.GetSpans().Snapshots() {
|
||||
if span.Name() != spanNameToolCall {
|
||||
continue
|
||||
}
|
||||
attrs := spanAttributes(span.Attributes())
|
||||
if attrs[AttrToolName] != "probe" {
|
||||
continue
|
||||
}
|
||||
if attrs[AttrToolAttempt] != "2" || attrs[AttrToolMaxAttempts] != "3" {
|
||||
t.Fatalf("tool retry span attempts = %#v", attrs)
|
||||
}
|
||||
if !spanEventHasAttr(span.Events(), "agent.tool", AttrToolAttempt, "2") || !spanEventHasAttr(span.Events(), "agent.tool", AttrToolMaxAttempts, "3") {
|
||||
t.Fatalf("tool retry event missing attempt attributes: %#v", span.Events())
|
||||
}
|
||||
sawToolSpan = true
|
||||
}
|
||||
if !sawToolSpan {
|
||||
t.Fatal("tool retry span not emitted")
|
||||
}
|
||||
|
||||
summaries, err := ListRunSummaries(st, "tool-retry-otel")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
events, err := LoadRunEvents(st, "tool-retry-otel", summaries[0].RunID)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, event := range events {
|
||||
if event.Kind == "tool" && event.Name == "probe" && event.Attempt == 2 && event.MaxAttempts == 3 {
|
||||
return
|
||||
}
|
||||
}
|
||||
t.Fatalf("persisted tool event missing retry attempts: %#v", events)
|
||||
}
|
||||
|
||||
func spanEventHasAttr(events []trace.Event, name, key, value string) bool {
|
||||
for _, event := range events {
|
||||
if event.Name != name {
|
||||
continue
|
||||
}
|
||||
attrs := spanAttributes(event.Attributes)
|
||||
if attrs[key] == value {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
func TestAgentRunObservabilityRedactsInputByDefault(t *testing.T) {
|
||||
secret := "deploy production with token sk-secret"
|
||||
exp := tracetest.NewInMemoryExporter()
|
||||
|
||||
@@ -83,62 +83,6 @@ func TestAskRetriesTransientErrorsThenSurfacesStructuredError(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestModelRetryDoesNotDuplicateCheckpointedToolSideEffects(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
cp := flow.StoreCheckpoint(store.NewMemoryStore(), "retry-tool-dedupe-agent")
|
||||
attempts := 0
|
||||
toolRuns := 0
|
||||
fakeGen = func(ctx context.Context, opts ai.Options, req *ai.Request) (*ai.Response, error) {
|
||||
attempts++
|
||||
if opts.ToolHandler == nil {
|
||||
t.Fatal("missing tool handler")
|
||||
}
|
||||
res := opts.ToolHandler(ctx, ai.ToolCall{ID: "create-1", Name: "external.create", Input: map[string]any{"title": "Retry safe"}})
|
||||
if res.Content != "created Retry safe" {
|
||||
t.Fatalf("tool result = %q, want cached create result", res.Content)
|
||||
}
|
||||
if attempts == 1 {
|
||||
return nil, testStatusError{code: 503}
|
||||
}
|
||||
return &ai.Response{Reply: "done", ToolCalls: []ai.ToolCall{{ID: "create-1", Name: "external.create", Input: map[string]any{"title": "Retry safe"}, Result: res.Content}}}, nil
|
||||
}
|
||||
defer func() { fakeGen = nil }()
|
||||
|
||||
a := newTestAgent(
|
||||
Name("retry-tool-dedupe-agent"),
|
||||
WithCheckpoint(cp),
|
||||
ModelRetry(2, time.Millisecond),
|
||||
WithTool("external.create", "create once", nil, func(context.Context, map[string]any) (string, error) {
|
||||
toolRuns++
|
||||
return "created Retry safe", nil
|
||||
}),
|
||||
)
|
||||
|
||||
resp, err := a.Ask(ctx, "create once despite a transient provider retry")
|
||||
if err != nil {
|
||||
t.Fatalf("Ask: %v", err)
|
||||
}
|
||||
if resp.Reply != "done" {
|
||||
t.Fatalf("reply = %q, want done", resp.Reply)
|
||||
}
|
||||
if attempts != 2 {
|
||||
t.Fatalf("model attempts = %d, want retry after transient provider failure", attempts)
|
||||
}
|
||||
if toolRuns != 1 {
|
||||
t.Fatalf("tool executions = %d, want checkpointed side effect reused across retry", toolRuns)
|
||||
}
|
||||
runs, err := cp.List(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("List: %v", err)
|
||||
}
|
||||
if len(runs) != 1 {
|
||||
t.Fatalf("checkpointed runs = %d, want 1", len(runs))
|
||||
}
|
||||
if _, ok := findStep(runs[0].Steps, `tool:external.create:{"title":"Retry safe"}`); !ok {
|
||||
t.Fatalf("checkpoint steps = %#v, want completed external.create step", runs[0].Steps)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAskRateLimitFailureSuggestsPreflightAndInspect(t *testing.T) {
|
||||
fakeGen = func(ctx context.Context, opts ai.Options, req *ai.Request) (*ai.Response, error) {
|
||||
return nil, testStatusError{code: 429}
|
||||
|
||||
+4
-10
@@ -71,15 +71,14 @@ func ResumeStreamAsk(ctx context.Context, ag Agent, runID string) (AgentStream,
|
||||
// StreamAsk runs tools like Ask, emits ToolStart/ToolEnd events as they execute,
|
||||
// then emits chunks of the final answer followed by a Done event.
|
||||
func (a *agentImpl) StreamAsk(ctx context.Context, message string) (AgentStream, error) {
|
||||
streamCtx, cancel := context.WithCancel(ctx)
|
||||
events := make(chan *StreamEvent, 16)
|
||||
done := make(chan struct{})
|
||||
s := &agentStream{events: events, done: done, cancel: cancel}
|
||||
s := &agentStream{events: events, done: done}
|
||||
|
||||
go func() {
|
||||
defer close(events)
|
||||
defer close(done)
|
||||
resp, err := a.askWithStreamEvents(streamCtx, message, events)
|
||||
resp, err := a.askWithStreamEvents(ctx, message, events)
|
||||
if err != nil {
|
||||
s.setErr(err)
|
||||
return
|
||||
@@ -95,15 +94,14 @@ func (a *agentImpl) StreamAsk(ctx context.Context, message string) (AgentStream,
|
||||
}
|
||||
|
||||
func (a *agentImpl) resumeStreamAsk(ctx context.Context, runID string) (AgentStream, error) {
|
||||
streamCtx, cancel := context.WithCancel(ctx)
|
||||
events := make(chan *StreamEvent, 16)
|
||||
done := make(chan struct{})
|
||||
s := &agentStream{events: events, done: done, cancel: cancel}
|
||||
s := &agentStream{events: events, done: done}
|
||||
|
||||
go func() {
|
||||
defer close(events)
|
||||
defer close(done)
|
||||
resp, err := a.resumeWithStreamEvents(streamCtx, runID, events)
|
||||
resp, err := a.resumeWithStreamEvents(ctx, runID, events)
|
||||
if err != nil {
|
||||
s.setErr(err)
|
||||
return
|
||||
@@ -262,7 +260,6 @@ func (a *agentImpl) streamAskAI(ctx context.Context, message string) (ai.Stream,
|
||||
type agentStream struct {
|
||||
events <-chan *StreamEvent
|
||||
done <-chan struct{}
|
||||
cancel context.CancelFunc
|
||||
mu sync.Mutex
|
||||
err error
|
||||
}
|
||||
@@ -281,9 +278,6 @@ func (s *agentStream) Recv() (*StreamEvent, error) {
|
||||
}
|
||||
|
||||
func (s *agentStream) Close() error {
|
||||
if s.cancel != nil {
|
||||
s.cancel()
|
||||
}
|
||||
<-s.done
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -5,7 +5,6 @@ import (
|
||||
"errors"
|
||||
"io"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"go-micro.dev/v6/ai"
|
||||
"go-micro.dev/v6/flow"
|
||||
@@ -76,39 +75,6 @@ func TestStreamAskEmitsToolEventsAndFinalTokens(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestStreamAskCloseCancelsInFlightModelCall(t *testing.T) {
|
||||
started := make(chan struct{})
|
||||
fakeGen = func(ctx context.Context, opts ai.Options, req *ai.Request) (*ai.Response, error) {
|
||||
close(started)
|
||||
<-ctx.Done()
|
||||
return nil, ctx.Err()
|
||||
}
|
||||
defer func() { fakeGen = nil }()
|
||||
|
||||
a := newTestAgent(Name("stream-cancel"))
|
||||
stream, err := a.StreamAsk(context.Background(), "cancel me")
|
||||
if err != nil {
|
||||
t.Fatalf("StreamAsk: %v", err)
|
||||
}
|
||||
|
||||
select {
|
||||
case <-started:
|
||||
case <-time.After(time.Second):
|
||||
t.Fatal("model call did not start")
|
||||
}
|
||||
|
||||
closed := make(chan error, 1)
|
||||
go func() { closed <- stream.Close() }()
|
||||
select {
|
||||
case err := <-closed:
|
||||
if err != nil {
|
||||
t.Fatalf("Close: %v", err)
|
||||
}
|
||||
case <-time.After(time.Second):
|
||||
t.Fatal("Close did not cancel the in-flight stream")
|
||||
}
|
||||
}
|
||||
|
||||
func TestStreamAskHelperRejectsUnsupportedAgent(t *testing.T) {
|
||||
_, err := StreamAsk(context.Background(), unsupportedAgent{}, "hello")
|
||||
if err == nil {
|
||||
|
||||
+15
-65
@@ -93,7 +93,20 @@ func (p *Provider) Options() ai.Options { return p.opts }
|
||||
func (p *Provider) String() string { return "atlascloud" }
|
||||
|
||||
func (p *Provider) Generate(ctx context.Context, req *ai.Request, opts ...ai.GenerateOption) (*ai.Response, error) {
|
||||
tools := atlascloudTools(req.Tools)
|
||||
var tools []map[string]any
|
||||
for _, t := range req.Tools {
|
||||
tools = append(tools, map[string]any{
|
||||
"type": "function",
|
||||
"function": map[string]any{
|
||||
"name": t.Name,
|
||||
"description": t.Description,
|
||||
"parameters": map[string]any{
|
||||
"type": "object",
|
||||
"properties": t.Properties,
|
||||
},
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
messages := []map[string]any{
|
||||
{"role": "system", "content": req.SystemPrompt},
|
||||
@@ -178,17 +191,7 @@ func (p *Provider) Generate(ctx context.Context, req *ai.Request, opts ...ai.Gen
|
||||
resp.ToolCalls = allToolCalls
|
||||
}
|
||||
if followUpResp.Reply != "" {
|
||||
if strings.Contains(followUpResp.Reply, "<tool_call") || strings.Contains(followUpResp.Reply, "function=") {
|
||||
// Preserve follow-up assistant content as Reply, not Answer, when
|
||||
// it may contain a text-encoded tool call. The agent harness
|
||||
// inspects Reply for text fallback calls after Generate returns,
|
||||
// which covers AtlasCloud/minimax turns that emit a second
|
||||
// required call (for example guarded delegate) as markup instead
|
||||
// of native tool_calls.
|
||||
resp.Reply = followUpResp.Reply
|
||||
} else {
|
||||
resp.Answer = followUpResp.Reply
|
||||
}
|
||||
resp.Answer = followUpResp.Reply
|
||||
} else if len(toolResults) > 0 {
|
||||
resp.Answer = strings.Join(toolResults, "\n")
|
||||
}
|
||||
@@ -376,59 +379,6 @@ func (p *Provider) callAPI(ctx context.Context, phase string, req map[string]any
|
||||
return response, rawMessage, nil
|
||||
}
|
||||
|
||||
func atlascloudTools(input []ai.Tool) []map[string]any {
|
||||
tools := make([]map[string]any, 0, len(input))
|
||||
for _, t := range input {
|
||||
tools = append(tools, map[string]any{
|
||||
"type": "function",
|
||||
"function": map[string]any{
|
||||
"name": t.Name,
|
||||
"description": t.Description,
|
||||
"parameters": map[string]any{
|
||||
"type": "object",
|
||||
"properties": normalizeAtlasCloudSchema(t.Properties),
|
||||
},
|
||||
},
|
||||
})
|
||||
}
|
||||
return tools
|
||||
}
|
||||
|
||||
func normalizeAtlasCloudSchema(schema map[string]any) map[string]any {
|
||||
if schema == nil {
|
||||
return nil
|
||||
}
|
||||
out := make(map[string]any, len(schema))
|
||||
for k, v := range schema {
|
||||
out[k] = normalizeAtlasCloudSchemaValue(v)
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
func normalizeAtlasCloudSchemaValue(v any) any {
|
||||
switch val := v.(type) {
|
||||
case map[string]any:
|
||||
out := make(map[string]any, len(val)+1)
|
||||
for k, nested := range val {
|
||||
out[k] = normalizeAtlasCloudSchemaValue(nested)
|
||||
}
|
||||
if typ, _ := out["type"].(string); typ == "array" {
|
||||
if _, ok := out["items"]; !ok {
|
||||
out["items"] = map[string]any{}
|
||||
}
|
||||
}
|
||||
return out
|
||||
case []any:
|
||||
out := make([]any, len(val))
|
||||
for i, nested := range val {
|
||||
out[i] = normalizeAtlasCloudSchemaValue(nested)
|
||||
}
|
||||
return out
|
||||
default:
|
||||
return v
|
||||
}
|
||||
}
|
||||
|
||||
func normalizeAtlasCloudToolCalls(toolCalls []atlasToolCall) []map[string]any {
|
||||
out := make([]map[string]any, 0, len(toolCalls))
|
||||
for _, tc := range toolCalls {
|
||||
|
||||
@@ -289,58 +289,6 @@ func TestProvider_GenerateMinimaxToolRequests(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestProvider_GenerateNormalizesBuiltInToolSchemas(t *testing.T) {
|
||||
var body map[string]any
|
||||
ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
if err := json.NewDecoder(r.Body).Decode(&body); err != nil {
|
||||
t.Fatalf("decode request: %v", err)
|
||||
}
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
_, _ = w.Write([]byte(`{"choices":[{"message":{"content":"ok"}}]}`))
|
||||
}))
|
||||
defer ts.Close()
|
||||
|
||||
planProperties := map[string]any{
|
||||
"steps": map[string]any{
|
||||
"type": "array",
|
||||
"description": "ordered plan steps",
|
||||
},
|
||||
}
|
||||
p := NewProvider(
|
||||
ai.WithAPIKey("test-key"),
|
||||
ai.WithBaseURL(ts.URL),
|
||||
ai.WithModel("minimaxai/minimax-m3"),
|
||||
)
|
||||
_, err := p.Generate(context.Background(), &ai.Request{
|
||||
Prompt: "plan and delegate",
|
||||
Tools: []ai.Tool{
|
||||
{Name: "task_TaskService_Add", Description: "add task", Properties: map[string]any{"title": map[string]any{"type": "string"}}},
|
||||
{Name: "plan", Description: "record a plan", Properties: planProperties},
|
||||
{Name: "request_input", Description: "request input", Properties: map[string]any{"prompt": map[string]any{"type": "string"}}},
|
||||
{Name: "delegate", Description: "delegate work", Properties: map[string]any{"task": map[string]any{"type": "string"}, "to": map[string]any{"type": "string"}}},
|
||||
},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("Generate returned error: %v", err)
|
||||
}
|
||||
|
||||
tools := body["tools"].([]any)
|
||||
if len(tools) != 4 {
|
||||
t.Fatalf("tools = %d, want custom tool plus built-ins", len(tools))
|
||||
}
|
||||
planTool := tools[1].(map[string]any)
|
||||
fn := planTool["function"].(map[string]any)
|
||||
params := fn["parameters"].(map[string]any)
|
||||
props := params["properties"].(map[string]any)
|
||||
steps := props["steps"].(map[string]any)
|
||||
if _, ok := steps["items"].(map[string]any); !ok {
|
||||
t.Fatalf("plan steps schema = %#v, want array items for AtlasCloud/minimax", steps)
|
||||
}
|
||||
if _, mutated := planProperties["steps"].(map[string]any)["items"]; mutated {
|
||||
t.Fatalf("Generate mutated caller tool schema: %#v", planProperties)
|
||||
}
|
||||
}
|
||||
|
||||
func TestProvider_GenerateExecutesFollowUpToolCall(t *testing.T) {
|
||||
var bodies []map[string]any
|
||||
ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
@@ -406,54 +354,6 @@ func TestProvider_GenerateExecutesFollowUpToolCall(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestProvider_GeneratePreservesFollowUpTextToolCallInReply(t *testing.T) {
|
||||
var bodies []map[string]any
|
||||
ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
var body map[string]any
|
||||
if err := json.NewDecoder(r.Body).Decode(&body); err != nil {
|
||||
t.Fatalf("decode request: %v", err)
|
||||
}
|
||||
bodies = append(bodies, body)
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
switch len(bodies) {
|
||||
case 1:
|
||||
_, _ = w.Write([]byte(`{"choices":[{"message":{"content":"","tool_calls":[{"id":"call-1","function":{"name":"conformance_echo","arguments":"{\"value\":\"agent-conformance\"}"}}]}}]}`))
|
||||
case 2:
|
||||
_, _ = w.Write([]byte(`{"choices":[{"message":{"content":"<tool_call name=\"delegate\">{\"task\":\"summarize the conformance marker\",\"to\":\"blocked-reviewer\"}</tool_call>"}}]}`))
|
||||
default:
|
||||
t.Fatalf("unexpected API call %d", len(bodies))
|
||||
}
|
||||
}))
|
||||
defer ts.Close()
|
||||
|
||||
p := NewProvider(
|
||||
ai.WithAPIKey("test-key"),
|
||||
ai.WithBaseURL(ts.URL),
|
||||
ai.WithToolHandler(func(ctx context.Context, call ai.ToolCall) ai.ToolResult {
|
||||
if call.Name != "conformance_echo" {
|
||||
t.Fatalf("unexpected structured tool call %+v", call)
|
||||
}
|
||||
return ai.ToolResult{ID: call.ID, Content: `{"marker":"agent-conformance-ok"}`}
|
||||
}),
|
||||
)
|
||||
resp, err := p.Generate(context.Background(), &ai.Request{
|
||||
Prompt: "run conformance",
|
||||
Tools: []ai.Tool{
|
||||
{Name: "conformance_echo", Description: "echo conformance marker", Properties: map[string]any{"value": map[string]any{"type": "string"}}},
|
||||
{Name: "delegate", Description: "delegate work", Properties: map[string]any{"task": map[string]any{"type": "string"}, "to": map[string]any{"type": "string"}}},
|
||||
},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("Generate returned error: %v", err)
|
||||
}
|
||||
if !strings.Contains(resp.Reply, `<tool_call name="delegate">`) {
|
||||
t.Fatalf("Reply = %q, want tagged delegate follow-up for agent text fallback", resp.Reply)
|
||||
}
|
||||
if resp.Answer != "" {
|
||||
t.Fatalf("Answer = %q, want follow-up text preserved only as Reply", resp.Answer)
|
||||
}
|
||||
}
|
||||
|
||||
func TestProvider_GenerateToolCallHTTPErrorIncludesRequestContext(t *testing.T) {
|
||||
ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
http.Error(w, `{"code":400,"msg":"bad request"}`, http.StatusBadRequest)
|
||||
|
||||
@@ -49,34 +49,6 @@ require_output() {
|
||||
fi
|
||||
}
|
||||
|
||||
require_ordered_output() {
|
||||
local description=$1
|
||||
shift
|
||||
local -a expected=()
|
||||
while [[ $# -gt 0 && "$1" != "--" ]]; do
|
||||
expected+=("$1")
|
||||
shift
|
||||
done
|
||||
shift
|
||||
|
||||
local output
|
||||
if ! output=$("$MICRO" "$@" 2>&1); then
|
||||
echo "micro $* failed while checking $description" >&2
|
||||
echo "$output" >&2
|
||||
exit 1
|
||||
fi
|
||||
|
||||
local remainder=$output
|
||||
for text in "${expected[@]}"; do
|
||||
if [[ "$remainder" != *"$text"* ]]; then
|
||||
echo "micro $* missing expected ordered text '$text' for $description" >&2
|
||||
echo "$output" >&2
|
||||
exit 1
|
||||
fi
|
||||
remainder=${remainder#*"$text"}
|
||||
done
|
||||
}
|
||||
|
||||
require_output "version" "micro version" --version
|
||||
require_output "root help" "COMMANDS" --help
|
||||
require_output "service scaffold" "micro new" new --help
|
||||
@@ -86,51 +58,4 @@ require_output "agent chat" "micro chat" chat --help
|
||||
require_output "agent inspection" "micro inspect agent" inspect agent --help
|
||||
require_output "flow inspection" "micro inspect flow" inspect flow --help
|
||||
|
||||
require_ordered_output "installed first-agent docs wayfinding" \
|
||||
"micro agent demo" \
|
||||
"no-secret-first-agent.html" \
|
||||
"your-first-agent.html" \
|
||||
"micro agent preflight # before micro run: prerequisites" \
|
||||
"micro run" \
|
||||
"micro chat" \
|
||||
"micro agent doctor # after micro run: chat/gateway/inspect recovery" \
|
||||
"debugging-agents.html" \
|
||||
"micro inspect agent <name>" \
|
||||
"zero-to-hero.html" \
|
||||
-- docs
|
||||
|
||||
require_ordered_output "installed provider-free examples wayfinding" \
|
||||
"go run ./examples/first-agent" \
|
||||
"go test ./internal/harness/zero-to-hero-ci -run TestNoSecretFirstAgentTranscript -count=1" \
|
||||
"go run ./examples/support" \
|
||||
"micro agent demo" \
|
||||
"micro docs" \
|
||||
"micro zero-to-hero" \
|
||||
"no-secret-first-agent.html" \
|
||||
"your-first-agent.html" \
|
||||
"debugging-agents.html" \
|
||||
"zero-to-hero.html" \
|
||||
-- examples
|
||||
|
||||
require_ordered_output "installed no-secret agent demo" \
|
||||
"provider-free" \
|
||||
"go test ./internal/harness/zero-to-hero-ci -run TestNoSecretFirstAgentTranscript -count=1" \
|
||||
"your-first-agent.html" \
|
||||
"debugging-agents.html" \
|
||||
"zero-to-hero.html" \
|
||||
"micro agent preflight # before micro run: prerequisites" \
|
||||
"micro run" \
|
||||
"micro chat" \
|
||||
"micro agent doctor # after micro run: chat/gateway/inspect recovery" \
|
||||
"micro inspect agent <name>" \
|
||||
-- agent demo
|
||||
|
||||
require_ordered_output "installed zero-to-hero lifecycle wayfinding" \
|
||||
"./internal/harness/zero-to-hero-ci/run.sh" \
|
||||
"go run ./examples/first-agent" \
|
||||
"go run ./examples/support" \
|
||||
"make harness" \
|
||||
"zero-to-hero.html" \
|
||||
-- zero-to-hero
|
||||
|
||||
echo "✓ install smoke path verified"
|
||||
|
||||
@@ -185,48 +185,22 @@ func (s *NotifyService) duplicateAttempts() int {
|
||||
}
|
||||
|
||||
func notifyDedupKey(to, message string) string {
|
||||
recipient := canonicalLaunchNotifyRecipient(normalizeNotifyText(to))
|
||||
recipient := strings.ToLower(strings.TrimSpace(to))
|
||||
body := normalizeNotifyText(message)
|
||||
if isLaunchReadinessNotify(body) {
|
||||
if recipient == "owner@acme.com" && isLaunchReadinessNotify(body) {
|
||||
body = "launch-readiness"
|
||||
}
|
||||
return recipient + "\x00" + body
|
||||
}
|
||||
|
||||
func canonicalLaunchNotifyRecipient(recipient string) string {
|
||||
switch recipient {
|
||||
case "owner", "launch owner", "plan owner", "owner acme com", "owner@acme com", "owner @ acme com":
|
||||
return "owner@acme.com"
|
||||
default:
|
||||
if strings.Contains(recipient, "owner") && strings.Contains(recipient, "acme") {
|
||||
return "owner@acme.com"
|
||||
}
|
||||
return recipient
|
||||
}
|
||||
}
|
||||
|
||||
func normalizeNotifyText(message string) string {
|
||||
message = strings.ToLower(strings.TrimSpace(message))
|
||||
message = strings.Map(func(r rune) rune {
|
||||
switch {
|
||||
case r >= 'a' && r <= 'z', r >= '0' && r <= '9':
|
||||
return r
|
||||
case r == '@':
|
||||
return r
|
||||
default:
|
||||
return ' '
|
||||
}
|
||||
}, message)
|
||||
return strings.Join(strings.Fields(message), " ")
|
||||
return strings.Join(strings.Fields(strings.ToLower(strings.TrimSpace(message))), " ")
|
||||
}
|
||||
|
||||
func isLaunchReadinessNotify(message string) bool {
|
||||
return strings.Contains(message, "launch") &&
|
||||
strings.Contains(message, "plan") &&
|
||||
(strings.Contains(message, "ready") ||
|
||||
strings.Contains(message, "readiness") ||
|
||||
strings.Contains(message, "prepared") ||
|
||||
strings.Contains(message, "complete"))
|
||||
(strings.Contains(message, "ready") || strings.Contains(message, "readiness"))
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -248,11 +222,6 @@ type mockModel struct {
|
||||
// duplicateNotify makes the comms mock replay the same notification call.
|
||||
// The notify service should collapse that replay to one durable side effect.
|
||||
duplicateNotify bool
|
||||
|
||||
// duplicateDelegate makes the conductor mock replay the same delegate call.
|
||||
// The delegate idempotency path should collapse that replay before it can
|
||||
// ask the delegated comms agent to notify twice.
|
||||
duplicateDelegate bool
|
||||
}
|
||||
|
||||
func newMock(opts ...ai.Option) ai.Model {
|
||||
@@ -273,12 +242,6 @@ func newMockDuplicateNotify(opts ...ai.Option) ai.Model {
|
||||
return m
|
||||
}
|
||||
|
||||
func newMockDuplicateDelegate(opts ...ai.Option) ai.Model {
|
||||
m := &mockModel{duplicateDelegate: true}
|
||||
_ = m.Init(opts...)
|
||||
return m
|
||||
}
|
||||
|
||||
func (m *mockModel) Init(opts ...ai.Option) error {
|
||||
for _, o := range opts {
|
||||
o(&m.opts)
|
||||
@@ -355,14 +318,10 @@ func (m *mockModel) Generate(ctx context.Context, req *ai.Request, _ ...ai.Gener
|
||||
"to": "comms",
|
||||
})
|
||||
} else {
|
||||
input := map[string]any{
|
||||
m.call("conductor", del, map[string]any{
|
||||
"task": delegatedNotifyTask,
|
||||
"to": "comms",
|
||||
}
|
||||
m.call("conductor", del, input)
|
||||
if m.duplicateDelegate {
|
||||
m.call("conductor", del, input)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
return &ai.Response{Answer: "Created Design, Build and Ship, and had comms notify the owner."}, nil
|
||||
@@ -398,8 +357,6 @@ func runPlanDelegate(provider string) error {
|
||||
ai.Register("mock-unknown-delegate", newMockUnknownDelegate)
|
||||
case "mock-duplicate-notify":
|
||||
ai.Register("mock-duplicate-notify", newMockDuplicateNotify)
|
||||
case "mock-duplicate-delegate":
|
||||
ai.Register("mock-duplicate-delegate", newMockDuplicateDelegate)
|
||||
default:
|
||||
apiKey = providerKey(provider)
|
||||
if apiKey == "" {
|
||||
@@ -615,7 +572,7 @@ func isClientTimeout(err error) bool {
|
||||
}
|
||||
|
||||
func main() {
|
||||
provider := flag.String("provider", "mock", "LLM provider: mock (default), mock-unknown-delegate, mock-duplicate-notify, mock-duplicate-delegate, anthropic, openai, gemini, groq, mistral, together, atlascloud")
|
||||
provider := flag.String("provider", "mock", "LLM provider: mock (default), mock-unknown-delegate, mock-duplicate-notify, anthropic, openai, gemini, groq, mistral, together, atlascloud")
|
||||
flag.Parse()
|
||||
|
||||
if err := runPlanDelegate(*provider); err != nil {
|
||||
|
||||
@@ -240,15 +240,6 @@ func TestPlanDelegateIdempotentDuplicateNotifyReplay(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestPlanDelegateIdempotentDuplicateDelegateReplay(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("0→hero harness boots an end-to-end system; skipped with -short")
|
||||
}
|
||||
if err := runPlanDelegate("mock-duplicate-delegate"); err != nil {
|
||||
t.Fatalf("0→hero harness with duplicate delegate replay: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestTaskServiceAddIsIdempotentForLaunchTitles(t *testing.T) {
|
||||
svc := new(TaskService)
|
||||
for _, title := range []string{"Design", "design task", "Build", "Build launch task", "Ship", "ship readiness"} {
|
||||
@@ -438,11 +429,7 @@ func TestNotifyServiceSendIsIdempotentForDuplicateDelivery(t *testing.T) {
|
||||
}
|
||||
for i, message := range messages {
|
||||
var rsp SendResponse
|
||||
to := "owner@acme.com"
|
||||
if i == len(messages)-1 {
|
||||
to = "owner"
|
||||
}
|
||||
if err := svc.Send(context.Background(), &SendRequest{To: to, Message: message}, &rsp); err != nil {
|
||||
if err := svc.Send(context.Background(), &SendRequest{To: "owner@acme.com", Message: message}, &rsp); err != nil {
|
||||
t.Fatalf("Send attempt %d: %v", i+1, err)
|
||||
}
|
||||
if !rsp.Sent {
|
||||
@@ -456,28 +443,3 @@ func TestNotifyServiceSendIsIdempotentForDuplicateDelivery(t *testing.T) {
|
||||
t.Fatalf("duplicate notify attempts = %d, want %d", got, len(messages)-1)
|
||||
}
|
||||
}
|
||||
|
||||
func TestNotifyServiceCollapsesProviderReadinessParaphrases(t *testing.T) {
|
||||
svc := new(NotifyService)
|
||||
requests := []SendRequest{
|
||||
{To: "owner@acme.com", Message: "The launch plan is ready"},
|
||||
{To: "owner @ acme.com", Message: "Launch plan ready."},
|
||||
{To: "launch owner", Message: "The launch readiness plan is prepared."},
|
||||
{To: "plan owner", Message: "Launch plan is complete!"},
|
||||
}
|
||||
for i, req := range requests {
|
||||
var rsp SendResponse
|
||||
if err := svc.Send(context.Background(), &req, &rsp); err != nil {
|
||||
t.Fatalf("Send attempt %d: %v", i+1, err)
|
||||
}
|
||||
if !rsp.Sent {
|
||||
t.Fatalf("Send attempt %d reported Sent=false", i+1)
|
||||
}
|
||||
}
|
||||
if got := svc.count(); got != 1 {
|
||||
t.Fatalf("notify count = %d, want 1 after provider paraphrase replays", got)
|
||||
}
|
||||
if got := svc.duplicateAttempts(); got != len(requests)-1 {
|
||||
t.Fatalf("duplicate notify attempts = %d, want %d", got, len(requests)-1)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -42,8 +42,6 @@ examples:
|
||||
guides:
|
||||
- title: Debugging your agent
|
||||
url: /docs/guides/debugging-agents.html
|
||||
- title: micro loop quickstart
|
||||
url: /docs/guides/micro-loop.html
|
||||
- title: Plan & Delegate
|
||||
url: /docs/guides/plan-delegate.html
|
||||
- title: Agent Guardrails
|
||||
@@ -88,7 +86,6 @@ search_order:
|
||||
- /docs/guides/your-first-agent.html
|
||||
- /docs/guides/zero-to-hero.html
|
||||
- /docs/guides/debugging-agents.html
|
||||
- /docs/guides/micro-loop.html
|
||||
- /docs/getting-started.html
|
||||
- /docs/mcp.html
|
||||
- /docs/architecture.html
|
||||
|
||||
@@ -259,5 +259,4 @@ The flow discovers all services as tools and lets the LLM decide which RPCs to c
|
||||
- [Agent Design](https://github.com/micro/go-micro/blob/master/internal/docs/AGENT_DESIGN.md) — the full agent interface specification
|
||||
- [MCP & AI Agents](mcp.html) — MCP gateway, tool discovery, and auth
|
||||
- [Data Model](model.html) — typed persistence with CRUD and queries
|
||||
- [`micro loop` quickstart](guides/micro-loop.html) — scaffold a CI-gated autonomous improvement loop for a repository
|
||||
- [Deployment](deployment.html) — deploy via SSH + systemd
|
||||
|
||||
@@ -1,96 +0,0 @@
|
||||
---
|
||||
layout: default
|
||||
---
|
||||
|
||||
# `micro loop` quickstart
|
||||
|
||||
`micro loop` scaffolds the autonomous improvement loop that Go Micro uses on
|
||||
this repository: GitHub Actions workflows for planning, building, evaluation
|
||||
feedback, coherence, security, and release. Use it when you want a repository to
|
||||
continuously turn a ranked queue into small PRs while CI remains the merge gate.
|
||||
|
||||
## 1. Initialize the loop
|
||||
|
||||
Run the default loop from the repository root:
|
||||
|
||||
```bash
|
||||
micro loop init
|
||||
```
|
||||
|
||||
For every role used by Go Micro itself, scaffold all workflows:
|
||||
|
||||
```bash
|
||||
micro loop init --roles all
|
||||
```
|
||||
|
||||
The command writes:
|
||||
|
||||
- `.github/loop/NORTH_STAR.md` — the direction every increment should optimize.
|
||||
- `.github/loop/PRIORITIES.md` — the ranked queue; the builder takes the top open issue.
|
||||
- `.github/loop/prompts/*.md` — editable policy for planner, builder, triage, coherence, and security roles.
|
||||
- `.github/workflows/loop-*.yml` — generated GitHub Actions mechanics.
|
||||
|
||||
Edit the files under `.github/loop/` to steer the loop. Re-run
|
||||
`micro loop init --roles all --force` only when you want to regenerate workflow
|
||||
mechanics from the installed CLI.
|
||||
|
||||
## 2. Configure the dispatch token
|
||||
|
||||
The scheduled builder needs a repository secret containing a token from a user
|
||||
account that the coding agent will answer. Go Micro names that secret
|
||||
`CODEX_TRIGGER_TOKEN` by default. If you use another secret name, pass it when
|
||||
you initialize the loop:
|
||||
|
||||
```bash
|
||||
micro loop init --agent @codex --token-secret LOOP_TOKEN --roles all
|
||||
```
|
||||
|
||||
The token needs enough repository permission to open issues, comment, push
|
||||
branches, create pull requests, and enable auto-merge. Run `gh auth setup-git` in
|
||||
the environment that will push branches so `git push` uses the same credentials
|
||||
as `gh`.
|
||||
|
||||
## 3. Make CI the gate
|
||||
|
||||
The loop should not be its own reviewer. Protect the default branch so PRs merge
|
||||
only after the required checks pass. At minimum, require the same commands the
|
||||
Go Micro loop verifies locally and in CI:
|
||||
|
||||
```bash
|
||||
go build ./...
|
||||
go test ./...
|
||||
golangci-lint run ./...
|
||||
```
|
||||
|
||||
If your repository has a harness or end-to-end grader, make that required too.
|
||||
Keep human approval requirements out of the autonomous path unless you intend the
|
||||
loop to pause for review.
|
||||
|
||||
## 4. Verify the wiring
|
||||
|
||||
After editing the North Star, queue, prompts, token secret, and branch
|
||||
protection, run:
|
||||
|
||||
```bash
|
||||
micro loop verify
|
||||
```
|
||||
|
||||
`micro loop verify` checks that the loop direction, queue, prompts, role
|
||||
workflows, and non-loop CI gate are present. Fix any reported missing items
|
||||
before relying on scheduled increments.
|
||||
|
||||
## 5. Operate the queue
|
||||
|
||||
Keep one ranked list in `.github/loop/PRIORITIES.md`. Each item should link a
|
||||
scoped issue and be small enough for one PR. The builder closes both the priority
|
||||
issue and the per-run tracker issue in the PR body, for example:
|
||||
|
||||
```text
|
||||
Closes #1234
|
||||
Closes #5678
|
||||
```
|
||||
|
||||
Use the North Star to keep the queue honest: favor small improvements that move
|
||||
developers through the services → agents → workflows lifecycle, and surface
|
||||
breaking API or brand/positioning decisions for humans instead of auto-merging
|
||||
them.
|
||||
@@ -30,7 +30,6 @@ Otherwise continue to read the docs for more information about the framework.
|
||||
- [Your First Agent](guides/your-first-agent.html) - Build a service-backed agent end to end
|
||||
- [MCP & AI Agents](mcp.html) - Turn services into AI-callable tools with the Model Context Protocol
|
||||
- [CLI & Gateway Guide](guides/cli-gateway.html) - Development vs Production modes
|
||||
- [`micro loop` quickstart](guides/micro-loop.html) - Scaffold an autonomous CI-gated improvement loop
|
||||
- [Quick Start](quickstart.html)
|
||||
- [Architecture](architecture.html)
|
||||
- [Configuration](config.html)
|
||||
@@ -65,6 +64,5 @@ Otherwise continue to read the docs for more information about the framework.
|
||||
- [Real-World Examples](examples/realworld/)
|
||||
- [Migration Guides](guides/migration/)
|
||||
- [Observability](observability.html)
|
||||
- [`micro loop` quickstart](guides/micro-loop.html)
|
||||
- [Contributing](contributing.html)
|
||||
- [Roadmap](roadmap.html)
|
||||
|
||||
@@ -220,13 +220,6 @@ func AgentResume(ctx context.Context, a Agent, runID string) (*AgentResponse, er
|
||||
return agent.Resume(ctx, a, runID)
|
||||
}
|
||||
|
||||
// AgentResumePending resumes every incomplete checkpointed agent run, oldest
|
||||
// first. It returns the first run id that fails again so startup recovery loops
|
||||
// can leave the durable backlog visible instead of swallowing the failure.
|
||||
func AgentResumePending(ctx context.Context, a Agent) (string, error) {
|
||||
return agent.ResumePending(ctx, a)
|
||||
}
|
||||
|
||||
// AgentResumeInput resumes a checkpointed agent run waiting for human input.
|
||||
func AgentResumeInput(ctx context.Context, a Agent, runID, input string) (*AgentResponse, error) {
|
||||
return agent.ResumeInput(ctx, a, runID, input)
|
||||
|
||||
Reference in New Issue
Block a user