package main import ( "context" "encoding/json" "fmt" "log/slog" "regexp" "strings" "time" "github.com/dtoro/oikos/internal/db/sqlcgen" "github.com/dtoro/oikos/internal/observability" "github.com/google/uuid" "github.com/jackc/pgx/v5/pgxpool" ) const maxToolResultSize = 4096 type store struct { pool *pgxpool.Pool } func newStore(ctx context.Context, databaseURL string) (*store, error) { if databaseURL == "" { return nil, nil } pool, err := pgxpool.New(ctx, databaseURL) if err != nil { return nil, fmt.Errorf("connect db: %w", err) } if err := pool.Ping(ctx); err != nil { pool.Close() return nil, fmt.Errorf("ping db: %w", err) } return &store{pool: pool}, nil } func (s *store) close() { if s.pool != nil { s.pool.Close() } } // session is a chat session elevated to a task: goal-structured work with a // lifecycle status and an outcome (see migration 018 / the task-board plan). // Outcome/Summary/EntityID are empty until set, hence omitempty. type session struct { ID string `json:"id"` Title string `json:"title"` Actor string `json:"actor"` Goal string `json:"goal"` Status string `json:"status"` Outcome string `json:"outcome,omitempty"` Summary string `json:"summary,omitempty"` EntityID string `json:"entity_id,omitempty"` CreatedAt time.Time `json:"created_at"` LastActiveAt time.Time `json:"last_active_at"` } type message struct { ID string `json:"id"` SessionID string `json:"session_id"` Role string `json:"role"` Content json.RawMessage `json:"content"` CreatedAt time.Time `json:"created_at"` } func (s *store) createSession(ctx context.Context, title string) (*session, error) { if s == nil { return &session{ID: "ephemeral", Title: title, Actor: "agent:nomos", Status: "active"}, nil } var id string err := s.pool.QueryRow(ctx, `INSERT INTO agent_sessions (title, actor) VALUES ($1, 'agent:nomos') RETURNING id`, title).Scan(&id) if err != nil { return nil, err } // Give the task its own entity so knowledge and involved-entity edges hang // off the existing relationships graph. Best-effort: a failure here must not // block the chat — the session is usable without a graph anchor. entityID := s.createTaskEntity(ctx, id, title) return &session{ID: id, Title: title, Actor: "agent:nomos", Status: "active", EntityID: entityID, CreatedAt: time.Now(), LastActiveAt: time.Now()}, nil } // createTaskEntity creates (or reuses) the task: entity that // anchors this task's knowledge and involved-entity relationships, and records // it on the session. Returns the entity id, or "" on failure — non-fatal, see // caller. Requires the 'task' entity type (seeds/ontology.yaml). func (s *store) createTaskEntity(ctx context.Context, sessionID, title string) string { entityID, _ := uuid.NewV7() slug := "task:" + sessionID // name is UNIQUE(type,name) and chat titles collide ("hi" ×6), so key the // name on the session id and keep the human title in attributes for display. name := "task " + sessionID attrs, _ := json.Marshal(map[string]any{"title": title}) if err := s.pool.QueryRow(ctx, ` INSERT INTO entities (id, slug, type, name, attributes) VALUES ($1, $2, 'task', $3, $4) ON CONFLICT (slug) DO UPDATE SET updated_at = now() RETURNING id`, entityID, slug, name, string(attrs)).Scan(&entityID); err != nil { slog.Warn("nomos: could not create task entity", "session", sessionID, "error", err) return "" } if _, err := s.pool.Exec(ctx, `UPDATE agent_sessions SET entity_id = $1 WHERE id = $2`, entityID, sessionID); err != nil { slog.Warn("nomos: could not link task entity", "session", sessionID, "error", err) } return entityID.String() } func (s *store) saveMessage(ctx context.Context, sessionID, role string, content json.RawMessage) error { if s == nil { return nil } _, err := s.pool.Exec(ctx, `INSERT INTO agent_messages (session_id, role, content) VALUES ($1, $2, $3)`, sessionID, role, truncateToolResults(content)) return err } // insertMessageReturningID and updateMessage exist for the auto-continuation // worker's live-progress persistence (see continue.go): rather than saving // one message only once the whole continuation finishes — which could be // several minutes of silence in the UI even though frontend polling exists — // the worker inserts a placeholder immediately and updates the SAME row as // each tool call completes, so a poller sees individual steps land, not just // a final rolled-up summary. func (s *store) insertMessageReturningID(ctx context.Context, sessionID, role string, content json.RawMessage) (uuid.UUID, error) { if s == nil { return uuid.Nil, nil } var id uuid.UUID err := s.pool.QueryRow(ctx, `INSERT INTO agent_messages (session_id, role, content) VALUES ($1, $2, $3) RETURNING id`, sessionID, role, truncateToolResults(content)).Scan(&id) return id, err } func (s *store) updateMessage(ctx context.Context, id uuid.UUID, content json.RawMessage) error { if s == nil || id == uuid.Nil { return nil } _, err := s.pool.Exec(ctx, `UPDATE agent_messages SET content = $2 WHERE id = $1`, id, truncateToolResults(content)) return err } func truncateToolResults(content json.RawMessage) json.RawMessage { var m map[string]any if err := json.Unmarshal(content, &m); err != nil { return content } toolCalls, ok := m["tool_calls"].([]any) if !ok || len(toolCalls) == 0 { return content } changed := false for i, raw := range toolCalls { tc, ok := raw.(map[string]any) if !ok { continue } if result, ok := tc["result"]; ok { resultJSON, _ := json.Marshal(result) if len(resultJSON) > maxToolResultSize { tc["result"] = string(resultJSON[:maxToolResultSize]) + fmt.Sprintf("...truncated (%d bytes total)", len(resultJSON)) toolCalls[i] = tc changed = true } } } if !changed { return content } m["tool_calls"] = toolCalls out, err := json.Marshal(m) if err != nil { return content } return out } func (s *store) touchSession(ctx context.Context, id string) { if s != nil { s.pool.Exec(ctx, `UPDATE agent_sessions SET last_active_at=now() WHERE id=$1`, id) } } func (s *store) listSessions(ctx context.Context) ([]session, error) { if s == nil { return nil, nil } rows, err := s.pool.Query(ctx, `SELECT id, title, actor, goal, status, COALESCE(outcome, ''), summary, COALESCE(entity_id::text, ''), created_at, last_active_at FROM agent_sessions ORDER BY last_active_at DESC LIMIT 50`) if err != nil { return nil, err } defer rows.Close() var out []session for rows.Next() { var sess session if err := rows.Scan(&sess.ID, &sess.Title, &sess.Actor, &sess.Goal, &sess.Status, &sess.Outcome, &sess.Summary, &sess.EntityID, &sess.CreatedAt, &sess.LastActiveAt); err != nil { return nil, err } out = append(out, sess) } return out, rows.Err() } // getMessages returns a session's ENTIRE message history, unbounded — used // for the UI's own transcript view (GET /sessions/{id}), where the operator // should be able to see everything a task has done regardless of how long // it's run. For LLM replay, see getRecentMessages: sending the operator's // full transcript is fine; sending the model's full transcript on every // single turn is not (see getRecentMessages's doc comment). func (s *store) getMessages(ctx context.Context, sessionID string) ([]message, error) { if s == nil { return nil, nil } rows, err := s.pool.Query(ctx, `SELECT id, session_id, role, content, created_at FROM agent_messages WHERE session_id=$1 ORDER BY created_at ASC`, sessionID) if err != nil { return nil, err } defer rows.Close() var out []message for rows.Next() { var m message if err := rows.Scan(&m.ID, &m.SessionID, &m.Role, &m.Content, &m.CreatedAt); err != nil { return nil, err } out = append(out, m) } return out, rows.Err() } // getRecentMessages returns the most recent `limit` messages for sessionID, // in chronological order, plus whether older messages exist beyond that // window. Used specifically for LLM replay (chatWith): without a bound, // every turn re-sent the ENTIRE session history into the model's context, // unconditionally growing with every turn — a real, observed-in-production // cost/latency/eventual-context-limit risk for exactly the long-running, // heavily-autonomous tasks (many auto-continuation cycles) this system is // built to run longest. Fetches limit+1 rows to detect "there's more" // without a separate COUNT query. func (s *store) getRecentMessages(ctx context.Context, sessionID string, limit int) (msgs []message, truncated bool, err error) { if s == nil { return nil, false, nil } rows, qerr := s.pool.Query(ctx, `SELECT id, session_id, role, content, created_at FROM agent_messages WHERE session_id=$1 ORDER BY created_at DESC LIMIT $2`, sessionID, limit+1) if qerr != nil { return nil, false, qerr } defer rows.Close() var out []message for rows.Next() { var m message if err := rows.Scan(&m.ID, &m.SessionID, &m.Role, &m.Content, &m.CreatedAt); err != nil { return nil, false, err } out = append(out, m) } if err := rows.Err(); err != nil { return nil, false, err } truncated = len(out) > limit if truncated { out = out[:limit] } // Rows came back newest-first (for the LIMIT to bound the right end); // reverse to chronological order for replay. for i, j := 0, len(out)-1; i < j; i, j = i+1, j-1 { out[i], out[j] = out[j], out[i] } return out, truncated, nil } func (s *store) deleteSession(ctx context.Context, id string) error { if s == nil { return nil } // Resolve the task entity so we can clean up its graph edges and events too // — otherwise deleting a session orphans its task: entity, its involves/ // documents relationships, and its task-scoped events. var entID uuid.UUID s.pool.QueryRow(ctx, `SELECT entity_id FROM agent_sessions WHERE id = $1`, id).Scan(&entID) if _, err := s.pool.Exec(ctx, `DELETE FROM agent_messages WHERE session_id = $1`, id); err != nil { return err } // task.status / entity.touched / knowledge.recorded are all correlated by // session id. s.pool.Exec(ctx, `DELETE FROM events WHERE correlation_id = $1`, id) if _, err := s.pool.Exec(ctx, `DELETE FROM agent_sessions WHERE id = $1`, id); err != nil { return err } if entID != uuid.Nil { // relationships FK is ON DELETE RESTRICT, so drop the task's edges first. s.pool.Exec(ctx, `DELETE FROM relationships WHERE source_id = $1 OR target_id = $1`, entID) s.pool.Exec(ctx, `DELETE FROM entities WHERE id = $1`, entID) } return nil } // taskEntityPtr returns the task entity id for a session, or nil — used as the // entity_id on task-scoped events so they anchor to the task in the graph. func (s *store) taskEntityPtr(ctx context.Context, sessionID string) *uuid.UUID { var id uuid.UUID if err := s.pool.QueryRow(ctx, `SELECT entity_id FROM agent_sessions WHERE id = $1`, sessionID).Scan(&id); err != nil || id == uuid.Nil { return nil } return &id } // setGoal records the task's goal and moves it into planning. Emits goal.set. func (s *store) setGoal(ctx context.Context, sessionID, goal string) error { if s == nil || sessionID == "" || sessionID == "ephemeral" { return nil } if _, err := s.pool.Exec(ctx, `UPDATE agent_sessions SET goal = $2, status = 'planning', last_active_at = now() WHERE id = $1`, sessionID, goal); err != nil { return err } _ = observability.Event(ctx, sqlcgen.New(s.pool), "goal.set", s.taskEntityPtr(ctx, sessionID), "info", "nomos", sessionID, map[string]any{"goal": goal}) return nil } // planStepInput is one step as the agent proposes it. type planStepInput struct { Title string Detail string TargetSlug string } // proposePlan sets the task's plan and moves it into executing. Emits // plan.proposed with the persisted steps (seq + id) so the panel can render // and later address them by id. // // Two modes, chosen by whether any existing step has left 'pending': // - Fresh/revise (no step started yet): full replace (delete + insert). This // covers the first call, and a genuine re-plan before any work began. // - Mid-flight (some step is running/done/failed/…): APPEND the new steps // after the current max seq instead of wiping. The model is instructed to // propose the whole plan in one call, but nothing stops it from calling // propose_plan again per-step as it goes — a destructive replace in that // case would erase every already-completed step, leaving the operator // seeing only the most recent single step ("1/1") instead of real // progress. Appending makes the panel's step history correct regardless // of how the model chooses to call the tool. func (s *store) proposePlan(ctx context.Context, sessionID string, steps []planStepInput) ([]map[string]any, error) { if s == nil || sessionID == "" || sessionID == "ephemeral" { return nil, nil } tx, err := s.pool.Begin(ctx) if err != nil { return nil, err } defer tx.Rollback(ctx) var startSeq int var anyStarted bool if err := tx.QueryRow(ctx, ` SELECT COALESCE(max(seq), 0), COALESCE(bool_or(status <> 'pending'), false) FROM session_plan_steps WHERE session_id = $1`, sessionID).Scan(&startSeq, &anyStarted); err != nil { return nil, err } if !anyStarted { if _, err := tx.Exec(ctx, `DELETE FROM session_plan_steps WHERE session_id = $1`, sessionID); err != nil { return nil, err } startSeq = 0 } out := make([]map[string]any, 0, len(steps)) for i, st := range steps { var targetSlug *string if st.TargetSlug != "" { targetSlug = &st.TargetSlug } seq := startSeq + i + 1 var id uuid.UUID if err := tx.QueryRow(ctx, ` INSERT INTO session_plan_steps (session_id, seq, title, detail, target_slug) VALUES ($1, $2, $3, $4, $5) RETURNING id`, sessionID, seq, st.Title, st.Detail, targetSlug).Scan(&id); err != nil { return nil, err } out = append(out, map[string]any{ "id": id.String(), "seq": seq, "title": st.Title, "detail": st.Detail, "target_slug": st.TargetSlug, }) } if _, err := tx.Exec(ctx, `UPDATE agent_sessions SET status = 'executing', last_active_at = now() WHERE id = $1`, sessionID); err != nil { return nil, err } if err := tx.Commit(ctx); err != nil { return nil, err } // Event after commit so subscribers only ever see a persisted plan. // appended=true tells the panel to add these steps to its existing list // rather than replace it (mirrors the mid-flight append above). _ = observability.Event(ctx, sqlcgen.New(s.pool), "plan.proposed", s.taskEntityPtr(ctx, sessionID), "info", "nomos", sessionID, map[string]any{"steps": out, "appended": anyStarted}) return out, nil } // updatePlanStep sets a step's status by seq, stamping started_at/finished_at // and linking an execution if given. Emits plan.step.started (running) or // plan.step.finished (terminal) so the panel advances live. The execution link // is also what lets the api auto-close the step when the execution finishes // (see closePlanStepForExecution). func (s *store) updatePlanStep(ctx context.Context, sessionID string, seq int, status, execID string) error { if s == nil || sessionID == "" || sessionID == "ephemeral" { return nil } stamp := "" switch status { case "running": stamp = ", started_at = COALESCE(started_at, now())" case "done", "failed", "skipped", "blocked": stamp = ", finished_at = now()" } var execPtr *uuid.UUID if id, err := uuid.Parse(execID); err == nil { execPtr = &id } var stepID uuid.UUID var targetSlug *string // stamp is a fixed literal from the switch above — never user input. if err := s.pool.QueryRow(ctx, ` UPDATE session_plan_steps SET status = $3, execution_id = COALESCE($4, execution_id)`+stamp+` WHERE session_id = $1 AND seq = $2 RETURNING id, target_slug`, sessionID, seq, status, execPtr).Scan(&stepID, &targetSlug); err != nil { return err } // Anchor the event to the step's target entity when it has one, else the task. entPtr := s.taskEntityPtr(ctx, sessionID) if targetSlug != nil && *targetSlug != "" { var tid uuid.UUID if s.pool.QueryRow(ctx, `SELECT id FROM entities WHERE slug = $1`, *targetSlug).Scan(&tid) == nil { entPtr = &tid } } evType := "plan.step.finished" if status == "running" { evType = "plan.step.started" } data := map[string]any{"step_id": stepID.String(), "seq": seq, "status": status} if execID != "" { data["execution_id"] = execID } _ = observability.Event(ctx, sqlcgen.New(s.pool), evType, entPtr, "info", "nomos", sessionID, data) return nil } // completeTask sets a task's terminal state, outcome, and one-line summary, // mirrors the outcome onto the task entity's attributes (so the board/graph // show it), and publishes task.status for the live context panel. outcome is // success|failure|partial; status is derived (failure → failed, else done). func (s *store) completeTask(ctx context.Context, sessionID, outcome, summary string) error { if s == nil || sessionID == "" || sessionID == "ephemeral" { return nil } status := "done" if outcome == "failure" { status = "failed" } if _, err := s.pool.Exec(ctx, ` UPDATE agent_sessions SET status = $2, outcome = $3, summary = $4, last_active_at = now() WHERE id = $1`, sessionID, status, outcome, summary); err != nil { return err } var entID uuid.UUID s.pool.QueryRow(ctx, `SELECT entity_id FROM agent_sessions WHERE id = $1`, sessionID).Scan(&entID) var entPtr *uuid.UUID if entID != uuid.Nil { attrs, _ := json.Marshal(map[string]any{"outcome": outcome, "status": status, "summary": summary}) s.pool.Exec(ctx, `UPDATE entities SET attributes = attributes || $2::jsonb, updated_at = now() WHERE id = $1`, entID, string(attrs)) entPtr = &entID } severity := "info" if outcome == "failure" { severity = "warning" } _ = observability.Event(ctx, sqlcgen.New(s.pool), "task.status", entPtr, severity, "nomos", sessionID, map[string]any{"status": status, "outcome": outcome, "summary": summary}) return nil } // staleGoalSession is a goal-bearing task that's gone idle without reaching // a terminal state — the idle-sweep worker's work list (fix 2+3 of // plans/2026-07-11-task-completion-safety-net.md). type staleGoalSession struct { ID string Goal string CompletionNudges int } // staleGoalSessions finds sessions that framed themselves as a real task // (goal != '', so the inline safety net in agent.go intentionally left them // alone) but have sat non-terminal past idleThreshold. completion_nudges // tells the caller whether to nudge (0) or give up and auto-close (>=1) — // see processIdleSweep in continue.go. func (s *store) staleGoalSessions(ctx context.Context, idleThreshold time.Duration, limit int) []staleGoalSession { if s == nil { return nil } rows, err := s.pool.Query(ctx, ` SELECT id, goal, completion_nudges FROM agent_sessions WHERE goal <> '' AND status IN ('active', 'planning', 'executing') AND last_active_at < now() - ($1 * interval '1 second') ORDER BY last_active_at LIMIT $2`, idleThreshold.Seconds(), limit) if err != nil { return nil } defer rows.Close() var out []staleGoalSession for rows.Next() { var s staleGoalSession if err := rows.Scan(&s.ID, &s.Goal, &s.CompletionNudges); err == nil { out = append(out, s) } } return out } // bumpCompletionNudge records that the idle sweep nudged a stalled session, // stamping last_active_at so it isn't picked up again until it's genuinely // idle again (a fresh nudge shouldn't fire every tick while the model is // mid-response to the previous one). func (s *store) bumpCompletionNudge(ctx context.Context, sessionID string) error { if s == nil { return nil } _, err := s.pool.Exec(ctx, ` UPDATE agent_sessions SET completion_nudges = completion_nudges + 1, last_active_at = now() WHERE id = $1`, sessionID) return err } // planStep is a persisted plan step, as returned to the frontend for hydration // (the panel otherwise only sees steps live via plan.proposed/plan.step.*). type planStep struct { ID string `json:"id"` Seq int `json:"seq"` Title string `json:"title"` Detail string `json:"detail"` Status string `json:"status"` ExecutionID *string `json:"execution_id,omitempty"` TargetSlug *string `json:"target_slug,omitempty"` StartedAt *string `json:"started_at,omitempty"` FinishedAt *string `json:"finished_at,omitempty"` } // getPlanSteps returns a task's plan in order — REST hydration for the context // panel when it first opens a task (live events only carry deltas from then on). func (s *store) getPlanSteps(ctx context.Context, sessionID string) ([]planStep, error) { if s == nil { return nil, nil } rows, err := s.pool.Query(ctx, ` SELECT id::text, seq, title, detail, status, execution_id::text, target_slug, started_at::text, finished_at::text FROM session_plan_steps WHERE session_id = $1 ORDER BY seq`, sessionID) if err != nil { return nil, err } defer rows.Close() var out []planStep for rows.Next() { var st planStep var execID, target, started, finished *string if err := rows.Scan(&st.ID, &st.Seq, &st.Title, &st.Detail, &st.Status, &execID, &target, &started, &finished); err != nil { return nil, err } st.ExecutionID, st.TargetSlug, st.StartedAt, st.FinishedAt = execID, target, started, finished out = append(out, st) } return out, rows.Err() } // sessionQuestion is a persisted question, as returned to the frontend. type sessionQuestion struct { ID string `json:"id"` Prompt string `json:"prompt"` Context map[string]any `json:"context"` Status string `json:"status"` Answer *string `json:"answer,omitempty"` CreatedAt string `json:"created_at"` AnsweredAt *string `json:"answered_at,omitempty"` } // getQuestions returns a task's questions (open and answered) newest-first — // REST hydration for the context panel's pinned question card and history. func (s *store) getQuestions(ctx context.Context, sessionID string) ([]sessionQuestion, error) { if s == nil { return nil, nil } rows, err := s.pool.Query(ctx, ` SELECT id::text, prompt, context, status, answer, created_at::text, answered_at::text FROM session_questions WHERE session_id = $1 ORDER BY created_at DESC`, sessionID) if err != nil { return nil, err } defer rows.Close() var out []sessionQuestion for rows.Next() { var q sessionQuestion var ctxJSON []byte var answer, answeredAt *string if err := rows.Scan(&q.ID, &q.Prompt, &ctxJSON, &q.Status, &answer, &q.CreatedAt, &answeredAt); err != nil { return nil, err } json.Unmarshal(ctxJSON, &q.Context) q.Answer, q.AnsweredAt = answer, answeredAt out = append(out, q) } return out, rows.Err() } // askOperator records a structured decision the agent needs from the operator, // moves the task to awaiting_input, and emits question.raised so the context // panel pins it. qctx carries {why, options, entities}. Returns the question id. func (s *store) askOperator(ctx context.Context, sessionID, prompt string, qctx map[string]any) (string, error) { if s == nil || sessionID == "" || sessionID == "ephemeral" { return "", nil } ctxJSON, _ := json.Marshal(qctx) var qid uuid.UUID if err := s.pool.QueryRow(ctx, ` INSERT INTO session_questions (session_id, prompt, context) VALUES ($1, $2, $3) RETURNING id`, sessionID, prompt, string(ctxJSON)).Scan(&qid); err != nil { return "", err } s.pool.Exec(ctx, `UPDATE agent_sessions SET status = 'awaiting_input', last_active_at = now() WHERE id = $1`, sessionID) data := map[string]any{"question_id": qid.String(), "prompt": prompt} for k, v := range qctx { data[k] = v } _ = observability.Event(ctx, sqlcgen.New(s.pool), "question.raised", s.taskEntityPtr(ctx, sessionID), "warning", "nomos", sessionID, data) return qid.String(), nil } // openQuestionID returns the id of the session's open question, or "". Used to // auto-close a pending question when the operator answers via a plain chat reply. func (s *store) openQuestionID(ctx context.Context, sessionID string) string { if s == nil || sessionID == "" || sessionID == "ephemeral" { return "" } var qid string s.pool.QueryRow(ctx, `SELECT id::text FROM session_questions WHERE session_id = $1 AND status = 'open' ORDER BY created_at DESC LIMIT 1`, sessionID).Scan(&qid) return qid } // getQuestion returns a question's prompt, answer, and session — used to build // the resume note when the operator answers via the panel. func (s *store) getQuestion(ctx context.Context, questionID string) (prompt, answer, sessionID string) { if s == nil || questionID == "" { return "", "", "" } qid, err := uuid.Parse(questionID) if err != nil { return "", "", "" } s.pool.QueryRow(ctx, `SELECT prompt, COALESCE(answer, ''), session_id::text FROM session_questions WHERE id = $1`, qid).Scan(&prompt, &answer, &sessionID) return } // answerQuestion records the operator's answer, returns the task to executing, // and emits question.answered. It does NOT itself resume the agent — the caller // decides: a chat reply IS the resuming turn, while a panel answer triggers a // continuation. func (s *store) answerQuestion(ctx context.Context, sessionID, questionID, answer string) error { if s == nil || sessionID == "" || sessionID == "ephemeral" || questionID == "" { return nil } qid, err := uuid.Parse(questionID) if err != nil { return err } if _, err := s.pool.Exec(ctx, ` UPDATE session_questions SET status = 'answered', answer = $2, answered_at = now() WHERE id = $1 AND status = 'open'`, qid, answer); err != nil { return err } s.pool.Exec(ctx, `UPDATE agent_sessions SET status = 'executing', last_active_at = now() WHERE id = $1`, sessionID) _ = observability.Event(ctx, sqlcgen.New(s.pool), "question.answered", s.taskEntityPtr(ctx, sessionID), "info", "nomos", sessionID, map[string]any{"question_id": questionID, "answer": answer}) return nil } // knowledgeSlugRe matches a nomos knowledge doc slug (:nomos/) as // printed in upsert_knowledge's result text. var knowledgeSlugRe = regexp.MustCompile(`[a-z]+:nomos/[a-z0-9-]+`) // linkKnowledgeToTask runs after a successful upsert_knowledge call within a // task: it links the created knowledge doc to the task entity (documents) so // get_relations(task) surfaces what the task learned, and publishes // knowledge.recorded for the live panel. Best-effort. The doc is ALSO linked to // the entity it's "about" by upsert_knowledge itself — that about-link is the // retrieval path future tasks use (get_entity_knowledge); this task-link is for // the task's own outcome/knowledge view. func (s *store) linkKnowledgeToTask(ctx context.Context, sessionID, resultText string) { if s == nil || sessionID == "" || sessionID == "ephemeral" { return } slug := knowledgeSlugRe.FindString(resultText) if slug == "" { return } var taskID, docID uuid.UUID if err := s.pool.QueryRow(ctx, `SELECT entity_id FROM agent_sessions WHERE id = $1`, sessionID).Scan(&taskID); err != nil || taskID == uuid.Nil { return } if err := s.pool.QueryRow(ctx, `SELECT id FROM entities WHERE slug = $1`, slug).Scan(&docID); err != nil { return } s.pool.Exec(ctx, ` INSERT INTO relationships (source_id, target_id, type, attributes, valid_from) SELECT $1, $2, 'documents', '{"by":"nomos"}'::jsonb, now() WHERE NOT EXISTS ( SELECT 1 FROM relationships WHERE source_id = $1 AND target_id = $2 AND type = 'documents' AND valid_to IS NULL)`, docID, taskID) _ = observability.Event(ctx, sqlcgen.New(s.pool), "knowledge.recorded", &docID, "info", "nomos", sessionID, map[string]any{"slug": slug}) } func (s *store) updateSessionTitle(ctx context.Context, id, title string) error { if s == nil { return nil } _, err := s.pool.Exec(ctx, `UPDATE agent_sessions SET title = $1 WHERE id = $2`, title, id) return err } // resolveAgentID looks up the UUID of the agent entity (e.g. "agent:nomos"). // Returns uuid.Nil if the store is absent or the slug is unknown. func (s *store) resolveAgentID(ctx context.Context, slug string) uuid.UUID { if s == nil { return uuid.Nil } var id uuid.UUID if err := s.pool.QueryRow(ctx, `SELECT id FROM entities WHERE slug = $1`, slug).Scan(&id); err != nil { return uuid.Nil } return id } // linkExecution records that a gated execution was initiated by a chat // session, so the auto-continuation worker can feed its result back to that // session when it finishes. Idempotent — the same execution may appear in // several tool results across a turn. func (s *store) linkExecution(ctx context.Context, execID uuid.UUID, sessionID string) { if s == nil || execID == uuid.Nil || sessionID == "" || sessionID == "ephemeral" { return } s.pool.Exec(ctx, ` INSERT INTO nomos_plan_executions (execution_id, session_id) VALUES ($1, $2) ON CONFLICT (execution_id) DO NOTHING`, execID, sessionID) } // taskSlugRe matches an entity slug: a lowercase type prefix then colon- // separated segments (host:hubris, lxc:caddy, check:ping:8cf). Mirrors the // frontend SessionGraph regex so the panel and the involves-graph agree on // what counts as an entity reference. var taskSlugRe = regexp.MustCompile(`[a-z][a-z-]*:[a-z0-9][a-z0-9._/-]*(?::[a-z0-9._/-]+)*`) // touchExcludedTypes are entity types too noisy to record as task involvement: // a health question names dozens of check:… slugs, executions/tasks are // bookkeeping, not things the task "worked on". var touchExcludedTypes = map[string]bool{"check": true, "execution": true, "task": true} // recordTouched links the task to every entity referenced in a tool call's // args (task —involves→ entity) and publishes one entity.touched event per // entity so the live context panel can pulse it. Best-effort: it never blocks // or fails the tool call. Only args are inspected — what the agent chose to act // on — never results, since a single bulk query result would otherwise pull the // whole fleet into the task's graph. func (s *store) recordTouched(ctx context.Context, sessionID, toolName string, args map[string]any) { if s == nil || sessionID == "" || sessionID == "ephemeral" || len(args) == 0 { return } slugs := map[string]struct{}{} collectTaskSlugs(args, slugs) if len(slugs) == 0 { return } var taskEntityID uuid.UUID if err := s.pool.QueryRow(ctx, `SELECT entity_id FROM agent_sessions WHERE id = $1`, sessionID).Scan(&taskEntityID); err != nil || taskEntityID == uuid.Nil { return // no task entity to anchor edges on } // One batched lookup instead of a SELECT per slug — a tool call naming // several entities (e.g. a multi-target comparison) used to issue N // round-trips here for N slugs found in its args. slugList := make([]string, 0, len(slugs)) for slug := range slugs { slugList = append(slugList, slug) } rows, err := s.pool.Query(ctx, `SELECT id, type, slug FROM entities WHERE slug = ANY($1)`, slugList) if err != nil { return } type found struct { id uuid.UUID etype string } matched := make(map[string]found, len(slugList)) for rows.Next() { var f found var slug string if rows.Scan(&f.id, &f.etype, &slug) == nil { matched[slug] = f } } rows.Close() if err := rows.Err(); err != nil { return } q := sqlcgen.New(s.pool) for slug, f := range matched { if touchExcludedTypes[f.etype] || f.id == taskEntityID { continue } // Idempotent involves edge (task → entity), same guard as upsert_knowledge. s.pool.Exec(ctx, ` INSERT INTO relationships (source_id, target_id, type, attributes, valid_from) SELECT $1, $2, 'involves', '{"by":"nomos"}'::jsonb, now() WHERE NOT EXISTS ( SELECT 1 FROM relationships WHERE source_id = $1 AND target_id = $2 AND type = 'involves' AND valid_to IS NULL)`, taskEntityID, f.id) // Live pulse for the panel. correlation_id = sessionID lets the frontend // filter to the active task. _ = observability.Event(ctx, q, "entity.touched", &f.id, "info", "nomos", sessionID, map[string]any{"slug": slug, "tool": toolName}) } } // collectTaskSlugs recursively pulls entity slugs out of tool-call args, // mirroring the frontend's collectSlugs so both sides see the same references. func collectTaskSlugs(v any, out map[string]struct{}) { switch t := v.(type) { case string: for _, m := range taskSlugRe.FindAllString(t, -1) { out[strings.TrimRight(m, ".,;)]")] = struct{}{} } case []any: for _, e := range t { collectTaskSlugs(e, out) } case map[string]any: for _, e := range t { collectTaskSlugs(e, out) } } } // pendingContinuation is one finished execution whose result hasn't yet been // fed back to its originating session. type pendingContinuation struct { ExecID uuid.UUID SessionID string Status string Result string Action string } // pendingContinuations returns executions that have reached a terminal state // but haven't been continued yet — the worker's work list. Bounded so one // tick can't fan out unboundedly. func (s *store) pendingContinuations(ctx context.Context, limit int) []pendingContinuation { if s == nil { return nil } rows, err := s.pool.Query(ctx, ` SELECT l.execution_id, l.session_id, e.status, COALESCE(e.result::text, ''), COALESCE(e.action, '') FROM nomos_plan_executions l JOIN executions e ON e.entity_id = l.execution_id WHERE l.continued_at IS NULL AND e.status IN ('completed', 'failed', 'cancelled', 'denied', 'revoked') ORDER BY l.created_at LIMIT $1`, limit) if err != nil { return nil } defer rows.Close() var out []pendingContinuation for rows.Next() { var p pendingContinuation if err := rows.Scan(&p.ExecID, &p.SessionID, &p.Status, &p.Result, &p.Action); err == nil { out = append(out, p) } } return out } // markContinued stamps an execution as fed-back so the worker won't process it // again (prevents an auto-continuation loop). func (s *store) markContinued(ctx context.Context, execID uuid.UUID) { if s == nil { return } s.pool.Exec(ctx, `UPDATE nomos_plan_executions SET continued_at = now() WHERE execution_id = $1`, execID) } // assentWindowActive reports whether THIS TASK currently has an open assent // window — the scope gate for auto-continuation. Scoped by session, not just // agent: with a single agent:nomos entity serving every concurrent task, an // agent-only key would let approving Task A's plan silently auto-run // unapproved config-mutation actions in a concurrently-running Task B. We // only auto-continue executions that are part of THIS session's approved // plan, never a stray action from another task riding the same window. func (s *store) assentWindowActive(ctx context.Context, agentID uuid.UUID, sessionID string) bool { if s == nil || agentID == uuid.Nil || sessionID == "" { return false // fail closed: no session to scope to means no window } var expires time.Time key := assentWindowKey(agentID, sessionID) if err := s.pool.QueryRow(ctx, `SELECT value::timestamptz FROM autonomy_settings WHERE key = $1`, key).Scan(&expires); err != nil { return false } return time.Now().Before(expires) } // assentWindowKey scopes the grant to one agent AND one session/task — see // assentWindowActive. Must match internal/mcp/server.go's copy (mirrored // there, not shared, since the two are separate Go packages/binaries reading // the same autonomy_settings row). func assentWindowKey(agentID uuid.UUID, sessionID string) string { return "assent_window.agent:" + agentID.String() + ".session:" + sessionID } // destructiveWindowDuration is intentionally shorter than the general assent // window (30 min): it's a narrow, scoped grant for a multi-step DESTRUCTIVE // recovery (e.g. "stop then destroy this specific half-provisioned // container"), not a standing license to destroy things. const destructiveWindowDuration = 15 * time.Minute // destructiveWindowKey scopes the grant to one agent, one target entity, AND // one session/task — an explicit typed confirmation ("I confirm") for a // destructive action on target X in task A must never be read as authorizing // a destructive action on target X from a DIFFERENT concurrently-running // task B, even though both share the same agent identity. func destructiveWindowKey(agentID uuid.UUID, targetSlug, sessionID string) string { return "destructive_window.agent:" + agentID.String() + ".target:" + targetSlug + ".session:" + sessionID } // openDestructiveWindow records a short, target-and-session-scoped grant // after an operator's EXPLICIT typed confirmation (never loose assent) // authorized a destructive action. Real case this exists for: recovering a // failed destroy took "stop" (destructive) then "destroy" (destructive) — // same container, two separate typed-confirmation round trips, because each // was gated independently. One explicit confirmation on a target should // cover the short follow-up sequence needed to finish what was just // confirmed — but only within the task that got the confirmation. func (s *store) openDestructiveWindow(ctx context.Context, agentID uuid.UUID, targetSlug, sessionID string) { if s == nil || agentID == uuid.Nil || targetSlug == "" || sessionID == "" { return } expires := time.Now().Add(destructiveWindowDuration).UTC().Format(time.RFC3339) s.pool.Exec(ctx, `INSERT INTO autonomy_settings (key, value) VALUES ($1, $2) ON CONFLICT (key) DO UPDATE SET value = $2`, destructiveWindowKey(agentID, targetSlug, sessionID), expires) } // destructiveWindowActive reports whether target has a live, explicitly- // confirmed destructive grant for this agent within this session/task. func (s *store) destructiveWindowActive(ctx context.Context, agentID uuid.UUID, targetSlug, sessionID string) bool { if s == nil || agentID == uuid.Nil || targetSlug == "" || sessionID == "" { return false } var expires time.Time if err := s.pool.QueryRow(ctx, `SELECT value::timestamptz FROM autonomy_settings WHERE key = $1`, destructiveWindowKey(agentID, targetSlug, sessionID)).Scan(&expires); err != nil { return false } return time.Now().Before(expires) } // executionTarget resolves the target entity slug for an execution — used to // scope the destructive window to the right entity when a chat-assent typed // confirmation grants a destructive execution. func (s *store) executionTarget(ctx context.Context, execID uuid.UUID) string { if s == nil { return "" } var slug string s.pool.QueryRow(ctx, ` SELECT e.slug FROM executions ex JOIN entities e ON e.id = ex.target_entity_id WHERE ex.entity_id = $1`, execID).Scan(&slug) return slug } // logActivity records a tool call. agent_id is the agent entity UUID and is // NOT NULL in the schema, so we skip logging when it can't be resolved. // The (nullable) session_id column carries the conversation id. func (s *store) logActivity(ctx context.Context, agentID uuid.UUID, sessionID, toolName, inputSummary, outputSummary string, durationMs int, success bool, correlationID string) { if s == nil || agentID == uuid.Nil { return } s.pool.Exec(ctx, ` INSERT INTO agent_activity (agent_id, session_id, activity_type, tool_name, input_summary, output_summary, duration_ms, success, correlation_id) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9)`, agentID, sessionID, "tool_call", toolName, inputSummary, outputSummary, durationMs, success, correlationID) }