Compare commits
2 Commits
4c29185c3b
...
48ab9d1dad
| Author | SHA1 | Date | |
|---|---|---|---|
| 48ab9d1dad | |||
| 377948708b |
@@ -85,24 +85,8 @@ func run(dir, skillDir, orgSlug string, dryRun bool, timeout time.Duration) erro
|
||||
return err
|
||||
}
|
||||
|
||||
// Parsed before anything is opened, so a malformed spec is a message rather
|
||||
// than a half-finished import. Every spec, not the first failure: an
|
||||
// operator fixing five typos should see five, not one per run.
|
||||
var problems []string
|
||||
for _, s := range specs {
|
||||
if len(s.parsed.Errors) > 0 {
|
||||
problems = append(problems, fmt.Sprintf(" %s: %s",
|
||||
s.name, strings.Join(s.parsed.Errors, "; ")))
|
||||
}
|
||||
if s.parsed.Status != "published" {
|
||||
problems = append(problems, fmt.Sprintf(
|
||||
" %s: status is %q; only a published spec can be imported",
|
||||
s.name, s.parsed.Status))
|
||||
}
|
||||
}
|
||||
if len(problems) > 0 {
|
||||
return fmt.Errorf("%d spec(s) will not import:\n%s",
|
||||
len(problems), strings.Join(problems, "\n"))
|
||||
if err := validateSpecs(specs); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
for _, s := range specs {
|
||||
@@ -116,6 +100,9 @@ func run(dir, skillDir, orgSlug string, dryRun bool, timeout time.Duration) erro
|
||||
// written. The runtime refuses to load an agent with a missing dependency,
|
||||
// so importing one without its skills produces an agent that exists and
|
||||
// cannot run — a failure that surfaces per request instead of here.
|
||||
if err := validateGraph(specs); err != nil {
|
||||
return err
|
||||
}
|
||||
if missing := missingSkills(specs, skills); len(missing) > 0 {
|
||||
return fmt.Errorf("%d skill(s) named by an agent are not in %s: %s",
|
||||
len(missing), skillDir, strings.Join(missing, ", "))
|
||||
@@ -160,73 +147,9 @@ func run(dir, skillDir, orgSlug string, dryRun bool, timeout time.Duration) erro
|
||||
}
|
||||
defer tx.Rollback(ctx) //nolint:errcheck // rolled back unless committed below
|
||||
|
||||
// Skills first. An agent row that lands before its dependencies exist is
|
||||
// briefly unloadable, and inside one transaction that is invisible — but
|
||||
// ordering them correctly costs nothing and means a future non-transactional
|
||||
// path is not silently broken.
|
||||
// Versions are recorded through the same transaction, so the history and
|
||||
// the definition it describes cannot disagree: either both land or neither
|
||||
// does.
|
||||
versions := repo.NewVersionsRepo(tx)
|
||||
ident := authctx.Identity{OrgID: orgID, UserID: author}
|
||||
|
||||
skillsWritten, skillVersions := 0, 0
|
||||
for _, sk := range skills {
|
||||
if err := upsertSkill(ctx, tx, orgID, author, sk); err != nil {
|
||||
return fmt.Errorf("%s: %w", sk.name, err)
|
||||
}
|
||||
skillsWritten++
|
||||
|
||||
// Skills are numbered by the server rather than by their author — they
|
||||
// have no `version:` to read. See snapshotSkill in internal/service,
|
||||
// which does the same for the authoring path.
|
||||
recorded, err := snapshotSkill(ctx, versions, ident, sk)
|
||||
if err != nil {
|
||||
return fmt.Errorf("%s: record version: %w", sk.name, err)
|
||||
}
|
||||
if recorded {
|
||||
skillVersions++
|
||||
}
|
||||
}
|
||||
|
||||
//
|
||||
// Snapshot is what refuses a spec that changed without raising its
|
||||
// `version:`. Every such spec is collected rather than the first one
|
||||
// returned, for the same reason the parse errors above are — an operator
|
||||
// who forgot to bump three files should see three. Collecting is safe here
|
||||
// because that refusal comes from comparing a row this code read, not from
|
||||
// a failed statement: the INSERT is ON CONFLICT DO NOTHING, so the
|
||||
// transaction is still healthy and the remaining specs can be checked.
|
||||
inserted, updated, versioned := 0, 0, 0
|
||||
var rewrites []string
|
||||
for _, s := range specs {
|
||||
wasNew, err := upsert(ctx, tx, orgID, author, s)
|
||||
if err != nil {
|
||||
return fmt.Errorf("%s: %w", s.name, err)
|
||||
}
|
||||
if wasNew {
|
||||
inserted++
|
||||
} else {
|
||||
updated++
|
||||
}
|
||||
|
||||
recorded, conflict, err := snapshotAgent(ctx, versions, ident, s)
|
||||
switch {
|
||||
case err != nil:
|
||||
return fmt.Errorf("%s: record version: %w", s.name, err)
|
||||
case conflict != "":
|
||||
rewrites = append(rewrites, fmt.Sprintf(" %s: %s", s.name, conflict))
|
||||
case recorded:
|
||||
versioned++
|
||||
}
|
||||
}
|
||||
|
||||
if len(rewrites) > 0 {
|
||||
return fmt.Errorf(
|
||||
"%d spec(s) would rewrite a version that is already published:\n%s\n\n"+
|
||||
"Nothing was written. Raise `version:` in the frontmatter of each, or "+
|
||||
"restore the published text.",
|
||||
len(rewrites), strings.Join(rewrites, "\n"))
|
||||
out, err := importInto(ctx, tx, orgID, author, specs, skills)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if err := tx.Commit(ctx); err != nil {
|
||||
@@ -235,7 +158,145 @@ func run(dir, skillDir, orgSlug string, dryRun bool, timeout time.Duration) erro
|
||||
|
||||
fmt.Printf("\n%d agent(s) published, %d updated, %d agent version(s) recorded, "+
|
||||
"%d skill(s) written, %d skill version(s) recorded, into %s\n",
|
||||
inserted, updated, versioned, skillsWritten, skillVersions, orgSlug)
|
||||
out.inserted, out.updated, out.versioned, out.skillsWritten,
|
||||
out.skillVersions, orgSlug)
|
||||
return nil
|
||||
}
|
||||
|
||||
// importCounts is what one import did.
|
||||
type importCounts struct {
|
||||
inserted int
|
||||
updated int
|
||||
versioned int
|
||||
skillsWritten int
|
||||
skillVersions int
|
||||
}
|
||||
|
||||
// importInto writes one validated set of specs and skills through a
|
||||
// transaction, and reports what it did.
|
||||
//
|
||||
// Separated from run() so it can be TESTED. run() loads configuration, opens
|
||||
// its own pool and resolves a tenant from a slug — none of which a test can
|
||||
// supply, which is why this command had no tests at all while carrying the
|
||||
// rules that decide whether a deploy is allowed to change a published agent.
|
||||
// Everything interesting lives here; run() is the wiring around it.
|
||||
//
|
||||
// The caller owns the transaction, and therefore the decision to commit. On any
|
||||
// error, including a refused rewrite, nothing here has committed and the
|
||||
// caller's deferred rollback undoes the writes that did happen.
|
||||
func importInto(ctx context.Context, tx pgx.Tx, orgID, author string,
|
||||
specs []spec, skills []skillSpec) (importCounts, error) {
|
||||
|
||||
var out importCounts
|
||||
|
||||
// Versions are recorded through the same transaction, so the history and
|
||||
// the definition it describes cannot disagree: either both land or neither
|
||||
// does.
|
||||
versions := repo.NewVersionsRepo(tx)
|
||||
ident := authctx.Identity{OrgID: orgID, UserID: author}
|
||||
|
||||
// Skills first. An agent row that lands before its dependencies exist is
|
||||
// briefly unloadable, and inside one transaction that is invisible — but
|
||||
// ordering them correctly costs nothing and means a future
|
||||
// non-transactional path is not silently broken.
|
||||
for _, sk := range skills {
|
||||
if err := upsertSkill(ctx, tx, orgID, author, sk); err != nil {
|
||||
return out, fmt.Errorf("%s: %w", sk.name, err)
|
||||
}
|
||||
out.skillsWritten++
|
||||
|
||||
// Skills are numbered by the server rather than by their author — they
|
||||
// have no `version:` to read. See snapshotSkill in internal/service,
|
||||
// which does the same for the authoring path.
|
||||
recorded, err := snapshotSkill(ctx, versions, ident, sk)
|
||||
if err != nil {
|
||||
return out, fmt.Errorf("%s: record version: %w", sk.name, err)
|
||||
}
|
||||
if recorded {
|
||||
out.skillVersions++
|
||||
}
|
||||
}
|
||||
|
||||
// snapshotAgent is what refuses a spec that changed without raising its
|
||||
// `version:`, or that lowers it. Every such spec is collected rather than
|
||||
// the first one returned, for the same reason the parse errors are — an
|
||||
// operator who forgot to bump three files should see three. Collecting is
|
||||
// safe because that refusal comes from comparing rows this code read, not
|
||||
// from a failed statement: the INSERT is ON CONFLICT DO NOTHING, so the
|
||||
// transaction is still healthy and the remaining specs can be checked.
|
||||
var rewrites []string
|
||||
for _, s := range specs {
|
||||
wasNew, err := upsert(ctx, tx, orgID, author, s)
|
||||
if err != nil {
|
||||
return out, fmt.Errorf("%s: %w", s.name, err)
|
||||
}
|
||||
if wasNew {
|
||||
out.inserted++
|
||||
} else {
|
||||
out.updated++
|
||||
}
|
||||
|
||||
recorded, conflict, err := snapshotAgent(ctx, versions, ident, s)
|
||||
switch {
|
||||
case err != nil:
|
||||
return out, fmt.Errorf("%s: record version: %w", s.name, err)
|
||||
case conflict != "":
|
||||
rewrites = append(rewrites, fmt.Sprintf(" %s: %s", s.name, conflict))
|
||||
case recorded:
|
||||
out.versioned++
|
||||
}
|
||||
}
|
||||
|
||||
if len(rewrites) > 0 {
|
||||
return out, fmt.Errorf(
|
||||
"%d spec(s) would rewrite a version that is already published:\n%s\n\n"+
|
||||
"Nothing was written. Raise `version:` in the frontmatter of each, or "+
|
||||
"restore the published text.",
|
||||
len(rewrites), strings.Join(rewrites, "\n"))
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// validateSpecs rejects specs that cannot be imported, reporting every one.
|
||||
//
|
||||
// Parsed before anything is opened, so a malformed spec is a message rather
|
||||
// than a half-finished import. Every spec, not the first failure: an operator
|
||||
// fixing five typos should see five, not one per run.
|
||||
func validateSpecs(specs []spec) error {
|
||||
var problems []string
|
||||
for _, s := range specs {
|
||||
if len(s.parsed.Errors) > 0 {
|
||||
problems = append(problems, fmt.Sprintf(" %s: %s",
|
||||
s.name, strings.Join(s.parsed.Errors, "; ")))
|
||||
}
|
||||
if s.parsed.Status != "published" {
|
||||
problems = append(problems, fmt.Sprintf(
|
||||
" %s: status is %q; only a published spec can be imported",
|
||||
s.name, s.parsed.Status))
|
||||
}
|
||||
}
|
||||
if len(problems) > 0 {
|
||||
return fmt.Errorf("%d spec(s) will not import:\n%s",
|
||||
len(problems), strings.Join(problems, "\n"))
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// validateGraph enforces §3's DAG requirement across the whole set.
|
||||
//
|
||||
// Every spec is in hand here, which is the only place that is cheaply true —
|
||||
// so this is where the check belongs. The runtime depth cap still bounds a
|
||||
// cycle that reaches run time by another route.
|
||||
func validateGraph(specs []spec) error {
|
||||
graph := make(map[string][]string, len(specs))
|
||||
for _, s := range specs {
|
||||
graph[s.parsed.ID] = s.parsed.Subagents
|
||||
}
|
||||
if cycle := definition.FindSubagentCycle(graph); cycle != "" {
|
||||
return fmt.Errorf("the subagent graph has a cycle: %s\n\n"+
|
||||
"Nothing was written. Delegation follows these edges, so a loop is a "+
|
||||
"run that delegates until it runs out of budget.", cycle)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -494,6 +555,19 @@ func snapshotSkill(ctx context.Context, versions *repo.VersionsRepo,
|
||||
func snapshotAgent(ctx context.Context, versions *repo.VersionsRepo,
|
||||
ident authctx.Identity, s spec) (recorded bool, conflict string, err error) {
|
||||
|
||||
// §3 calls the version monotonic and nothing enforced it. A spec edited
|
||||
// from an older copy republishes an older number whose content still
|
||||
// matches what was published under it — no conflict, no complaint, and the
|
||||
// deployed agent quietly goes backwards.
|
||||
latest, err := versions.LatestVersion(ctx, ident, repo.KindAgent, s.parsed.ID)
|
||||
if err != nil {
|
||||
return false, "", err
|
||||
}
|
||||
if latest > 0 && s.parsed.Version < latest {
|
||||
return false, definition.ErrVersionWentBackwards(
|
||||
s.parsed.ID, latest, s.parsed.Version).Error(), nil
|
||||
}
|
||||
|
||||
stored, err := versions.Load(ctx, ident, repo.KindAgent, s.parsed.ID, s.parsed.Version)
|
||||
if err != nil {
|
||||
var apiErr *domain.Error
|
||||
|
||||
244
go-api/cmd/importagents/main_test.go
Normal file
244
go-api/cmd/importagents/main_test.go
Normal file
@@ -0,0 +1,244 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/krow/krow-backend/go-api/internal/definition"
|
||||
"github.com/krow/krow-backend/go-api/internal/testutil"
|
||||
)
|
||||
|
||||
// specFor builds a parsed spec the way loadSpecs would, without a file.
|
||||
func specFor(t *testing.T, id string, version int, subagents ...string) spec {
|
||||
t.Helper()
|
||||
var sub string
|
||||
if len(subagents) > 0 {
|
||||
sub = "subagents:\n"
|
||||
for _, s := range subagents {
|
||||
sub += " - " + s + "\n"
|
||||
}
|
||||
}
|
||||
raw := fmt.Sprintf(`---
|
||||
id: %s
|
||||
name: %s
|
||||
description: a spec built for a test
|
||||
icon: layers
|
||||
status: published
|
||||
version: %d
|
||||
reasoning: balanced
|
||||
pages:
|
||||
- talent-pool
|
||||
%s---
|
||||
|
||||
## Instructions
|
||||
Answer the question, version %d.
|
||||
`, id, strings.ToUpper(id[:1])+id[1:], version, sub, version)
|
||||
|
||||
parsed, err := definition.ParseAgent(raw, definition.Options{})
|
||||
if err != nil {
|
||||
t.Fatalf("fixture %q does not parse: %v", id, err)
|
||||
}
|
||||
return spec{name: id + ".md", raw: raw, parsed: parsed}
|
||||
}
|
||||
|
||||
/* ── The pure checks, which need no database ─────────────────────────────── */
|
||||
|
||||
func TestValidateSpecsReportsEveryProblem(t *testing.T) {
|
||||
draft := specFor(t, "draft-agent", 1)
|
||||
draft.parsed.Status = "draft"
|
||||
broken := specFor(t, "broken-agent", 1)
|
||||
broken.parsed.Errors = []string{"something is wrong"}
|
||||
|
||||
err := validateSpecs([]spec{draft, broken, specFor(t, "fine-agent", 1)})
|
||||
if err == nil {
|
||||
t.Fatal("two bad specs were accepted")
|
||||
}
|
||||
// Both, not the first: an operator fixing two problems should see two.
|
||||
for _, want := range []string{"draft-agent", "broken-agent"} {
|
||||
if !strings.Contains(err.Error(), want) {
|
||||
t.Errorf("the report does not mention %q:\n%s", want, err)
|
||||
}
|
||||
}
|
||||
if strings.Contains(err.Error(), "fine-agent") {
|
||||
t.Errorf("a valid spec was reported as a problem:\n%s", err)
|
||||
}
|
||||
|
||||
if err := validateSpecs([]spec{specFor(t, "fine-agent", 1)}); err != nil {
|
||||
t.Errorf("a valid spec was refused: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestValidateGraphRefusesACycle(t *testing.T) {
|
||||
acyclic := []spec{
|
||||
specFor(t, "a-agent", 1, "b-agent"),
|
||||
specFor(t, "b-agent", 1),
|
||||
}
|
||||
if err := validateGraph(acyclic); err != nil {
|
||||
t.Errorf("a chain was called a cycle: %v", err)
|
||||
}
|
||||
|
||||
cyclic := []spec{
|
||||
specFor(t, "a-agent", 1, "b-agent"),
|
||||
specFor(t, "b-agent", 1, "a-agent"),
|
||||
}
|
||||
err := validateGraph(cyclic)
|
||||
if err == nil {
|
||||
t.Fatal("a cycle was accepted")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "a-agent") || !strings.Contains(err.Error(), "b-agent") {
|
||||
t.Errorf("the message does not name the edge to cut:\n%s", err)
|
||||
}
|
||||
}
|
||||
|
||||
/* ── The write phase, against a real database ────────────────────────────── */
|
||||
|
||||
// importOnce runs one import in its own transaction and commits it, the way
|
||||
// run() does.
|
||||
//
|
||||
// The author comes from resolveAuthor rather than a literal, so this exercises
|
||||
// the production path and fails loudly if an organization has nobody to
|
||||
// attribute specs to — which is a real deployment condition, not a test
|
||||
// detail.
|
||||
func importOnce(t *testing.T, h *testutil.Harness, specs []spec) (importCounts, error) {
|
||||
t.Helper()
|
||||
ctx := context.Background()
|
||||
|
||||
author, err := resolveAuthor(ctx, h.Pool, h.OrgID)
|
||||
if err != nil {
|
||||
t.Fatalf("resolve author: %v", err)
|
||||
}
|
||||
|
||||
tx, err := h.Pool.Begin(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("begin: %v", err)
|
||||
}
|
||||
defer tx.Rollback(ctx) //nolint:errcheck
|
||||
|
||||
out, err := importInto(ctx, tx, h.OrgID, author, specs, nil)
|
||||
if err != nil {
|
||||
return out, err
|
||||
}
|
||||
if err := tx.Commit(ctx); err != nil {
|
||||
t.Fatalf("commit: %v", err)
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
func TestImportRecordsVersionsAndIsIdempotent(t *testing.T) {
|
||||
h := testutil.New(t)
|
||||
specs := []spec{specFor(t, "import-a", 1), specFor(t, "import-b", 1)}
|
||||
|
||||
out, err := importOnce(t, h, specs)
|
||||
if err != nil {
|
||||
t.Fatalf("first import: %v", err)
|
||||
}
|
||||
if out.inserted != 2 || out.versioned != 2 {
|
||||
t.Errorf("first import: inserted=%d versioned=%d, want 2 and 2", out.inserted, out.versioned)
|
||||
}
|
||||
|
||||
// Again, unchanged. Nothing new is recorded — the number must not report
|
||||
// nine every deploy, which is what it used to do.
|
||||
out, err = importOnce(t, h, specs)
|
||||
if err != nil {
|
||||
t.Fatalf("second import: %v", err)
|
||||
}
|
||||
if out.versioned != 0 {
|
||||
t.Errorf("re-importing unchanged specs recorded %d version(s), want 0", out.versioned)
|
||||
}
|
||||
if out.updated != 2 {
|
||||
t.Errorf("second import: updated=%d, want 2", out.updated)
|
||||
}
|
||||
}
|
||||
|
||||
func TestImportRefusesRewritingAPublishedVersion(t *testing.T) {
|
||||
h := testutil.New(t)
|
||||
if _, err := importOnce(t, h, []spec{specFor(t, "rewrite-me", 1)}); err != nil {
|
||||
t.Fatalf("first import: %v", err)
|
||||
}
|
||||
|
||||
// Same version, different body.
|
||||
changed := specFor(t, "rewrite-me", 1)
|
||||
changed.raw = strings.Replace(changed.raw, "version 1.", "something else entirely.", 1)
|
||||
reparsed, err := definition.ParseAgent(changed.raw, definition.Options{})
|
||||
if err != nil {
|
||||
t.Fatalf("fixture does not parse: %v", err)
|
||||
}
|
||||
changed.parsed = reparsed
|
||||
|
||||
_, err = importOnce(t, h, []spec{changed})
|
||||
if err == nil {
|
||||
t.Fatal("a changed spec republished at the same version was accepted")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "rewrite") {
|
||||
t.Errorf("unexpected error: %v", err)
|
||||
}
|
||||
|
||||
// And nothing landed: the live row still says what v1 said.
|
||||
var live string
|
||||
if err := h.Pool.QueryRow(context.Background(),
|
||||
`SELECT markdown FROM agent_definitions WHERE org_id = $1::uuid AND definition_id = 'rewrite-me'`,
|
||||
h.OrgID).Scan(&live); err != nil {
|
||||
t.Fatalf("read back: %v", err)
|
||||
}
|
||||
if strings.Contains(live, "something else entirely") {
|
||||
t.Error("the refused import was committed anyway")
|
||||
}
|
||||
}
|
||||
|
||||
func TestImportRefusesAVersionGoingBackwards(t *testing.T) {
|
||||
h := testutil.New(t)
|
||||
if _, err := importOnce(t, h, []spec{specFor(t, "backwards", 1)}); err != nil {
|
||||
t.Fatalf("v1: %v", err)
|
||||
}
|
||||
if _, err := importOnce(t, h, []spec{specFor(t, "backwards", 2)}); err != nil {
|
||||
t.Fatalf("v2: %v", err)
|
||||
}
|
||||
|
||||
// Back to v1, byte-for-byte what v1 said. Nothing conflicts, which is why
|
||||
// this used to succeed and silently revert the deployed agent.
|
||||
_, err := importOnce(t, h, []spec{specFor(t, "backwards", 1)})
|
||||
if err == nil {
|
||||
t.Fatal("a lowered version was accepted")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "monotonic") {
|
||||
t.Errorf("the message does not explain why: %v", err)
|
||||
}
|
||||
|
||||
var version int
|
||||
if err := h.Pool.QueryRow(context.Background(),
|
||||
`SELECT version FROM agent_definitions WHERE org_id = $1::uuid AND definition_id = 'backwards'`,
|
||||
h.OrgID).Scan(&version); err != nil {
|
||||
t.Fatalf("read back: %v", err)
|
||||
}
|
||||
if version != 2 {
|
||||
t.Errorf("live version = %d, want 2 — the refused import rolled it back", version)
|
||||
}
|
||||
}
|
||||
|
||||
// Raising the version is the supported way to change a published spec.
|
||||
func TestImportAcceptsARaisedVersion(t *testing.T) {
|
||||
h := testutil.New(t)
|
||||
if _, err := importOnce(t, h, []spec{specFor(t, "raised", 1)}); err != nil {
|
||||
t.Fatalf("v1: %v", err)
|
||||
}
|
||||
out, err := importOnce(t, h, []spec{specFor(t, "raised", 2)})
|
||||
if err != nil {
|
||||
t.Fatalf("v2 was refused: %v", err)
|
||||
}
|
||||
if out.versioned != 1 {
|
||||
t.Errorf("recorded %d version(s) for a raised version, want 1", out.versioned)
|
||||
}
|
||||
|
||||
var n int
|
||||
if err := h.Pool.QueryRow(context.Background(),
|
||||
`SELECT count(*) FROM definition_versions
|
||||
WHERE org_id = $1::uuid AND kind = 'agent' AND definition_id = 'raised'`,
|
||||
h.OrgID).Scan(&n); err != nil {
|
||||
t.Fatalf("count: %v", err)
|
||||
}
|
||||
if n != 2 {
|
||||
t.Errorf("history holds %d versions, want 2 (v1 and v2)", n)
|
||||
}
|
||||
}
|
||||
103
go-api/internal/definition/graph.go
Normal file
103
go-api/internal/definition/graph.go
Normal file
@@ -0,0 +1,103 @@
|
||||
package definition
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"sort"
|
||||
"strings"
|
||||
)
|
||||
|
||||
// FindSubagentCycle reports the first delegation cycle in a set of agents, or
|
||||
// "" if the graph is acyclic.
|
||||
//
|
||||
// §3: "subagents must form a DAG. Cycle detection runs at publish." This is the
|
||||
// publish-time half. The runtime half is runtime.MaxDelegationDepth, which
|
||||
// bounds a cycle that reaches run time anyway — because a graph can only be
|
||||
// checked against the agents the checker was GIVEN, and an agent published
|
||||
// while another is being edited can complete a loop neither publish saw.
|
||||
//
|
||||
// The returned string names the cycle in the order it was walked, so an
|
||||
// operator can see which edge to cut:
|
||||
//
|
||||
// a -> b -> c -> a
|
||||
//
|
||||
// Edges pointing at agents not in the set are ignored rather than treated as
|
||||
// missing. Resolving those is a different check with a different message
|
||||
// (runtime.unknown_subagent), and conflating the two produces "cycle detected"
|
||||
// for what is actually a typo.
|
||||
func FindSubagentCycle(subagents map[string][]string) string {
|
||||
// Depth-first search tracking the path, so the cycle can be REPORTED
|
||||
// rather than merely detected — "there is a cycle" leaves an operator to
|
||||
// find it by hand across a set of specs.
|
||||
//
|
||||
// Recursive, and deliberately: a goroutine stack grows on demand, so depth
|
||||
// here costs memory rather than a crash, and a 5000-long chain is covered
|
||||
// by a test. An explicit stack would buy nothing and lose the path
|
||||
// bookkeeping that makes the message useful.
|
||||
const (
|
||||
unvisited = 0
|
||||
onPath = 1
|
||||
done = 2
|
||||
)
|
||||
state := make(map[string]int, len(subagents))
|
||||
|
||||
// Sorted, so the same set of agents always reports the same cycle. An
|
||||
// error message that changes between runs on identical input is one
|
||||
// nobody trusts.
|
||||
roots := make([]string, 0, len(subagents))
|
||||
for id := range subagents {
|
||||
roots = append(roots, id)
|
||||
}
|
||||
sort.Strings(roots)
|
||||
|
||||
var path []string
|
||||
var walk func(id string) string
|
||||
walk = func(id string) string {
|
||||
switch state[id] {
|
||||
case done:
|
||||
return ""
|
||||
case onPath:
|
||||
// Found it. Report from the first occurrence of this id, so the
|
||||
// message is the cycle itself and not the walk that reached it.
|
||||
for i, seen := range path {
|
||||
if seen == id {
|
||||
return strings.Join(append(append([]string{}, path[i:]...), id), " -> ")
|
||||
}
|
||||
}
|
||||
return id + " -> " + id
|
||||
}
|
||||
|
||||
state[id] = onPath
|
||||
path = append(path, id)
|
||||
for _, next := range subagents[id] {
|
||||
if _, known := subagents[next]; !known {
|
||||
continue // not ours to judge; see the doc comment
|
||||
}
|
||||
if cycle := walk(next); cycle != "" {
|
||||
return cycle
|
||||
}
|
||||
}
|
||||
path = path[:len(path)-1]
|
||||
state[id] = done
|
||||
return ""
|
||||
}
|
||||
|
||||
for _, id := range roots {
|
||||
if cycle := walk(id); cycle != "" {
|
||||
return cycle
|
||||
}
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
// ErrVersionWentBackwards describes a publish that lowers a version.
|
||||
//
|
||||
// §3 calls the version monotonic. Nothing enforced it: the upsert wrote
|
||||
// whatever the frontmatter said, so a spec edited from an older copy silently
|
||||
// rolled a deployed agent backwards — no conflict, because the older version's
|
||||
// content still matched what was published under that number.
|
||||
func ErrVersionWentBackwards(id string, from, to int) error {
|
||||
return fmt.Errorf(
|
||||
"%q is published at version %d and this publishes version %d; "+
|
||||
"a version is monotonic, so raise it above %d rather than lowering it",
|
||||
id, from, to, from)
|
||||
}
|
||||
101
go-api/internal/definition/graph_test.go
Normal file
101
go-api/internal/definition/graph_test.go
Normal file
@@ -0,0 +1,101 @@
|
||||
package definition_test
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/krow/krow-backend/go-api/internal/definition"
|
||||
)
|
||||
|
||||
func TestFindSubagentCycle(t *testing.T) {
|
||||
for _, tc := range []struct {
|
||||
name string
|
||||
graph map[string][]string
|
||||
want string // "" means acyclic; otherwise a substring the report must contain
|
||||
}{
|
||||
{"empty", map[string][]string{}, ""},
|
||||
{"no edges", map[string][]string{"a": nil, "b": nil}, ""},
|
||||
{"a chain is not a cycle", map[string][]string{
|
||||
"a": {"b"}, "b": {"c"}, "c": nil,
|
||||
}, ""},
|
||||
{"a diamond is not a cycle", map[string][]string{
|
||||
"a": {"b", "c"}, "b": {"d"}, "c": {"d"}, "d": nil,
|
||||
}, ""},
|
||||
{"self reference", map[string][]string{"a": {"a"}}, "a -> a"},
|
||||
{"two-agent loop", map[string][]string{
|
||||
"a": {"b"}, "b": {"a"},
|
||||
}, "a -> b -> a"},
|
||||
{"longer loop", map[string][]string{
|
||||
"a": {"b"}, "b": {"c"}, "c": {"a"},
|
||||
}, "a -> b -> c -> a"},
|
||||
{"cycle not involving the first agent walked", map[string][]string{
|
||||
"a": {"b"}, "b": {"c"}, "c": {"b"},
|
||||
}, "b -> c -> b"},
|
||||
{"an edge to an unknown agent is not a cycle", map[string][]string{
|
||||
"a": {"nowhere"},
|
||||
}, ""},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
got := definition.FindSubagentCycle(tc.graph)
|
||||
switch {
|
||||
case tc.want == "" && got != "":
|
||||
t.Errorf("reported a cycle %q in an acyclic graph", got)
|
||||
case tc.want != "" && got == "":
|
||||
t.Errorf("missed the cycle; want something containing %q", tc.want)
|
||||
case tc.want != "" && !strings.Contains(got, tc.want):
|
||||
t.Errorf("cycle = %q, want it to contain %q", got, tc.want)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// The report must be stable: the same graph reported differently on different
|
||||
// runs is an error message nobody trusts, and map iteration order in Go is
|
||||
// deliberately random.
|
||||
func TestFindSubagentCycleIsDeterministic(t *testing.T) {
|
||||
graph := map[string][]string{
|
||||
"e": {"f"}, "f": {"e"},
|
||||
"a": {"b"}, "b": {"c"}, "c": {"a"},
|
||||
"z": nil, "y": {"z"},
|
||||
}
|
||||
first := definition.FindSubagentCycle(graph)
|
||||
if first == "" {
|
||||
t.Fatal("no cycle found in a graph with two")
|
||||
}
|
||||
for i := 0; i < 50; i++ {
|
||||
if got := definition.FindSubagentCycle(graph); got != first {
|
||||
t.Fatalf("run %d reported %q, first run reported %q — the report "+
|
||||
"depends on map iteration order", i, got, first)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// A deep chain must not overflow the stack. An author supplies this graph.
|
||||
func TestFindSubagentCycleHandlesADeepChain(t *testing.T) {
|
||||
graph := map[string][]string{}
|
||||
const n = 5000
|
||||
for i := 0; i < n; i++ {
|
||||
graph[itoa(i)] = []string{itoa(i + 1)}
|
||||
}
|
||||
graph[itoa(n)] = nil
|
||||
if got := definition.FindSubagentCycle(graph); got != "" {
|
||||
t.Errorf("reported a cycle %q in a %d-long chain", got, n)
|
||||
}
|
||||
// And the same chain closed into a loop is found.
|
||||
graph[itoa(n)] = []string{itoa(0)}
|
||||
if definition.FindSubagentCycle(graph) == "" {
|
||||
t.Error("missed a cycle closing a long chain")
|
||||
}
|
||||
}
|
||||
|
||||
func itoa(i int) string {
|
||||
if i == 0 {
|
||||
return "0"
|
||||
}
|
||||
var b []byte
|
||||
for i > 0 {
|
||||
b = append([]byte{byte('0' + i%10)}, b...)
|
||||
i /= 10
|
||||
}
|
||||
return string(b)
|
||||
}
|
||||
@@ -1175,3 +1175,155 @@ The instructions, unchanged throughout.
|
||||
res.code, res.body)
|
||||
}
|
||||
}
|
||||
|
||||
// TestPublishedVersionCannotGoBackwards covers §3's "monotonic".
|
||||
//
|
||||
// The rewrite guard only compares content at ONE version number, so an older
|
||||
// number republished with the text that was originally published under it
|
||||
// looked like a no-op: no conflict, nothing to refuse, and the live row
|
||||
// silently reverted. The agent in the UI then reads v1 while the newest thing
|
||||
// anybody approved was v2.
|
||||
func TestPublishedVersionCannotGoBackwards(t *testing.T) {
|
||||
r := newRBAC(t)
|
||||
|
||||
const v1 = `---
|
||||
id: monotonic-agent
|
||||
name: Monotonic Agent
|
||||
description: published twice, then rolled back
|
||||
status: published
|
||||
version: 1
|
||||
pages:
|
||||
- candidates
|
||||
---
|
||||
|
||||
## Instructions
|
||||
The first version.
|
||||
`
|
||||
res := r.as(r.admin, "POST", "/api/v1/agent-definitions", map[string]any{
|
||||
"markdown": v1, "visibility": "personal",
|
||||
})
|
||||
if res.code != http.StatusCreated {
|
||||
t.Fatalf("create v1: status %d (%v)", res.code, res.body)
|
||||
}
|
||||
id, _ := res.record(t)["id"].(string)
|
||||
|
||||
v2 := strings.Replace(strings.Replace(v1, "version: 1", "version: 2", 1),
|
||||
"The first version.", "The second version.", 1)
|
||||
res = r.as(r.admin, "PATCH", "/api/v1/agent-definitions/"+id, map[string]any{"markdown": v2})
|
||||
if res.code != http.StatusOK {
|
||||
t.Fatalf("publish v2: status %d (%v)", res.code, res.body)
|
||||
}
|
||||
|
||||
// Back to v1, byte-for-byte what v1 said. Nothing here conflicts — which
|
||||
// is exactly why it used to succeed.
|
||||
res = r.as(r.admin, "PATCH", "/api/v1/agent-definitions/"+id, map[string]any{"markdown": v1})
|
||||
if res.code != http.StatusConflict {
|
||||
t.Fatalf("republishing v1 after v2: status %d, want 409 (%v)", res.code, res.body)
|
||||
}
|
||||
|
||||
// And the live definition is still v2, not silently reverted.
|
||||
res = r.as(r.admin, "GET", "/api/v1/agent-definitions/"+id, nil)
|
||||
if got := fmt.Sprint(res.record(t)["version"]); got != "2" {
|
||||
t.Errorf("live version = %s, want 2 — the refused publish rolled it back anyway", got)
|
||||
}
|
||||
}
|
||||
|
||||
// TestSubagentCycleIsRefusedAtPublish covers §3's DAG requirement.
|
||||
//
|
||||
// runtime.MaxDelegationDepth bounds a cycle that reaches run time, so this is
|
||||
// not a safety hole — it is a budget one. Every run that entered the loop would
|
||||
// spend its whole allowance delegating in a circle before terminating, and the
|
||||
// person who wrote the loop would learn about it from a bill rather than from
|
||||
// the publish that created it.
|
||||
func TestSubagentCycleIsRefusedAtPublish(t *testing.T) {
|
||||
r := newRBAC(t)
|
||||
|
||||
// version is a parameter so the loop-closing edit can BUMP it. Otherwise
|
||||
// the rewrite guard refuses that edit for changing published text, the
|
||||
// test passes for the wrong reason, and it would keep passing with cycle
|
||||
// detection removed entirely.
|
||||
agent := func(id, name string, version int, subagents ...string) string {
|
||||
var sub string
|
||||
if len(subagents) > 0 {
|
||||
sub = "subagents:\n"
|
||||
for _, s := range subagents {
|
||||
sub += " - " + s + "\n"
|
||||
}
|
||||
}
|
||||
return fmt.Sprintf(`---
|
||||
id: %s
|
||||
name: %s
|
||||
description: part of a delegation graph
|
||||
status: published
|
||||
version: %d
|
||||
pages:
|
||||
- candidates
|
||||
%s---
|
||||
|
||||
## Instructions
|
||||
Delegate.
|
||||
`, id, name, version, sub)
|
||||
}
|
||||
|
||||
// A, with no subagents yet.
|
||||
res := r.as(r.admin, "POST", "/api/v1/agent-definitions", map[string]any{
|
||||
"markdown": agent("cycle-a", "Cycle A", 1), "visibility": "organization",
|
||||
})
|
||||
if res.code != http.StatusCreated {
|
||||
t.Fatalf("create A: status %d (%v)", res.code, res.body)
|
||||
}
|
||||
idA, _ := res.record(t)["id"].(string)
|
||||
|
||||
// B delegates to A. Still a DAG.
|
||||
res = r.as(r.admin, "POST", "/api/v1/agent-definitions", map[string]any{
|
||||
"markdown": agent("cycle-b", "Cycle B", 1, "cycle-a"), "visibility": "organization",
|
||||
})
|
||||
if res.code != http.StatusCreated {
|
||||
t.Fatalf("create B pointing at A: status %d, want 201 — a chain is not a cycle (%v)",
|
||||
res.code, res.body)
|
||||
}
|
||||
|
||||
// Now close the loop: A delegates to B.
|
||||
res = r.as(r.admin, "PATCH", "/api/v1/agent-definitions/"+idA, map[string]any{
|
||||
"markdown": agent("cycle-a", "Cycle A", 2, "cycle-b"),
|
||||
})
|
||||
if res.code != http.StatusUnprocessableEntity && res.code != http.StatusBadRequest {
|
||||
t.Fatalf("closing the loop: status %d, want a validation failure (%v)", res.code, res.body)
|
||||
}
|
||||
if body := fmt.Sprint(res.body); !strings.Contains(body, "cycle") {
|
||||
t.Errorf("the refusal did not mention a cycle: %v", res.body)
|
||||
}
|
||||
|
||||
// A must be unchanged — refused, not half-applied.
|
||||
res = r.as(r.admin, "GET", "/api/v1/agent-definitions/"+idA, nil)
|
||||
if md, _ := res.record(t)["markdown"].(string); strings.Contains(md, "cycle-b") {
|
||||
t.Error("the refused edit was applied anyway")
|
||||
}
|
||||
}
|
||||
|
||||
// A self-reference is the shortest cycle and the easiest to write by accident.
|
||||
func TestSelfReferencingSubagentIsRefused(t *testing.T) {
|
||||
r := newRBAC(t)
|
||||
|
||||
const md = `---
|
||||
id: narcissus-agent
|
||||
name: Narcissus Agent
|
||||
description: names itself
|
||||
status: published
|
||||
version: 1
|
||||
pages:
|
||||
- candidates
|
||||
subagents:
|
||||
- narcissus-agent
|
||||
---
|
||||
|
||||
## Instructions
|
||||
Ask myself.
|
||||
`
|
||||
res := r.as(r.admin, "POST", "/api/v1/agent-definitions", map[string]any{
|
||||
"markdown": md, "visibility": "organization",
|
||||
})
|
||||
if res.code == http.StatusCreated {
|
||||
t.Fatal("an agent naming itself as its own subagent was published")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -234,6 +234,13 @@ func (s *DefinitionsService) CreateAgent(ctx context.Context, ident authctx.Iden
|
||||
input.OwnerUserID = &ident.UserID
|
||||
}
|
||||
|
||||
if err := s.refuseVersionGoingBackwards(ctx, ident, repo.KindAgent,
|
||||
agent.ID, agent.Status, agent.Version); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if err := s.refuseSubagentCycle(ctx, ident, agent.ID, agent.Subagents); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if err := s.refusePublishedRewrite(ctx, ident, repo.KindAgent,
|
||||
agent.ID, markdown, agent.Status, agent.Version); err != nil {
|
||||
return nil, err
|
||||
@@ -302,6 +309,13 @@ func (s *DefinitionsService) UpdateAgent(ctx context.Context, ident authctx.Iden
|
||||
input.Version = &agent.Version
|
||||
input.Pages = agent.Pages
|
||||
|
||||
if err := s.refuseVersionGoingBackwards(ctx, ident, repo.KindAgent,
|
||||
agent.ID, agent.Status, agent.Version); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if err := s.refuseSubagentCycle(ctx, ident, agent.ID, agent.Subagents); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if err := s.refusePublishedRewrite(ctx, ident, repo.KindAgent,
|
||||
agent.ID, markdown, agent.Status, agent.Version); err != nil {
|
||||
return nil, err
|
||||
@@ -535,6 +549,80 @@ func (s *DefinitionsService) DeleteSkill(ctx context.Context, ident authctx.Iden
|
||||
|
||||
/* ── Publishing ─────────────────────────────────────────────────────────── */
|
||||
|
||||
// refuseVersionGoingBackwards enforces §3's "monotonic".
|
||||
//
|
||||
// Nothing enforced it before. A spec edited from an older copy republishes an
|
||||
// older number whose content still matches what was published under it, so the
|
||||
// rewrite guard sees no conflict and the deployed agent quietly goes
|
||||
// backwards — the version in the UI reads 2 while the newest thing anybody
|
||||
// approved was 3.
|
||||
func (s *DefinitionsService) refuseVersionGoingBackwards(ctx context.Context,
|
||||
ident authctx.Identity, kind repo.VersionKind, definitionID, status string, version int) error {
|
||||
|
||||
if s.versions == nil || status != "published" || definitionID == "" || version < 1 {
|
||||
return nil
|
||||
}
|
||||
latest, err := s.versions.LatestVersion(ctx, ident, kind, definitionID)
|
||||
if err != nil || latest == 0 {
|
||||
// Unreadable, or nothing published yet. Neither is grounds to refuse a
|
||||
// save: the rewrite guard is what protects published text, and this
|
||||
// only orders the numbers.
|
||||
return nil
|
||||
}
|
||||
if version < latest {
|
||||
return domain.Conflict(definition.ErrVersionWentBackwards(definitionID, latest, version).Error())
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// refuseSubagentCycle enforces §3's DAG at publish.
|
||||
//
|
||||
// The graph is every organization-visible agent plus the one being published,
|
||||
// with the incoming definition standing in for its stored self — otherwise an
|
||||
// edit that CREATES a cycle is checked against the version that did not have
|
||||
// one, and passes.
|
||||
//
|
||||
// Personal agents are not included. They are invisible to everyone else, so
|
||||
// they cannot complete a loop for anybody else, and loading them would mean
|
||||
// reading other people's drafts to validate your own.
|
||||
func (s *DefinitionsService) refuseSubagentCycle(ctx context.Context,
|
||||
ident authctx.Identity, definitionID string, subagents []string) error {
|
||||
|
||||
if len(subagents) == 0 {
|
||||
return nil
|
||||
}
|
||||
rows, _, err := s.repo.ListAgents(ctx, ident, repo.DefinitionListParams{
|
||||
Visibility: "organization", Limit: 500,
|
||||
})
|
||||
if err != nil {
|
||||
// A graph we could not read is not a graph we can call cyclic. The
|
||||
// runtime depth cap is what holds when this cannot run.
|
||||
return nil
|
||||
}
|
||||
|
||||
graph := make(map[string][]string, len(rows)+1)
|
||||
for _, rec := range rows {
|
||||
id, _ := rec["definition_id"].(string)
|
||||
markdown, _ := rec["markdown"].(string)
|
||||
if id == "" || markdown == "" || id == definitionID {
|
||||
continue // the incoming definition replaces its stored self, below
|
||||
}
|
||||
if parsed, err := definition.ParseAgent(markdown, definition.Options{}); err == nil && parsed != nil {
|
||||
graph[id] = parsed.Subagents
|
||||
}
|
||||
}
|
||||
graph[definitionID] = subagents
|
||||
|
||||
if cycle := definition.FindSubagentCycle(graph); cycle != "" {
|
||||
return domain.Validation(
|
||||
"that would create a delegation cycle: "+cycle+
|
||||
". Delegation follows these edges, so a loop is a run that delegates "+
|
||||
"until it runs out of budget.",
|
||||
map[string]string{"subagents": "cycle"})
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// refusePublishedRewrite fails a publish that would change a version already
|
||||
// published, BEFORE anything is written.
|
||||
//
|
||||
|
||||
Reference in New Issue
Block a user