From f9d8f1a1c97ce88b25cd84dc17e310897f427e50 Mon Sep 17 00:00:00 2001 From: Algis Dumbris Date: Tue, 14 Apr 2026 20:27:25 +0300 Subject: [PATCH] feat(018): /goals page, budget cascade, quarantine, secrets loop - New /api/goals + /api/goals/{id} endpoints serving list + task tree + cost rollup + billing breakdown + spawned agents + timeline. - New Svelte /goals and /goals/[id] pages with sidebar link. - goals.Service.EvaluateBudget returns a soft/hard verdict; agent runner posts the 80% warning once and auto-pauses at 100%. - Auto-quarantine: after each reputation append the agent runner checks rolling score < 0.3 and writes quarantined_at; reactor refuses new reactive dispatches to quarantined agents. - Reactor exposes SetSecretProvider; main.go wires secrets.Store so reactive subprocess runs inherit user/agent-scoped env vars. - cli-verifier demonstrates the resource-request protocol: checks MCPPROXY_API_KEY, posts to #requests + resource_requests row if missing. New `synapbus secrets set/list` CLI (direct-DB) closes the loop. Co-Authored-By: Claude Opus 4.6 (1M context) --- .gitignore | 1 + cmd/docgardener/agent.go | 96 ++++++++ cmd/docgardener/flow.go | 12 +- cmd/synapbus/goals_adapter.go | 36 +++ cmd/synapbus/main.go | 26 ++ cmd/synapbus/secrets_cli.go | 173 ++++++++++++++ internal/api/goals_handler.go | 317 +++++++++++++++++++++++++ internal/api/router.go | 15 ++ internal/goals/service.go | 38 +++ internal/reactor/reactor.go | 41 ++++ internal/web/dist/index.html | 22 +- web/src/lib/api/client.ts | 10 + web/src/lib/components/Sidebar.svelte | 1 + web/src/routes/goals/+page.svelte | 135 +++++++++++ web/src/routes/goals/[id]/+page.svelte | 225 ++++++++++++++++++ 15 files changed, 1136 insertions(+), 12 deletions(-) create mode 100644 cmd/synapbus/goals_adapter.go create mode 100644 cmd/synapbus/secrets_cli.go create mode 100644 internal/api/goals_handler.go create mode 100644 web/src/routes/goals/+page.svelte create mode 100644 web/src/routes/goals/[id]/+page.svelte diff --git a/.gitignore b/.gitignore index a40e593..7a70434 100644 --- a/.gitignore +++ b/.gitignore @@ -45,6 +45,7 @@ __pycache__/ benchmark/data/ benchmark/results/ .venv-bench/ +.venv-kimi/ .venv/ # Debug diff --git a/cmd/docgardener/agent.go b/cmd/docgardener/agent.go index 7425e73..2a2c515 100644 --- a/cmd/docgardener/agent.go +++ b/cmd/docgardener/agent.go @@ -393,6 +393,29 @@ func (a *agentRunner) handleSpecialist(ctx context.Context) (string, error) { tokensOut := int64(400 + 100*(taskID%3)) costCents := int64(25 + 10*(taskID%3)) _ = a.tasks.AddSpend(ctx, taskID, tokensIn+tokensOut, costCents) + + // 4a. Budget cascade — re-roll up the goal's total cents and + // check the thresholds. On first crossing of 80% we post a + // warning; at 100% the goal is auto-paused and new claims + // would bounce. + if g.RootTaskID != nil { + _, rollupCents, _, _ := a.tasks.RollupCosts(ctx, *g.RootTaskID) + verdict, err := a.goals.EvaluateBudget(ctx, g.ID, rollupCents) + if err == nil && verdict != nil { + if verdict.TriggerSoftAlert { + a.postSystemMessage(ctx, g.ChannelID, + fmt.Sprintf("⚠️ Budget soft alert: goal has consumed %.0f%% of its dollar budget.", + verdict.PercentBudget)) + _ = a.goals.MarkSoftAlertPosted(ctx, g.ID) + } + if verdict.TriggerHardPause { + a.postSystemMessage(ctx, g.ChannelID, + fmt.Sprintf("🛑 Budget hard cap: goal at %.0f%% → auto-paused.", verdict.PercentBudget)) + _ = a.goals.TransitionStatus(ctx, g.ID, goals.StatusPaused) + } + } + } + _ = a.tasks.Transition(ctx, taskID, goaltasks.StatusAwaitingVerification, goaltasks.Extras{CompletionMessageID: &artifactMsgID}) extras := goaltasks.Extras{} @@ -412,12 +435,32 @@ func (a *agentRunner) handleSpecialist(ctx context.Context) (string, error) { a.logger.Warn("append evidence failed", "err", err) } + // 5a. Quarantine check — if the rolling reputation dropped below + // 0.3 after this evidence, flag the agent as quarantined so + // future reactive runs refuse to spawn. + if score, _, err := a.ledger.RollingScore(ctx, hash, "default", 30); err == nil && score < 0.3 { + _, _ = a.db.ExecContext(ctx, + `UPDATE agents SET quarantined_at = ?, quarantine_reason = ? WHERE id = ? AND quarantined_at IS NULL`, + time.Now().UTC(), fmt.Sprintf("reputation=%.2f", score), specialistID) + a.postSystemMessage(ctx, g.ChannelID, + fmt.Sprintf("⛔ Agent %s quarantined — reputation %.2f below 0.3.", a.agentName, score)) + } + // 6. Post a system summary line with real telemetry. a.postSystemMessage(ctx, g.ChannelID, fmt.Sprintf("Task %d %q %s by %s — tokens_in=%d tokens_out=%d cost=$%.2f duration=%dms Δrep=%+.2f", taskID, t.Title, verdict, role, tokensIn, tokensOut, float64(costCents)/100, duration.Milliseconds(), scoreDelta)) + // 6a. Resource-request protocol demo: the cli-verifier needs + // MCPPROXY_API_KEY. If the injected env doesn't carry it, it + // posts a structured request to the #requests channel so a + // human can set it via `synapbus secrets set`. + if role == "cli-verifier" && os.Getenv("MCPPROXY_API_KEY") == "" { + a.postResourceRequest(ctx, role, taskID, "MCPPROXY_API_KEY", + "Need the mcpproxy admin API key to re-run live CLI verification against a remote proxy; set it with `synapbus secrets set MCPPROXY_API_KEY --scope agent:cli-verifier`.") + } + // 7. DM coordinator with DONE or FAIL. reply := fmt.Sprintf("%s task=%d role=%s tokens_in=%d tokens_out=%d cost_cents=%d duration_ms=%d", strings.ToUpper(string(verdict)), taskID, role, tokensIn, tokensOut, costCents, duration.Milliseconds()) @@ -562,6 +605,59 @@ func (a *agentRunner) spawnSpecialist(ctx context.Context, ownerID, parentID int return id, nil } +// postResourceRequest writes a structured resource_requests row AND +// posts a #requests channel message describing the missing secret. +// The human reads it, runs `synapbus secrets set` to provision it, +// and the next reactive run picks up the injected env var. +func (a *agentRunner) postResourceRequest(ctx context.Context, role string, taskID int64, resourceName, reason string) { + // Ensure #requests channel exists. Admin CLI creates it in + // start.sh; we re-check defensively here and create if missing. + var reqChannelID int64 + err := a.db.QueryRowContext(ctx, + `SELECT id FROM channels WHERE name='requests' LIMIT 1`).Scan(&reqChannelID) + if err != nil { + // Channel missing — create it inline (no CreatedBy enforcement in the demo). + res, cerr := a.db.ExecContext(ctx, + `INSERT INTO channels (name, description, type, created_by, is_private, is_system) + VALUES ('requests','Resource requests','blackboard', ?, 0, 1)`, + a.agentName) + if cerr != nil { + a.logger.Warn("could not create #requests channel", "err", cerr) + return + } + reqChannelID, _ = res.LastInsertId() + } + + body := fmt.Sprintf("#resource-request agent=%s task=%d resource=%s type=env_var\nreason: %s", + a.agentName, taskID, resourceName, reason) + + // Insert the message directly (the #requests channel is not + // reactive so bypassing the reactor dispatcher is fine here). + convID, cerr := a.ensureConversation(ctx, reqChannelID) + if cerr != nil { + a.logger.Warn("could not ensure #requests conversation", "err", cerr) + return + } + now := time.Now().UTC() + _, err = a.db.ExecContext(ctx, ` + INSERT INTO messages (conversation_id, from_agent, to_agent, channel_id, body, priority, status, metadata, created_at, updated_at) + VALUES (?, ?, NULL, ?, ?, 7, 'done', '{"kind":"resource-request"}', ?, ?)`, + convID, a.agentName, reqChannelID, body, now, now) + if err != nil { + a.logger.Warn("could not post resource-request message", "err", err) + return + } + + // Also write a resource_requests row (feature 018) so the /goals + // page + /api could display it later. + _, _ = a.db.ExecContext(ctx, ` + INSERT INTO resource_requests (requester_agent_id, task_id, resource_name, resource_type, reason, status) + SELECT id, ?, ?, 'env_var', ?, 'pending' FROM agents WHERE name=?`, + taskID, resourceName, reason, a.agentName) + + a.logger.Info("resource request posted", "resource", resourceName, "task_id", taskID) +} + // postRealArtifactDirect is a copy of flow.go's postRealArtifact that // does not rely on a shared conversation — it creates a fresh // conversation scoped to this single message write if one doesn't diff --git a/cmd/docgardener/flow.go b/cmd/docgardener/flow.go index ac9061f..1e4fc02 100644 --- a/cmd/docgardener/flow.go +++ b/cmd/docgardener/flow.go @@ -493,11 +493,21 @@ flags=--port --config --socket --data-dir --log-format --log-level --otel-endpoi timestamp=$(date -u +%Y-%m-%dT%H:%M:%SZ) EOF` case "cli-verifier": - return `cat <, agent:, task: + if !strings.Contains(s, ":") { + return "", 0, fmt.Errorf("scope must be user:NAME, agent:NAME, or task:ID") + } + typ, ident, _ := strings.Cut(s, ":") + db, _, err := openDirect() + if err != nil { + return "", 0, err + } + defer db.Close() + switch typ { + case "user": + var id int64 + err := db.QueryRowContext(context.Background(), `SELECT id FROM users WHERE username=?`, ident).Scan(&id) + if err != nil { + return "", 0, fmt.Errorf("user %q not found: %w", ident, err) + } + return secrets.ScopeUser, id, nil + case "agent": + var id int64 + err := db.QueryRowContext(context.Background(), `SELECT id FROM agents WHERE name=?`, ident).Scan(&id) + if err != nil { + return "", 0, fmt.Errorf("agent %q not found: %w", ident, err) + } + return secrets.ScopeAgent, id, nil + case "task": + id, err := strconv.ParseInt(ident, 10, 64) + if err != nil { + return "", 0, fmt.Errorf("task scope id must be an integer") + } + return secrets.ScopeTask, id, nil + } + return "", 0, fmt.Errorf("unknown scope type %q", typ) + } + + root := &cobra.Command{ + Use: "secrets", + Short: "Manage encrypted scoped secrets (resource-request protocol)", + } + root.PersistentFlags().StringVar(&dbPath, "db", "", "Path to synapbus.db (defaults to ./data or SYNAPBUS_DATA_DIR)") + + setCmd := &cobra.Command{ + Use: "set NAME VALUE", + Short: "Store a secret under a scope", + Args: cobra.ExactArgs(2), + RunE: func(cmd *cobra.Command, args []string) error { + name, value := args[0], args[1] + scopeType, scopeID, err := parseScope(scope) + if err != nil { + return err + } + db, dataDir, err := openDirect() + if err != nil { + return err + } + defer db.Close() + store, err := secrets.NewStore(db, dataDir, slog.Default()) + if err != nil { + return err + } + s, err := store.Set(cmd.Context(), name, scopeType, scopeID, 0, value) + if err != nil { + return err + } + fmt.Printf("stored secret id=%d name=%s scope=%s:%d\n", s.ID, s.Name, scopeType, scopeID) + return nil + }, + } + setCmd.Flags().StringVar(&scope, "scope", "", "Scope (user:NAME, agent:NAME, task:ID)") + _ = setCmd.MarkFlagRequired("scope") + + listCmd := &cobra.Command{ + Use: "list", + Short: "List secrets visible to a scope (names only — never values)", + RunE: func(cmd *cobra.Command, args []string) error { + scopeType, scopeID, err := parseScope(scope) + if err != nil { + return err + } + db, dataDir, err := openDirect() + if err != nil { + return err + } + defer db.Close() + store, err := secrets.NewStore(db, dataDir, slog.Default()) + if err != nil { + return err + } + infos, err := store.List(cmd.Context(), []secrets.Scope{{Type: scopeType, ID: scopeID}}) + if err != nil { + return err + } + tw := tabwriter.NewWriter(os.Stdout, 0, 0, 2, ' ', 0) + fmt.Fprintln(tw, "NAME\tSCOPE\tLAST USED") + for _, i := range infos { + last := "—" + if i.LastUsedAt != nil { + last = i.LastUsedAt.Format("2006-01-02 15:04") + } + fmt.Fprintf(tw, "%s\t%s:%d\t%s\n", i.Name, i.ScopeType, i.ScopeID, last) + } + return tw.Flush() + }, + } + listCmd.Flags().StringVar(&scope, "scope", "", "Scope (user:NAME, agent:NAME, task:ID)") + _ = listCmd.MarkFlagRequired("scope") + + root.AddCommand(setCmd, listCmd) + return root +} diff --git a/internal/api/goals_handler.go b/internal/api/goals_handler.go new file mode 100644 index 0000000..a828215 --- /dev/null +++ b/internal/api/goals_handler.go @@ -0,0 +1,317 @@ +package api + +import ( + "database/sql" + "encoding/json" + "net/http" + "strconv" + + "github.com/go-chi/chi/v5" + + "github.com/synapbus/synapbus/internal/goals" + "github.com/synapbus/synapbus/internal/goaltasks" +) + +// GoalsHandler serves the /api/goals endpoints used by the Web UI /goals +// page: list goals, show a single goal's full task tree with cost +// rollup, billing-code breakdown, and the spawned agents attached. +type GoalsHandler struct { + goals *goals.Service + tasks *goaltasks.Service + db *sql.DB +} + +func NewGoalsHandler(g *goals.Service, t *goaltasks.Service, db *sql.DB) *GoalsHandler { + return &GoalsHandler{goals: g, tasks: t, db: db} +} + +// ListGoals returns recent goals with basic metadata + total spend. +func (h *GoalsHandler) ListGoals(w http.ResponseWriter, r *http.Request) { + limit := 50 + if l := r.URL.Query().Get("limit"); l != "" { + if v, err := strconv.Atoi(l); err == nil && v > 0 && v <= 200 { + limit = v + } + } + + gs, err := h.goals.ListGoals(r.Context(), nil, limit) + if err != nil { + writeJSON(w, http.StatusInternalServerError, errorBody("internal_error", err.Error())) + return + } + + type goalSummary struct { + ID int64 `json:"id"` + Slug string `json:"slug"` + Title string `json:"title"` + Status string `json:"status"` + ChannelID int64 `json:"channel_id"` + OwnerUsername string `json:"owner_username"` + RootTaskID *int64 `json:"root_task_id"` + SpentTokens int64 `json:"spent_tokens"` + SpentDollarsCents int64 `json:"spent_dollars_cents"` + TaskCount int `json:"task_count"` + BudgetTokens *int64 `json:"budget_tokens"` + BudgetDollarsCents *int64 `json:"budget_dollars_cents"` + PercentBudget float64 `json:"percent_budget"` + CreatedAt string `json:"created_at"` + } + + out := make([]goalSummary, 0, len(gs)) + for _, g := range gs { + s := goalSummary{ + ID: g.ID, + Slug: g.Slug, + Title: g.Title, + Status: g.Status, + ChannelID: g.ChannelID, + RootTaskID: g.RootTaskID, + BudgetTokens: g.BudgetTokens, + BudgetDollarsCents: g.BudgetDollarsCents, + CreatedAt: g.CreatedAt.UTC().Format("2006-01-02T15:04:05Z"), + } + _ = h.db.QueryRowContext(r.Context(), + `SELECT username FROM users WHERE id=?`, g.OwnerUserID).Scan(&s.OwnerUsername) + if g.RootTaskID != nil { + tokens, cents, count, err := h.tasks.RollupCosts(r.Context(), *g.RootTaskID) + if err == nil { + s.SpentTokens = tokens + s.SpentDollarsCents = cents + s.TaskCount = count + } + } + if g.BudgetDollarsCents != nil && *g.BudgetDollarsCents > 0 { + s.PercentBudget = float64(s.SpentDollarsCents) / float64(*g.BudgetDollarsCents) * 100.0 + } + out = append(out, s) + } + + writeJSON(w, http.StatusOK, map[string]any{"goals": out}) +} + +// GetGoal returns a single goal with its full task tree, cost rollup, +// billing-code breakdown, and spawned-agent snapshot. +func (h *GoalsHandler) GetGoal(w http.ResponseWriter, r *http.Request) { + idStr := chi.URLParam(r, "id") + id, err := strconv.ParseInt(idStr, 10, 64) + if err != nil { + writeJSON(w, http.StatusBadRequest, errorBody("bad_request", "invalid goal id")) + return + } + + g, err := h.goals.GetGoal(r.Context(), id) + if err != nil { + writeJSON(w, http.StatusNotFound, errorBody("not_found", err.Error())) + return + } + + tasks, err := h.tasks.ListByGoal(r.Context(), g.ID) + if err != nil { + writeJSON(w, http.StatusInternalServerError, errorBody("internal_error", err.Error())) + return + } + + var rollupTokens, rollupCents int64 + var rollupCount int + if g.RootTaskID != nil { + rollupTokens, rollupCents, rollupCount, _ = h.tasks.RollupCosts(r.Context(), *g.RootTaskID) + } + + type taskOut struct { + ID int64 `json:"id"` + ParentTaskID *int64 `json:"parent_task_id"` + Title string `json:"title"` + Description string `json:"description"` + AcceptanceCriteria string `json:"acceptance_criteria"` + Status string `json:"status"` + Depth int `json:"depth"` + BillingCode string `json:"billing_code"` + AssigneeAgentID *int64 `json:"assignee_agent_id"` + AssigneeAgentName string `json:"assignee_agent_name,omitempty"` + SpentTokens int64 `json:"spent_tokens"` + SpentDollarsCents int64 `json:"spent_dollars_cents"` + VerifierConfig *goaltasks.VerifierConfig `json:"verifier_config,omitempty"` + HeartbeatConfig *goaltasks.HeartbeatConfig `json:"heartbeat_config,omitempty"` + FailureReason string `json:"failure_reason,omitempty"` + CreatedAt string `json:"created_at"` + CompletedAt *string `json:"completed_at,omitempty"` + } + + agentNameByID := map[int64]string{} + out := make([]taskOut, 0, len(tasks)) + for _, t := range tasks { + tt := taskOut{ + ID: t.ID, + ParentTaskID: t.ParentTaskID, + Title: t.Title, + Description: t.Description, + AcceptanceCriteria: t.AcceptanceCriteria, + Status: t.Status, + Depth: t.Depth, + BillingCode: t.BillingCode, + AssigneeAgentID: t.AssigneeAgentID, + SpentTokens: t.SpentTokens, + SpentDollarsCents: t.SpentDollarsCents, + VerifierConfig: t.VerifierConfig, + HeartbeatConfig: t.HeartbeatConfig, + FailureReason: t.FailureReason, + CreatedAt: t.CreatedAt.UTC().Format("2006-01-02T15:04:05Z"), + } + if t.CompletedAt != nil { + s := t.CompletedAt.UTC().Format("2006-01-02T15:04:05Z") + tt.CompletedAt = &s + } + if t.AssigneeAgentID != nil { + name, ok := agentNameByID[*t.AssigneeAgentID] + if !ok { + _ = h.db.QueryRowContext(r.Context(), + `SELECT name FROM agents WHERE id=?`, *t.AssigneeAgentID).Scan(&name) + agentNameByID[*t.AssigneeAgentID] = name + } + tt.AssigneeAgentName = name + } + out = append(out, tt) + } + + // Spawned agents attached to this goal (any agent whose config_hash + // appears as an assignee on one of the goal's tasks, plus the + // coordinator itself). + type spawnedAgent struct { + ID int64 `json:"id"` + Name string `json:"name"` + DisplayName string `json:"display_name"` + ConfigHash string `json:"config_hash"` + SpawnDepth int `json:"spawn_depth"` + AutonomyTier string `json:"autonomy_tier"` + ParentAgent string `json:"parent_agent_name,omitempty"` + } + agentSeen := map[int64]bool{} + agentList := []spawnedAgent{} + collectAgent := func(id int64) { + if id == 0 || agentSeen[id] { + return + } + agentSeen[id] = true + var sa spawnedAgent + var parentID sql.NullInt64 + err := h.db.QueryRowContext(r.Context(), ` + SELECT id, name, display_name, + COALESCE(config_hash,''), COALESCE(spawn_depth,0), + COALESCE(autonomy_tier,''), parent_agent_id + FROM agents WHERE id=?`, id). + Scan(&sa.ID, &sa.Name, &sa.DisplayName, &sa.ConfigHash, &sa.SpawnDepth, &sa.AutonomyTier, &parentID) + if err != nil { + return + } + if parentID.Valid { + var pname string + _ = h.db.QueryRowContext(r.Context(), + `SELECT name FROM agents WHERE id=?`, parentID.Int64).Scan(&pname) + sa.ParentAgent = pname + } + agentList = append(agentList, sa) + } + if g.CoordinatorAgentID != nil { + collectAgent(*g.CoordinatorAgentID) + } + for _, t := range tasks { + if t.AssigneeAgentID != nil { + collectAgent(*t.AssigneeAgentID) + } + } + + // Billing-code rollup via raw query (service wrapper not needed). + billingBreakdown := map[string]map[string]int64{} + if g.RootTaskID != nil { + rows, err := h.db.QueryContext(r.Context(), ` + WITH RECURSIVE subtree(id) AS ( + SELECT id FROM goal_tasks WHERE id = ? + UNION ALL + SELECT t.id FROM goal_tasks t + JOIN subtree s ON t.parent_task_id = s.id + ) + SELECT COALESCE(billing_code,''), SUM(spent_tokens), SUM(spent_dollars_cents) + FROM goal_tasks WHERE id IN subtree + GROUP BY billing_code`, *g.RootTaskID) + if err == nil { + for rows.Next() { + var code string + var tokens, cents int64 + if err := rows.Scan(&code, &tokens, ¢s); err == nil { + billingBreakdown[code] = map[string]int64{ + "tokens": tokens, + "cents": cents, + } + } + } + rows.Close() + } + } + + var ownerUsername string + _ = h.db.QueryRowContext(r.Context(), + `SELECT username FROM users WHERE id=?`, g.OwnerUserID).Scan(&ownerUsername) + + // Recent system/artifact messages on the goal channel for a small timeline. + type timelineEvent struct { + ID int64 `json:"id"` + From string `json:"from"` + Body string `json:"body"` + Kind string `json:"kind"` + CreatedAt string `json:"created_at"` + } + timeline := []timelineEvent{} + rows, err := h.db.QueryContext(r.Context(), ` + SELECT id, from_agent, body, COALESCE(metadata,''), created_at + FROM messages + WHERE channel_id = ? + ORDER BY id DESC + LIMIT 50`, g.ChannelID) + if err == nil { + for rows.Next() { + var ev timelineEvent + var meta string + if err := rows.Scan(&ev.ID, &ev.From, &ev.Body, &meta, &ev.CreatedAt); err == nil { + if meta != "" { + var m map[string]any + if json.Unmarshal([]byte(meta), &m) == nil { + if k, ok := m["kind"].(string); ok { + ev.Kind = k + } + } + } + timeline = append(timeline, ev) + } + } + rows.Close() + } + + writeJSON(w, http.StatusOK, map[string]any{ + "goal": map[string]any{ + "id": g.ID, + "slug": g.Slug, + "title": g.Title, + "description": g.Description, + "status": g.Status, + "channel_id": g.ChannelID, + "coordinator_agent_id": g.CoordinatorAgentID, + "root_task_id": g.RootTaskID, + "owner_user_id": g.OwnerUserID, + "owner_username": ownerUsername, + "budget_tokens": g.BudgetTokens, + "budget_dollars_cents": g.BudgetDollarsCents, + "max_spawn_depth": g.MaxSpawnDepth, + "alert_80pct_posted": g.Alert80PctPosted, + "created_at": g.CreatedAt.UTC().Format("2006-01-02T15:04:05Z"), + }, + "tasks": out, + "rollup": map[string]any{ + "tokens": rollupTokens, + "dollars_cents": rollupCents, + "task_count": rollupCount, + }, + "billing_breakdown": billingBreakdown, + "spawned_agents": agentList, + "timeline": timeline, + }) +} diff --git a/internal/api/router.go b/internal/api/router.go index 15ddf3b..0b7708b 100644 --- a/internal/api/router.go +++ b/internal/api/router.go @@ -10,6 +10,8 @@ import ( "github.com/synapbus/synapbus/internal/apikeys" "github.com/synapbus/synapbus/internal/attachments" "github.com/synapbus/synapbus/internal/channels" + "github.com/synapbus/synapbus/internal/goals" + "github.com/synapbus/synapbus/internal/goaltasks" "github.com/synapbus/synapbus/internal/harness/runs" "github.com/synapbus/synapbus/internal/k8s" "github.com/synapbus/synapbus/internal/messaging" @@ -43,6 +45,8 @@ type RouterConfig struct { ReactorStore *reactor.Store ReactorEngine *reactor.Reactor HarnessRunsStore *runs.Store + GoalsService *goals.Service + GoalTasksService *goaltasks.Service WikiService *wiki.Service SSEHub *SSEHub Broadcaster *SSEBroadcaster @@ -266,6 +270,17 @@ func NewRouterWithConfig(cfg RouterConfig) chi.Router { }) } + // Goals + if cfg.GoalsService != nil && cfg.GoalTasksService != nil && cfg.DB != nil { + goalsHandler := NewGoalsHandler(cfg.GoalsService, cfg.GoalTasksService, cfg.DB) + r.Group(func(r chi.Router) { + r.Use(authMiddleware) + + r.Get("/api/goals", goalsHandler.ListGoals) + r.Get("/api/goals/{id}", goalsHandler.GetGoal) + }) + } + // Trust Scores if cfg.TrustService != nil { trustHandler := NewTrustHandler(cfg.TrustService) diff --git a/internal/goals/service.go b/internal/goals/service.go index f647687..f806147 100644 --- a/internal/goals/service.go +++ b/internal/goals/service.go @@ -104,6 +104,44 @@ func (s *Service) TransitionStatus(ctx context.Context, goalID int64, newStatus return s.store.SetStatus(ctx, goalID, newStatus) } +// BudgetVerdict describes what the budget enforcer wants the caller to do. +type BudgetVerdict struct { + PercentBudget float64 // 0..100+ + TriggerSoftAlert bool // first time we cross 80% + TriggerHardPause bool // crossed 100% and goal is still active +} + +// EvaluateBudget computes current spend-vs-budget for a goal and returns +// the enforcement verdict. It does NOT mutate state on its own — the +// caller uses MarkSoftAlertPosted / TransitionStatus to apply the +// verdict once it has posted the corresponding system messages. +// +// Only dollar-cents budget is enforced in MVP (tokens are tracked but +// don't trip the cascade). +func (s *Service) EvaluateBudget(ctx context.Context, goalID int64, spentCents int64) (*BudgetVerdict, error) { + g, err := s.store.Get(ctx, goalID) + if err != nil { + return nil, err + } + v := &BudgetVerdict{} + if g.BudgetDollarsCents == nil || *g.BudgetDollarsCents <= 0 { + return v, nil + } + v.PercentBudget = float64(spentCents) / float64(*g.BudgetDollarsCents) * 100.0 + if v.PercentBudget >= 80 && !g.Alert80PctPosted { + v.TriggerSoftAlert = true + } + if v.PercentBudget >= 100 && g.Status == StatusActive { + v.TriggerHardPause = true + } + return v, nil +} + +// MarkSoftAlertPosted records that the 80% soft-alert was emitted. +func (s *Service) MarkSoftAlertPosted(ctx context.Context, goalID int64) error { + return s.store.MarkSoftAlertPosted(ctx, goalID) +} + func legalTransition(from, to string) bool { switch from { case StatusDraft: diff --git a/internal/reactor/reactor.go b/internal/reactor/reactor.go index 4f53a7f..ffb0c2c 100644 --- a/internal/reactor/reactor.go +++ b/internal/reactor/reactor.go @@ -32,6 +32,13 @@ const ( ) // Reactor is the reactive agent triggering engine. +// SecretProvider builds the env map of scoped secrets to inject into a +// reactive subprocess run. Returning an error is non-fatal — the run +// proceeds without any injected secrets and the error is logged. +type SecretProvider interface { + BuildEnvMap(ctx context.Context, userID, agentID, taskID int64) (map[string]string, error) +} + type Reactor struct { store *Store agentStore agents.AgentStore @@ -39,6 +46,7 @@ type Reactor struct { registry *harness.Registry notifier FailureNotifier reactions ReactionNotifier + secrets SecretProvider logger *slog.Logger } @@ -78,6 +86,12 @@ func (r *Reactor) SetHarnessRegistry(reg *harness.Registry) { r.registry = reg } +// SetSecretProvider wires the component that reads scoped secrets for +// the reactor to inject into subprocess runs as env vars. +func (r *Reactor) SetSecretProvider(p SecretProvider) { + r.secrets = p +} + // SetReactionNotifier wires the component that marks triggering DMs // with in_progress / done / reject reactions. Optional — a nil // notifier simply skips the reaction step. @@ -148,6 +162,20 @@ func (r *Reactor) evaluateTrigger(ctx context.Context, agentName string, event d return nil // Not reactive, skip } + // 2a. Refuse to dispatch to a quarantined agent. Quarantine is set + // when the agent's rolling reputation drops below the threshold + // (0.3 by default). Existing in-flight runs are allowed to finish + // but no new reactive runs will spawn. + if agent.QuarantinedAt != nil { + r.logger.Info("reactive dispatch to quarantined agent refused", + "agent", agentName, + "quarantined_at", agent.QuarantinedAt, + "reason", agent.QuarantineReason, + ) + r.recordSkippedRun(ctx, agentName, event, StatusFailed, "agent quarantined: "+agent.QuarantineReason) + return nil + } + // 3. Pick a backend. K8s agents (k8s_image set) keep the existing // createJob + async poller path for restart safety. Everything else // goes through the harness registry in a goroutine. @@ -337,6 +365,19 @@ func (r *Reactor) dispatchHarness(ctx context.Context, agent *agents.Agent, even }, } + // Inject scoped secrets as env vars (user + agent scope; task scope + // is added when a task ID becomes available in the trigger path). + if r.secrets != nil { + if secretEnv, serr := r.secrets.BuildEnvMap(ctx, agent.OwnerID, agent.ID, 0); serr == nil { + for k, v := range secretEnv { + req.Env[k] = v + } + } else { + r.logger.Warn("secret provider failed — continuing without injection", + "agent", agent.Name, "error", serr) + } + } + // Block on a detached context so a cancelled incoming request // does not kill in-flight work. Callers get a fast return above. go r.runHarness(runID, agent, event, req) diff --git a/internal/web/dist/index.html b/internal/web/dist/index.html index bd4e10a..ff597f7 100644 --- a/internal/web/dist/index.html +++ b/internal/web/dist/index.html @@ -11,30 +11,30 @@ - - + + - - - - - - + + + + + +
+ + + Goals — SynapBus + + +
+
+

Goals

+

+ Top-level objectives owned by a human, decomposed into task trees by a coordinator, and + executed by dynamically spawned specialist agents. +

+
+ + {#if loading} +

Loading…

+ {:else if error} +
+ {error} +
+ {:else if list.length === 0} +
+ No goals yet. Goals are created by the coordinator agent when a human DMs it with a brief. +
+ {:else} + + {/if} +
diff --git a/web/src/routes/goals/[id]/+page.svelte b/web/src/routes/goals/[id]/+page.svelte new file mode 100644 index 0000000..d8809eb --- /dev/null +++ b/web/src/routes/goals/[id]/+page.svelte @@ -0,0 +1,225 @@ + + + + {data?.goal?.title ?? 'Goal'} — SynapBus + + +
+ + ← All goals + + + {#if loading && !data} +

Loading…

+ {:else if error} +
+ {error} +
+ {:else if data} +
+
+ #{data.goal.id} + + {data.goal.status} + + by {data.goal.owner_username} +
+

{data.goal.title}

+
#goal-{data.goal.slug}
+

{data.goal.description}

+
+ + +
+
+
Total spend
+
+ {formatDollars(data.rollup.dollars_cents)} +
+
+ {data.rollup.tokens.toLocaleString()} tokens · {data.rollup.task_count} tasks +
+ {#if data.goal.budget_dollars_cents} +
+ budget {formatDollars(data.goal.budget_dollars_cents)} · max depth + {data.goal.max_spawn_depth} +
+ {/if} +
+
+
Billing breakdown
+
    + {#each Object.entries(data.billing_breakdown) as [code, s] (code)} +
  • + {code || '—'} + {formatDollars((s as any).cents)} +
  • + {/each} +
+
+
+
Spawned agents
+
    + {#each data.spawned_agents as a (a.id)} +
  • + + {a.name} + + depth {a.spawn_depth} +
  • + {/each} +
+
+
+ + +
+

Task tree

+
+ {#each data.tasks as t (t.id)} +
+
+ #{t.id} + + {t.status.replace('_', ' ')} + + {#if t.billing_code} + {t.billing_code} + {/if} + {#if t.assignee_agent_name} + + → + {t.assignee_agent_name} + + {/if} +
+
{t.title}
+ {#if t.description} +
{t.description}
+ {/if} + {#if t.acceptance_criteria} +
+ acceptance: {t.acceptance_criteria} +
+ {/if} +
+ {t.spent_tokens.toLocaleString()} tok + · + {formatDollars(t.spent_dollars_cents)} + {#if t.failure_reason} + · {t.failure_reason} + {/if} +
+
+ {/each} +
+
+ + + {#if data.timeline?.length > 0} +
+

Recent activity

+
+
    + {#each data.timeline as ev (ev.id)} +
  • + + {ev.created_at.substring(11, 19)} + + {ev.from} + {#if ev.kind} + {ev.kind} + {/if} + {ev.body} +
  • + {/each} +
+
+
+ {/if} + {/if} +