Compare commits

..

21 Commits

Author SHA1 Message Date
Codex e4c8e8cf5e Correlate agent trace spans
Harness (E2E) / Harnesses (mock LLM) (push) Waiting to run
Harness (E2E) / Provider harnesses (live LLM conformance) (push) Waiting to run
Lint / golangci-lint (push) Waiting to run
Run Tests / Unit Tests (push) Waiting to run
Run Tests / Etcd Integration Tests (push) Waiting to run
2026-06-25 22:09:21 +00:00
Asim Aslam f1f987cb8d provider conformance: print capability matrix (#3092)
Co-authored-by: Codex <codex@openai.com>
2026-06-25 22:18:34 +01:00
Asim Aslam 7b9c583dcf blog: "How Go Micro Builds Itself" — the autonomous Codex loop (#3090)
A post on how the framework is increasingly built by an autonomous loop of
two AI agents (Codex implementing scoped increments, Claude Code
orchestrating, the human setting direction): the dispatch→build→auto-merge
mechanism, the three altitudes (architect/increment/DevRel), the failure
modes we hit wiring it up, and the concrete increments it produces.


Claude-Session: https://claude.ai/code/session_01CmdEY7pYmV5zzwCjNJ4ykL

Co-authored-by: Claude <noreply@anthropic.com>
2026-06-25 21:59:00 +01:00
Asim Aslam b7f1ac6d6d Add flow run correlation to step contexts (#3089)
Co-authored-by: Codex <codex@openai.com>
2026-06-25 21:20:38 +01:00
Asim Aslam fb14ccbc17 Add deterministic AI capability rows (#3087)
Co-authored-by: Codex <codex@openai.com>
2026-06-25 20:21:00 +01:00
Asim Aslam 5f40cae7af ci: add DevRel + Architect overseer passes to the loop (#3085)
The hourly loop ships increments but nothing watches the whole. Add two
periodic high-altitude passes, same dispatch mechanism (fresh issue → @codex):

- devrel-review.yml (daily): audits README, website, docs, blog for coherence
  with the North Star, README crispness, and blog-worthy material. Safe
  alignment/crispness fixes auto-merge; brand/positioning copy and blog drafts
  are surfaced in a report for the human, never auto-merged.
- architecture-review.yml (every ~3 days): reviews the framework/harness against
  the thesis and files scoped follow-up issues that feed the increment loop. It
  does not make breaking/architectural changes itself.

Documented both in CONTINUOUS_IMPROVEMENT.md (Overseer passes).


Claude-Session: https://claude.ai/code/session_01CmdEY7pYmV5zzwCjNJ4ykL

Co-authored-by: Claude <noreply@anthropic.com>
2026-06-25 19:42:58 +01:00
Asim Aslam 496d5b7644 Show parent agent runs in CLI history (#3084)
Co-authored-by: Codex <codex@openai.com>
2026-06-25 19:28:32 +01:00
Asim Aslam 44e62f5e8b flow: honor canceled checkpoint contexts (#3082)
Co-authored-by: Codex <codex@openai.com>
2026-06-25 18:24:35 +01:00
Asim Aslam 18ae00dad0 Harden agent harness contracts (#3080)
Co-authored-by: Codex <codex@openai.com>
2026-06-25 17:36:36 +01:00
Asim Aslam 6f9adfd375 ai: expose unsupported streaming sentinel (#3078)
Co-authored-by: Codex <codex@openai.com>
2026-06-25 16:38:14 +01:00
Asim Aslam 8fc4f58a17 Add AI provider capability matrix (#3076)
Co-authored-by: Codex <codex@openai.com>
2026-06-25 15:26:49 +01:00
Asim Aslam 19ad28e04b Propagate A2A request cancellation (#3074)
Co-authored-by: Codex <codex@openai.com>
2026-06-25 14:32:06 +01:00
Asim Aslam 319ff1972b ci: require live provider conformance keys (#3072)
Co-authored-by: Codex <codex@openai.com>
2026-06-25 13:28:57 +01:00
Asim Aslam 03791759ef Add filtered agent run summaries (#3070)
Co-authored-by: Codex <codex@openai.com>
2026-06-25 12:35:39 +01:00
Asim Aslam 5e3cef7cd9 Add agent run summary status and duration (#3068)
Co-authored-by: Codex <codex@openai.com>
2026-06-25 10:50:49 +01:00
Asim Aslam 751530aaec website: refresh Developer Experience copy for the harness DX (#3066)
The section only mentioned micro run (services) and micro deploy. Update it
to the current lifecycle DX: micro new scaffolds a service or agent (every
endpoint an MCP tool), micro run starts everything with hot reload + gateway
+ console, and micro chat talks to agents and services. Drops micro deploy.


Claude-Session: https://claude.ai/code/session_01CmdEY7pYmV5zzwCjNJ4ykL

Co-authored-by: Claude <noreply@anthropic.com>
2026-06-25 10:33:34 +01:00
Asim Aslam 3a9f45750b docs + website: loop mechanics doc; landing "Features" + subtitle trim (#3065)
* ci: self-merge Codex PRs via native auto-merge; retire the sweep

With branch protection + "Allow auto-merge" now enabled on master, Codex
enables GitHub auto-merge on its own PR (gh pr merge --squash --auto) right
after opening it, so the PR lands the moment the required CI checks pass —
no polling sweep, and the green-CI gate is enforced by GitHub instead of by
gh pr checks in a cron. Removes auto-merge-codex.yml and updates the dispatch
and AGENTS.md accordingly.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01CmdEY7pYmV5zzwCjNJ4ykL

* docs: document the durable loop mechanics (stub→gh, branch, auto-merge)

Capture the hard-won wiring of the autonomous loop so it isn't re-derived:
fresh issue per increment, user-PAT dispatch (Codex ignores the bot), Codex
opening the PR via gh (make_pr is a no-op stub), unique codex/ branch + label,
and native auto-merge gated by branch protection with 0 approvals. Adds a
"do-not-break" list (don't re-add approvals, don't reuse one tracker issue,
don't use make_pr, don't re-implement during the summary→PR lag).

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01CmdEY7pYmV5zzwCjNJ4ykL

* website: rename section to "Features"; trim subtitle

Rename the feature-grid heading from "The Runtime Around the Agent" to
"Features", and drop "once they leave the demo." from the subtitle.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01CmdEY7pYmV5zzwCjNJ4ykL

---------

Co-authored-by: Claude <noreply@anthropic.com>
2026-06-25 10:22:51 +01:00
Asim Aslam 2def581408 ci: self-merge Codex PRs via native auto-merge; retire the sweep (#3064)
With branch protection + "Allow auto-merge" now enabled on master, Codex
enables GitHub auto-merge on its own PR (gh pr merge --squash --auto) right
after opening it, so the PR lands the moment the required CI checks pass —
no polling sweep, and the green-CI gate is enforced by GitHub instead of by
gh pr checks in a cron. Removes auto-merge-codex.yml and updates the dispatch
and AGENTS.md accordingly.


Claude-Session: https://claude.ai/code/session_01CmdEY7pYmV5zzwCjNJ4ykL

Co-authored-by: Claude <noreply@anthropic.com>
2026-06-25 10:13:24 +01:00
Asim Aslam 95ee402c8b ci: pin Codex PRs to a unique codex/ branch + codex label (#3063)
Codex was pushing to a generic branch (e.g. "work"), which the auto-merge
sweep ignores (it only matches codex/* branches with the codex label) and
which collides across runs. The dispatch and AGENTS.md now instruct Codex to
create a unique codex/increment-<issue> branch and pass --label codex to
gh pr create, so every increment is isolated and actually swept.


Claude-Session: https://claude.ai/code/session_01CmdEY7pYmV5zzwCjNJ4ykL

Co-authored-by: Claude <noreply@anthropic.com>
2026-06-25 09:48:25 +01:00
Asim Aslam 8aeba56839 Harden provider conformance selection (#3061)
Co-authored-by: Codex <codex@openai.com>
2026-06-25 09:47:55 +01:00
Asim Aslam 42fdad439c docs: AGENTS.md — open PRs via gh, not the make_pr stub (#3062)
The initial AGENTS.md instructed Codex tasks to use the make_pr tool and
said a shell GitHub token "is not a substitute" — but make_pr in the Codex
sandbox is a no-op stub that only records metadata and never opens a PR, so
that guidance steers every run back into the broken path. Replace it with
the working flow: push the branch and open the PR with the gh CLI
(git push + gh pr create), which is what actually creates PRs now.


Claude-Session: https://claude.ai/code/session_01CmdEY7pYmV5zzwCjNJ4ykL

Co-authored-by: Claude <noreply@anthropic.com>
2026-06-25 09:32:50 +01:00
45 changed files with 1018 additions and 159 deletions
+48
View File
@@ -0,0 +1,48 @@
name: Architecture Review
# Periodic high-altitude review of the whole framework and harness against the
# North Star (internal/docs/THESIS.md) — part of the autonomous loop
# (internal/docs/CONTINUOUS_IMPROVEMENT.md). Where DevRel watches the public
# story, the architect watches the system: API coherence, lifecycle gaps, and
# whether recent increments are converging on the thesis or sprawling.
#
# The architect's OUTPUT is an assessment plus scoped follow-up issues that feed
# the hourly increment loop — NOT large refactors. Breaking public-API and
# architectural changes stay with the human (see CONTINUOUS_IMPROVEMENT.md).
#
# Opens a fresh issue and dispatches Codex via CODEX_TRIGGER_TOKEN.
on:
workflow_dispatch: {}
schedule:
- cron: "0 8 */3 * *" # roughly every 3 days, 08:00 UTC (tunable)
permissions:
issues: write
concurrency:
group: architecture-review
cancel-in-progress: false
jobs:
dispatch:
runs-on: ubuntu-latest
steps:
- name: Open an architecture review issue and dispatch Codex
env:
GH_TOKEN: ${{ secrets.CODEX_TRIGGER_TOKEN || github.token }}
HAS_TRIGGER_TOKEN: ${{ secrets.CODEX_TRIGGER_TOKEN != '' }}
REPO: ${{ github.repository }}
RUN_NUMBER: ${{ github.run_number }}
run: |
if [ "$HAS_TRIGGER_TOKEN" != "true" ]; then
echo "CODEX_TRIGGER_TOKEN is not set — skipping (Codex ignores Actions-bot comments)."
exit 0
fi
ISSUE_URL=$(gh issue create --repo "$REPO" \
--title "Architecture review #$RUN_NUMBER" \
--body "Periodic architecture / harness review against the North Star in internal/docs/THESIS.md. Output: an assessment plus scoped follow-up issues for the increment loop.")
ISSUE_NUM="${ISSUE_URL##*/}"
echo "Opened issue #$ISSUE_NUM — dispatching Codex (Architect)."
gh issue comment "$ISSUE_NUM" --repo "$REPO" --body \
"@codex Act as the architect for go-micro. Review the overall framework and harness against the North Star in internal/docs/THESIS.md (services → agents → workflows as one runtime) and the roadmap in ROADMAP.md. Assess: API coherence and consistency across the core packages (agent, ai, flow, gateway/mcp, gateway/a2a, model, server, store, registry), gaps or missing pieces in the services → agents → workflows lifecycle, duplication or drift, and whether recent increments (scan recently merged PRs) are converging on the thesis or sprawling. Then: (A) post a concise architectural assessment as a comment on this issue (#$ISSUE_NUM) — strengths, the top risks/gaps, and a recommended direction for the next increments; (B) file concrete, scoped follow-up issues for the highest-value gaps so the hourly increment loop can pick them up — \`gh issue create --label codex --label enhancement --title \"<scoped task>\" --body \"<goal, scope, acceptance criteria>\"\` (each must be a single, self-contained, CI-verifiable chunk). Do NOT make breaking public-API or architectural changes yourself — your output is the assessment and the issues. A small, safe doc/comment correction may be a PR (\`git switch -c codex/architect-$ISSUE_NUM\` … \`gh pr create --base master --label codex …\` … \`gh pr merge --squash --auto --delete-branch\`). Do not use the make_pr tool (it is a no-op stub)."
-42
View File
@@ -1,42 +0,0 @@
name: Auto-merge Codex PRs
# Part of the autonomous improvement loop (internal/docs/CONTINUOUS_IMPROVEMENT.md).
# Merges Codex's PRs once CI is green — no human involvement; CI (build, test,
# golangci-lint, harnesses) is the only gate. Scoped to PRs that are BOTH
# codex-labelled AND from a codex/* branch, so nothing else can auto-merge.
on:
schedule:
- cron: "*/15 * * * *" # sweep every 15 min
workflow_dispatch: {}
permissions:
contents: write
pull-requests: write
concurrency:
group: auto-merge-codex
cancel-in-progress: false
jobs:
merge:
runs-on: ubuntu-latest
steps:
- name: Merge green Codex PRs
env:
GH_TOKEN: ${{ github.token }}
REPO: ${{ github.repository }}
run: |
gh pr list --repo "$REPO" --label codex --state open \
--json number,headRefName \
--jq '.[] | select(.headRefName | startswith("codex/")) | .number' \
| while read -r pr; do
[ -z "$pr" ] && continue
if gh pr checks "$pr" --repo "$REPO" >/dev/null 2>&1; then
echo "Checks green on #$pr — merging."
gh pr merge "$pr" --repo "$REPO" --squash --delete-branch \
|| echo "skip #$pr (not mergeable — conflicts?)"
else
echo "skip #$pr (checks pending/failing)"
fi
done
+7 -4
View File
@@ -5,8 +5,11 @@ name: Continuous Improvement
#
# A Claude Max subscription provides no API key for CI, so the loop is driven by
# Codex rather than Claude Code: on a cadence this opens a fresh tracking issue and
# posts an @codex instruction on it, and Codex runs one improvement increment + opens
# a PR. (See the per-issue rationale below.)
# posts an @codex instruction on it, and Codex runs one improvement increment, opens
# a PR (git push + gh pr create — the make_pr tool is a no-op stub), and enables
# GitHub auto-merge (gh pr merge --auto) so the PR lands once the required CI checks
# pass. No separate merge sweep — branch protection + native auto-merge is the gate.
# (See the per-issue rationale below.)
#
# Codex does NOT respond to comments authored by the github-actions bot, so the
# dispatch is GATED on a CODEX_TRIGGER_TOKEN secret (a PAT for a user account Codex
@@ -52,11 +55,11 @@ jobs:
echo "a user account Codex follows) to activate the loop."
exit 0
fi
# A unique issue per run → Codex derives a unique branch → no collisions.
# A unique issue per run → unique codex/ branch → no collisions.
ISSUE_URL=$(gh issue create --repo "$REPO" \
--title "Continuous improvement increment #$RUN_NUMBER" \
--body "Autonomous continuous-improvement increment. North Star: internal/docs/THESIS.md; charter: internal/docs/CONTINUOUS_IMPROVEMENT.md. Tracker: #3024.")
ISSUE_NUM="${ISSUE_URL##*/}"
echo "Opened issue #$ISSUE_NUM — dispatching Codex."
gh issue comment "$ISSUE_NUM" --repo "$REPO" --body \
"@codex Run one continuous-improvement increment per internal/docs/CONTINUOUS_IMPROVEMENT.md, aligned to the North Star in internal/docs/THESIS.md (the holistic services → agents → workflows lifecycle). Pick the single highest-value roadmap/issue/improvement-radar item that advances that thesis, implement it, and verify \`go build ./...\`, \`go test ./...\`, and \`golangci-lint run ./...\`. Then open the PR YOURSELF from the shell — do NOT use the make_pr tool (in this environment it only records metadata and never creates a PR). Run \`git push -u origin HEAD\` then \`gh pr create --base master --title \"<title>\" --body \"<body, including 'Closes #$ISSUE_NUM'>\"\`; the gh CLI is installed and authenticated and origin points to $REPO. One concern per PR; stay out of brand/positioning copy and breaking public API."
"@codex Run one continuous-improvement increment per internal/docs/CONTINUOUS_IMPROVEMENT.md, aligned to the North Star in internal/docs/THESIS.md (the holistic services → agents → workflows lifecycle). Pick the single highest-value roadmap/issue/improvement-radar item that advances that thesis, implement it, and verify \`go build ./...\`, \`go test ./...\`, and \`golangci-lint run ./...\`. Then open the PR YOURSELF from the shell — do NOT use the make_pr tool (in this environment it only records metadata and never creates a PR). Create a uniquely-named branch under the codex/ prefix and open the PR from it: \`git switch -c codex/increment-$ISSUE_NUM\`, then \`git push -u origin codex/increment-$ISSUE_NUM\`, then \`gh pr create --base master --label codex --title \"<title>\" --body \"<body, including 'Closes #$ISSUE_NUM'>\"\`. Finally enable auto-merge so GitHub merges it once CI is green: \`gh pr merge --squash --auto --delete-branch\`. The gh CLI is installed and authenticated and origin points to $REPO. One concern per PR; stay out of brand/positioning copy and breaking public API."
+47
View File
@@ -0,0 +1,47 @@
name: DevRel Review
# Daily higher-altitude coherence pass over the PUBLIC surface — README,
# website (landing + docs), and blog — part of the autonomous loop
# (internal/docs/CONTINUOUS_IMPROVEMENT.md). The hourly increment loop ships
# code; this keeps the story coherent: docs/website aligned, README crisp, and
# a steady supply of things worth blogging about.
#
# Like the increment loop it opens a fresh issue and dispatches Codex via
# CODEX_TRIGGER_TOKEN (Codex ignores Actions-bot comments). Autonomy boundary:
# SAFE factual-alignment and crispness fixes auto-merge; brand/positioning copy
# and blog drafts are surfaced in the report for the human, never auto-merged.
on:
workflow_dispatch: {}
schedule:
- cron: "0 7 * * *" # daily, 07:00 UTC (tunable)
permissions:
issues: write
concurrency:
group: devrel-review
cancel-in-progress: false
jobs:
dispatch:
runs-on: ubuntu-latest
steps:
- name: Open a DevRel review issue and dispatch Codex
env:
GH_TOKEN: ${{ secrets.CODEX_TRIGGER_TOKEN || github.token }}
HAS_TRIGGER_TOKEN: ${{ secrets.CODEX_TRIGGER_TOKEN != '' }}
REPO: ${{ github.repository }}
RUN_NUMBER: ${{ github.run_number }}
run: |
if [ "$HAS_TRIGGER_TOKEN" != "true" ]; then
echo "CODEX_TRIGGER_TOKEN is not set — skipping (Codex ignores Actions-bot comments)."
exit 0
fi
ISSUE_URL=$(gh issue create --repo "$REPO" \
--title "DevRel coherence review #$RUN_NUMBER" \
--body "Daily DevRel / coherence pass over README, website (landing + docs), and the blog. North Star: internal/docs/THESIS.md.")
ISSUE_NUM="${ISSUE_URL##*/}"
echo "Opened issue #$ISSUE_NUM — dispatching Codex (DevRel)."
gh issue comment "$ISSUE_NUM" --repo "$REPO" --body \
"@codex Act as DevRel for go-micro. Audit the PUBLIC surface — \`README.md\`, \`internal/website/\` (landing \`index.html\` + \`docs/\`), and the blog under \`internal/website/blog/\` — for coherence with the North Star in internal/docs/THESIS.md (an agent harness and service framework; the services → agents → workflows lifecycle). Look for: (1) places where README / website / docs contradict each other, are stale, or describe behavior that has since changed (cross-check against the code and recent merged PRs / CHANGELOG.md); (2) whether the README is crisp and leads with the harness positioning; (3) one to three genuinely blog-worthy items from recently shipped work. Then do BOTH of these: (A) post a concise findings report as a comment on this issue (#$ISSUE_NUM) — what is aligned, what drifted, what you fixed, and the blog ideas; (B) for SAFE factual-alignment and crispness fixes only (NOT brand/marketing/positioning rewrites), open one PR: \`git switch -c codex/devrel-$ISSUE_NUM\`, \`git push -u origin codex/devrel-$ISSUE_NUM\`, \`gh pr create --base master --label codex --title \"<title>\" --body \"<summary, including 'Closes #$ISSUE_NUM'>\"\`, then \`gh pr merge --squash --auto --delete-branch\`. Leave brand/positioning copy and blog drafts for the human — describe them in the report, do NOT open auto-merging PRs for them. Do not use the make_pr tool (it is a no-op stub). If you touch code, verify go build/test/golangci-lint. Stay out of breaking public-API changes."
+6 -5
View File
@@ -2,8 +2,9 @@ name: Harness (E2E)
# Runs the end-to-end harnesses for agents, services, flows, and provider
# conformance. The default job uses deterministic mock LLMs and needs no
# secrets. A second job runs the same harnesses against any live providers
# whose API key secrets are configured.
# secrets. A second job runs the same harnesses against the live provider set and
# fails if any selected provider secret is missing, so scheduled conformance
# cannot silently degrade.
on:
push:
@@ -34,7 +35,7 @@ jobs:
run: go run ./internal/harness/plan-delegate
harness-live:
name: Provider harnesses (live LLM, if keys present)
name: Provider harnesses (live LLM conformance)
runs-on: ubuntu-latest
# Only on the daily schedule or a manual run — never automatically on
# every push/PR, so changes don't quietly burn API credits. Trigger it
@@ -47,7 +48,7 @@ jobs:
with:
go-version: stable
cache: true
- name: Provider conformance against configured live models
- name: Provider conformance against live models
env:
ANTHROPIC_API_KEY: ${{ secrets.ANTHROPIC_API_KEY }}
OPENAI_API_KEY: ${{ secrets.OPENAI_API_KEY }}
@@ -56,4 +57,4 @@ jobs:
MISTRAL_API_KEY: ${{ secrets.MISTRAL_API_KEY }}
TOGETHER_API_KEY: ${{ secrets.TOGETHER_API_KEY }}
ATLASCLOUD_API_KEY: ${{ secrets.ATLASCLOUD_API_KEY }}
run: go run ./internal/harness/provider-conformance
run: go run ./internal/harness/provider-conformance -require-configured
+28 -4
View File
@@ -7,9 +7,33 @@ These instructions apply to the entire repository.
When a Codex task makes repository changes and the requested outcome is a PR:
1. Keep the change focused on the assigned issue or prompt.
2. Run the relevant verification commands and capture their results.
2. Run the relevant verification commands and capture their results
(`go build ./...`, `go test ./...`, `golangci-lint run ./...`).
3. Check `git status --short` and review the diff before finishing.
4. Stage the intended files and create a local git commit on the current branch.
5. Use the Codex `make_pr` tool to open the pull request with a concise title and a body that summarizes the change and testing.
4. Create a uniquely-named branch under the `codex/` prefix (do not work on
`master`, and do not use a generic name like `work`):
Do not just say that a PR was opened. If local changes exist, the task is not complete until the changes are committed and the `make_pr` tool has been called. A GitHub token in the shell environment is not a substitute for the Codex `make_pr` tool in this environment.
```sh
git switch -c codex/<issue-number>-<short-slug>
```
5. Stage the intended files and commit on that branch.
6. Open the pull request yourself with the GitHub CLI, which is installed in the
environment and whose `origin` points at this repository, then enable
auto-merge so GitHub merges it once the required CI checks pass:
```sh
git push -u origin HEAD
gh pr create --base master --label codex \
--title "<concise title>" \
--body "<summary of the change and testing, including 'Closes #<issue>'>"
gh pr merge --squash --auto --delete-branch
```
The branch should start with `codex/` and the PR should carry the `codex`
label. Auto-merge waits for the required status checks (build, tests,
golangci-lint) — never merge a PR manually before CI is green.
Do not just say that a PR was opened, and do **not** rely on the `make_pr` tool:
in this environment `make_pr` only records the title/body and never pushes a
branch or creates a PR. The task is not complete until `gh pr create` has opened
a real pull request and printed its URL.
+1
View File
@@ -261,6 +261,7 @@ func (a *agentImpl) Run() error {
a.server = server.NewServer(
server.Name(a.opts.Name),
server.Address(a.opts.Address),
server.Registry(a.opts.Registry),
server.Metadata(map[string]string{
"type": "agent",
+8
View File
@@ -39,6 +39,7 @@ type Options struct {
Provider string
Model string
APIKey string
Address string
Registry registry.Registry
Client client.Client
Store store.Store
@@ -135,6 +136,13 @@ func APIKey(k string) Option {
return func(o *Options) { o.APIKey = k }
}
// Address sets the network address for the agent's service endpoint.
// Use "127.0.0.1:0" in local harnesses/tests to bind an ephemeral loopback
// port and avoid advertising the default service address.
func Address(addr string) Option {
return func(o *Options) { o.Address = addr }
}
// WithRegistry sets the service registry.
func WithRegistry(r registry.Registry) Option {
return func(o *Options) { o.Registry = r }
+90 -22
View File
@@ -56,18 +56,31 @@ type RunEvent struct {
type Usage = ai.Usage
// RunListOptions controls how recorded agent run summaries are returned.
// Zero values preserve the full deterministic run list.
type RunListOptions struct {
// Status, when set, keeps only runs with the matching status
// (for example "running", "done", "error", or "refused").
Status string
// Limit, when positive, returns the most recently updated runs up to
// the limit. Limited results are ordered newest first.
Limit int
}
// RunSummary is a compact index entry for a recorded agent run.
type RunSummary struct {
RunID string `json:"run_id"`
Agent string `json:"agent"`
ParentID string `json:"parent_id,omitempty"`
TraceID string `json:"trace_id,omitempty"`
SpanID string `json:"span_id,omitempty"`
StartedAt time.Time `json:"started_at"`
UpdatedAt time.Time `json:"updated_at"`
Events int `json:"events"`
LastKind string `json:"last_kind,omitempty"`
LastError string `json:"last_error,omitempty"`
RunID string `json:"run_id"`
Agent string `json:"agent"`
ParentID string `json:"parent_id,omitempty"`
TraceID string `json:"trace_id,omitempty"`
SpanID string `json:"span_id,omitempty"`
StartedAt time.Time `json:"started_at"`
UpdatedAt time.Time `json:"updated_at"`
DurationMS int64 `json:"duration_ms,omitempty"`
Events int `json:"events"`
Status string `json:"status,omitempty"`
LastKind string `json:"last_kind,omitempty"`
LastError string `json:"last_error,omitempty"`
}
func (a *agentImpl) tracer() trace.Tracer {
@@ -111,7 +124,13 @@ func (m *tracedModel) Generate(ctx context.Context, req *ai.Request, opts ...ai.
info, _ := ai.RunInfoFrom(ctx)
provider := m.String()
model := m.Options().Model
ctx, span := m.a.tracer().Start(ctx, spanNameModelCall, trace.WithAttributes(attribute.String(AttrProvider, provider), attribute.String(AttrModel, model)))
ctx, span := m.a.tracer().Start(ctx, spanNameModelCall, trace.WithAttributes(
attribute.String(AttrRunID, info.RunID),
attribute.String(AttrParentRunID, info.ParentID),
attribute.String(AttrAgentName, info.Agent),
attribute.String(AttrProvider, provider),
attribute.String(AttrModel, model),
))
start := time.Now()
resp, err := m.Model.Generate(ctx, req, opts...)
dur := time.Since(start).Milliseconds()
@@ -156,7 +175,13 @@ func (a *agentImpl) traceTool(next ai.ToolHandler) ai.ToolHandler {
}
return func(ctx context.Context, call ai.ToolCall) ai.ToolResult {
info, _ := ai.RunInfoFrom(ctx)
ctx, span := a.tracer().Start(ctx, spanNameToolCall, trace.WithAttributes(attribute.String(AttrToolName, call.Name), attribute.Bool(AttrDelegate, call.Name == toolDelegate)))
ctx, span := a.tracer().Start(ctx, spanNameToolCall, trace.WithAttributes(
attribute.String(AttrRunID, info.RunID),
attribute.String(AttrParentRunID, info.ParentID),
attribute.String(AttrAgentName, info.Agent),
attribute.String(AttrToolName, call.Name),
attribute.Bool(AttrDelegate, call.Name == toolDelegate),
))
start := time.Now()
res := next(ctx, call)
dur := time.Since(start).Milliseconds()
@@ -210,6 +235,12 @@ func (a *agentImpl) recordRunEvent(e RunEvent) {
// ListRunSummaries returns a deterministic summary of recorded runs for agentName.
func ListRunSummaries(s store.Store, agentName string) ([]RunSummary, error) {
return ListRunSummariesWithOptions(s, agentName, RunListOptions{})
}
// ListRunSummariesWithOptions returns summaries of recorded runs for agentName,
// optionally filtered by status and limited to the most recently updated runs.
func ListRunSummariesWithOptions(s store.Store, agentName string, opts RunListOptions) ([]RunSummary, error) {
st := store.Scope(s, "agent", agentName)
keys, err := st.List(store.ListPrefix("runs/"))
if err != nil {
@@ -240,16 +271,18 @@ func ListRunSummaries(s store.Store, agentName string) ([]RunSummary, error) {
first := events[0]
last := events[len(events)-1]
summary := RunSummary{
RunID: id,
Agent: first.Agent,
ParentID: first.ParentID,
TraceID: first.TraceID,
SpanID: first.SpanID,
StartedAt: first.Time,
UpdatedAt: last.Time,
Events: len(events),
LastKind: last.Kind,
LastError: last.Error,
RunID: id,
Agent: first.Agent,
ParentID: first.ParentID,
TraceID: first.TraceID,
SpanID: first.SpanID,
StartedAt: first.Time,
UpdatedAt: last.Time,
DurationMS: last.Time.Sub(first.Time).Milliseconds(),
Events: len(events),
Status: runStatus(events),
LastKind: last.Kind,
LastError: last.Error,
}
for _, e := range events {
if e.Agent != "" {
@@ -268,11 +301,46 @@ func ListRunSummaries(s store.Store, agentName string) ([]RunSummary, error) {
summary.LastError = e.Error
}
}
if opts.Status != "" && summary.Status != opts.Status {
continue
}
summaries = append(summaries, summary)
}
if opts.Limit > 0 {
sort.SliceStable(summaries, func(i, j int) bool {
return summaries[i].UpdatedAt.After(summaries[j].UpdatedAt)
})
if len(summaries) > opts.Limit {
summaries = summaries[:opts.Limit]
}
}
return summaries, nil
}
func runStatus(events []RunEvent) string {
if len(events) == 0 {
return ""
}
status := "running"
for _, e := range events {
if e.Error != "" {
status = "error"
}
if e.Refused != "" && status != "error" {
status = "refused"
}
switch e.Kind {
case "error":
status = "error"
case "done":
if status == "running" {
status = "done"
}
}
}
return status
}
func LoadRunEvents(s store.Store, agentName, runID string) ([]RunEvent, error) {
st := store.Scope(s, "agent", agentName)
keys, err := st.List(store.ListPrefix("runs/" + runID + "/"))
+62 -2
View File
@@ -8,6 +8,7 @@ import (
"go-micro.dev/v6/ai"
"go-micro.dev/v6/store"
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/sdk/trace"
"go.opentelemetry.io/otel/sdk/trace/tracetest"
)
@@ -46,16 +47,33 @@ func TestAgentOpenTelemetrySpans(t *testing.T) {
}
spans := exp.GetSpans().Snapshots()
want := map[string]bool{spanNameRun: false, spanNameModelCall: false, spanNameToolCall: false}
var runID string
for _, s := range spans {
if _, ok := want[s.Name()]; ok {
want[s.Name()] = true
}
attrs := spanAttributes(s.Attributes())
if s.Name() == spanNameRun {
runID = attrs[AttrRunID]
}
}
for name, seen := range want {
if !seen {
t.Fatalf("span %s not emitted; got %d spans", name, len(spans))
}
}
if runID == "" {
t.Fatal("run span missing run id attribute")
}
for _, s := range spans {
if s.Name() != spanNameModelCall && s.Name() != spanNameToolCall {
continue
}
attrs := spanAttributes(s.Attributes())
if attrs[AttrRunID] != runID || attrs[AttrAgentName] != "runner" {
t.Fatalf("%s missing run correlation attributes: %#v", s.Name(), attrs)
}
}
keys, err := store.Scope(st, "agent", "runner").List(store.ListPrefix("runs/"))
if err != nil {
t.Fatal(err)
@@ -73,6 +91,12 @@ func TestAgentOpenTelemetrySpans(t *testing.T) {
if summaries[0].LastKind != "done" {
t.Fatalf("LastKind = %q, want done", summaries[0].LastKind)
}
if summaries[0].Status != "done" {
t.Fatalf("Status = %q, want done", summaries[0].Status)
}
if summaries[0].DurationMS < 0 {
t.Fatalf("DurationMS = %d, want non-negative", summaries[0].DurationMS)
}
if summaries[0].TraceID == "" || summaries[0].SpanID == "" {
t.Fatalf("summary missing trace correlation: %#v", summaries[0])
}
@@ -85,6 +109,14 @@ func TestAgentOpenTelemetrySpans(t *testing.T) {
}
}
func spanAttributes(attrs []attribute.KeyValue) map[string]string {
out := make(map[string]string, len(attrs))
for _, attr := range attrs {
out[string(attr.Key)] = attr.Value.AsString()
}
return out
}
func TestAgentOpenTelemetryNoopWhenUnconfigured(t *testing.T) {
st := store.NewMemoryStore()
a := New(Name("runner-noop"), Provider("oteltest"), WithStore(st), WithTool("probe", "probe", nil, func(context.Context, map[string]any) (string, error) { return "ok", nil }))
@@ -164,10 +196,38 @@ func TestListRunSummaries(t *testing.T) {
if len(got) != 2 {
t.Fatalf("got %d summaries, want 2: %#v", len(got), got)
}
if got[0].RunID != "run-a" || got[0].TraceID != "trace-a" || got[0].SpanID != "span-a" || got[0].Events != 2 || got[0].LastKind != "tool" || !got[0].UpdatedAt.Equal(time.Unix(0, 2)) {
if got[0].RunID != "run-a" || got[0].TraceID != "trace-a" || got[0].SpanID != "span-a" || got[0].Events != 2 || got[0].Status != "running" || got[0].DurationMS != 0 || got[0].LastKind != "tool" || !got[0].UpdatedAt.Equal(time.Unix(0, 2)) {
t.Fatalf("unexpected run-a summary: %#v", got[0])
}
if got[1].RunID != "run-b" || got[1].ParentID != "parent" || got[1].Events != 2 || got[1].LastKind != "error" || got[1].LastError != "boom" {
if got[1].RunID != "run-b" || got[1].ParentID != "parent" || got[1].Events != 2 || got[1].Status != "error" || got[1].DurationMS != 0 || got[1].LastKind != "error" || got[1].LastError != "boom" {
t.Fatalf("unexpected run-b summary: %#v", got[1])
}
}
func TestListRunSummariesWithOptionsFiltersAndLimits(t *testing.T) {
st := store.NewMemoryStore()
scoped := store.Scope(st, "agent", "runner")
events := []RunEvent{
{Time: time.Unix(0, 1), RunID: "run-old", Agent: "runner", Kind: "run"},
{Time: time.Unix(0, 2), RunID: "run-old", Agent: "runner", Kind: "done"},
{Time: time.Unix(0, 3), RunID: "run-new", Agent: "runner", Kind: "run"},
{Time: time.Unix(0, 4), RunID: "run-new", Agent: "runner", Kind: "error", Error: "boom"},
}
for _, e := range events {
b, err := json.Marshal(e)
if err != nil {
t.Fatal(err)
}
if err := scoped.Write(&store.Record{Key: "runs/" + e.RunID + "/" + e.Time.Format("20060102150405.000000000") + "-" + e.Kind, Value: b}); err != nil {
t.Fatal(err)
}
}
got, err := ListRunSummariesWithOptions(st, "runner", RunListOptions{Status: "error", Limit: 1})
if err != nil {
t.Fatal(err)
}
if len(got) != 1 || got[0].RunID != "run-new" || got[0].Status != "error" {
t.Fatalf("filtered summaries = %#v", got)
}
}
+1 -1
View File
@@ -161,7 +161,7 @@ func (p *Provider) Generate(ctx context.Context, req *ai.Request, opts ...ai.Gen
// Stream generates a streaming response (not yet implemented)
func (p *Provider) Stream(ctx context.Context, req *ai.Request, opts ...ai.GenerateOption) (ai.Stream, error) {
return nil, fmt.Errorf("streaming not yet implemented for anthropic provider")
return nil, fmt.Errorf("%w: anthropic provider", ai.ErrStreamingUnsupported)
}
// callAPI makes an HTTP request to the Anthropic API
+3 -2
View File
@@ -2,6 +2,7 @@ package anthropic
import (
"context"
"errors"
"testing"
"go-micro.dev/v6/ai"
@@ -88,7 +89,7 @@ func TestProvider_Stream_NotImplemented(t *testing.T) {
}
_, err := p.Stream(context.Background(), req)
if err == nil {
t.Error("Expected error for unimplemented streaming, got nil")
if !errors.Is(err, ai.ErrStreamingUnsupported) {
t.Fatalf("Stream error = %v, want ErrStreamingUnsupported", err)
}
}
+1 -1
View File
@@ -143,7 +143,7 @@ func (p *Provider) Generate(ctx context.Context, req *ai.Request, opts ...ai.Gen
}
func (p *Provider) Stream(ctx context.Context, req *ai.Request, opts ...ai.GenerateOption) (ai.Stream, error) {
return nil, fmt.Errorf("streaming not yet implemented for atlascloud provider")
return nil, fmt.Errorf("%w: atlascloud provider", ai.ErrStreamingUnsupported)
}
func (p *Provider) callAPI(ctx context.Context, req map[string]any) (*ai.Response, map[string]any, error) {
+3 -2
View File
@@ -2,6 +2,7 @@ package atlascloud
import (
"context"
"errors"
"testing"
"go-micro.dev/v6/ai"
@@ -88,8 +89,8 @@ func TestProvider_Stream_NotImplemented(t *testing.T) {
}
_, err := p.Stream(context.Background(), req)
if err == nil {
t.Error("Expected error for unimplemented streaming, got nil")
if !errors.Is(err, ai.ErrStreamingUnsupported) {
t.Fatalf("Stream error = %v, want ErrStreamingUnsupported", err)
}
}
+116
View File
@@ -0,0 +1,116 @@
package ai
import "sort"
// CapabilityRow is one deterministic row in a provider capability matrix.
type CapabilityRow struct {
// Provider is the registered provider name.
Provider string
Capabilities
}
// Capabilities describes the AI interfaces a provider has registered.
// It is intentionally based on package registration rather than external
// provider marketing claims, so it reflects what this build can actually use.
type Capabilities struct {
// Model reports whether ai.New can construct a chat/text model provider.
Model bool
// Image reports whether ai.NewImage can construct an image model provider.
Image bool
// Video reports whether ai.NewVideo can construct a video model provider.
Video bool
}
// ProviderCapabilities reports the capabilities registered for provider.
func ProviderCapabilities(provider string) Capabilities {
_, hasModel := providers[provider]
_, hasImage := imageProviders[provider]
_, hasVideo := videoProviders[provider]
return Capabilities{
Model: hasModel,
Image: hasImage,
Video: hasVideo,
}
}
// CapabilityMatrix returns a snapshot of all registered AI providers and the
// interfaces they support. The returned map is a copy and can be modified by
// callers without mutating the registry. Use CapabilityRows when rendering a
// deterministic table or report.
func CapabilityMatrix() map[string]Capabilities {
names := map[string]struct{}{}
for name := range providers {
names[name] = struct{}{}
}
for name := range imageProviders {
names[name] = struct{}{}
}
for name := range videoProviders {
names[name] = struct{}{}
}
matrix := make(map[string]Capabilities, len(names))
for name := range names {
matrix[name] = ProviderCapabilities(name)
}
return matrix
}
// CapabilityRows returns a deterministic capability support matrix for every
// registered AI provider. It is the ordered form of CapabilityMatrix, intended
// for CLIs, docs generators, and conformance reports that need stable output.
func CapabilityRows() []CapabilityRow {
names := RegisteredProviders("")
rows := make([]CapabilityRow, 0, len(names))
for _, name := range names {
rows = append(rows, CapabilityRow{
Provider: name,
Capabilities: ProviderCapabilities(name),
})
}
return rows
}
// RegisteredProviders returns the registered provider names in sorted order.
// kind may be "model", "image", "video", or empty for the union of all
// provider registries.
func RegisteredProviders(kind string) []string {
names := map[string]struct{}{}
add := func(registry any) {
switch r := registry.(type) {
case map[string]NewFunc:
for name := range r {
names[name] = struct{}{}
}
case map[string]NewImageFunc:
for name := range r {
names[name] = struct{}{}
}
case map[string]NewVideoFunc:
for name := range r {
names[name] = struct{}{}
}
}
}
switch kind {
case "model":
add(providers)
case "image":
add(imageProviders)
case "video":
add(videoProviders)
default:
add(providers)
add(imageProviders)
add(videoProviders)
}
out := make([]string, 0, len(names))
for name := range names {
out = append(out, name)
}
sort.Strings(out)
return out
}
+75
View File
@@ -0,0 +1,75 @@
package ai_test
import (
"reflect"
"testing"
"go-micro.dev/v6/ai"
_ "go-micro.dev/v6/ai/anthropic"
_ "go-micro.dev/v6/ai/atlascloud"
_ "go-micro.dev/v6/ai/gemini"
_ "go-micro.dev/v6/ai/groq"
_ "go-micro.dev/v6/ai/mistral"
_ "go-micro.dev/v6/ai/openai"
_ "go-micro.dev/v6/ai/together"
)
func TestRegisteredProviders(t *testing.T) {
got := ai.RegisteredProviders("")
want := []string{"anthropic", "atlascloud", "gemini", "groq", "mistral", "openai", "together"}
if !reflect.DeepEqual(got, want) {
t.Fatalf("RegisteredProviders() = %#v, want %#v", got, want)
}
got = ai.RegisteredProviders("image")
want = []string{"atlascloud", "openai"}
if !reflect.DeepEqual(got, want) {
t.Fatalf("RegisteredProviders(image) = %#v, want %#v", got, want)
}
got = ai.RegisteredProviders("video")
want = []string{"atlascloud"}
if !reflect.DeepEqual(got, want) {
t.Fatalf("RegisteredProviders(video) = %#v, want %#v", got, want)
}
}
func TestCapabilityRows(t *testing.T) {
got := ai.CapabilityRows()
want := []ai.CapabilityRow{
{Provider: "anthropic", Capabilities: ai.Capabilities{Model: true}},
{Provider: "atlascloud", Capabilities: ai.Capabilities{Model: true, Image: true, Video: true}},
{Provider: "gemini", Capabilities: ai.Capabilities{Model: true}},
{Provider: "groq", Capabilities: ai.Capabilities{Model: true}},
{Provider: "mistral", Capabilities: ai.Capabilities{Model: true}},
{Provider: "openai", Capabilities: ai.Capabilities{Model: true, Image: true}},
{Provider: "together", Capabilities: ai.Capabilities{Model: true}},
}
if !reflect.DeepEqual(got, want) {
t.Fatalf("CapabilityRows() = %#v, want %#v", got, want)
}
}
func TestCapabilityMatrix(t *testing.T) {
matrix := ai.CapabilityMatrix()
for _, provider := range []string{"anthropic", "atlascloud", "gemini", "groq", "mistral", "openai", "together"} {
caps, ok := matrix[provider]
if !ok {
t.Fatalf("CapabilityMatrix missing %q", provider)
}
if !caps.Model {
t.Fatalf("CapabilityMatrix(%s).Model = false, want true", provider)
}
}
if caps := ai.ProviderCapabilities("openai"); caps != (ai.Capabilities{Model: true, Image: true}) {
t.Fatalf("ProviderCapabilities(openai) = %#v", caps)
}
if caps := ai.ProviderCapabilities("atlascloud"); caps != (ai.Capabilities{Model: true, Image: true, Video: true}) {
t.Fatalf("ProviderCapabilities(atlascloud) = %#v", caps)
}
if caps := ai.ProviderCapabilities("missing"); caps != (ai.Capabilities{}) {
t.Fatalf("ProviderCapabilities(missing) = %#v", caps)
}
}
+1 -1
View File
@@ -135,7 +135,7 @@ func (p *Provider) Generate(ctx context.Context, req *ai.Request, opts ...ai.Gen
}
func (p *Provider) Stream(ctx context.Context, req *ai.Request, opts ...ai.GenerateOption) (ai.Stream, error) {
return nil, fmt.Errorf("streaming not yet implemented for gemini provider")
return nil, fmt.Errorf("%w: gemini provider", ai.ErrStreamingUnsupported)
}
func (p *Provider) callAPI(ctx context.Context, req map[string]any) (*ai.Response, []map[string]any, error) {
+3 -2
View File
@@ -2,6 +2,7 @@ package gemini
import (
"context"
"errors"
"testing"
"go-micro.dev/v6/ai"
@@ -88,8 +89,8 @@ func TestProvider_Stream_NotImplemented(t *testing.T) {
}
_, err := p.Stream(context.Background(), req)
if err == nil {
t.Error("Expected error for unimplemented streaming, got nil")
if !errors.Is(err, ai.ErrStreamingUnsupported) {
t.Fatalf("Stream error = %v, want ErrStreamingUnsupported", err)
}
}
+1 -1
View File
@@ -119,7 +119,7 @@ func (p *Provider) Generate(ctx context.Context, req *ai.Request, opts ...ai.Gen
}
func (p *Provider) Stream(ctx context.Context, req *ai.Request, opts ...ai.GenerateOption) (ai.Stream, error) {
return nil, fmt.Errorf("streaming not yet implemented for groq provider")
return nil, fmt.Errorf("%w: groq provider", ai.ErrStreamingUnsupported)
}
func (p *Provider) callAPI(ctx context.Context, req map[string]any) (*ai.Response, map[string]any, error) {
+3 -2
View File
@@ -2,6 +2,7 @@ package groq
import (
"context"
"errors"
"testing"
"go-micro.dev/v6/ai"
@@ -40,8 +41,8 @@ func TestProvider_Generate_NoAPIKey(t *testing.T) {
}
func TestProvider_Stream_NotImplemented(t *testing.T) {
if _, err := NewProvider().Stream(context.Background(), &ai.Request{Prompt: "hi"}); err == nil {
t.Error("expected error")
if _, err := NewProvider().Stream(context.Background(), &ai.Request{Prompt: "hi"}); !errors.Is(err, ai.ErrStreamingUnsupported) {
t.Fatalf("Stream error = %v, want ErrStreamingUnsupported", err)
}
}
+1 -1
View File
@@ -119,7 +119,7 @@ func (p *Provider) Generate(ctx context.Context, req *ai.Request, opts ...ai.Gen
}
func (p *Provider) Stream(ctx context.Context, req *ai.Request, opts ...ai.GenerateOption) (ai.Stream, error) {
return nil, fmt.Errorf("streaming not yet implemented for mistral provider")
return nil, fmt.Errorf("%w: mistral provider", ai.ErrStreamingUnsupported)
}
func (p *Provider) callAPI(ctx context.Context, req map[string]any) (*ai.Response, map[string]any, error) {
+3 -2
View File
@@ -2,6 +2,7 @@ package mistral
import (
"context"
"errors"
"testing"
"go-micro.dev/v6/ai"
@@ -40,8 +41,8 @@ func TestProvider_Generate_NoAPIKey(t *testing.T) {
}
func TestProvider_Stream_NotImplemented(t *testing.T) {
if _, err := NewProvider().Stream(context.Background(), &ai.Request{Prompt: "hi"}); err == nil {
t.Error("expected error")
if _, err := NewProvider().Stream(context.Background(), &ai.Request{Prompt: "hi"}); !errors.Is(err, ai.ErrStreamingUnsupported) {
t.Fatalf("Stream error = %v, want ErrStreamingUnsupported", err)
}
}
+7 -1
View File
@@ -4,6 +4,7 @@ package ai
import (
"context"
"encoding/json"
"errors"
"strings"
)
@@ -133,7 +134,12 @@ func RunInfoFrom(ctx context.Context) (RunInfo, bool) {
return r, ok
}
// Stream is the interface for streaming responses (future implementation)
// ErrStreamingUnsupported is returned by providers that implement the Model
// interface but do not yet support token streaming. Use errors.Is so callers
// can distinguish an unsupported capability from transient provider failures.
var ErrStreamingUnsupported = errors.New("ai: streaming unsupported")
// Stream is the interface for streaming responses.
type Stream interface {
// Recv receives the next chunk of the response
Recv() (*Response, error)
+1 -1
View File
@@ -142,7 +142,7 @@ func (p *Provider) Generate(ctx context.Context, req *ai.Request, opts ...ai.Gen
// Stream generates a streaming response (not yet implemented)
func (p *Provider) Stream(ctx context.Context, req *ai.Request, opts ...ai.GenerateOption) (ai.Stream, error) {
return nil, fmt.Errorf("streaming not yet implemented for openai provider")
return nil, fmt.Errorf("%w: openai provider", ai.ErrStreamingUnsupported)
}
// callAPI makes an HTTP request to the OpenAI API
+3 -2
View File
@@ -2,6 +2,7 @@ package openai
import (
"context"
"errors"
"testing"
"go-micro.dev/v6/ai"
@@ -88,8 +89,8 @@ func TestProvider_Stream_NotImplemented(t *testing.T) {
}
_, err := p.Stream(context.Background(), req)
if err == nil {
t.Error("Expected error for unimplemented streaming, got nil")
if !errors.Is(err, ai.ErrStreamingUnsupported) {
t.Fatalf("Stream error = %v, want ErrStreamingUnsupported", err)
}
}
+1 -1
View File
@@ -119,7 +119,7 @@ func (p *Provider) Generate(ctx context.Context, req *ai.Request, opts ...ai.Gen
}
func (p *Provider) Stream(ctx context.Context, req *ai.Request, opts ...ai.GenerateOption) (ai.Stream, error) {
return nil, fmt.Errorf("streaming not yet implemented for together provider")
return nil, fmt.Errorf("%w: together provider", ai.ErrStreamingUnsupported)
}
func (p *Provider) callAPI(ctx context.Context, req map[string]any) (*ai.Response, map[string]any, error) {
+3 -2
View File
@@ -2,6 +2,7 @@ package together
import (
"context"
"errors"
"testing"
"go-micro.dev/v6/ai"
@@ -40,8 +41,8 @@ func TestProvider_Generate_NoAPIKey(t *testing.T) {
}
func TestProvider_Stream_NotImplemented(t *testing.T) {
if _, err := NewProvider().Stream(context.Background(), &ai.Request{Prompt: "hi"}); err == nil {
t.Error("expected error")
if _, err := NewProvider().Stream(context.Background(), &ai.Request{Prompt: "hi"}); !errors.Is(err, ai.ErrStreamingUnsupported) {
t.Fatalf("Stream error = %v, want ErrStreamingUnsupported", err)
}
}
+36 -12
View File
@@ -19,9 +19,7 @@ func init() {
Name: "runs",
Usage: "Show recorded agent runs",
ArgsUsage: "[agent] [run-id]",
Flags: []cli.Flag{
&cli.BoolFlag{Name: "json", Usage: "Print run data as JSON for automation"},
},
Flags: runFlags(),
Action: func(c *cli.Context) error {
name := c.Args().First()
if name == "" {
@@ -30,7 +28,7 @@ func init() {
if runID := c.Args().Get(1); runID != "" {
return printRunHistory(name, runID, c.Bool("json"))
}
return printRunIndex(name, c.Bool("json"))
return printRunIndex(name, runOptions(c), c.Bool("json"))
},
})
@@ -102,9 +100,7 @@ func init() {
Name: "history",
Usage: "Show an agent's stored conversation and run history",
ArgsUsage: "[name] [run-id]",
Flags: []cli.Flag{
&cli.BoolFlag{Name: "json", Usage: "Print run data as JSON for automation"},
},
Flags: runFlags(),
Action: func(c *cli.Context) error {
name := c.Args().First()
if name == "" {
@@ -114,7 +110,7 @@ func init() {
return printRunHistory(name, runID, c.Bool("json"))
}
if c.Bool("json") {
return printRunIndex(name, true)
return printRunIndex(name, runOptions(c), true)
}
// Read from the agent's scoped state store (database
// "agent", table = name) — available whether or not the
@@ -128,15 +124,27 @@ func init() {
fmt.Printf(" \033[2m%s:\033[0m %v\n", m.Role, m.Content)
}
}
return printRunIndex(name, c.Bool("json"))
return printRunIndex(name, runOptions(c), c.Bool("json"))
},
},
},
})
}
func printRunIndex(name string, asJSON bool) error {
runs, err := goagent.ListRunSummaries(store.DefaultStore, name)
func runFlags() []cli.Flag {
return []cli.Flag{
&cli.BoolFlag{Name: "json", Usage: "Print run data as JSON for automation"},
&cli.StringFlag{Name: "status", Usage: "Only show runs with this status (running, done, error, refused)"},
&cli.IntFlag{Name: "limit", Usage: "Show the most recently updated N runs"},
}
}
func runOptions(c *cli.Context) goagent.RunListOptions {
return goagent.RunListOptions{Status: c.String("status"), Limit: c.Int("limit")}
}
func printRunIndex(name string, opts goagent.RunListOptions, asJSON bool) error {
runs, err := goagent.ListRunSummariesWithOptions(store.DefaultStore, name, opts)
if err != nil {
return err
}
@@ -155,7 +163,10 @@ func writeRunIndex(w io.Writer, name string, runs []goagent.RunSummary, asJSON b
}
fmt.Fprintln(w, " Runs:")
for _, run := range runs {
line := fmt.Sprintf(" %s events=%d last=%s updated=%s", run.RunID, run.Events, run.LastKind, run.UpdatedAt.Format("2006-01-02 15:04:05"))
line := fmt.Sprintf(" %s status=%s events=%d duration=%s last=%s updated=%s", run.RunID, run.Status, run.Events, formatDurationMS(run.DurationMS), run.LastKind, run.UpdatedAt.Format("2006-01-02 15:04:05"))
if run.ParentID != "" {
line += " parent=" + run.ParentID
}
if run.TraceID != "" {
line += " trace=" + shortTraceID(run.TraceID)
}
@@ -199,6 +210,9 @@ func writeRunHistory(w io.Writer, name, runID string, events []goagent.RunEvent,
if e.Tokens.TotalTokens > 0 {
line += fmt.Sprintf(" tokens=%d", e.Tokens.TotalTokens)
}
if e.ParentID != "" {
line += " parent=" + e.ParentID
}
if e.TraceID != "" {
line += " trace=" + shortTraceID(e.TraceID)
}
@@ -213,6 +227,16 @@ func writeRunHistory(w io.Writer, name, runID string, events []goagent.RunEvent,
return nil
}
func formatDurationMS(ms int64) string {
if ms <= 0 {
return "0ms"
}
if ms < 1000 {
return fmt.Sprintf("%dms", ms)
}
return fmt.Sprintf("%.1fs", float64(ms)/1000)
}
func shortTraceID(id string) string {
if len(id) <= 12 {
return id
+34 -8
View File
@@ -13,13 +13,15 @@ import (
func TestWriteRunIndexJSON(t *testing.T) {
runs := []goagent.RunSummary{{
RunID: "run-1",
Agent: "runner",
StartedAt: time.Unix(0, 1),
UpdatedAt: time.Unix(0, 2),
Events: 2,
LastKind: "tool",
TraceID: "1234567890abcdef",
RunID: "run-1",
Agent: "runner",
StartedAt: time.Unix(0, 1),
UpdatedAt: time.Unix(0, 2),
DurationMS: 1234,
Events: 2,
Status: "done",
LastKind: "tool",
TraceID: "1234567890abcdef",
}}
var out bytes.Buffer
if err := writeRunIndex(&out, "runner", runs, true); err != nil {
@@ -34,6 +36,29 @@ func TestWriteRunIndexJSON(t *testing.T) {
}
}
func TestWriteRunIndexHumanIncludesStatusAndDuration(t *testing.T) {
runs := []goagent.RunSummary{{
RunID: "run-1",
Agent: "runner",
UpdatedAt: time.Date(2026, 6, 25, 12, 34, 56, 0, time.UTC),
DurationMS: 1234,
Events: 2,
Status: "done",
LastKind: "tool",
ParentID: "parent-run",
}}
var out bytes.Buffer
if err := writeRunIndex(&out, "runner", runs, false); err != nil {
t.Fatal(err)
}
line := out.String()
for _, want := range []string{"run-1", "status=done", "events=2", "duration=1.2s", "last=tool", "parent=parent-run"} {
if !strings.Contains(line, want) {
t.Fatalf("human output %q missing %q", line, want)
}
}
}
func TestWriteRunHistoryHumanAndJSON(t *testing.T) {
events := []goagent.RunEvent{{
Time: time.Date(2026, 6, 25, 12, 34, 56, 7_000_000, time.UTC),
@@ -46,6 +71,7 @@ func TestWriteRunHistoryHumanAndJSON(t *testing.T) {
LatencyMS: 42,
Tokens: ai.Usage{TotalTokens: 5},
TraceID: "1234567890abcdef",
ParentID: "parent-run",
}}
var human bytes.Buffer
@@ -53,7 +79,7 @@ func TestWriteRunHistoryHumanAndJSON(t *testing.T) {
t.Fatal(err)
}
line := human.String()
for _, want := range []string{"12:34:56.007 tool", "probe", "oteltest/unit-model", "42ms", "tokens=5", "trace=1234567890ab"} {
for _, want := range []string{"12:34:56.007 tool", "probe", "oteltest/unit-model", "42ms", "tokens=5", "parent=parent-run", "trace=1234567890ab"} {
if !strings.Contains(line, want) {
t.Fatalf("human output %q missing %q", line, want)
}
+19 -3
View File
@@ -111,16 +111,25 @@ func StoreCheckpoint(s store.Store, scope string) Checkpoint {
return &storeCheckpoint{store: store.Scope(s, "flow", scope)}
}
func (c *storeCheckpoint) Save(_ context.Context, run Run) error {
func (c *storeCheckpoint) Save(ctx context.Context, run Run) error {
if err := ctx.Err(); err != nil {
return err
}
run.Updated = time.Now()
b, err := json.Marshal(run)
if err != nil {
return err
}
if err := ctx.Err(); err != nil {
return err
}
return c.store.Write(&store.Record{Key: run.ID, Value: b})
}
func (c *storeCheckpoint) Load(_ context.Context, runID string) (Run, bool, error) {
func (c *storeCheckpoint) Load(ctx context.Context, runID string) (Run, bool, error) {
if err := ctx.Err(); err != nil {
return Run{}, false, err
}
recs, err := c.store.Read(runID)
if err == store.ErrNotFound || len(recs) == 0 {
return Run{}, false, nil
@@ -135,11 +144,17 @@ func (c *storeCheckpoint) Load(_ context.Context, runID string) (Run, bool, erro
return run, true, nil
}
func (c *storeCheckpoint) Delete(_ context.Context, runID string) error {
func (c *storeCheckpoint) Delete(ctx context.Context, runID string) error {
if err := ctx.Err(); err != nil {
return err
}
return c.store.Delete(runID)
}
func (c *storeCheckpoint) List(ctx context.Context) ([]Run, error) {
if err := ctx.Err(); err != nil {
return nil, err
}
keys, err := c.store.List()
if err != nil {
return nil, err
@@ -368,6 +383,7 @@ func (f *Flow) Pending(ctx context.Context) ([]Run, error) {
func (f *Flow) runFrom(ctx context.Context, run Run) (Run, error) {
steps := f.opts.Steps
ctx = withDeps(ctx, &runDeps{client: f.client, model: f.model, tools: f.toolSet})
ctx = ai.WithRunInfo(ctx, ai.RunInfo{RunID: run.ID, Agent: f.name})
start := stepIndex(steps, run.State.Stage)
if start < 0 {
+49
View File
@@ -6,6 +6,7 @@ import (
"testing"
"time"
"go-micro.dev/v6/ai"
"go-micro.dev/v6/store"
)
@@ -86,6 +87,33 @@ func TestFlowCheckpointResume(t *testing.T) {
}
}
func TestFlowStepContextIncludesRunInfo(t *testing.T) {
var got ai.RunInfo
step := Step{Name: "inspect", Run: func(ctx context.Context, in State) (State, error) {
var ok bool
got, ok = ai.RunInfoFrom(ctx)
if !ok {
t.Fatal("RunInfo missing from step context")
}
in.Data = []byte("ok")
return in, nil
}}
f := New("correlated",
WithCheckpoint(StoreCheckpoint(store.NewMemoryStore(), "correlated")),
Steps(step),
)
if err := f.Execute(context.Background(), "start"); err != nil {
t.Fatalf("Execute: %v", err)
}
if got.Agent != "correlated" {
t.Fatalf("RunInfo.Agent = %q, want correlated", got.Agent)
}
if got.RunID == "" {
t.Fatal("RunInfo.RunID is empty")
}
}
func TestFlowResumePendingResumesOldestRunsUntilFailure(t *testing.T) {
mem := store.NewMemoryStore()
ctx := context.Background()
@@ -333,6 +361,27 @@ func TestStoreCheckpointListReturnsRunsInStartedOrder(t *testing.T) {
}
}
func TestStoreCheckpointHonorsCanceledContext(t *testing.T) {
cp := StoreCheckpoint(store.NewMemoryStore(), "canceled")
ctx, cancel := context.WithCancel(context.Background())
cancel()
run := Run{ID: "canceled", Started: time.Now()}
if err := cp.Save(ctx, run); !errors.Is(err, context.Canceled) {
t.Fatalf("Save error = %v, want context.Canceled", err)
}
if _, ok, err := cp.Load(ctx, run.ID); !errors.Is(err, context.Canceled) || ok {
t.Fatalf("Load ok, error = %v, %v; want false, context.Canceled", ok, err)
}
if err := cp.Delete(ctx, run.ID); !errors.Is(err, context.Canceled) {
t.Fatalf("Delete error = %v, want context.Canceled", err)
}
if _, err := cp.List(ctx); !errors.Is(err, context.Canceled) {
t.Fatalf("List error = %v, want context.Canceled", err)
}
}
func TestStateSetScan(t *testing.T) {
var s State
type payload struct {
+26 -7
View File
@@ -415,7 +415,7 @@ func (d *dispatcher) serve(w http.ResponseWriter, r *http.Request, invoke Invoke
switch req.Method {
case "message/send":
d.send(w, req, invoke)
d.send(requestContext(r.Context()), w, req, invoke)
case "tasks/get":
d.get(w, req)
case "tasks/cancel":
@@ -432,7 +432,7 @@ type sendParams struct {
Message Message `json:"message"`
}
func (d *dispatcher) send(w http.ResponseWriter, req rpcRequest, invoke Invoke) {
func (d *dispatcher) send(ctx context.Context, w http.ResponseWriter, req rpcRequest, invoke Invoke) {
var p sendParams
if err := json.Unmarshal(req.Params, &p); err != nil {
writeRPC(w, req.ID, nil, &rpcError{Code: errInvalidParams, Message: "invalid params"})
@@ -444,7 +444,7 @@ func (d *dispatcher) send(w http.ResponseWriter, req rpcRequest, invoke Invoke)
return
}
reply, err := invoke(r2ctx(), text)
reply, err := invoke(ctx, text)
contextID := p.Message.ContextID
if contextID == "" {
contextID = uuid.New().String()
@@ -542,6 +542,29 @@ func textArtifact(text string) Artifact {
}
}
// requestContext carries request cancellation and deadlines into the downstream
// agent call without leaking HTTP transport context values into the go-micro
// client stack.
func requestContext(parent context.Context) context.Context {
if err := parent.Err(); err != nil {
ctx, cancel := context.WithCancel(context.Background())
cancel()
return ctx
}
ctx := context.Background()
var cancel context.CancelFunc
if deadline, ok := parent.Deadline(); ok {
ctx, cancel = context.WithDeadline(ctx, deadline)
} else {
ctx, cancel = context.WithCancel(ctx)
}
go func() {
<-parent.Done()
cancel()
}()
return ctx
}
func writeJSON(w http.ResponseWriter, status int, v any) {
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(status)
@@ -554,7 +577,3 @@ func writeRPC(w http.ResponseWriter, id json.RawMessage, result any, e *rpcError
}
writeJSON(w, http.StatusOK, rpcResponse{JSONRPC: "2.0", ID: id, Result: result, Error: e})
}
// r2ctx returns a background context for the agent call. (Kept as a seam
// so request-scoped context/deadlines can be threaded later.)
func r2ctx() context.Context { return context.Background() }
+36
View File
@@ -109,6 +109,42 @@ func TestMessageSendAndGet(t *testing.T) {
}
}
func TestMessageSendUsesRequestContext(t *testing.T) {
d := newDispatcher()
ctx, cancel := context.WithCancel(context.Background())
cancel()
req := httptest.NewRequest(http.MethodPost, "/", bytes.NewBufferString(`{
"jsonrpc":"2.0","id":1,"method":"message/send",
"params":{"message":{"role":"user","kind":"message","messageId":"m1",
"parts":[{"kind":"text","text":"ping"}]}}}`))
req = req.WithContext(ctx)
rr := httptest.NewRecorder()
d.serve(rr, req, func(ctx context.Context, text string) (string, error) {
if err := ctx.Err(); err != nil {
return "", err
}
return "unexpected success", nil
})
var resp struct {
Result Task `json:"result"`
Error *rpcError `json:"error"`
}
if err := json.NewDecoder(rr.Result().Body).Decode(&resp); err != nil {
t.Fatalf("decode: %v", err)
}
if resp.Error != nil {
t.Fatalf("rpc error: %+v", resp.Error)
}
if resp.Result.Status.State != stateFailed {
t.Fatalf("task state = %q, want failed", resp.Result.Status.State)
}
if len(resp.Result.Artifacts) != 1 || textOf(resp.Result.Artifacts[0].Parts) != "error: context canceled" {
t.Fatalf("artifact = %+v, want context cancellation", resp.Result.Artifacts)
}
}
func TestUnknownMethod(t *testing.T) {
ts, cleanup := newGatewayWithAgent(t)
defer cleanup()
+73 -4
View File
@@ -69,11 +69,80 @@ Grounded in real signal, never speculative rewrites. Each cycle draws from:
recurring jobs expire after 7 days, so it is **not** a durable scheduler.
- **GitHub Actions (durable)** — a scheduled workflow that runs the loop
independently of any session. This is the real backbone; it opens a fresh
tracking issue for each increment and dispatches Codex there, so every run gets
a unique Codex-derived branch and a PR that closes its tracking issue. It needs
a `CODEX_TRIGGER_TOKEN` repo secret from a user account Codex responds to;
tracking issue for each increment and dispatches Codex there. It needs a
`CODEX_TRIGGER_TOKEN` repo secret from a user account Codex responds to;
without that secret the workflow deliberately no-ops to avoid ignored bot
comments. See `.github/workflows/continuous-improvement.yml`.
comments. See `.github/workflows/continuous-improvement.yml` and the mechanics
below.
## How the durable loop works (mechanics)
Hard-won wiring — change any one piece and the loop silently stops producing
merged PRs. Each scheduled run:
1. **Opens a fresh issue per increment** (`Continuous improvement increment #N`)
and posts the `@codex` instruction on it. *Why a fresh issue:* Codex derives
its branch name from the triggering issue's context, so re-using one tracker
issue collapses every run onto one branch name and only the first PR opens —
the rest collide and silently fail.
2. **Posts as a user, not the Actions bot.** Codex ignores `@codex` comments
authored by `github-actions[bot]`, so the dispatch uses `CODEX_TRIGGER_TOKEN`
(a PAT for a user account Codex follows). No token → the step no-ops.
3. **Codex opens the PR itself with `gh` — never `make_pr`.** In the Codex Cloud
sandbox the `make_pr` tool is a **no-op stub**: it records the PR title/body
for the manual "Create PR" button and never pushes a branch or calls the API.
So the dispatch and [`AGENTS.md`](../../AGENTS.md) tell Codex to do it by hand:
```sh
git switch -c codex/increment-<issue> # unique branch, codex/ prefix
git push -u origin codex/increment-<issue>
gh pr create --base master --label codex --title "…" --body "… Closes #<issue>"
gh pr merge --squash --auto --delete-branch
```
This requires the Codex setup script to install `gh` and run `gh auth
setup-git` (so `git push` is authenticated) with a write-scoped token.
4. **Merges via GitHub native auto-merge, gated by branch protection.** `master`
requires the CI status checks (build, tests, golangci-lint) and **0 approving
reviews**. `gh pr merge --auto` enables auto-merge; GitHub lands the PR the
moment checks pass and deletes the branch. `Closes #<issue>` auto-closes the
tracking issue. There is **no merge sweep workflow** — branch protection is
the gate.
### Do-not-break list
- **Don't re-add required approvals** to `master` — it blocks every autonomous
merge. The intended gate is **green CI only**.
- **Don't point the dispatch at one standing tracker issue** — one issue per run.
- **Don't tell Codex to use `make_pr`** (or imply a token "isn't a substitute"):
it cannot open a PR. `gh` is the only path.
- **Don't manually re-implement a Codex increment during the summary→PR lag**
(Codex posts an optimistic "opened a PR" comment ~3045 min before the PR
actually appears). Re-doing it creates duplicate PRs and stale branches that
then block the next run. Wait for the PR, or let it ride.
## Overseer passes (DevRel + Architect)
The hourly loop ships increments; two periodic passes keep the *whole* heading in
the right direction. Both use the same mechanism (fresh issue → `@codex` →
output) but produce direction and coherence, not just code.
- **DevRel — daily** (`.github/workflows/devrel-review.yml`). Audits the public
surface (README, website landing + docs, blog) for coherence with the North
Star, README crispness, and blog-worthy material. **Autonomy boundary:** safe
factual-alignment and crispness fixes auto-merge like any increment;
brand/positioning copy and blog drafts are *surfaced in a report* for the
human, never auto-merged.
- **Architect — every few days** (`.github/workflows/architecture-review.yml`).
Reviews the framework/harness against the thesis: API coherence, lifecycle
gaps, drift/sprawl. **Its output is an assessment plus scoped follow-up
issues** that feed the hourly increment loop — it does **not** make breaking or
architectural changes itself (those stay with the human).
Together they close the loop: the architect decides *what* should change and files
issues, the increment loop *builds* them, and DevRel keeps the public story
honest. Cadence is tunable in each workflow's `cron`. Codex is serial, so these
passes queue behind any in-flight increment rather than running concurrently.
## Stop / redirect
+3 -2
View File
@@ -207,18 +207,19 @@ func main() {
mem := store.NewMemoryStore()
wsSvc := new(WorkspaceService)
ws := service.New(service.Name("workspace"), service.Registry(reg), service.Client(cl))
ws := service.New(service.Name("workspace"), service.Address("127.0.0.1:0"), service.Registry(reg), service.Client(cl))
ws.Handle(wsSvc)
go ws.Run()
ntSvc := new(NotifyService)
nt := service.New(service.Name("notify"), service.Registry(reg), service.Client(cl))
nt := service.New(service.Name("notify"), service.Address("127.0.0.1:0"), service.Registry(reg), service.Client(cl))
nt.Handle(ntSvc)
go nt.Run()
// The onboarder agent, registered so the flow can reach it over RPC.
onboarder := agent.New(
agent.Name("onboarder"),
agent.Address("127.0.0.1:0"),
agent.Services("workspace", "notify"),
agent.Prompt("You onboard new users. Create their workspace and send a welcome message."),
agent.Provider(*provider),
+3 -2
View File
@@ -36,14 +36,14 @@ func TestEventTriggersAgentNoPrompt(t *testing.T) {
mem := store.NewMemoryStore()
wsSvc := new(WorkspaceService)
ws := service.New(service.Name("workspace"), service.Registry(reg), service.Client(cl))
ws := service.New(service.Name("workspace"), service.Address("127.0.0.1:0"), service.Registry(reg), service.Client(cl))
if err := ws.Handle(wsSvc); err != nil {
t.Fatalf("handle workspace: %v", err)
}
go ws.Run()
ntSvc := new(NotifyService)
nt := service.New(service.Name("notify"), service.Registry(reg), service.Client(cl))
nt := service.New(service.Name("notify"), service.Address("127.0.0.1:0"), service.Registry(reg), service.Client(cl))
if err := nt.Handle(ntSvc); err != nil {
t.Fatalf("handle notify: %v", err)
}
@@ -51,6 +51,7 @@ func TestEventTriggersAgentNoPrompt(t *testing.T) {
onboarder := agent.New(
agent.Name("onboarder"),
agent.Address("127.0.0.1:0"),
agent.Services("workspace", "notify"),
agent.Prompt("You onboard new users. Create their workspace and send a welcome message."),
agent.Provider("mock"),
+8 -4
View File
@@ -47,14 +47,14 @@ func TestPlanDelegateEndToEnd(t *testing.T) {
// Real services on the shared registry/client.
taskSvc := new(TaskService)
task := service.New(service.Name("task"), service.Registry(reg), service.Client(cl))
task := service.New(service.Name("task"), service.Address("127.0.0.1:0"), service.Registry(reg), service.Client(cl))
if err := task.Handle(taskSvc); err != nil {
t.Fatalf("handle task: %v", err)
}
go task.Run()
notifySvc := new(NotifyService)
notify := service.New(service.Name("notify"), service.Registry(reg), service.Client(cl))
notify := service.New(service.Name("notify"), service.Address("127.0.0.1:0"), service.Registry(reg), service.Client(cl))
if err := notify.Handle(notifySvc); err != nil {
t.Fatalf("handle notify: %v", err)
}
@@ -63,6 +63,7 @@ func TestPlanDelegateEndToEnd(t *testing.T) {
// Real comms agent (owns notify), registered so delegate reaches it over RPC.
comms := agent.New(
agent.Name("comms"),
agent.Address("127.0.0.1:0"),
agent.Services("notify"),
agent.Prompt("You handle outbound notifications."),
agent.Provider("mock"),
@@ -80,6 +81,7 @@ func TestPlanDelegateEndToEnd(t *testing.T) {
// Real conductor agent (owns task), driven programmatically.
conductor := agent.New(
agent.Name("conductor"),
agent.Address("127.0.0.1:0"),
agent.Services("task"),
agent.Prompt("Plan first, create tasks, delegate notifications to the comms agent."),
agent.Provider("mock"),
@@ -128,14 +130,14 @@ func TestFlowDispatchesToAgentEndToEnd(t *testing.T) {
mem := store.NewMemoryStore()
taskSvc := new(TaskService)
task := service.New(service.Name("task"), service.Registry(reg), service.Client(cl))
task := service.New(service.Name("task"), service.Address("127.0.0.1:0"), service.Registry(reg), service.Client(cl))
if err := task.Handle(taskSvc); err != nil {
t.Fatalf("handle task: %v", err)
}
go task.Run()
notifySvc := new(NotifyService)
notify := service.New(service.Name("notify"), service.Registry(reg), service.Client(cl))
notify := service.New(service.Name("notify"), service.Address("127.0.0.1:0"), service.Registry(reg), service.Client(cl))
if err := notify.Handle(notifySvc); err != nil {
t.Fatalf("handle notify: %v", err)
}
@@ -143,6 +145,7 @@ func TestFlowDispatchesToAgentEndToEnd(t *testing.T) {
comms := agent.New(
agent.Name("comms"),
agent.Address("127.0.0.1:0"),
agent.Services("notify"),
agent.Prompt("You handle outbound notifications."),
agent.Provider("mock"),
@@ -157,6 +160,7 @@ func TestFlowDispatchesToAgentEndToEnd(t *testing.T) {
// so the flow can reach it over RPC.
conductor := agent.New(
agent.Name("conductor"),
agent.Address("127.0.0.1:0"),
agent.Services("task"),
agent.Prompt("Plan first, create tasks, delegate notifications to the comms agent."),
agent.Provider("mock"),
+39 -2
View File
@@ -22,6 +22,15 @@ import (
"slices"
"strings"
"time"
"go-micro.dev/v6/ai"
_ "go-micro.dev/v6/ai/anthropic"
_ "go-micro.dev/v6/ai/atlascloud"
_ "go-micro.dev/v6/ai/gemini"
_ "go-micro.dev/v6/ai/groq"
_ "go-micro.dev/v6/ai/mistral"
_ "go-micro.dev/v6/ai/openai"
_ "go-micro.dev/v6/ai/together"
)
var providerEnv = map[string]string{
@@ -38,6 +47,8 @@ func main() {
providersFlag := flag.String("providers", "anthropic,openai,gemini,groq,mistral,together,atlascloud", "comma-separated providers to check; use mock for deterministic local checks")
harnessesFlag := flag.String("harnesses", "universe,agent-flow,plan-delegate", "comma-separated harness names under internal/harness")
timeoutFlag := flag.Duration("timeout", 10*time.Minute, "timeout per provider/harness run")
requireConfiguredFlag := flag.Bool("require-configured", false, "fail when a selected live provider is missing an API key")
capabilitiesFlag := flag.Bool("capabilities", true, "print the registered provider capability matrix before running conformance")
flag.Parse()
providers := splitCSV(*providersFlag)
@@ -47,11 +58,21 @@ func main() {
os.Exit(2)
}
if *capabilitiesFlag {
printCapabilityMatrix()
}
var ran, skipped, failed int
for _, provider := range providers {
if provider != "mock" && providerKey(provider) == "" {
fmt.Printf("- %s: skipped (set MICRO_AI_API_KEY or %s)\n", provider, providerEnv[provider])
skipped++
msg := fmt.Sprintf("set MICRO_AI_API_KEY or %s", providerEnv[provider])
if *requireConfiguredFlag {
fmt.Printf("FAIL %s: missing API key (%s)\n", provider, msg)
failed++
} else {
fmt.Printf("- %s: skipped (%s)\n", provider, msg)
skipped++
}
continue
}
@@ -72,6 +93,22 @@ func main() {
}
}
func printCapabilityMatrix() {
fmt.Println("Provider capability matrix:")
fmt.Println("provider model image video")
for _, row := range ai.CapabilityRows() {
fmt.Printf("%-12s %-5s %-5s %-5s\n", row.Provider, yesNo(row.Model), yesNo(row.Image), yesNo(row.Video))
}
fmt.Println()
}
func yesNo(ok bool) string {
if ok {
return "yes"
}
return "no"
}
func validateSelection(providers, harnesses []string) error {
if len(providers) == 0 {
return fmt.Errorf("no providers selected")
@@ -3,6 +3,8 @@ package main
import (
"strings"
"testing"
"go-micro.dev/v6/ai"
)
func TestValidateSelectionAcceptsKnownProviderAndHarness(t *testing.T) {
@@ -30,3 +32,23 @@ func TestValidateSelectionRejectsUnsafeHarnessName(t *testing.T) {
t.Fatalf("validateSelection error = %q, want invalid harness message", err)
}
}
func TestCapabilityMatrixHasRegisteredProviders(t *testing.T) {
rows := ai.CapabilityRows()
if len(rows) == 0 {
t.Fatal("CapabilityRows returned no providers")
}
var foundOpenAI bool
for _, row := range rows {
if row.Provider == "openai" {
foundOpenAI = true
if !row.Model || !row.Image || row.Video {
t.Fatalf("openai capabilities = %#v, want model+image only", row.Capabilities)
}
}
}
if !foundOpenAI {
t.Fatalf("CapabilityRows = %#v, want openai row", rows)
}
}
+17 -11
View File
@@ -207,23 +207,27 @@ func waitFor(reg registry.Registry, names ...string) {
func main() {
provider := flag.String("provider", "mock", "LLM provider: mock (default), anthropic, openai, ...")
flag.Parse()
os.Exit(runUniverse(*provider))
}
func runUniverse(provider string) int {
failures = 0
apiKey := ""
if *provider == "mock" {
if provider == "mock" {
ai.Register("mock", newMock)
} else if apiKey = providerKey(*provider); apiKey == "" {
fmt.Printf("no API key for provider %q — set MICRO_AI_API_KEY or the provider's key env\n", *provider)
os.Exit(2)
} else if apiKey = providerKey(provider); apiKey == "" {
fmt.Printf("no API key for provider %q — set MICRO_AI_API_KEY or the provider's key env\n", provider)
return 2
}
fmt.Printf("\n\033[1mUNIVERSE — booting a mini go-micro world (provider: %s)\033[0m\n\n", *provider)
fmt.Printf("\n\033[1mUNIVERSE — booting a mini go-micro world (provider: %s)\033[0m\n\n", provider)
// Infrastructure — all in-memory, all real.
reg := registry.NewMemoryRegistry()
br := broker.NewMemoryBroker()
if err := br.Connect(); err != nil {
fmt.Println("broker connect:", err)
os.Exit(2)
return 2
}
cl := client.NewClient(client.Registry(reg), client.Selector(selector.NewSelector(selector.Registry(reg))))
st := store.NewMemoryStore()
@@ -231,7 +235,7 @@ func main() {
// Services.
inv, pay, ord, ntf := new(Inventory), new(Payment), new(Orders), new(Notify)
for name, h := range map[string]any{"inventory": inv, "payment": pay, "orders": ord, "notify": ntf} {
svc := service.New(service.Name(name), service.Registry(reg), service.Client(cl))
svc := service.New(service.Name(name), service.Address("127.0.0.1:0"), service.Registry(reg), service.Client(cl))
svc.Handle(h)
go svc.Run()
}
@@ -243,7 +247,8 @@ func main() {
agent.Name("concierge"),
agent.Services("notify"),
agent.Prompt("You notify buyers when their order is confirmed."),
agent.Provider(*provider), agent.APIKey(apiKey),
agent.Address("127.0.0.1:0"),
agent.Provider(provider), agent.APIKey(apiKey),
agent.MaxSteps(5),
agent.WrapTool(func(next ai.ToolHandler) ai.ToolHandler {
return func(ctx context.Context, call ai.ToolCall) ai.ToolResult {
@@ -272,7 +277,7 @@ func main() {
)
if err := checkout.Register(reg, br, cl); err != nil {
fmt.Println("flow register:", err)
os.Exit(2)
return 2
}
defer checkout.Stop()
@@ -282,7 +287,7 @@ func main() {
fmt.Println("\033[1m> event:\033[0m events.order.placed {\"order\":\"order-1\"}")
if err := br.Publish("events.order.placed", &broker.Message{Body: []byte(`{"order":"order-1"}`)}); err != nil {
fmt.Println("publish:", err)
os.Exit(2)
return 2
}
// Wait for the run to be checkpointed as failed at "charge".
@@ -352,7 +357,8 @@ func main() {
if failures > 0 {
fmt.Printf("\n\033[31m✗ universe failed: %d assertion(s)\033[0m\n", failures)
os.Exit(1)
return 1
}
fmt.Println("\n\033[32m✓ universe: booted, survived a crash, resumed, and shut down cleanly\033[0m")
return 0
}
+20
View File
@@ -0,0 +1,20 @@
package main
import (
"testing"
)
// TestUniverseHarnessContract makes the 0→hero harness part of the ordinary
// Go test contract. The harness boots real services, a durable workflow, an
// agent, scoped state, and the A2A gateway with only the LLM mocked; running it
// here prevents the full services → agents → workflows lifecycle from silently
// drifting while developers rely on `go test ./...`.
func TestUniverseHarnessContract(t *testing.T) {
if testing.Short() {
t.Skip("universe harness boots an end-to-end system; skipped with -short")
}
if code := runUniverse("mock"); code != 0 {
t.Fatalf("universe harness exited with code %d", code)
}
}
+75
View File
@@ -0,0 +1,75 @@
---
layout: blog
title: "How Go Micro Builds Itself"
permalink: /blog/31
description: "Go Micro is increasingly built by an autonomous loop of two AI agents — Codex implementing scoped increments, Claude Code orchestrating, the human setting direction. Here's how the loop actually works, including the parts that broke."
---
# How Go Micro Builds Itself
*June 25, 2026 • By the Go Micro Team*
Go Micro is an agent harness. The most honest test of that claim is to use agents to build it — so increasingly, we do. A scheduled loop of two AI agents now opens issues, writes increments, and merges its own pull requests against this repo, on a cadence, with a human setting direction rather than typing the code.
This isn't a stunt. If a harness is good enough to operate a loop that builds itself, that's evidence it's good enough to operate the loop that builds *your* software. So we pointed the thesis at itself and wired up the loop. This post is how it actually works — including the parts that didn't, because the failure modes are the interesting part.
## Two agents and a human
The work splits across three roles:
- **Codex** is the serial builder. It takes one scoped task at a time, implements it, runs the build, tests, and linter, and opens a pull request.
- **Claude Code** is the orchestrator: it sets up the machinery, reviews, integrates, and handles the judgment calls Codex shouldn't make alone.
- **The human** sets direction and owns taste — brand and positioning copy, breaking public API changes, architectural decisions. Those never merge autonomously.
Everything else is automated.
## The mechanism
Each cycle is deliberately boring, which is the point:
1. A scheduled workflow opens a **fresh tracking issue** and dispatches Codex on it with a single instruction: pick the highest-value improvement that advances the [North Star](/docs/guides/agent-harness.html), implement it, verify it builds and passes tests and lint, and open a PR.
2. Codex does the work on its own branch and opens the PR.
3. GitHub **native auto-merge** lands it the moment the required CI checks go green — build, tests, golangci-lint. There is no human approval step. **CI is the only gate**, and that's not an approval, it's just a refusal to ship broken code.
Every increment is small, single-concern, and reversible. Nothing clever survives that can't pass the same checks a human contributor's PR would.
## Three altitudes
One loop only produces increments; it doesn't know whether they're adding up to anything. So there are three passes at different altitudes:
- An **architect** pass reviews the whole framework against the thesis every few days — API coherence, gaps in the services → agents → workflows lifecycle, drift — and files scoped issues. It decides *what* to build.
- The **hourly increment loop** builds those issues.
- A **DevRel** pass audits the README, website, docs, and blog each day for coherence, and surfaces things worth writing about. (This post is the kind of thing it's meant to catch.)
The architect points, the loop builds, DevRel keeps the story honest. Direction flows down; code flows up.
## The parts that broke
Wiring an autonomous loop is mostly plumbing and failure modes, which is exactly why it's a good test of a harness. A few we hit:
- The agent's "open a pull request" tool turned out to be a **stub** — it recorded the PR's title and body and returned them for a downstream step, but never pushed a branch or called the API. The agent cheerfully reported "opened a PR" every time, and no PR ever appeared. The fix was to stop trusting the tool and have the agent push and open the PR itself.
- Dispatching every run from a single tracking issue made the agent derive the **same branch name** each time, so the first increment opened a PR and the rest silently collided. One fresh issue per run fixed it.
- At one point an increment helpfully rewrote the repo's own agent instructions to point *back* at the broken tool. An autonomous loop will faithfully encode its own mistakes, so the guardrails have to be explicit.
None of these are exotic. They're the ordinary reality of operating an agent loop: tools that lie, state that collides, instructions that drift. The things the harness ships — observability, durable runs, resilience, guardrails — are the things you reach for the moment you try to run a loop like this in earnest.
## What it produces
The increments are unglamorous and real. Recent ones hardened the agent run loop with **OpenTelemetry run timelines** and a `micro runs` command to inspect them, correlated those timelines with trace spans, added **retry backoff** to durable flows, and made flow steps **cancellation-safe** so a canceled run stops retrying instead of burning its budget. Each one landed as a small PR that passed CI on its own.
That's the texture of the work: not a model writing a framework in one shot, but a loop making it a little better, continuously, under a gate that keeps it honest.
## The loop is the proof
We think the future of agentic software is scheduled, looping, work-performing agents — not chat. Go Micro is built by exactly that, against its own repo. The human still sets direction and owns the calls that need taste; CI is the gate; everything is reversible. Within those bounds, the harness builds itself.
If it can do that, it can build yours.
---
*Go Micro is an open source agent harness and service framework for Go. [Star us on GitHub](https://github.com/micro/go-micro).*
<div class="post-nav">
<div><a href="/blog/30">&larr; Go Micro is an Agent Harness</a></div>
<div><a href="/blog/">All Posts</a></div>
</div>
+7
View File
@@ -11,6 +11,13 @@ permalink: /blog/
<div class="posts">
<article style="margin-bottom: 2rem; padding-bottom: 1.5rem; border-bottom: 1px solid #e5e5e5;">
<h2 style="margin: 0 0 0.5rem;"><a href="/blog/31">How Go Micro Builds Itself</a></h2>
<p class="meta" style="color: #666; font-size: 0.85rem;">June 25, 2026</p>
<p>Go Micro is increasingly built by an autonomous loop of two AI agents — Codex writing scoped increments, Claude Code orchestrating, the human setting direction. Here's how the loop actually works, including the parts that broke.</p>
<a href="/blog/31">Read more &rarr;</a>
</article>
<article style="margin-bottom: 2rem; padding-bottom: 1.5rem; border-bottom: 1px solid #e5e5e5;">
<h2 style="margin: 0 0 0.5rem;"><a href="/blog/30">Go Micro is an Agent Harness</a></h2>
<p class="meta" style="color: #666; font-size: 0.85rem;">June 24, 2026</p>
@@ -21,6 +21,31 @@ ai/
└── yourprovider_test.go # Unit tests
```
## Discover registered provider capabilities
Go Micro exposes the provider interfaces registered in the current build, so
runtime tooling and docs can report what is actually available after blank
imports are linked in:
```go
for _, row := range ai.CapabilityRows() {
fmt.Printf("%s: chat=%t image=%t video=%t\n", row.Provider, row.Model, row.Image, row.Video)
}
```
The built-in providers currently register these capability interfaces:
| Provider | Chat/text (`ai.Model`) | Image (`ai.ImageModel`) | Video (`ai.VideoModel`) |
| --- | --- | --- | --- |
| `anthropic` | Yes | No | No |
| `atlascloud` | Yes | Yes | Yes |
| `gemini` | Yes | No | No |
| `groq` | Yes | No | No |
| `mistral` | Yes | No | No |
| `openai` | Yes | Yes | No |
| `together` | Yes | No | No |
## Step 1: Implement the `ai.Model` Interface
Every provider must satisfy `ai.Model`:
+3 -3
View File
@@ -183,8 +183,8 @@
<section class="section-alt">
<div class="section">
<h2>The Runtime Around the Agent</h2>
<p class="subtitle">An agent harness and a service framework in one — agents, services, and flows on the same runtime, with the production pieces agents need once they leave the demo.</p>
<h2>Features</h2>
<p class="subtitle">An agent harness and a service framework in one — agents, services, and flows on the same runtime, with the production pieces agents need.</p>
<div class="features">
<div class="feature">
<strong>Agent Harness</strong>
@@ -247,7 +247,7 @@
<div class="two-col">
<div class="two-col-text">
<h2>Developer Experience</h2>
<p><code>micro run</code> starts your services with hot reload, an API gateway, and an interactive console for talking to them. <code>micro deploy</code> pushes to production via SSH.</p>
<p><code>micro new</code> scaffolds a service or an agent — every endpoint is automatically an MCP tool. <code>micro run</code> starts everything with hot reload, an API gateway, and an interactive console, and <code>micro chat</code> talks to your agents and services from the terminal.</p>
</div>
<div>
<img src="/images/generated/developer-experience.jpg" alt="Terminal showing micro run and micro chat" />