package service import ( "context" "fmt" "net/url" "strconv" "strings" "github.com/krow/krow-backend/go-api/internal/authctx" "github.com/krow/krow-backend/go-api/internal/definition" "github.com/krow/krow-backend/go-api/internal/domain" "github.com/krow/krow-backend/go-api/internal/repo" ) // allowedDefinitionFilters names the accepted query parameters for definition collections. var allowedDefinitionFilters = map[string]bool{ "visibility": true, "status": true, "definition_id": true, "sort": true, "limit": true, "offset": true, } // DefinitionsService manages authored Agent and Skill definitions. type DefinitionsService struct { repo *repo.DefinitionsRepo // versions is the append-only history beside the editable definition. // Written on publish, never on save — see snapshotIfPublished. versions *repo.VersionsRepo // unknownTools reports which of a spec's tool names are not registered. // nil means no check — the shipped importer path, which has its own. unknownTools func([]string) []string } // NewDefinitions builds a definitions service over a repository. func NewDefinitions(db repo.Querier) *DefinitionsService { return &DefinitionsService{ repo: repo.NewDefinitionsRepo(db), versions: repo.NewVersionsRepo(db), } } // WithToolCheck teaches the service which tool names exist. // // §3 requires an unknown tool name to fail validation at PUBLISH. Without it // the runtime records the name and drops it, so a typo produces an agent that // is silently missing a capability its author believes it has — and the author // finds out by watching it fail to answer. // // Injected rather than imported so this package does not depend on the tool // registry, and so a test can supply its own vocabulary. func (s *DefinitionsService) WithToolCheck(unknown func([]string) []string) *DefinitionsService { s.unknownTools = unknown return s } // rejectUnknownTools fails a definition that names a tool that does not exist. func (s *DefinitionsService) rejectUnknownTools(markdown string) error { if s.unknownTools == nil { return nil } agent, err := definition.ParseAgent(markdown, definition.Options{}) if err != nil || agent == nil { // ValidateAgent has already run and reported anything real; a parse // failure here is not a second opinion worth raising. return nil } if bad := s.unknownTools(agent.Tools); len(bad) > 0 { return domain.Validation( fmt.Sprintf("unknown tool(s): %s", strings.Join(bad, ", ")), map[string]string{"tools": strings.Join(bad, ", ")}) } return nil } // ParseListParams validates query parameters for listing definitions. func (s *DefinitionsService) ParseListParams(q url.Values) (repo.DefinitionListParams, error) { p := repo.DefinitionListParams{ Limit: 100, Sort: "created_date", Desc: true, } for name := range q { if !allowedDefinitionFilters[name] { return p, domain.Invalid(fmt.Sprintf("unknown filter field %q", name)) } } if raw := q.Get("visibility"); raw != "" { if raw != "personal" && raw != "organization" { return p, domain.Invalid("visibility must be one of: personal, organization") } p.Visibility = raw } if raw := q.Get("status"); raw != "" { p.Status = raw } if raw := q.Get("definition_id"); raw != "" { p.DefinitionID = raw } if raw := q.Get("sort"); raw != "" { field := raw desc := false if strings.HasPrefix(field, "-") { desc = true field = field[1:] } switch field { case "created_date", "updated_date", "name", "definition_id", "status", "version": p.Sort = field p.Desc = desc default: return p, domain.Invalid(fmt.Sprintf("cannot sort by %q", field)) } } if raw := q.Get("limit"); raw != "" { n, err := strconv.Atoi(raw) if err != nil || n < 0 { return p, domain.Invalid("limit must be a non-negative integer") } if n > MaxLimit { n = MaxLimit } p.Limit = n } if raw := q.Get("offset"); raw != "" { n, err := strconv.Atoi(raw) if err != nil || n < 0 { return p, domain.Invalid("offset must be a non-negative integer") } p.Offset = n } return p, nil } /* ── Agents ─────────────────────────────────────────────────────────────── */ // ListAgents returns a page of authored agent definitions. func (s *DefinitionsService) ListAgents(ctx context.Context, ident authctx.Identity, p repo.DefinitionListParams) (*domain.Page, error) { records, total, err := s.repo.ListAgents(ctx, ident, p) if err != nil { return nil, err } if records == nil { records = []domain.Record{} } return &domain.Page{ Records: records, Total: total, Limit: p.Limit, Offset: p.Offset, }, nil } // GetAgent returns one agent definition by id within the caller's tenant and ownership scope. func (s *DefinitionsService) GetAgent(ctx context.Context, ident authctx.Identity, id string) (domain.Record, error) { if !isUUID(id) { return nil, domain.NotFound("AgentDefinition", id) } rec, err := s.repo.GetAgent(ctx, ident, id) if err != nil { return nil, err } if rec == nil { return nil, domain.NotFound("AgentDefinition", id) } return rec, nil } // CreateAgent validates, parses and persists a new authored agent definition. func (s *DefinitionsService) CreateAgent(ctx context.Context, ident authctx.Identity, body domain.Record) (domain.Record, error) { mdRaw, ok := body["markdown"] if !ok || mdRaw == nil { return nil, domain.Validation("Paste or upload a Markdown definition.", nil) } markdown, isStr := mdRaw.(string) if !isStr { return nil, domain.Validation("markdown must be a string", nil) } if err := definition.ValidateAgent(markdown); err != nil { return nil, domain.Validation(err.Error(), nil) } if err := s.rejectUnknownTools(markdown); err != nil { return nil, err } visibility := "personal" if visRaw, ok := body["visibility"]; ok && visRaw != nil { v, isStr := visRaw.(string) if !isStr || (v != "personal" && v != "organization") { return nil, domain.Validation("visibility must be one of: personal, organization", map[string]string{"visibility": "invalid"}) } visibility = v } if visibility == "organization" { role, known := domain.ParseRole(ident.Role) if !known || role == domain.RoleTalent { return nil, domain.Forbidden() } } agent, err := definition.ParseAgent(markdown, definition.Options{}) if err != nil { return nil, domain.Validation("That definition could not be parsed. "+err.Error(), nil) } input := repo.AgentInsertInput{ DefinitionID: agent.ID, OrgID: ident.OrgID, Visibility: visibility, CreatedBy: &ident.UserID, Markdown: markdown, Status: agent.Status, Version: agent.Version, Name: agent.Name, Description: agent.Description, Pages: agent.Pages, } if visibility == "personal" { 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 } rec, err := s.repo.InsertAgent(ctx, ident, input) if err != nil { return nil, err } // Recorded after the row exists, and never allowed to fail the save: the // author's work is already stored, and losing it to protect a record of it // would be the wrong trade. _ = s.snapshotIfPublished(ctx, ident, repo.KindAgent, rec) return rec, nil } // UpdateAgent validates and applies updates to an authored agent definition. func (s *DefinitionsService) UpdateAgent(ctx context.Context, ident authctx.Identity, id string, patch domain.Record) (domain.Record, error) { if !isUUID(id) { return nil, domain.NotFound("AgentDefinition", id) } existing, err := s.repo.GetAgent(ctx, ident, id) if err != nil { return nil, err } if existing == nil { return nil, domain.NotFound("AgentDefinition", id) } if existing["visibility"] == "organization" { role, known := domain.ParseRole(ident.Role) if !known || role == domain.RoleTalent { return nil, domain.Forbidden() } } if visRaw, ok := patch["visibility"]; ok && visRaw != nil { if v, isStr := visRaw.(string); isStr && v != existing["visibility"] { return nil, domain.Validation("visibility cannot be modified after creation", map[string]string{"visibility": "immutable"}) } } var input repo.AgentUpdateInput if mdRaw, ok := patch["markdown"]; ok && mdRaw != nil { markdown, isStr := mdRaw.(string) if !isStr { return nil, domain.Validation("markdown must be a string", nil) } if err := definition.ValidateAgent(markdown); err != nil { return nil, domain.Validation(err.Error(), nil) } if err := s.rejectUnknownTools(markdown); err != nil { return nil, err } agent, err := definition.ParseAgent(markdown, definition.Options{}) if err != nil { return nil, domain.Validation("That definition could not be parsed. "+err.Error(), nil) } input.Markdown = &markdown input.DefinitionID = &agent.ID input.Name = &agent.Name input.Description = &agent.Description input.Status = &agent.Status 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 } } else if statusRaw, ok := patch["status"]; ok && statusRaw != nil { status, isStr := statusRaw.(string) if !isStr || (status != "draft" && status != "published" && status != "archived") { return nil, domain.Validation("status must be one of: draft, published, archived", map[string]string{"status": "invalid"}) } input.Status = &status } rec, err := s.repo.UpdateAgent(ctx, ident, id, input) if err != nil { return nil, err } _ = s.snapshotIfPublished(ctx, ident, repo.KindAgent, rec) return rec, nil } // DeleteAgent removes an agent definition following idempotent delete semantics. func (s *DefinitionsService) DeleteAgent(ctx context.Context, ident authctx.Identity, id string) (domain.Record, error) { if !isUUID(id) { return domain.Record{"id": id}, nil } existing, err := s.repo.GetAgent(ctx, ident, id) if err != nil { return nil, err } if existing == nil { return domain.Record{"id": id}, nil } if existing["visibility"] == "organization" { role, known := domain.ParseRole(ident.Role) if !known || role == domain.RoleTalent { return nil, domain.Forbidden() } } if _, err := s.repo.DeleteAgent(ctx, ident, id); err != nil { return nil, err } return domain.Record{"id": id}, nil } /* ── Skills ─────────────────────────────────────────────────────────────── */ // ListSkills returns a page of authored skill definitions. func (s *DefinitionsService) ListSkills(ctx context.Context, ident authctx.Identity, p repo.DefinitionListParams) (*domain.Page, error) { records, total, err := s.repo.ListSkills(ctx, ident, p) if err != nil { return nil, err } if records == nil { records = []domain.Record{} } return &domain.Page{ Records: records, Total: total, Limit: p.Limit, Offset: p.Offset, }, nil } // GetSkill returns one skill definition by id within the caller's tenant and ownership scope. func (s *DefinitionsService) GetSkill(ctx context.Context, ident authctx.Identity, id string) (domain.Record, error) { if !isUUID(id) { return nil, domain.NotFound("SkillDefinition", id) } rec, err := s.repo.GetSkill(ctx, ident, id) if err != nil { return nil, err } if rec == nil { return nil, domain.NotFound("SkillDefinition", id) } return rec, nil } // CreateSkill validates, parses and persists a new authored skill definition. func (s *DefinitionsService) CreateSkill(ctx context.Context, ident authctx.Identity, body domain.Record) (domain.Record, error) { mdRaw, ok := body["markdown"] if !ok || mdRaw == nil { return nil, domain.Validation("Paste or upload a Markdown definition.", nil) } markdown, isStr := mdRaw.(string) if !isStr { return nil, domain.Validation("markdown must be a string", nil) } if err := definition.ValidateSkill(markdown); err != nil { return nil, domain.Validation(err.Error(), nil) } visibility := "personal" if visRaw, ok := body["visibility"]; ok && visRaw != nil { v, isStr := visRaw.(string) if !isStr || (v != "personal" && v != "organization") { return nil, domain.Validation("visibility must be one of: personal, organization", map[string]string{"visibility": "invalid"}) } visibility = v } if visibility == "organization" { role, known := domain.ParseRole(ident.Role) if !known || role == domain.RoleTalent { return nil, domain.Forbidden() } } skill, err := definition.ParseSkill(markdown, definition.Options{}) if err != nil { return nil, domain.Validation("That definition could not be parsed. "+err.Error(), nil) } input := repo.SkillInsertInput{ DefinitionID: skill.ID, OrgID: ident.OrgID, Visibility: visibility, CreatedBy: &ident.UserID, Markdown: markdown, Status: skill.Status, Name: skill.Name, Description: skill.Description, Pages: skill.Pages, } if visibility == "personal" { input.OwnerUserID = &ident.UserID } rec, err := s.repo.InsertSkill(ctx, ident, input) if err != nil { return nil, err } _ = s.snapshotSkill(ctx, ident, rec) return rec, nil } // UpdateSkill validates and applies updates to an authored skill definition. func (s *DefinitionsService) UpdateSkill(ctx context.Context, ident authctx.Identity, id string, patch domain.Record) (domain.Record, error) { if !isUUID(id) { return nil, domain.NotFound("SkillDefinition", id) } existing, err := s.repo.GetSkill(ctx, ident, id) if err != nil { return nil, err } if existing == nil { return nil, domain.NotFound("SkillDefinition", id) } if existing["visibility"] == "organization" { role, known := domain.ParseRole(ident.Role) if !known || role == domain.RoleTalent { return nil, domain.Forbidden() } } if visRaw, ok := patch["visibility"]; ok && visRaw != nil { if v, isStr := visRaw.(string); isStr && v != existing["visibility"] { return nil, domain.Validation("visibility cannot be modified after creation", map[string]string{"visibility": "immutable"}) } } var input repo.SkillUpdateInput if mdRaw, ok := patch["markdown"]; ok && mdRaw != nil { markdown, isStr := mdRaw.(string) if !isStr { return nil, domain.Validation("markdown must be a string", nil) } if err := definition.ValidateSkill(markdown); err != nil { return nil, domain.Validation(err.Error(), nil) } skill, err := definition.ParseSkill(markdown, definition.Options{}) if err != nil { return nil, domain.Validation("That definition could not be parsed. "+err.Error(), nil) } input.Markdown = &markdown input.DefinitionID = &skill.ID input.Name = &skill.Name input.Description = &skill.Description input.Status = &skill.Status input.Pages = skill.Pages } else if statusRaw, ok := patch["status"]; ok && statusRaw != nil { status, isStr := statusRaw.(string) if !isStr || (status != "active" && status != "inactive") { return nil, domain.Validation("status must be one of: active, inactive", map[string]string{"status": "invalid"}) } input.Status = &status } rec, err := s.repo.UpdateSkill(ctx, ident, id, input) if err != nil { return nil, err } _ = s.snapshotSkill(ctx, ident, rec) return rec, nil } // DeleteSkill removes a skill definition following idempotent delete semantics. func (s *DefinitionsService) DeleteSkill(ctx context.Context, ident authctx.Identity, id string) (domain.Record, error) { if !isUUID(id) { return domain.Record{"id": id}, nil } existing, err := s.repo.GetSkill(ctx, ident, id) if err != nil { return nil, err } if existing == nil { return domain.Record{"id": id}, nil } if existing["visibility"] == "organization" { role, known := domain.ParseRole(ident.Role) if !known || role == domain.RoleTalent { return nil, domain.Forbidden() } } if _, err := s.repo.DeleteSkill(ctx, ident, id); err != nil { return nil, err } return domain.Record{"id": id}, nil } /* ── 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. // // snapshotIfPublished below deliberately never fails a save: the author's work // is already stored and losing it to protect a record of it is the wrong trade. // That is right for a recording failure — the disk, the pool, the network — and // wrong for exactly one case. When repo.VersionsRepo.Snapshot refuses because // the version already says something different, that is not the history failing // to record; it is §3 firing. Swallowing it leaves two different definitions // both called v2: the live row the runtime serves, and the snapshot the history // shows. runtime.LoadAgentVersion resolves a pin by returning the CURRENT // definition whenever the pinned number equals the current one, so the run gets // the changed text while the audit trail says otherwise. // // So the conflict is detected here instead, before the write, where refusing // costs the author nothing but a version bump. The post-write snapshot keeps // its original contract for every other kind of failure. // // A concurrent publish of the same number with different content can still slip // past this check and be caught by the unique index afterwards, where it is // swallowed as before. That leaves the live row ahead of its snapshot, which is // the pre-existing behaviour and not something this guard makes worse. func (s *DefinitionsService) refusePublishedRewrite(ctx context.Context, ident authctx.Identity, kind repo.VersionKind, definitionID, markdown, status string, version int) error { if s.versions == nil || status != "published" || definitionID == "" || version < 1 { return nil } stored, err := s.versions.Load(ctx, ident, kind, definitionID, version) if err != nil || stored == nil { // Absent (the ordinary case for a new version) or unreadable. Either // way there is no published text to contradict, so this is not the // place to fail the save. return nil } // Semantic, not textual. The authoring UI re-serialises a definition when // it is saved, so a byte comparison refuses a publish over frontmatter key // order and a defaulted value written out in full — see // definition.SameAgent, which is deliberately conservative about what it // treats as inert. if definition.SameAgent(stored.Markdown, markdown) { return nil // republishing the same version unchanged is a no-op } return domain.Conflict(fmt.Sprintf( "version %d of %q is already published and says something different; "+ "raise the version in the frontmatter to publish a change", version, definitionID)) } // snapshotIfPublished records an immutable copy when a definition is published. // // §3: editing publishes a NEW version, and a published version never changes. // This is the half that records it. The half that enforces it is a trigger on // the table, because the repository is not the only thing that can reach it. // // ONLY ON PUBLISH. A draft is a work in progress and snapshotting every save // would fill the history with keystrokes — the version number would stop // meaning "a thing somebody decided to ship" and start meaning "a time somebody // pressed save", which is the number a run records and a person has to // recognise. // // A failure here does NOT fail the save. The definition is already written; // refusing the whole operation because its history could not be recorded would // lose the author's work to protect a record of it. It is returned so the // caller can log it, and the missing version shows up as a gap rather than as // a wrong answer. func (s *DefinitionsService) snapshotIfPublished(ctx context.Context, ident authctx.Identity, kind repo.VersionKind, rec domain.Record) error { if s.versions == nil || rec == nil { return nil } status, _ := rec["status"].(string) if status != "published" { return nil } markdown, _ := rec["markdown"].(string) definitionID, _ := rec["definition_id"].(string) if markdown == "" || definitionID == "" { return nil } version := 1 switch v := rec["version"].(type) { case int: version = v case int32: version = int(v) case int64: version = int(v) case float64: version = int(v) } name, _ := rec["name"].(string) description, _ := rec["description"].(string) return s.versions.Snapshot(ctx, ident, repo.SnapshotInput{ Kind: kind, DefinitionID: definitionID, Version: version, Markdown: markdown, Name: name, Description: description, Pages: recordPages(rec), }) } // recordPages reads a record's pages, which arrive as []string from the // repository and as []any when they have been through JSON. func recordPages(rec domain.Record) []string { if raw, ok := rec["pages"].([]string); ok { return raw } raw, ok := rec["pages"].([]any) if !ok { return nil } pages := make([]string, 0, len(raw)) for _, p := range raw { if str, isStr := p.(string); isStr { pages = append(pages, str) } } return pages } // snapshotSkill records an immutable copy of a skill, numbered by the server. // // Skills carry no version. An agent's frontmatter names one, so its author // decides when a change is a new version and can be refused for rewriting an // old one. A skill has no such field, and giving it one would mean a migration, // a parser change on BOTH sides of the conformance test in // internal/definition, and an edit to all 23 shipped skills — a feature, not // the fix this is. // // So the number is the server's: one after whatever was last published. That // is what repo.VersionsRepo.LatestVersion was written for ("the next published // version has to follow what was actually published rather than what somebody // wrote in the frontmatter") and what migration 000010 means by "agents and // skills version identically". It was built and never wired to anything. // // Because the author never names a version, there is nothing here to refuse: // an edit is always a NEW version, and a save that changed nothing is not a // version at all. The comparison against the last published copy is what keeps // the history from filling with keystrokes. // // Two simultaneous edits can both compute the same next number; one wins and // the other's snapshot is dropped, leaving a version unrecorded. That is the // same narrow race the agent path has, and the same reason it is tolerated // here: failing an author's save to record it is the wrong trade. func (s *DefinitionsService) snapshotSkill(ctx context.Context, ident authctx.Identity, rec domain.Record) error { if s.versions == nil || rec == nil { return nil } // Only a skill that is in service. "inactive" is the skill vocabulary's // equivalent of a draft — see definition.SkillStatuses. if status, _ := rec["status"].(string); status != "active" { return nil } markdown, _ := rec["markdown"].(string) definitionID, _ := rec["definition_id"].(string) if markdown == "" || definitionID == "" { return nil } latest, err := s.versions.LatestVersion(ctx, ident, repo.KindSkill, definitionID) if err != nil { return err } if latest > 0 { stored, err := s.versions.Load(ctx, ident, repo.KindSkill, definitionID, latest) if err == nil && stored != nil && definition.SameSkill(stored.Markdown, markdown) { return nil // unchanged since the last publish } } name, _ := rec["name"].(string) description, _ := rec["description"].(string) return s.versions.Snapshot(ctx, ident, repo.SnapshotInput{ Kind: repo.KindSkill, DefinitionID: definitionID, Version: latest + 1, Markdown: markdown, Name: name, Description: description, Pages: recordPages(rec), }) } // AgentHistory lists an agent's published versions, newest first. func (s *DefinitionsService) AgentHistory(ctx context.Context, ident authctx.Identity, definitionID string, limit int) ([]repo.Version, error) { return s.versions.History(ctx, ident, repo.KindAgent, definitionID, limit) } // AgentVersion loads one published version of an agent, as it was. func (s *DefinitionsService) AgentVersion(ctx context.Context, ident authctx.Identity, definitionID string, version int) (*repo.Version, error) { return s.versions.Load(ctx, ident, repo.KindAgent, definitionID, version) }