Compare commits
31 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 80102ae526 | |||
| 47ff1bc6f5 | |||
| 6c6ce0e2a7 | |||
| d91ac00da9 | |||
| ff94e8b316 | |||
| ab0bf29c79 | |||
| 96e822d656 | |||
| c60c03d485 | |||
| 7963342999 | |||
| 593846093f | |||
| 5b27c5b239 | |||
| 0a10655231 | |||
| cc3502d156 | |||
| 014a6431f2 | |||
| 3725a53bf9 | |||
| b24b19f191 | |||
| bc732bb513 | |||
| 522786d21d | |||
| 84cc4c6e8d | |||
| f96b1014a2 | |||
| 3cafff8789 | |||
| 69e23eb1f6 | |||
| 1841599980 | |||
| 6f9994b3ea | |||
| 55070f1597 | |||
| b3940860b6 | |||
| 5c38a23fe4 | |||
| 9e7273dde1 | |||
| 42080b6d6d | |||
| c3df8b833d | |||
| 8694e3a20c |
@@ -21,11 +21,9 @@ changes, architectural rewrites. Those go to the human.
|
||||
|
||||
## Work queue (ranked)
|
||||
|
||||
1. **Isolate file-store tests from shared default directory** ([#3751](https://github.com/micro/go-micro/issues/3751)) — with the atlascloud delegated-notify live harness blocker closed by #3816, repeated `go test -race -cover ./...` failures around shared file-store directories are now the highest-value reliability gap. Keep this first because a flaky evaluator undermines every autonomous adoption and hardening increment that follows.
|
||||
2. **Stabilize file-store suffix expiry test timing** ([#3780](https://github.com/micro/go-micro/issues/3780)) — the remaining 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.
|
||||
3. **Propagate agent run cancellation and deadlines through model and tool calls** ([#3544](https://github.com/micro/go-micro/issues/3544)) — once red CI blockers are cleared, the highest-value remaining Now-phase 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.
|
||||
4. **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.
|
||||
5. **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.
|
||||
1. **Add cross-provider agent conformance** ([#3901](https://github.com/micro/go-micro/issues/3901)) — #3895 closed the provider deadline/retry guidance gap and #3899 closed the AP2 foundation, leaving no open `codex` PRs or issues in flight. The highest-value Now-phase gap is therefore battle-testing the same first-agent behavior across providers: tool calling, multi-step runs, plan/delegate, and guardrails should fail loudly by provider instead of drifting behind a single happy-path backend. Keep default CI no-secret friendly with mock coverage and key-gated live providers.
|
||||
2. **Checkpoint and resume durable agent runs** ([#3902](https://github.com/micro/go-micro/issues/3902)) — adoption now has install, scaffold, run, chat, inspect/history, and 0→hero coverage, but the lived harness story still has an asymmetry: flows checkpoint and resume while long agent runs do not. This is the top Next-phase agentic-depth gap because it makes services → agents → workflows feel like one runtime under interruption and deploy/restart conditions. Keep it additive; do not break the public agent API.
|
||||
3. **Broaden provider streaming and keep chat/A2A streaming end to end** ([#3903](https://github.com/micro/go-micro/issues/3903)) — after conformance and durability, streaming is the next developer-visible seam in the inner loop. Real chat and long-running A2A tasks need token streaming to stay coherent from provider → `ai.Stream` → `micro chat` → A2A `message/stream`, with mock/default CI coverage plus key-gated live provider checks and safe fallback for non-streaming providers.
|
||||
|
||||
_Seeded by Claude Code from the roadmap + open issues; thereafter maintained by the
|
||||
architecture-review pass._
|
||||
|
||||
@@ -0,0 +1,30 @@
|
||||
<!--
|
||||
The SECURITY prompt — go-micro's security audit. Editable policy; the workflow
|
||||
prepends the agent @mention and substitutes __ISSUE__ before posting. Keep
|
||||
__ISSUE__ literal.
|
||||
|
||||
Deliberately conservative: it does NOT auto-merge fixes, and it does NOT publish
|
||||
exploit details in public issues (responsible disclosure).
|
||||
-->
|
||||
Act as the security reviewer for go-micro. Audit for real, exploitable vulnerabilities — skip theoretical or lint-style noise.
|
||||
|
||||
GO-MICRO ATTACK SURFACE — weight these:
|
||||
- **MCP gateway** (`gateway/mcp`) and **A2A gateway** (`gateway/a2a`) — untrusted input from agents/tools: auth/scope enforcement, injection into downstream RPC, SSRF via tool/agent URLs, rate-limit/circuit-breaker bypass, info leak in errors.
|
||||
- **x402 payments** (`wrapper/x402`) — payment verification and settlement: signature/mandate validation, replay, budget-reservation races, facilitator auth (CDP bearer) handling, amount/network confusion.
|
||||
- **Auth** (`auth/jwt`, `wrapper/auth`) — token validation, algorithm confusion, scope/priority rule bypass, missing checks on endpoints.
|
||||
- **AI providers** (`ai/*`) — base-URL and endpoint handling: SSRF via config-controlled `BaseURL`, API keys leaking into logs/errors, TLS verification.
|
||||
- **Agent tool loop** (`agent/`) — prompt injection reaching real tool calls, guardrail (`MaxSteps`/`LoopLimit`/`ApproveTool`) bypass, delegate/plan side effects.
|
||||
- **Trust boundaries** — `server` RPC handlers, `broker` consumers, `store`/`registry` inputs, `transport` TLS defaults (v6 verifies by default — confirm nothing regressed).
|
||||
- **The loop itself** — `.github/workflows/loop-*.yml`: the `CODEX_TRIGGER_TOKEN` PAT must never be echoed/leaked; workflow inputs must not enable script injection.
|
||||
- **Dependencies** — run `govulncheck ./...` (install if needed) and inspect `go.mod` for known CVEs.
|
||||
|
||||
DEDUPE against open issues first.
|
||||
|
||||
HOW TO REPORT:
|
||||
- **Known/public dependency CVEs**: file a `security` issue referencing the CVE + module; you MAY open a PR bumping to the patched version. Do NOT enable auto-merge.
|
||||
- **Novel, exploitable vulnerabilities in this code** (not yet public): do NOT post an exploit or PoC in a public issue. File a CONCISE `security` + `needs-human` issue naming the class, location (file/function), and impact only — and note it should go through GitHub private vulnerability reporting. Do NOT open a public fix PR that reveals it.
|
||||
- **Low-risk hardening**: a normal `security` issue is fine.
|
||||
|
||||
NEVER auto-merge a security change. Never weaken a control to make a test pass. Architectural/breaking fixes → `needs-human` with the tradeoff.
|
||||
|
||||
Post a summary as a comment on this issue (#__ISSUE__) — findings by severity, what you filed, what needs a human — then close it (`gh issue close __ISSUE__`). If you open a dependency-bump PR: `git switch -c loop/security-__ISSUE__`, `git push -u origin loop/security-__ISSUE__`, `gh pr create --base master --label codex --label security --title "<title>" --body "<summary, Closes #__ISSUE__>"` — then STOP, do NOT run `gh pr merge --auto`. Do not use the make_pr tool.
|
||||
@@ -0,0 +1,60 @@
|
||||
name: "Loop: Security"
|
||||
|
||||
# Generated by `micro loop init`. A dispatch role of the autonomous loop: on a
|
||||
# cadence it opens a fresh tracking issue and posts the instruction in
|
||||
# .github/loop/prompts/security.md to the agent (@codex).
|
||||
#
|
||||
# The workflow is the MECHANISM; that prompt file is the editable POLICY —
|
||||
# change what this role does by editing the prompt, not this YAML. A FRESH
|
||||
# issue per run is deliberate: agents derive the PR branch name from the
|
||||
# triggering issue, so reusing one tracker collapses every run onto one branch.
|
||||
#
|
||||
# Gated on CODEX_TRIGGER_TOKEN: the agent ignores @mentions from the
|
||||
# github-actions bot, so dispatch posts as a real user (a PAT). No token → no-op.
|
||||
|
||||
on:
|
||||
workflow_dispatch: {}
|
||||
schedule:
|
||||
- cron: "0 6 * * 1"
|
||||
|
||||
permissions:
|
||||
issues: write
|
||||
|
||||
concurrency:
|
||||
group: loop-security
|
||||
cancel-in-progress: false
|
||||
|
||||
jobs:
|
||||
dispatch:
|
||||
runs-on: ubuntu-latest
|
||||
steps:
|
||||
- uses: actions/checkout@v4 # needed to read the prompt file
|
||||
- name: Dispatch security
|
||||
env:
|
||||
GH_TOKEN: ${{ secrets.CODEX_TRIGGER_TOKEN || github.token }}
|
||||
HAS_TOKEN: ${{ secrets.CODEX_TRIGGER_TOKEN != '' }}
|
||||
REPO: ${{ github.repository }}
|
||||
RUN_NUMBER: ${{ github.run_number }}
|
||||
run: |
|
||||
if [ "$HAS_TOKEN" != "true" ]; then
|
||||
echo "CODEX_TRIGGER_TOKEN is not set — skipping (the agent ignores bot @mentions)."
|
||||
exit 0
|
||||
fi
|
||||
PROMPT=".github/loop/prompts/security.md"
|
||||
if [ ! -f "$PROMPT" ]; then
|
||||
echo "missing $PROMPT — run 'micro loop init'." >&2
|
||||
exit 1
|
||||
fi
|
||||
ISSUE_URL=$(gh issue create --repo "$REPO" \
|
||||
--title "Loop: security review #$RUN_NUMBER" \
|
||||
--body "Autonomous security pass. Direction: .github/loop/NORTH_STAR.md; queue: .github/loop/PRIORITIES.md.")
|
||||
ISSUE_NUM="${ISSUE_URL##*/}"
|
||||
echo "Opened issue #$ISSUE_NUM — dispatching security."
|
||||
# The prompt file is the policy; strip its editorial <!-- --> header and
|
||||
# substitute the tracking issue number (__ISSUE__) at runtime.
|
||||
{
|
||||
echo "@codex"
|
||||
echo
|
||||
sed -e '/<!--/,/-->/d' -e "s/__ISSUE__/$ISSUE_NUM/g" "$PROMPT"
|
||||
} > "$RUNNER_TEMP/loop-body.md"
|
||||
gh issue comment "$ISSUE_NUM" --repo "$REPO" --body-file "$RUNNER_TEMP/loop-body.md"
|
||||
@@ -16,12 +16,25 @@ next version when it ships.
|
||||
|
||||
## [Unreleased]
|
||||
|
||||
## [6.3.14] - July 2026
|
||||
|
||||
### Added
|
||||
- **MiniMax provider** — run agents against MiniMax's `MiniMax-M3` model via its OpenAI-compatible endpoint, with tool calling and streaming; auto-detected from the base URL. (`ai/minimax/`)
|
||||
- **`micro loop` security role** — a new opt-in loop role (`--roles …,security`) that periodically audits a repo for vulnerabilities and files `security` issues. It is deliberately conservative: it never auto-merges fixes and never publishes exploit detail in public issues (responsible disclosure), and risky fixes are marked `needs-human`. go-micro now runs it against its own attack surface (MCP/A2A gateways, x402, auth, provider URLs, agent tool loop, deps). (`cmd/micro/loop/`)
|
||||
- **Agent run tracing** — agent model streaming and run-event kinds now emit richer trace detail for debugging agent execution. (`agent/`)
|
||||
|
||||
### Changed
|
||||
- **Agent memory** — streamed agent replies are persisted in conversation memory so later turns can reference streamed responses. (`agent/`)
|
||||
|
||||
### Fixed
|
||||
- **Plan/delegate completion** — agents now continue unfinished plan steps more reliably, fail checkpointed runs that leave delegated plans unfinished, recover from unknown plan-delegate tool calls, avoid duplicate side effects, and complete timeout paths deterministically. (`agent/`)
|
||||
- **AtlasCloud tool calls** — streaming and request fallback handling now recovers tool-call results from provider responses that omit the expected structured fields. (`ai/atlascloud/`)
|
||||
- **Agent preflight diagnostics** — provider setup failures now surface more actionable errors before an agent run starts. (`agent/`)
|
||||
- **A2A fallback streams** — fallback stream validation is stricter for malformed or incomplete A2A streaming responses. (`gateway/a2a/`)
|
||||
- **File-store test isolation** — file-store expiry and table tests are less timing-sensitive and isolate their state more reliably. (`store/file/`)
|
||||
|
||||
### Documentation
|
||||
- **First-agent debugging path** — docs now include no-secret transcript checkpoints, durable resume examples, and clearer CLI/website wayfinding for first-agent debugging. (`README.md`, `internal/website/docs/`, `examples/agent-durable/`)
|
||||
|
||||
---
|
||||
|
||||
|
||||
@@ -443,6 +443,7 @@ 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).
|
||||
|
||||
+6
-1
@@ -222,12 +222,16 @@ func (a *agentImpl) Stream(ctx context.Context, message string) (ai.Stream, erro
|
||||
return nil, fmt.Errorf("discover tools: %w", err)
|
||||
}
|
||||
a.mem.Add("user", message)
|
||||
return a.model.Stream(ctx, &ai.Request{
|
||||
stream, err := a.model.Stream(ctx, &ai.Request{
|
||||
Prompt: message,
|
||||
SystemPrompt: a.buildPrompt(),
|
||||
Tools: toolList,
|
||||
Messages: a.mem.Messages(),
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &memoryRecordingStream{stream: stream, memory: a.mem}, nil
|
||||
}
|
||||
|
||||
// Pending returns checkpointed agent runs that have not completed. It mirrors
|
||||
@@ -308,6 +312,7 @@ func (a *agentImpl) askLocked(ctx context.Context, runID, message, parentRunID s
|
||||
})
|
||||
if err != nil {
|
||||
run.Status = agentRunFailureStatus(err)
|
||||
err = agentOperationalError(err)
|
||||
if a.currentRun != nil {
|
||||
run.Steps = a.currentRun.Steps
|
||||
}
|
||||
|
||||
@@ -191,6 +191,43 @@ func agentRunFailureStatus(err error) string {
|
||||
}
|
||||
}
|
||||
|
||||
type operationalError struct {
|
||||
err error
|
||||
hint string
|
||||
}
|
||||
|
||||
func (e *operationalError) Error() string {
|
||||
if e == nil {
|
||||
return ""
|
||||
}
|
||||
return e.err.Error() + "; " + e.hint
|
||||
}
|
||||
|
||||
func (e *operationalError) Unwrap() error {
|
||||
if e == nil {
|
||||
return nil
|
||||
}
|
||||
return e.err
|
||||
}
|
||||
|
||||
func agentOperationalError(err error) error {
|
||||
if err == nil {
|
||||
return nil
|
||||
}
|
||||
switch ai.ClassifyError(err) {
|
||||
case ai.ErrorKindCanceled:
|
||||
return &operationalError{err: err, hint: "agent run canceled; inspect run history with `micro inspect agent <name> --status canceled` or see docs/guides/debugging-agents.md"}
|
||||
case ai.ErrorKindTimeout:
|
||||
return &operationalError{err: err, hint: "agent provider call timed out; inspect run history with `micro inspect agent <name> --status timeout`, then adjust AgentModelCallTimeout/AgentModelRetry or see docs/guides/debugging-agents.md"}
|
||||
case ai.ErrorKindRateLimited:
|
||||
return &operationalError{err: err, hint: "agent provider was rate limited; inspect run history with `micro inspect agent <name> --status rate_limited`, check provider keys with `micro agent preflight`, or see docs/guides/debugging-agents.md"}
|
||||
case ai.ErrorKindUnavailable:
|
||||
return &operationalError{err: err, hint: "agent provider appears temporarily unavailable; retry with bounded AgentModelRetry and verify provider setup with `micro agent preflight` or docs/guides/debugging-agents.md"}
|
||||
default:
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
func (a *agentImpl) checkpointToolWrap(next ai.ToolHandler) ai.ToolHandler {
|
||||
return func(ctx context.Context, call ai.ToolCall) ai.ToolResult {
|
||||
if a.opts.Checkpoint == nil || a.currentRun == nil {
|
||||
|
||||
@@ -2,6 +2,7 @@ package agent
|
||||
|
||||
import (
|
||||
"context"
|
||||
"io"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
@@ -16,6 +17,7 @@ import (
|
||||
// it with a deferred cleanup. Tests in this package are not parallel,
|
||||
// so a package-level hook is safe.
|
||||
var fakeGen func(ctx context.Context, opts ai.Options, req *ai.Request) (*ai.Response, error)
|
||||
var fakeStream func(ctx context.Context, opts ai.Options, req *ai.Request) (ai.Stream, error)
|
||||
|
||||
type fakeModel struct{ opts ai.Options }
|
||||
|
||||
@@ -33,10 +35,33 @@ func (m *fakeModel) Generate(ctx context.Context, req *ai.Request, _ ...ai.Gener
|
||||
return &ai.Response{Reply: "ok"}, nil
|
||||
}
|
||||
func (m *fakeModel) Stream(ctx context.Context, req *ai.Request, _ ...ai.GenerateOption) (ai.Stream, error) {
|
||||
return nil, nil
|
||||
if fakeStream != nil {
|
||||
return fakeStream(ctx, m.opts, req)
|
||||
}
|
||||
return &sliceStream{chunks: []string{"ok"}}, nil
|
||||
}
|
||||
func (m *fakeModel) String() string { return "fake" }
|
||||
|
||||
type sliceStream struct {
|
||||
chunks []string
|
||||
idx int
|
||||
closed bool
|
||||
}
|
||||
|
||||
func (s *sliceStream) Recv() (*ai.Response, error) {
|
||||
if s.idx >= len(s.chunks) {
|
||||
return nil, io.EOF
|
||||
}
|
||||
chunk := s.chunks[s.idx]
|
||||
s.idx++
|
||||
return &ai.Response{Reply: chunk}, nil
|
||||
}
|
||||
|
||||
func (s *sliceStream) Close() error {
|
||||
s.closed = true
|
||||
return nil
|
||||
}
|
||||
|
||||
func init() {
|
||||
ai.Register("fake", func(opts ...ai.Option) ai.Model {
|
||||
m := &fakeModel{}
|
||||
|
||||
+125
-3
@@ -3,7 +3,9 @@ package agent
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"sort"
|
||||
"strings"
|
||||
"time"
|
||||
@@ -18,9 +20,10 @@ import (
|
||||
const agentInstrumentationName = "go-micro.dev/v6/agent"
|
||||
|
||||
const (
|
||||
spanNameRun = "agent.run"
|
||||
spanNameModelCall = "agent.model.call"
|
||||
spanNameToolCall = "agent.tool.call"
|
||||
spanNameRun = "agent.run"
|
||||
spanNameModelCall = "agent.model.call"
|
||||
spanNameModelStream = "agent.model.stream"
|
||||
spanNameToolCall = "agent.tool.call"
|
||||
|
||||
AttrRunID = "agent.run.id"
|
||||
AttrParentRunID = "agent.run.parent_id"
|
||||
@@ -45,6 +48,7 @@ const (
|
||||
AttrFlowStep = "agent.flow.step"
|
||||
AttrDispatch = "agent.dispatch"
|
||||
AttrTrigger = "agent.trigger"
|
||||
AttrRunEventKind = "agent.event.kind"
|
||||
)
|
||||
|
||||
type RunEvent struct {
|
||||
@@ -219,6 +223,123 @@ func (m *tracedModel) Generate(ctx context.Context, req *ai.Request, opts ...ai.
|
||||
return resp, err
|
||||
}
|
||||
|
||||
func (m *tracedModel) Stream(ctx context.Context, req *ai.Request, opts ...ai.GenerateOption) (ai.Stream, error) {
|
||||
info, _ := ai.RunInfoFrom(ctx)
|
||||
provider := m.String()
|
||||
model := m.Options().Model
|
||||
start := time.Now()
|
||||
|
||||
if m.a.opts.TraceProvider == nil {
|
||||
stream, err := m.Model.Stream(ctx, req, opts...)
|
||||
if err != nil {
|
||||
m.a.recordRunEvent(RunEvent{Time: time.Now(), RunID: info.RunID, ParentID: info.ParentID, Agent: info.Agent, Kind: "stream", Provider: provider, Model: model, Attempt: info.Attempt, MaxAttempts: info.MaxAttempts, LatencyMS: time.Since(start).Milliseconds(), Error: err.Error(), ErrorKind: string(ai.ClassifyError(err))})
|
||||
return nil, err
|
||||
}
|
||||
return &tracedStream{Stream: stream, a: m.a, info: info, provider: provider, model: model, start: start}, nil
|
||||
}
|
||||
|
||||
attrs := appendRunInfoAttributes([]attribute.KeyValue{
|
||||
attribute.String(AttrRunID, info.RunID),
|
||||
attribute.String(AttrParentRunID, info.ParentID),
|
||||
attribute.String(AttrAgentName, info.Agent),
|
||||
attribute.String(AttrProvider, provider),
|
||||
attribute.String(AttrModel, model),
|
||||
}, info)
|
||||
ctx, span := m.a.tracer().Start(ctx, spanNameModelStream, trace.WithAttributes(attrs...))
|
||||
stream, err := m.Model.Stream(ctx, req, opts...)
|
||||
if err != nil {
|
||||
dur := time.Since(start).Milliseconds()
|
||||
span.SetAttributes(attribute.Int64(AttrLatencyMS, dur), attribute.String(AttrErrorKind, string(ai.ClassifyError(err))))
|
||||
span.RecordError(err)
|
||||
span.SetStatus(codes.Error, err.Error())
|
||||
e := RunEvent{Time: time.Now(), RunID: info.RunID, ParentID: info.ParentID, Agent: info.Agent, Kind: "stream", Provider: provider, Model: model, Attempt: info.Attempt, MaxAttempts: info.MaxAttempts, LatencyMS: dur, Error: err.Error(), ErrorKind: string(ai.ClassifyError(err))}
|
||||
m.a.recordSpanEvent(span, e)
|
||||
span.End()
|
||||
return nil, err
|
||||
}
|
||||
return &tracedStream{Stream: stream, a: m.a, info: info, provider: provider, model: model, start: start, span: span}, nil
|
||||
}
|
||||
|
||||
type tracedStream struct {
|
||||
ai.Stream
|
||||
a *agentImpl
|
||||
info ai.RunInfo
|
||||
provider string
|
||||
model string
|
||||
start time.Time
|
||||
span trace.Span
|
||||
usage ai.Usage
|
||||
closed bool
|
||||
}
|
||||
|
||||
func (s *tracedStream) Recv() (*ai.Response, error) {
|
||||
resp, err := s.Stream.Recv()
|
||||
if resp != nil {
|
||||
s.usage = mergeUsage(s.usage, resp.Usage)
|
||||
}
|
||||
if err != nil {
|
||||
if errors.Is(err, io.EOF) {
|
||||
s.finish(nil)
|
||||
} else {
|
||||
s.finish(err)
|
||||
}
|
||||
}
|
||||
return resp, err
|
||||
}
|
||||
|
||||
func (s *tracedStream) Close() error {
|
||||
err := s.Stream.Close()
|
||||
s.finish(err)
|
||||
return err
|
||||
}
|
||||
|
||||
func (s *tracedStream) finish(err error) {
|
||||
if s.closed {
|
||||
return
|
||||
}
|
||||
s.closed = true
|
||||
dur := time.Since(s.start).Milliseconds()
|
||||
e := RunEvent{Time: time.Now(), RunID: s.info.RunID, ParentID: s.info.ParentID, Agent: s.info.Agent, Kind: "stream", Provider: s.provider, Model: s.model, Attempt: s.info.Attempt, MaxAttempts: s.info.MaxAttempts, LatencyMS: dur, Tokens: s.usage}
|
||||
if err != nil {
|
||||
e.Error = err.Error()
|
||||
e.ErrorKind = string(ai.ClassifyError(err))
|
||||
}
|
||||
if s.span == nil {
|
||||
s.a.recordRunEvent(e)
|
||||
return
|
||||
}
|
||||
attrs := appendUsage([]attribute.KeyValue{attribute.Int64(AttrLatencyMS, dur)}, s.usage)
|
||||
if s.info.Attempt > 0 {
|
||||
attrs = append(attrs, attribute.Int(AttrAttempt, s.info.Attempt))
|
||||
}
|
||||
if s.info.MaxAttempts > 0 {
|
||||
attrs = append(attrs, attribute.Int(AttrMaxAttempts, s.info.MaxAttempts))
|
||||
}
|
||||
if err != nil {
|
||||
attrs = append(attrs, attribute.String(AttrErrorKind, e.ErrorKind))
|
||||
s.span.RecordError(err)
|
||||
s.span.SetStatus(codes.Error, err.Error())
|
||||
} else {
|
||||
s.span.SetStatus(codes.Ok, "")
|
||||
}
|
||||
s.span.SetAttributes(attrs...)
|
||||
s.a.recordSpanEvent(s.span, e)
|
||||
s.span.End()
|
||||
}
|
||||
|
||||
func mergeUsage(current, next ai.Usage) ai.Usage {
|
||||
if next.InputTokens > current.InputTokens {
|
||||
current.InputTokens = next.InputTokens
|
||||
}
|
||||
if next.OutputTokens > current.OutputTokens {
|
||||
current.OutputTokens = next.OutputTokens
|
||||
}
|
||||
if next.TotalTokens > current.TotalTokens {
|
||||
current.TotalTokens = next.TotalTokens
|
||||
}
|
||||
return current
|
||||
}
|
||||
|
||||
func appendUsage(attrs []attribute.KeyValue, u ai.Usage) []attribute.KeyValue {
|
||||
if u.InputTokens > 0 {
|
||||
attrs = append(attrs, attribute.Int(AttrInputTokens, u.InputTokens))
|
||||
@@ -323,6 +444,7 @@ 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))
|
||||
|
||||
+92
-1
@@ -5,6 +5,7 @@ import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
@@ -279,7 +280,8 @@ func spanEventHasRunInfo(events []trace.Event, name, runID, agentName string) bo
|
||||
continue
|
||||
}
|
||||
attrs := spanAttributes(event.Attributes)
|
||||
if attrs[AttrRunID] == runID && attrs[AttrAgentName] == agentName {
|
||||
wantKind := strings.TrimPrefix(name, "agent.")
|
||||
if attrs[AttrRunID] == runID && attrs[AttrAgentName] == agentName && attrs[AttrRunEventKind] == wantKind {
|
||||
return true
|
||||
}
|
||||
}
|
||||
@@ -577,3 +579,92 @@ func TestListRunSummariesWithOptionsFiltersAndLimits(t *testing.T) {
|
||||
t.Fatalf("filtered summaries = %#v", got)
|
||||
}
|
||||
}
|
||||
|
||||
type otelStreamModel struct{ opts ai.Options }
|
||||
|
||||
func (m *otelStreamModel) Init(opts ...ai.Option) error {
|
||||
for _, o := range opts {
|
||||
o(&m.opts)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
func (m *otelStreamModel) Options() ai.Options { return m.opts }
|
||||
func (m *otelStreamModel) String() string { return "otelstream" }
|
||||
func (m *otelStreamModel) Generate(context.Context, *ai.Request, ...ai.GenerateOption) (*ai.Response, error) {
|
||||
return &ai.Response{Reply: "unused"}, nil
|
||||
}
|
||||
func (m *otelStreamModel) Stream(context.Context, *ai.Request, ...ai.GenerateOption) (ai.Stream, error) {
|
||||
return &otelTestStream{chunks: []*ai.Response{{Reply: "one", Usage: ai.Usage{InputTokens: 1, OutputTokens: 2, TotalTokens: 3}}, {Reply: "two", Usage: ai.Usage{InputTokens: 1, OutputTokens: 4, TotalTokens: 5}}}}, nil
|
||||
}
|
||||
|
||||
type otelTestStream struct {
|
||||
chunks []*ai.Response
|
||||
idx int
|
||||
}
|
||||
|
||||
func (s *otelTestStream) Recv() (*ai.Response, error) {
|
||||
if s.idx >= len(s.chunks) {
|
||||
return nil, io.EOF
|
||||
}
|
||||
resp := s.chunks[s.idx]
|
||||
s.idx++
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
func (s *otelTestStream) Close() error { return nil }
|
||||
|
||||
func TestAgentOpenTelemetrySpansModelStream(t *testing.T) {
|
||||
exp := tracetest.NewInMemoryExporter()
|
||||
tp := trace.NewTracerProvider(trace.WithSyncer(exp))
|
||||
st := store.NewMemoryStore()
|
||||
a := New(Name("stream-runner"), Provider("oteltest"), Model("stream-model"), WithStore(st), TraceProvider(tp))
|
||||
m := a.(*agentImpl).tracedModel(&otelStreamModel{opts: ai.Options{Model: "stream-model"}})
|
||||
ctx := ai.WithRunInfo(context.Background(), ai.RunInfo{RunID: "stream-run-1", ParentID: "parent-run", Agent: "stream-runner", Attempt: 2, MaxAttempts: 3, Flow: "deploy", Step: "plan"})
|
||||
|
||||
stream, err := m.Stream(ctx, &ai.Request{Prompt: "stream"})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for {
|
||||
_, err := stream.Recv()
|
||||
if errors.Is(err, io.EOF) {
|
||||
break
|
||||
}
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
if err := stream.Close(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
spans := exp.GetSpans().Snapshots()
|
||||
var sawStream bool
|
||||
for _, s := range spans {
|
||||
if s.Name() != spanNameModelStream {
|
||||
continue
|
||||
}
|
||||
attrs := spanAttributes(s.Attributes())
|
||||
if attrs[AttrRunID] != "stream-run-1" || attrs[AttrParentRunID] != "parent-run" || attrs[AttrAgentName] != "stream-runner" {
|
||||
t.Fatalf("stream span missing run lineage: %#v", attrs)
|
||||
}
|
||||
if attrs[AttrFlowName] != "deploy" || attrs[AttrFlowStep] != "plan" {
|
||||
t.Fatalf("stream span missing workflow attributes: %#v", attrs)
|
||||
}
|
||||
if attrs[AttrAttempt] != "2" || attrs[AttrMaxAttempts] != "3" || attrs[AttrTotalTokens] != "5" {
|
||||
t.Fatalf("stream span missing attempt/usage attributes: %#v", attrs)
|
||||
}
|
||||
sawStream = true
|
||||
}
|
||||
if !sawStream {
|
||||
t.Fatalf("stream span not emitted; got %d spans", len(spans))
|
||||
}
|
||||
|
||||
events, err := LoadRunEvents(st, "stream-runner", "stream-run-1")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(events) != 1 || events[0].Kind != "stream" || events[0].TraceID == "" || events[0].SpanID == "" || events[0].Tokens.TotalTokens != 5 {
|
||||
t.Fatalf("unexpected stream run event: %#v", events)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -77,6 +77,30 @@ func TestAskRetriesTransientErrorsThenSurfacesStructuredError(t *testing.T) {
|
||||
if attempts != 2 {
|
||||
t.Fatalf("model attempts = %d, want 2", attempts)
|
||||
}
|
||||
if !strings.Contains(err.Error(), "micro inspect agent <name> --status timeout") ||
|
||||
!strings.Contains(err.Error(), "docs/guides/debugging-agents.md") {
|
||||
t.Fatalf("Ask error = %q, want actionable timeout/debugging guidance", err.Error())
|
||||
}
|
||||
}
|
||||
|
||||
func TestAskRateLimitFailureSuggestsPreflightAndInspect(t *testing.T) {
|
||||
fakeGen = func(ctx context.Context, opts ai.Options, req *ai.Request) (*ai.Response, error) {
|
||||
return nil, testStatusError{code: 429}
|
||||
}
|
||||
defer func() { fakeGen = nil }()
|
||||
|
||||
a := newTestAgent(Name("rate-limit-guidance"), ModelRetry(1, time.Millisecond))
|
||||
_, err := a.Ask(context.Background(), "hello")
|
||||
if err == nil {
|
||||
t.Fatal("Ask succeeded, want rate-limit failure")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "micro inspect agent <name> --status rate_limited") ||
|
||||
!strings.Contains(err.Error(), "micro agent preflight") {
|
||||
t.Fatalf("Ask error = %q, want inspect and preflight guidance", err.Error())
|
||||
}
|
||||
if ai.ClassifyError(err) != ai.ErrorKindRateLimited {
|
||||
t.Fatalf("ClassifyError(wrapped error) = %q, want rate_limited", ai.ClassifyError(err))
|
||||
}
|
||||
}
|
||||
|
||||
func TestCanceledAskContextSkipsToolExecution(t *testing.T) {
|
||||
@@ -132,6 +156,34 @@ 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
|
||||
|
||||
@@ -185,6 +185,45 @@ type agentStreamAdapter struct {
|
||||
stream AgentStream
|
||||
}
|
||||
|
||||
type memoryRecordingStream struct {
|
||||
stream ai.Stream
|
||||
memory Memory
|
||||
|
||||
mu sync.Mutex
|
||||
chunks []string
|
||||
closed bool
|
||||
}
|
||||
|
||||
func (s *memoryRecordingStream) Recv() (*ai.Response, error) {
|
||||
resp, err := s.stream.Recv()
|
||||
if resp != nil && resp.Reply != "" {
|
||||
s.mu.Lock()
|
||||
s.chunks = append(s.chunks, resp.Reply)
|
||||
s.mu.Unlock()
|
||||
}
|
||||
if errors.Is(err, io.EOF) {
|
||||
s.recordAssistant()
|
||||
}
|
||||
return resp, err
|
||||
}
|
||||
|
||||
func (s *memoryRecordingStream) Close() error {
|
||||
s.recordAssistant()
|
||||
return s.stream.Close()
|
||||
}
|
||||
|
||||
func (s *memoryRecordingStream) recordAssistant() {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
if s.closed {
|
||||
return
|
||||
}
|
||||
s.closed = true
|
||||
if reply := strings.Join(s.chunks, ""); reply != "" {
|
||||
s.memory.Add("assistant", reply)
|
||||
}
|
||||
}
|
||||
|
||||
func (s *agentStreamAdapter) Recv() (*ai.Response, error) {
|
||||
for {
|
||||
event, err := s.stream.Recv()
|
||||
|
||||
@@ -82,6 +82,52 @@ func TestStreamAskHelperRejectsUnsupportedAgent(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestAgentStreamUsesProviderStreamingAndRecordsAssistantMemory(t *testing.T) {
|
||||
var sawRequest bool
|
||||
fakeStream = func(ctx context.Context, opts ai.Options, req *ai.Request) (ai.Stream, error) {
|
||||
sawRequest = true
|
||||
if req.Prompt != "stream the answer" {
|
||||
t.Fatalf("Prompt = %q, want stream the answer", req.Prompt)
|
||||
}
|
||||
if len(req.Messages) != 1 || req.Messages[0].Role != "user" || req.Messages[0].Content != "stream the answer" {
|
||||
t.Fatalf("Messages = %#v, want current user turn in memory", req.Messages)
|
||||
}
|
||||
return &sliceStream{chunks: []string{"hel", "lo"}}, nil
|
||||
}
|
||||
defer func() { fakeStream = nil }()
|
||||
|
||||
a := newTestAgent(Name("provider-stream"))
|
||||
stream, err := a.Stream(context.Background(), "stream the answer")
|
||||
if err != nil {
|
||||
t.Fatalf("Stream: %v", err)
|
||||
}
|
||||
|
||||
var reply string
|
||||
for {
|
||||
chunk, err := stream.Recv()
|
||||
if errors.Is(err, io.EOF) {
|
||||
break
|
||||
}
|
||||
if err != nil {
|
||||
t.Fatalf("Recv: %v", err)
|
||||
}
|
||||
reply += chunk.Reply
|
||||
}
|
||||
if err := stream.Close(); err != nil {
|
||||
t.Fatalf("Close: %v", err)
|
||||
}
|
||||
if !sawRequest {
|
||||
t.Fatal("provider Stream was not called")
|
||||
}
|
||||
if reply != "hello" {
|
||||
t.Fatalf("reply = %q, want hello", reply)
|
||||
}
|
||||
got := a.mem.Messages()
|
||||
if len(got) != 2 || got[0].Role != "user" || got[0].Content != "stream the answer" || got[1].Role != "assistant" || got[1].Content != "hello" {
|
||||
t.Fatalf("memory = %#v, want user turn and streamed assistant reply", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestResumeStreamAskDoesNotReplayCompletedTool(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
cp := flow.StoreCheckpoint(store.NewStore(), "stream-resume-agent")
|
||||
|
||||
+8
-4
@@ -99,15 +99,19 @@ 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 {
|
||||
|
||||
@@ -646,6 +646,7 @@ micro loop verify # check a repo is wired correctly
|
||||
| Builder | `loop-builder.yml` | Builds the top open item as a single-concern PR, auto-merged on green CI |
|
||||
| Triage | `loop-triage.yml` | Turns CI failures into scoped fix issues, back into the queue |
|
||||
| Coherence | `loop-coherence.yml` | Keeps README/docs/CHANGELOG aligned with the North Star *(opt-in)* |
|
||||
| Security | `loop-security.yml` | Audits for vulnerabilities and files them; never auto-merges fixes, never publishes exploit detail *(opt-in)* |
|
||||
| Release | `loop-release.yml` | Cuts the next patch tag when the branch has new commits *(opt-in)* |
|
||||
|
||||
The workflows are the **mechanism**; each dispatch role's instruction is an editable file in `.github/loop/prompts/` — the **policy**. Edit those prompts (and `.github/loop/NORTH_STAR.md`) to steer the loop without touching the CLI. That split is what lets go-micro itself use `micro loop` while keeping its own richer prompts.
|
||||
|
||||
@@ -16,6 +16,7 @@ type preflightCheck struct {
|
||||
OK bool
|
||||
Detail string
|
||||
Fix string
|
||||
Next string
|
||||
}
|
||||
|
||||
type preflightDeps struct {
|
||||
@@ -52,6 +53,9 @@ 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)
|
||||
@@ -87,19 +91,23 @@ func agentPreflightChecks(deps preflightDeps) []preflightCheck {
|
||||
func checkGoToolchain(deps preflightDeps) preflightCheck {
|
||||
path, err := deps.lookPath("go")
|
||||
if err != nil {
|
||||
return preflightCheck{Name: "Go toolchain", Fix: "Install Go 1.24 or newer and ensure go is on PATH."}
|
||||
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."}
|
||||
}
|
||||
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."}
|
||||
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", OK: true, Detail: fmt.Sprintf("%s (%s)", firstLine(out), path)}
|
||||
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)}
|
||||
}
|
||||
|
||||
func checkMicroBinary(deps preflightDeps) preflightCheck {
|
||||
exe, err := deps.executable()
|
||||
if err != nil || exe == "" {
|
||||
return preflightCheck{Name: "micro binary", Fix: "Install the micro CLI or run this check through go run ./cmd/micro agent preflight."}
|
||||
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."}
|
||||
}
|
||||
version := deps.version()
|
||||
if version == "" {
|
||||
@@ -117,7 +125,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."}
|
||||
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", OK: true, Detail: "found " + strings.Join(found, ", ")}
|
||||
}
|
||||
@@ -125,7 +133,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 + " or run micro run --address with a free port."}
|
||||
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."}
|
||||
}
|
||||
_ = ln.Close()
|
||||
return preflightCheck{Name: "local port " + addr, OK: true, Detail: "available for " + use}
|
||||
@@ -138,3 +146,18 @@ 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,13 +61,59 @@ func TestRunAgentPreflightReportsActionableFailures(t *testing.T) {
|
||||
t.Fatal("runAgentPreflight() error = nil")
|
||||
}
|
||||
got := out.String()
|
||||
for _, want := range []string{"✗ Go toolchain", "Install Go 1.24", "✗ micro binary", "✗ provider API key", "ANTHROPIC_API_KEY", "✗ local port :8080", "micro run --address"} {
|
||||
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"} {
|
||||
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)
|
||||
|
||||
@@ -71,10 +71,11 @@ var dispatchRoles = map[string]dispatchRole{
|
||||
"planner": {"Loop: Planner", "Loop: planning review", "loop-planner", "planner-cron", "0 * * * *"},
|
||||
"builder": {"Loop: Builder", "Loop: build increment", "loop-builder", "builder-cron", "30 * * * *"},
|
||||
"coherence": {"Loop: Coherence", "Loop: coherence review", "loop-coherence", "coherence-cron", "0 7 * * *"},
|
||||
"security": {"Loop: Security", "Loop: security review", "loop-security", "security-cron", "0 6 * * 1"},
|
||||
}
|
||||
|
||||
// allRoles is the full set, in a stable order, for --roles=all and help text.
|
||||
var allRoles = []string{"planner", "builder", "triage", "coherence", "release"}
|
||||
var allRoles = []string{"planner", "builder", "triage", "coherence", "security", "release"}
|
||||
|
||||
const (
|
||||
promptDir = ".github/loop/prompts"
|
||||
@@ -95,6 +96,7 @@ Roles (choose with --roles, default: planner,builder,triage):
|
||||
builder builds the top open item as a single-concern PR (auto-merged on green CI)
|
||||
triage turns CI failures into scoped fix issues back into the queue
|
||||
coherence keeps README/docs/CHANGELOG aligned with the North Star
|
||||
security audits for vulnerabilities and files them (fixes stay human-reviewed)
|
||||
release cuts the next patch tag when the branch has new commits
|
||||
|
||||
Each dispatch role's instruction is an editable file in .github/loop/prompts/ —
|
||||
@@ -127,6 +129,7 @@ Examples:
|
||||
&cli.StringFlag{Name: "planner-cron", Usage: "Cron schedule for the planner", Value: "0 * * * *"},
|
||||
&cli.StringFlag{Name: "builder-cron", Usage: "Cron schedule for the builder", Value: "30 * * * *"},
|
||||
&cli.StringFlag{Name: "coherence-cron", Usage: "Cron schedule for the coherence role", Value: "0 7 * * *"},
|
||||
&cli.StringFlag{Name: "security-cron", Usage: "Cron schedule for the security role", Value: "0 6 * * 1"},
|
||||
&cli.StringFlag{Name: "release-cron", Usage: "Cron schedule for the release role", Value: "0 23 * * *"},
|
||||
&cli.StringFlag{Name: "tag-prefix", Usage: "Tag prefix the release role matches and bumps", Value: "v"},
|
||||
&cli.BoolFlag{Name: "force", Usage: "Overwrite existing loop files"},
|
||||
@@ -171,6 +174,7 @@ func runInit(c *cli.Context) error {
|
||||
"planner": c.String("planner-cron"),
|
||||
"builder": c.String("builder-cron"),
|
||||
"coherence": c.String("coherence-cron"),
|
||||
"security": c.String("security-cron"),
|
||||
}
|
||||
|
||||
if err := scaffold(dir, cfg, roles, crons, c.Bool("force")); err != nil {
|
||||
|
||||
@@ -29,6 +29,7 @@ func renderCases() map[string]config {
|
||||
"templates/prompts/planner.md.tmpl": testCfg,
|
||||
"templates/prompts/builder.md.tmpl": testCfg,
|
||||
"templates/prompts/coherence.md.tmpl": testCfg,
|
||||
"templates/prompts/security.md.tmpl": testCfg,
|
||||
}
|
||||
for role, d := range dispatchRoles {
|
||||
rc := testCfg
|
||||
@@ -59,7 +60,7 @@ func TestRenderIsPlaceholderFreeAndKeepsGHAExpressions(t *testing.T) {
|
||||
|
||||
func TestBaseBranchSubstitutedIntoPrompts(t *testing.T) {
|
||||
// The base branch appears in the PR-opening instructions of these prompts.
|
||||
for _, p := range []string{"planner", "builder", "coherence"} {
|
||||
for _, p := range []string{"planner", "builder", "coherence", "security"} {
|
||||
s := mustRender(t, "templates/prompts/"+p+".md.tmpl", testCfg)
|
||||
if !strings.Contains(s, "--base main") {
|
||||
t.Errorf("%s prompt missing substituted base branch", p)
|
||||
@@ -117,7 +118,7 @@ func TestDispatchWorkflowsStripPromptComments(t *testing.T) {
|
||||
|
||||
func TestPromptsLeaveRuntimeTokensLiteral(t *testing.T) {
|
||||
// __ISSUE__ must survive render (the workflow substitutes it at runtime).
|
||||
for _, p := range []string{"planner", "builder", "coherence", "triage"} {
|
||||
for _, p := range []string{"planner", "builder", "coherence", "triage", "security"} {
|
||||
s := mustRender(t, "templates/prompts/"+p+".md.tmpl", testCfg)
|
||||
if !strings.Contains(s, "__ISSUE__") {
|
||||
t.Errorf("%s prompt lost its __ISSUE__ runtime token", p)
|
||||
@@ -133,19 +134,19 @@ func TestScaffoldAllRolesWritesEverything(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
mustWrite(t, filepath.Join(dir, wfDir, "ci.yml"), "name: CI\n")
|
||||
|
||||
roles := []string{"planner", "builder", "triage", "coherence", "release"}
|
||||
roles := []string{"planner", "builder", "triage", "coherence", "security", "release"}
|
||||
if err := scaffold(dir, testCfg, roles, testCrons, false); err != nil {
|
||||
t.Fatalf("scaffold: %v", err)
|
||||
}
|
||||
|
||||
wantWorkflows := []string{"loop-planner.yml", "loop-builder.yml", "loop-triage.yml", "loop-coherence.yml", "loop-release.yml"}
|
||||
wantWorkflows := []string{"loop-planner.yml", "loop-builder.yml", "loop-triage.yml", "loop-coherence.yml", "loop-security.yml", "loop-release.yml"}
|
||||
for _, w := range wantWorkflows {
|
||||
if !fileExists(filepath.Join(dir, wfDir, w)) {
|
||||
t.Errorf("expected %s", w)
|
||||
}
|
||||
}
|
||||
// Dispatch + triage roles have prompts; release does not.
|
||||
for _, p := range []string{"planner.md", "builder.md", "triage.md", "coherence.md"} {
|
||||
for _, p := range []string{"planner.md", "builder.md", "triage.md", "coherence.md", "security.md"} {
|
||||
if !fileExists(filepath.Join(dir, promptDir, p)) {
|
||||
t.Errorf("expected prompt %s", p)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,22 @@
|
||||
<!--
|
||||
The SECURITY prompt — the editable policy for the security role. The workflow
|
||||
prepends the agent @mention and substitutes __ISSUE__ before posting. Keep
|
||||
__ISSUE__ literal.
|
||||
|
||||
Security is deliberately more conservative than the other roles: it does NOT
|
||||
auto-merge fixes, and it does NOT publish exploit details in public issues.
|
||||
-->
|
||||
Act as the security reviewer for this repository. Audit for real, exploitable vulnerabilities — do not pad the report with theoretical or low-value lint-style noise.
|
||||
|
||||
WHAT TO LOOK FOR: injection (SQL/command/template), authentication and authorization bypass, credential/secret/token exposure (in code, logs, or error messages), SSRF and unsafe outbound requests (especially user- or config-controlled URLs), path traversal, unsafe deserialization, missing or incorrect input validation on trust boundaries (HTTP handlers, RPC endpoints, message consumers), insecure defaults (TLS, auth, permissions), unsafe use of `crypto`/randomness, and known-vulnerable dependencies (run `govulncheck ./...` if available, or inspect `go.mod`).
|
||||
|
||||
DEDUPE against open issues before filing anything.
|
||||
|
||||
HOW TO REPORT — this matters:
|
||||
- **Known/public dependency CVEs** (already disclosed): file an issue labeled `security` referencing the CVE and the affected module, and you MAY open a PR that bumps the dependency to the patched version. Do **NOT** enable auto-merge — leave it for human review.
|
||||
- **Novel, exploitable vulnerabilities in this codebase** (not yet public): do **NOT** post a working exploit, proof-of-concept, or step-by-step reproduction in a public issue — that is irresponsible disclosure. File a CONCISE issue labeled `security` and `needs-human` that names the vulnerability *class*, the *location* (file/function), and the *impact*, with only enough detail for a maintainer to find it — and note it should be handled via the repository's private vulnerability reporting if the repo is public. Do NOT open a public fix PR that reveals the vulnerability; leave the fix to a human.
|
||||
- **Low-risk hardening** (defense-in-depth, missing validation with no proven exploit): a normal `security` issue is fine.
|
||||
|
||||
NEVER auto-merge a security change. Never weaken a control to make a test pass. Anything requiring an architectural or breaking change: label it `needs-human` and describe the tradeoff.
|
||||
|
||||
Post a summary as a comment on this issue (#__ISSUE__) — how many findings by severity, what you filed, and what needs a human — then close it (`gh issue close __ISSUE__`). If you open a dependency-bump PR, do it yourself from the shell: `git switch -c loop/security-__ISSUE__`, `git push -u origin loop/security-__ISSUE__`, `gh pr create --base << .DefaultBranch >> --title "<title>" --body "<summary, Closes #__ISSUE__>"` — then STOP; do NOT run `gh pr merge --auto`. Do not use a make_pr tool.
|
||||
+22
-18
@@ -201,12 +201,13 @@ type Part struct {
|
||||
|
||||
// Message is a turn in an A2A conversation.
|
||||
type Message struct {
|
||||
Role string `json:"role"` // "user" | "agent"
|
||||
Parts []Part `json:"parts"`
|
||||
MessageID string `json:"messageId,omitempty"`
|
||||
TaskID string `json:"taskId,omitempty"`
|
||||
ContextID string `json:"contextId,omitempty"`
|
||||
Kind string `json:"kind"` // "message"
|
||||
Role string `json:"role"` // "user" | "agent"
|
||||
Parts []Part `json:"parts"`
|
||||
MessageID string `json:"messageId,omitempty"`
|
||||
TaskID string `json:"taskId,omitempty"`
|
||||
ContextID string `json:"contextId,omitempty"`
|
||||
Kind string `json:"kind"` // "message"
|
||||
AP2Mandates []AP2SignedMandate `json:"ap2Mandates,omitempty"`
|
||||
}
|
||||
|
||||
// TaskStatus is a task's lifecycle state.
|
||||
@@ -223,12 +224,14 @@ type Artifact struct {
|
||||
|
||||
// Task is the unit of work returned by message/send and tasks/get.
|
||||
type Task struct {
|
||||
ID string `json:"id"`
|
||||
ContextID string `json:"contextId"`
|
||||
Status TaskStatus `json:"status"`
|
||||
Artifacts []Artifact `json:"artifacts,omitempty"`
|
||||
History []Message `json:"history,omitempty"`
|
||||
Kind string `json:"kind"` // "task"
|
||||
ID string `json:"id"`
|
||||
ContextID string `json:"contextId"`
|
||||
Status TaskStatus `json:"status"`
|
||||
Artifacts []Artifact `json:"artifacts,omitempty"`
|
||||
History []Message `json:"history,omitempty"`
|
||||
Kind string `json:"kind"` // "task"
|
||||
AP2Mandates []AP2SignedMandate `json:"ap2Mandates,omitempty"`
|
||||
AP2Verifications []AP2Verification `json:"ap2Verifications,omitempty"`
|
||||
}
|
||||
|
||||
// PushNotificationConfig tells the gateway where to POST task updates for a
|
||||
@@ -848,12 +851,13 @@ func taskFromReplyWithIDsAndHistory(input Message, reply, state, taskID, context
|
||||
input.Kind = "message"
|
||||
}
|
||||
task := &Task{
|
||||
ID: taskID,
|
||||
ContextID: contextID,
|
||||
Kind: "task",
|
||||
History: append(append([]Message{}, history...), input),
|
||||
Status: TaskStatus{State: state, Timestamp: time.Now().UTC().Format(time.RFC3339)},
|
||||
Artifacts: []Artifact{textArtifact(reply)},
|
||||
ID: taskID,
|
||||
ContextID: contextID,
|
||||
Kind: "task",
|
||||
History: append(append([]Message{}, history...), input),
|
||||
Status: TaskStatus{State: state, Timestamp: time.Now().UTC().Format(time.RFC3339)},
|
||||
Artifacts: []Artifact{textArtifact(reply)},
|
||||
AP2Mandates: append([]AP2SignedMandate{}, input.AP2Mandates...),
|
||||
}
|
||||
task.History = append(task.History, Message{
|
||||
Role: "agent",
|
||||
|
||||
@@ -0,0 +1,150 @@
|
||||
package a2a
|
||||
|
||||
import (
|
||||
"crypto/ed25519"
|
||||
"crypto/rand"
|
||||
"crypto/sha256"
|
||||
"encoding/base64"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"time"
|
||||
)
|
||||
|
||||
// AP2MandateKind identifies the AP2 mandate stage represented by a credential.
|
||||
type AP2MandateKind string
|
||||
|
||||
const (
|
||||
AP2CheckoutMandate AP2MandateKind = "checkout"
|
||||
AP2PaymentMandate AP2MandateKind = "payment"
|
||||
)
|
||||
|
||||
// AP2RailRef names the settlement rail authorized by an AP2 payment mandate.
|
||||
type AP2RailRef struct {
|
||||
Type string `json:"type"`
|
||||
Reference string `json:"reference"`
|
||||
}
|
||||
|
||||
// AP2Mandate is a small verifiable AP2 credential. It is deliberately rail-neutral:
|
||||
// x402 is represented as one possible rail reference under a payment mandate.
|
||||
type AP2Mandate struct {
|
||||
ID string `json:"id"`
|
||||
Kind AP2MandateKind `json:"kind"`
|
||||
Subject string `json:"subject,omitempty"`
|
||||
Merchant string `json:"merchant,omitempty"`
|
||||
Amount string `json:"amount,omitempty"`
|
||||
Currency string `json:"currency,omitempty"`
|
||||
Description string `json:"description,omitempty"`
|
||||
TaskID string `json:"taskId,omitempty"`
|
||||
ContextID string `json:"contextId,omitempty"`
|
||||
Rail *AP2RailRef `json:"rail,omitempty"`
|
||||
IssuedAt time.Time `json:"issuedAt"`
|
||||
}
|
||||
|
||||
// AP2SignedMandate is an AP2 mandate plus an Ed25519 signature over its canonical JSON.
|
||||
type AP2SignedMandate struct {
|
||||
Mandate AP2Mandate `json:"mandate"`
|
||||
KeyID string `json:"keyId,omitempty"`
|
||||
Signature string `json:"signature"`
|
||||
}
|
||||
|
||||
// AP2Verification records mandate verification on an A2A task without mixing in
|
||||
// payment-settlement state.
|
||||
type AP2Verification struct {
|
||||
MandateID string `json:"mandateId"`
|
||||
Kind string `json:"kind"`
|
||||
Verified bool `json:"verified"`
|
||||
Error string `json:"error,omitempty"`
|
||||
}
|
||||
|
||||
// NewAP2Keypair returns an Ed25519 keypair suitable for tests or local demos.
|
||||
func NewAP2Keypair() (ed25519.PublicKey, ed25519.PrivateKey, error) {
|
||||
return ed25519.GenerateKey(rand.Reader)
|
||||
}
|
||||
|
||||
// SignAP2Mandate signs a mandate as a verifiable AP2 credential.
|
||||
func SignAP2Mandate(m AP2Mandate, keyID string, private ed25519.PrivateKey) (AP2SignedMandate, error) {
|
||||
if m.ID == "" {
|
||||
return AP2SignedMandate{}, errors.New("ap2: mandate id is required")
|
||||
}
|
||||
if m.Kind == "" {
|
||||
return AP2SignedMandate{}, errors.New("ap2: mandate kind is required")
|
||||
}
|
||||
if m.IssuedAt.IsZero() {
|
||||
m.IssuedAt = time.Now().UTC()
|
||||
}
|
||||
payload, err := ap2Payload(m)
|
||||
if err != nil {
|
||||
return AP2SignedMandate{}, err
|
||||
}
|
||||
return AP2SignedMandate{Mandate: m, KeyID: keyID, Signature: base64.RawURLEncoding.EncodeToString(ed25519.Sign(private, payload))}, nil
|
||||
}
|
||||
|
||||
// VerifyAP2Mandate verifies a signed mandate credential.
|
||||
func VerifyAP2Mandate(s AP2SignedMandate, public ed25519.PublicKey) error {
|
||||
sig, err := base64.RawURLEncoding.DecodeString(s.Signature)
|
||||
if err != nil {
|
||||
return fmt.Errorf("ap2: invalid signature encoding: %w", err)
|
||||
}
|
||||
payload, err := ap2Payload(s.Mandate)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if !ed25519.Verify(public, payload, sig) {
|
||||
return errors.New("ap2: mandate signature verification failed")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// AP2BindMandateToMessage returns a copy of m bound to the A2A message's task/context.
|
||||
func AP2BindMandateToMessage(m AP2Mandate, msg Message) AP2Mandate {
|
||||
m.TaskID = msg.TaskID
|
||||
m.ContextID = msg.ContextID
|
||||
return m
|
||||
}
|
||||
|
||||
// AP2AttachMandate returns a copy of msg carrying the signed AP2 mandate.
|
||||
func AP2AttachMandate(msg Message, mandate AP2SignedMandate) Message {
|
||||
msg.AP2Mandates = append(append([]AP2SignedMandate{}, msg.AP2Mandates...), mandate)
|
||||
return msg
|
||||
}
|
||||
|
||||
// VerifyAP2ForTask verifies signature, task binding, and optional settlement rail reference.
|
||||
func VerifyAP2ForTask(s AP2SignedMandate, public ed25519.PublicKey, task Task, rail *AP2RailRef) AP2Verification {
|
||||
out := AP2Verification{MandateID: s.Mandate.ID, Kind: string(s.Mandate.Kind), Verified: true}
|
||||
if err := VerifyAP2Mandate(s, public); err != nil {
|
||||
out.Verified = false
|
||||
out.Error = err.Error()
|
||||
return out
|
||||
}
|
||||
if s.Mandate.TaskID != "" && s.Mandate.TaskID != task.ID {
|
||||
out.Verified = false
|
||||
out.Error = "ap2: mandate task binding mismatch"
|
||||
return out
|
||||
}
|
||||
if s.Mandate.ContextID != "" && s.Mandate.ContextID != task.ContextID {
|
||||
out.Verified = false
|
||||
out.Error = "ap2: mandate context binding mismatch"
|
||||
return out
|
||||
}
|
||||
if s.Mandate.Kind == AP2PaymentMandate && rail != nil {
|
||||
if s.Mandate.Rail == nil || *s.Mandate.Rail != *rail {
|
||||
out.Verified = false
|
||||
out.Error = "ap2: settlement rail reference mismatch"
|
||||
return out
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// X402AP2Rail builds the x402 settlement rail reference carried under a payment mandate.
|
||||
func X402AP2Rail(reference string) AP2RailRef { return AP2RailRef{Type: "x402", Reference: reference} }
|
||||
|
||||
func ap2Payload(m AP2Mandate) ([]byte, error) {
|
||||
b, err := json.Marshal(m)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("ap2: marshal mandate: %w", err)
|
||||
}
|
||||
sum := sha256.Sum256(b)
|
||||
return sum[:], nil
|
||||
}
|
||||
@@ -0,0 +1,78 @@
|
||||
package a2a
|
||||
|
||||
import (
|
||||
"crypto/ed25519"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
func testAP2Key(t *testing.T) (ed25519.PublicKey, ed25519.PrivateKey) {
|
||||
t.Helper()
|
||||
pub, priv, err := NewAP2Keypair()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return pub, priv
|
||||
}
|
||||
|
||||
func TestAP2CheckoutMandateSignAttachAndVerify(t *testing.T) {
|
||||
pub, priv := testAP2Key(t)
|
||||
msg := Message{Role: "user", Kind: "message", TaskID: "task-1", ContextID: "ctx-1", Parts: []Part{{Kind: "text", Text: "buy"}}}
|
||||
mandate := AP2BindMandateToMessage(AP2Mandate{ID: "checkout-1", Kind: AP2CheckoutMandate, Subject: "alice", Merchant: "store", Amount: "10.00", Currency: "USD", Description: "demo", IssuedAt: time.Unix(1, 0).UTC()}, msg)
|
||||
signed, err := SignAP2Mandate(mandate, "test-key", priv)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
msg = AP2AttachMandate(msg, signed)
|
||||
task := taskFromReplyWithIDs(msg, "ok", stateCompleted, msg.TaskID, msg.ContextID)
|
||||
|
||||
if len(task.AP2Mandates) != 1 {
|
||||
t.Fatalf("expected mandate carried on task, got %d", len(task.AP2Mandates))
|
||||
}
|
||||
got := VerifyAP2ForTask(task.AP2Mandates[0], pub, *task, nil)
|
||||
if !got.Verified || got.Error != "" {
|
||||
t.Fatalf("expected verified mandate, got %+v", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAP2PaymentMandateX402RailReference(t *testing.T) {
|
||||
pub, priv := testAP2Key(t)
|
||||
rail := X402AP2Rail("payreq_123")
|
||||
task := Task{ID: "task-2", ContextID: "ctx-2"}
|
||||
signed, err := SignAP2Mandate(AP2Mandate{ID: "payment-1", Kind: AP2PaymentMandate, TaskID: task.ID, ContextID: task.ContextID, Rail: &rail, IssuedAt: time.Unix(1, 0).UTC()}, "test-key", priv)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
got := VerifyAP2ForTask(signed, pub, task, &rail)
|
||||
if !got.Verified {
|
||||
t.Fatalf("expected x402 rail to verify, got %+v", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAP2TamperCasesFailDistinctly(t *testing.T) {
|
||||
pub, priv := testAP2Key(t)
|
||||
rail := X402AP2Rail("payreq_123")
|
||||
task := Task{ID: "task-3", ContextID: "ctx-3"}
|
||||
signed, err := SignAP2Mandate(AP2Mandate{ID: "payment-2", Kind: AP2PaymentMandate, TaskID: task.ID, ContextID: task.ContextID, Rail: &rail, IssuedAt: time.Unix(1, 0).UTC()}, "test-key", priv)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
tampered := signed
|
||||
tampered.Mandate.Amount = "999.00"
|
||||
if got := VerifyAP2ForTask(tampered, pub, task, &rail); got.Verified || !strings.Contains(got.Error, "signature") {
|
||||
t.Fatalf("expected signature failure, got %+v", got)
|
||||
}
|
||||
|
||||
wrongTask := task
|
||||
wrongTask.ID = "other-task"
|
||||
if got := VerifyAP2ForTask(signed, pub, wrongTask, &rail); got.Verified || !strings.Contains(got.Error, "task binding") {
|
||||
t.Fatalf("expected task binding failure, got %+v", got)
|
||||
}
|
||||
|
||||
otherRail := X402AP2Rail("payreq_other")
|
||||
if got := VerifyAP2ForTask(signed, pub, task, &otherRail); got.Verified || !strings.Contains(got.Error, "rail reference") {
|
||||
t.Fatalf("expected rail reference failure, got %+v", got)
|
||||
}
|
||||
}
|
||||
@@ -107,12 +107,13 @@ func (c *Client) SendMessage(ctx context.Context, message Message) (*Task, error
|
||||
return nil, err
|
||||
}
|
||||
return &Task{
|
||||
ID: m.TaskID,
|
||||
ContextID: m.ContextID,
|
||||
Kind: "task",
|
||||
Status: TaskStatus{State: stateCompleted, Timestamp: time.Now().UTC().Format(time.RFC3339)},
|
||||
Artifacts: []Artifact{textArtifact(textOf(m.Parts))},
|
||||
History: []Message{m},
|
||||
ID: m.TaskID,
|
||||
ContextID: m.ContextID,
|
||||
Kind: "task",
|
||||
Status: TaskStatus{State: stateCompleted, Timestamp: time.Now().UTC().Format(time.RFC3339)},
|
||||
Artifacts: []Artifact{textArtifact(textOf(m.Parts))},
|
||||
History: []Message{m},
|
||||
AP2Mandates: append([]AP2SignedMandate{}, m.AP2Mandates...),
|
||||
}, nil
|
||||
}
|
||||
|
||||
|
||||
@@ -24,6 +24,7 @@ Actions instead of subagents. Each role is a workflow:
|
||||
| **Evaluator** | `harness.yml` — *Harness (E2E)*, plus the CI gate (`tests.yaml`, `lint.yaml`) | Grades every change: the mock harness + unit/lint on each push/PR, and real-model conformance hourly. A *separate* grader — never the generator judging itself. |
|
||||
| **Evaluator → feedback** | `loop-triage.yml` — *Loop: Triage (Evaluator feedback)* | When a gate workflow (Lint, Run Tests, or the harness) fails on a non-PR run, root-causes, dedupes, and files scoped fix issues back into the planner's queue. The hill-climbing feedback path. |
|
||||
| **Coherence** | `loop-coherence.yml` — *Loop: Coherence* | Keeps README/website/docs/blog aligned with the North Star, keeps `CHANGELOG.md` living (reconciling `[Unreleased]` against merged PRs and rolling it into version headings as tags cut), and drafts the changelog blog post. |
|
||||
| **Security** | `loop-security.yml` — *Loop: Security* | Weekly vulnerability audit of the attack surface (MCP/A2A gateways, x402, auth, provider URLs, agent tool loop, deps via `govulncheck`). Files `security` issues; **never auto-merges** fixes and **never publishes exploit detail** in public issues (responsible disclosure); risky fixes are `needs-human`. |
|
||||
| **Release** | `loop-release.yml` — *Loop: Release (daily patch)* | Cuts a daily patch tag when master has new commits, so the *installable* framework tracks the loop's improvements (triggers `release.yml`/goreleaser). Minor/major bumps stay with the human. |
|
||||
|
||||
Generation is separated from evaluation on purpose: an agent grading its own work
|
||||
|
||||
@@ -148,13 +148,17 @@ func main() {
|
||||
fmt.Fprintf(os.Stderr, "content-type = %q, want text/event-stream\n", ct)
|
||||
os.Exit(1)
|
||||
}
|
||||
payload, err := readSSEData(res.Body)
|
||||
summary, err := readSSESummary(res.Body)
|
||||
if err != nil {
|
||||
fmt.Fprintln(os.Stderr, err)
|
||||
os.Exit(1)
|
||||
}
|
||||
if !strings.Contains(payload, "a2a-fallback-ok") {
|
||||
fmt.Fprintf(os.Stderr, "stream payload missing marker: %s\n", payload)
|
||||
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)
|
||||
os.Exit(1)
|
||||
}
|
||||
if !sawTool || !sawRunInfo {
|
||||
@@ -164,10 +168,16 @@ func main() {
|
||||
fmt.Println("\n\033[32m✓ A2A message/stream fell back to Ask and preserved tool/run metadata\033[0m")
|
||||
}
|
||||
|
||||
func readSSEData(r io.Reader) (string, error) {
|
||||
type streamSummary struct {
|
||||
Payload string
|
||||
State string
|
||||
HasArtifactText bool
|
||||
}
|
||||
|
||||
func readSSESummary(r io.Reader) (streamSummary, error) {
|
||||
scanner := bufio.NewScanner(r)
|
||||
var event strings.Builder
|
||||
var payload strings.Builder
|
||||
var summary streamSummary
|
||||
seen := false
|
||||
flush := func() error {
|
||||
data := strings.TrimSpace(event.String())
|
||||
@@ -175,19 +185,40 @@ func readSSEData(r io.Reader) (string, error) {
|
||||
if data == "" {
|
||||
return nil
|
||||
}
|
||||
if !json.Valid([]byte(data)) {
|
||||
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 {
|
||||
return fmt.Errorf("SSE data event is not JSON: %s", data)
|
||||
}
|
||||
seen = true
|
||||
payload.WriteString(data)
|
||||
payload.WriteByte('\n')
|
||||
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
|
||||
}
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
for scanner.Scan() {
|
||||
line := scanner.Text()
|
||||
if strings.TrimSpace(line) == "" {
|
||||
if err := flush(); err != nil {
|
||||
return "", err
|
||||
return streamSummary{}, err
|
||||
}
|
||||
continue
|
||||
}
|
||||
@@ -201,13 +232,13 @@ func readSSEData(r io.Reader) (string, error) {
|
||||
event.WriteString(strings.TrimSpace(data))
|
||||
}
|
||||
if err := scanner.Err(); err != nil {
|
||||
return "", err
|
||||
return streamSummary{}, err
|
||||
}
|
||||
if err := flush(); err != nil {
|
||||
return "", err
|
||||
return streamSummary{}, err
|
||||
}
|
||||
if !seen {
|
||||
return "", errors.New("no SSE data received")
|
||||
return streamSummary{}, errors.New("no SSE data received")
|
||||
}
|
||||
return payload.String(), nil
|
||||
return summary, nil
|
||||
}
|
||||
|
||||
@@ -5,19 +5,26 @@ import (
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestReadSSEDataAcceptsMultipleJSONEvents(t *testing.T) {
|
||||
payload, err := readSSEData(strings.NewReader("event: status\ndata: {\"phase\":\"started\"}\n\ndata:{\"result\":\"a2a-fallback-ok\"}\n\n"))
|
||||
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"))
|
||||
if err != nil {
|
||||
t.Fatalf("readSSEData returned error: %v", err)
|
||||
t.Fatalf("readSSESummary() error = %v", err)
|
||||
}
|
||||
if !strings.Contains(payload, "started") || !strings.Contains(payload, "a2a-fallback-ok") {
|
||||
t.Fatalf("payload = %q, want both event payloads", payload)
|
||||
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)
|
||||
}
|
||||
}
|
||||
|
||||
func TestReadSSEDataRejectsInvalidJSONEvent(t *testing.T) {
|
||||
_, err := readSSEData(strings.NewReader("data: {\"ok\":true}\n\ndata: {bad json}\n\n"))
|
||||
func TestReadSSESummaryRejectsNonJSONData(t *testing.T) {
|
||||
_, err := readSSESummary(strings.NewReader("data: not-json\n\n"))
|
||||
if err == nil {
|
||||
t.Fatal("readSSEData returned nil error for invalid JSON event")
|
||||
t.Fatal("readSSESummary() error = nil, want non-JSON error")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -451,9 +451,12 @@ func waitForPlanDelegateExecution(done <-chan error, taskSvc *TaskService, notif
|
||||
tasks := taskSvc.count()
|
||||
notify := notifySvc.count()
|
||||
if err != nil {
|
||||
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
|
||||
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)
|
||||
}
|
||||
return fmt.Errorf("flow execute after side effects tasks=%d notify=%d: %w", tasks, notify, err)
|
||||
}
|
||||
@@ -461,12 +464,18 @@ func waitForPlanDelegateExecution(done <-chan error, taskSvc *TaskService, notif
|
||||
if recoverMissingNotify == nil || tasks != 3 || notify != 0 {
|
||||
return fmt.Errorf("delegation completed without required notify side effect: notify=%d, want 1", notify)
|
||||
}
|
||||
fmt.Print("\n\033[33mwarning:\033[0m flow completed before delegated notify; retrying the missing comms handoff once.\n")
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
||||
retryErr := recoverMissingNotify(ctx)
|
||||
cancel()
|
||||
if retryErr != nil {
|
||||
return fmt.Errorf("delegation completed without required notify side effect and recovery failed: notify=%d, want 1: %w", notify, retryErr)
|
||||
settled, err := waitForNotifySideEffect(notifySvc, 2*time.Second)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if !settled {
|
||||
fmt.Print("\n\033[33mwarning:\033[0m flow completed before delegated notify; retrying the missing comms handoff once.\n")
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
||||
retryErr := recoverMissingNotify(ctx)
|
||||
cancel()
|
||||
if retryErr != nil {
|
||||
return fmt.Errorf("delegation completed without required notify side effect and recovery failed: notify=%d, want 1: %w", notify, retryErr)
|
||||
}
|
||||
}
|
||||
if notify = notifySvc.count(); notify != 1 {
|
||||
return fmt.Errorf("delegation recovery completed without required notify side effect: notify=%d, want 1", notify)
|
||||
@@ -481,6 +490,26 @@ func waitForPlanDelegateExecution(done <-chan error, taskSvc *TaskService, notif
|
||||
}
|
||||
}
|
||||
|
||||
func waitForNotifySideEffect(notifySvc *NotifyService, timeout time.Duration) (bool, error) {
|
||||
deadline := time.Now().Add(timeout)
|
||||
for {
|
||||
if notifySvc.count() == 1 {
|
||||
return true, nil
|
||||
}
|
||||
if dup := notifySvc.duplicateAttempts(); dup > 0 {
|
||||
return false, fmt.Errorf("duplicate notify attempts: got %d duplicate replay(s), want 0", dup)
|
||||
}
|
||||
if !time.Now().Before(deadline) {
|
||||
return false, nil
|
||||
}
|
||||
time.Sleep(50 * time.Millisecond)
|
||||
}
|
||||
}
|
||||
|
||||
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")
|
||||
|
||||
@@ -316,6 +316,47 @@ func TestPlanDelegateExecutionRecoversMissingNotifyOnce(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestPlanDelegateExecutionWaitsForInFlightNotifyAfterFlowCompletion(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)
|
||||
}
|
||||
}
|
||||
notifySvc := new(NotifyService)
|
||||
done := make(chan error, 1)
|
||||
done <- nil
|
||||
|
||||
go func() {
|
||||
time.Sleep(100 * time.Millisecond)
|
||||
var rsp SendResponse
|
||||
_ = notifySvc.Send(context.Background(), &SendRequest{To: "owner@acme.com", Message: "The launch plan is ready"}, &rsp)
|
||||
}()
|
||||
|
||||
recovered := false
|
||||
err := waitForPlanDelegateExecution(done, taskSvc, notifySvc, func(ctx context.Context) error {
|
||||
recovered = true
|
||||
var rsp SendResponse
|
||||
return notifySvc.Send(ctx, &SendRequest{To: "owner@acme.com", Message: "The launch plan is ready"}, &rsp)
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("waitForPlanDelegateExecution returned %v, want in-flight notify success", err)
|
||||
}
|
||||
if recovered {
|
||||
t.Fatal("missing notify recovery ran while delegated notify was still in flight")
|
||||
}
|
||||
if got := taskSvc.count(); got != 3 {
|
||||
t.Fatalf("task count = %d, want 3 after in-flight notify settles", got)
|
||||
}
|
||||
if got := notifySvc.count(); got != 1 {
|
||||
t.Fatalf("notify count = %d, want 1 after in-flight notify settles", got)
|
||||
}
|
||||
if got := notifySvc.duplicateAttempts(); got != 0 {
|
||||
t.Fatalf("duplicate notify attempts = %d, want 0", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestPlanDelegateExecutionAcceptsClientTimeoutAfterSideEffects(t *testing.T) {
|
||||
taskSvc := new(TaskService)
|
||||
for _, title := range []string{"Design", "Build", "Ship"} {
|
||||
@@ -338,7 +379,7 @@ func TestPlanDelegateExecutionAcceptsClientTimeoutAfterSideEffects(t *testing.T)
|
||||
}
|
||||
}
|
||||
|
||||
func TestPlanDelegateExecutionRejectsClientTimeoutBeforeSideEffects(t *testing.T) {
|
||||
func TestPlanDelegateExecutionClassifiesClientTimeoutBeforeSideEffects(t *testing.T) {
|
||||
done := make(chan error, 1)
|
||||
done <- errors.New(`{"id":"go.micro.client","code":408,"detail":"<nil>","status":"Request Timeout"}`)
|
||||
|
||||
@@ -346,8 +387,35 @@ func TestPlanDelegateExecutionRejectsClientTimeoutBeforeSideEffects(t *testing.T
|
||||
if err == nil {
|
||||
t.Fatal("waitForPlanDelegateExecution returned nil, want timeout before side effects to fail")
|
||||
}
|
||||
if got := err.Error(); !strings.Contains(got, "tasks=0 notify=0") {
|
||||
t.Fatalf("error = %q, want side-effect counts", got)
|
||||
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)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -74,6 +74,56 @@ func TestGuidesNavigationLeadsWithDoing(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestFirstAgentWayfindingDocs(t *testing.T) {
|
||||
root := filepath.Clean(filepath.Join("..", "..", ".."))
|
||||
checks := []struct {
|
||||
name string
|
||||
file string
|
||||
heading string
|
||||
links []string
|
||||
}{
|
||||
{
|
||||
name: "README first-agent on-ramp",
|
||||
file: filepath.Join(root, "README.md"),
|
||||
heading: "### First agent on-ramp",
|
||||
links: []string{
|
||||
"internal/website/docs/guides/no-secret-first-agent.md",
|
||||
"internal/website/docs/guides/your-first-agent.md",
|
||||
"internal/website/docs/guides/debugging-agents.md",
|
||||
"internal/website/docs/guides/zero-to-hero.md",
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "website getting-started on-ramp",
|
||||
file: filepath.Join(root, "internal", "website", "docs", "getting-started.md"),
|
||||
heading: "### First-agent on-ramp",
|
||||
links: []string{
|
||||
"guides/no-secret-first-agent.html",
|
||||
"guides/your-first-agent.html",
|
||||
"guides/debugging-agents.html",
|
||||
"guides/zero-to-hero.html",
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
for _, check := range checks {
|
||||
t.Run(check.name, func(t *testing.T) {
|
||||
doc := firstMarkdownSection(t, readFile(t, check.file), check.heading)
|
||||
last := -1
|
||||
for _, link := range check.links {
|
||||
idx := strings.Index(doc, link)
|
||||
if idx == -1 {
|
||||
t.Fatalf("%s missing first-agent wayfinding link %q; keep the no-secret → first-agent → debugging → 0→hero path discoverable", check.name, link)
|
||||
}
|
||||
if idx < last {
|
||||
t.Fatalf("%s link %q appeared out of order; expected no-secret → first-agent → debugging → 0→hero", check.name, link)
|
||||
}
|
||||
last = idx
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestNoSecretFirstAgentTranscript(t *testing.T) {
|
||||
root := filepath.Clean(filepath.Join("..", "..", ".."))
|
||||
guide := readFile(t, filepath.Join(root, "internal", "website", "docs", "guides", "no-secret-first-agent.md"))
|
||||
@@ -94,6 +144,20 @@ func TestNoSecretFirstAgentTranscript(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
debugCheckpoint := firstMarkdownSection(t, guide, "## Debug transcript checkpoint")
|
||||
for _, want := range []string{
|
||||
`micro chat assistant --prompt "Triage ticket-1 for Alice"`,
|
||||
"micro inspect agent assistant --limit 1",
|
||||
"micro agent history assistant",
|
||||
"status, event count, last event",
|
||||
"Debugging your agent",
|
||||
"debugging-agents.html",
|
||||
} {
|
||||
if !strings.Contains(debugCheckpoint, want) {
|
||||
t.Fatalf("no-secret debug transcript checkpoint missing %q", want)
|
||||
}
|
||||
}
|
||||
|
||||
readme := readFile(t, filepath.Join(root, "README.md"))
|
||||
if !strings.Contains(readme, "internal/website/docs/guides/no-secret-first-agent.md") {
|
||||
t.Fatal("README does not point to the no-secret first-agent transcript")
|
||||
@@ -105,6 +169,19 @@ func TestNoSecretFirstAgentTranscript(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func firstMarkdownSection(t *testing.T, doc, heading string) string {
|
||||
t.Helper()
|
||||
start := strings.Index(doc, heading)
|
||||
if start == -1 {
|
||||
t.Fatalf("missing %q section", heading)
|
||||
}
|
||||
section := doc[start+len(heading):]
|
||||
if next := strings.Index(section, "\n##"); next != -1 {
|
||||
section = section[:next]
|
||||
}
|
||||
return section
|
||||
}
|
||||
|
||||
func readFile(t *testing.T, name string) string {
|
||||
t.Helper()
|
||||
data, err := os.ReadFile(name)
|
||||
|
||||
@@ -8,6 +8,7 @@ cd "$ROOT"
|
||||
# without secrets or long-running daemons.
|
||||
go test ./cmd/micro -run 'TestFirstAgentWalkthroughCLIBoundaries|TestZeroToHeroCLIBoundaries' -count=1
|
||||
go test ./cmd/micro/cli/deploy -run TestDeployDryRun -count=1
|
||||
go test ./internal/harness/zero-to-hero-ci -run 'TestNoSecretFirstAgentTranscript|TestZeroToHeroReferenceDocs' -count=1
|
||||
|
||||
# Deterministic no-secret reference scenarios. These use the real Go Micro
|
||||
# runtime and mock only the LLM provider. The support example is the maintained
|
||||
|
||||
@@ -16,6 +16,7 @@ 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. |
|
||||
|
||||
@@ -35,7 +36,11 @@ 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)
|
||||
and [`examples/support`](https://github.com/micro/go-micro/tree/master/examples/support).
|
||||
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.
|
||||
|
||||
## 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, the `micro` binary, provider-key setup, and whether the default `micro run` gateway port is free, without calling an LLM provider.
|
||||
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.
|
||||
|
||||
## Install
|
||||
|
||||
@@ -86,9 +86,14 @@ Created project Launch and added task 'Write docs' to it.
|
||||
|
||||
The console discovers services from the registry and orchestrates across them via the agent. Use `micro run -d` for detached mode without the console, or `micro chat` as a standalone command.
|
||||
|
||||
If the agent surprises you while iterating, use the [Debugging your agent](guides/debugging-agents.html) guide to inspect service registration, tool calls, run history, memory, provider failures, and flow handoffs.
|
||||
### First-agent on-ramp
|
||||
|
||||
When you are ready to prove the whole path end to end, follow the [0→hero reference path](guides/zero-to-hero.html). It is the canonical handoff from this quick start: scaffold a service, run it locally, chat with an agent, inspect durable agent/flow history, and finish with `micro deploy --dry-run` using the same commands exercised by `make harness`.
|
||||
After this quick start, follow the agent path in order:
|
||||
|
||||
1. [No-secret first-agent transcript](guides/no-secret-first-agent.html) — run a useful support agent with a mock model before setting up a provider key.
|
||||
2. [Your First Agent](guides/your-first-agent.html) — build a service-backed agent and talk to it with `micro chat`.
|
||||
3. [Debugging your agent](guides/debugging-agents.html) — inspect service registration, tool calls, run history, memory, provider failures, and flow handoffs when the agent surprises you.
|
||||
4. [0→hero reference path](guides/zero-to-hero.html) — prove the full scaffold → run → chat → inspect → deploy dry-run lifecycle with commands exercised by `make harness`.
|
||||
|
||||
## Quick Start: Write a Service
|
||||
|
||||
|
||||
@@ -186,6 +186,16 @@ This is the JSON-RPC binding for task execution:
|
||||
|
||||
Both directions work: the gateway exposes your agents, and `a2a.Client` (via `flow.A2A` or `delegate` to a URL) calls external ones. The task binding is what makes a Go Micro agent both reachable from, and able to reach, the A2A ecosystem today.
|
||||
|
||||
## AP2 mandate layer (opt-in)
|
||||
|
||||
AP2 sits above A2A as a verifiable-intent and audit layer. Go Micro keeps the
|
||||
A2A envelope separate from payment settlement: an A2A message can carry signed
|
||||
AP2 checkout or payment mandates, and the resulting task can retain the stable
|
||||
mandate reference plus verification result. Payment settlement state remains in
|
||||
the payment rail. For x402, use an AP2 payment mandate with an `x402` rail
|
||||
reference to name the payment requirement; the existing x402 facilitator still
|
||||
performs verification and settlement.
|
||||
|
||||
## See also
|
||||
|
||||
- [MCP & AI Agents](../mcp.html) — exposing services as tools
|
||||
|
||||
@@ -17,7 +17,9 @@ 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.
|
||||
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.
|
||||
|
||||
## 1. Reproduce one small turn
|
||||
|
||||
|
||||
@@ -89,6 +89,26 @@ CI keeps those CLI boundaries present with:
|
||||
go test ./cmd/micro -run TestFirstAgentWalkthroughCLIBoundaries -count=1
|
||||
```
|
||||
|
||||
If chat behaves unexpectedly, continue to
|
||||
## Debug transcript checkpoint
|
||||
|
||||
A successful first chat turn should always leave an inspectable trail. After the
|
||||
chat command finishes, continue the same terminal transcript with the inspection
|
||||
and history commands before changing prompts or provider settings:
|
||||
|
||||
```sh
|
||||
micro chat assistant --prompt "Triage ticket-1 for Alice"
|
||||
micro inspect agent assistant --limit 1
|
||||
micro agent history assistant
|
||||
```
|
||||
|
||||
The inspection output is the checkpoint that the runnable loop did not stop at
|
||||
chat: it should show a recent agent run with a status, event count, last event,
|
||||
and trace breadcrumb when tracing is configured. `micro agent history assistant`
|
||||
then confirms the conversation memory that future turns will reuse. If either
|
||||
command is empty after a successful chat turn, keep the failing transcript and
|
||||
use [Debugging your agent](debugging-agents.html) to check provider failures, run
|
||||
history, memory, and tool-call inspection before changing application code.
|
||||
|
||||
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
|
||||
[Debugging your agent](debugging-agents.html) for provider checks, run history,
|
||||
memory, and tool-call inspection.
|
||||
|
||||
@@ -128,3 +128,11 @@ Leave those variables unset in normal CI; the live test skips unless the facilit
|
||||
- [Building Effective Agents — Agents and Workflows](agents-and-workflows.html)
|
||||
- [MCP & AI Agents](../mcp.html)
|
||||
- [x402 — Coinbase Developer Docs](https://docs.cdp.coinbase.com/x402/welcome) · [x402 on Solana](https://solana.com/x402/what-is-x402)
|
||||
|
||||
## AP2 payment mandates
|
||||
|
||||
AP2 can authorize an x402 payment without making A2A carry settlement state. A
|
||||
payment mandate records the buyer intent and names an `x402` rail reference; the
|
||||
existing x402 facilitator remains responsible for payment verification and
|
||||
settlement. This keeps AP2 as the signed mandate/audit layer while x402 stays the
|
||||
pluggable payment rail.
|
||||
|
||||
@@ -53,7 +53,7 @@ Run the read-only first-agent preflight before starting the walkthrough. The sam
|
||||
micro agent preflight
|
||||
```
|
||||
|
||||
It checks Go, the `micro` binary, provider-key setup, and the default local gateway port without contacting a provider.
|
||||
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.
|
||||
|
||||
## 1. Create a workspace
|
||||
|
||||
|
||||
@@ -18,6 +18,7 @@ cloud credentials?"
|
||||
| Boundary | Contract | CI check |
|
||||
| --- | --- | --- |
|
||||
| Scaffold | `micro new` generates a runnable service with and without MCP support. | `go test ./cmd/micro/cli/new -run TestZeroToOne -count=1` |
|
||||
| First-agent wayfinding | README and the website getting-started docs keep the no-secret → first-agent → debugging → 0→hero links present and in order. | `go test ./internal/harness/zero-to-hero-ci -run TestFirstAgentWayfindingDocs -count=1` |
|
||||
| First agent | `micro new`, `micro agent preflight`, `micro run`, `micro chat`, and `micro inspect agent` stay available for the documented first-agent walkthrough. | `go test ./cmd/micro -run TestFirstAgentWalkthroughCLIBoundaries -count=1` |
|
||||
| Run | `micro run` remains the local development entry point. | `go test ./cmd/micro -run TestZeroToHeroCLIBoundaries -count=1` |
|
||||
| Chat | `micro chat` remains the interactive agent entry point. | `go test ./cmd/micro -run TestZeroToHeroCLIBoundaries -count=1` |
|
||||
|
||||
@@ -56,7 +56,11 @@ After that first-agent path, branch out to:
|
||||
```go
|
||||
package main
|
||||
|
||||
import "go-micro.dev/v6"
|
||||
import (
|
||||
"context"
|
||||
|
||||
"go-micro.dev/v6"
|
||||
)
|
||||
|
||||
type Greeter struct{}
|
||||
|
||||
@@ -74,7 +78,11 @@ func main() {
|
||||
|
||||
### Pub/Sub Event Handler
|
||||
```go
|
||||
import "go-micro.dev/v6"
|
||||
import (
|
||||
"context"
|
||||
|
||||
"go-micro.dev/v6"
|
||||
)
|
||||
|
||||
func main() {
|
||||
service := micro.NewService("subscriber")
|
||||
@@ -82,7 +90,7 @@ func main() {
|
||||
// Subscribe to events
|
||||
micro.RegisterSubscriber("user.created", service.Server(),
|
||||
func(ctx context.Context, event *UserCreatedEvent) error {
|
||||
log.Infof("User created: %s", event.Email)
|
||||
// Handle the event here.
|
||||
return nil
|
||||
},
|
||||
)
|
||||
|
||||
+17
-32
@@ -135,22 +135,22 @@ func fileTest(s Store, t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// Write 3 records with various expiry and get with Suffix
|
||||
// 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.
|
||||
records = []*Record{
|
||||
{
|
||||
Key: "foo",
|
||||
Value: []byte("foofoo"),
|
||||
},
|
||||
{
|
||||
Key: "barfoo",
|
||||
Value: []byte("barfoobarfoo"),
|
||||
|
||||
Expiry: time.Millisecond * 100,
|
||||
Key: "barfoo",
|
||||
Value: []byte("barfoobarfoo"),
|
||||
Expiry: -time.Second,
|
||||
},
|
||||
{
|
||||
Key: "bazbarfoo",
|
||||
Value: []byte("bazbarfoobazbarfoo"),
|
||||
Expiry: 2 * time.Millisecond * 100,
|
||||
Key: "bazbarfoo",
|
||||
Value: []byte("bazbarfoobazbarfoo"),
|
||||
},
|
||||
}
|
||||
for _, r := range records {
|
||||
@@ -160,39 +160,24 @@ 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) != 3 {
|
||||
t.Errorf("Expected 3 items, got %d", len(results))
|
||||
// t.Logf("Table test: %v\n", spew.Sdump(results))
|
||||
}
|
||||
} 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)
|
||||
}
|
||||
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) != 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))
|
||||
}
|
||||
} else if len(results) != 1 {
|
||||
t.Errorf("Expected 1 unexpired suffix item, got %d (%# v)", len(results), 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