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
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
// Parsed before anything is opened, so a malformed spec is a message rather
|
if err := validateSpecs(specs); err != nil {
|
||||||
// than a half-finished import. Every spec, not the first failure: an
|
return err
|
||||||
// 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"))
|
|
||||||
}
|
}
|
||||||
|
|
||||||
for _, s := range specs {
|
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,
|
// written. The runtime refuses to load an agent with a missing dependency,
|
||||||
// so importing one without its skills produces an agent that exists and
|
// so importing one without its skills produces an agent that exists and
|
||||||
// cannot run — a failure that surfaces per request instead of here.
|
// 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 {
|
if missing := missingSkills(specs, skills); len(missing) > 0 {
|
||||||
return fmt.Errorf("%d skill(s) named by an agent are not in %s: %s",
|
return fmt.Errorf("%d skill(s) named by an agent are not in %s: %s",
|
||||||
len(missing), skillDir, strings.Join(missing, ", "))
|
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
|
defer tx.Rollback(ctx) //nolint:errcheck // rolled back unless committed below
|
||||||
|
|
||||||
// Skills first. An agent row that lands before its dependencies exist is
|
out, err := importInto(ctx, tx, orgID, author, specs, skills)
|
||||||
// briefly unloadable, and inside one transaction that is invisible — but
|
if err != nil {
|
||||||
// ordering them correctly costs nothing and means a future non-transactional
|
return err
|
||||||
// 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"))
|
|
||||||
}
|
}
|
||||||
|
|
||||||
if err := tx.Commit(ctx); err != nil {
|
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, "+
|
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",
|
"%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
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -494,6 +555,19 @@ func snapshotSkill(ctx context.Context, versions *repo.VersionsRepo,
|
|||||||
func snapshotAgent(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) {
|
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)
|
stored, err := versions.Load(ctx, ident, repo.KindAgent, s.parsed.ID, s.parsed.Version)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
var apiErr *domain.Error
|
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)
|
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
|
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,
|
if err := s.refusePublishedRewrite(ctx, ident, repo.KindAgent,
|
||||||
agent.ID, markdown, agent.Status, agent.Version); err != nil {
|
agent.ID, markdown, agent.Status, agent.Version); err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
@@ -302,6 +309,13 @@ func (s *DefinitionsService) UpdateAgent(ctx context.Context, ident authctx.Iden
|
|||||||
input.Version = &agent.Version
|
input.Version = &agent.Version
|
||||||
input.Pages = agent.Pages
|
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,
|
if err := s.refusePublishedRewrite(ctx, ident, repo.KindAgent,
|
||||||
agent.ID, markdown, agent.Status, agent.Version); err != nil {
|
agent.ID, markdown, agent.Status, agent.Version); err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
@@ -535,6 +549,80 @@ func (s *DefinitionsService) DeleteSkill(ctx context.Context, ident authctx.Iden
|
|||||||
|
|
||||||
/* ── Publishing ─────────────────────────────────────────────────────────── */
|
/* ── 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
|
// refusePublishedRewrite fails a publish that would change a version already
|
||||||
// published, BEFORE anything is written.
|
// published, BEFORE anything is written.
|
||||||
//
|
//
|
||||||
|
|||||||
Reference in New Issue
Block a user