Phase 1 — crash recovery: SSE auto-reconnect + backoff, polling gate during disconnect, connection banner with retry button, empty-response retry 3x, non-terminal resume on empty response, persistent error cards. Phase 2/4 — visibility + continuation: custom ExecutionStatus renderer, approvals extracted on every tool_result (not just done), activity bar with status/goal, SessionDigest live polling, Continue button. Phase 3 — cleanup: complete_task auto-cancels orphaned approvals, deletes assent/destructive window keys, propose_plan marks pending steps as replaced, plan step seq-order enforcement. Phase 5 — knowledge loop: list_lxcs state filter (active/destroyed), SOUL.md unmissable writeback section, propose_plan validation nudge, complete_task writeback check, upsert_knowledge about array support, plan generation grouping in frontend, session approval count badge. Retire request_execution — all mutations now route through run. Updated SOUL.md, AGENTS.md, CLIENTS.md, skills, and agent system notes. Migration 020: plan step generation column, audit_log session_id index, nomos_plan_executions pending-approval index.
307 lines
13 KiB
Go
307 lines
13 KiB
Go
package main
|
||
|
||
import (
|
||
"context"
|
||
"fmt"
|
||
"log/slog"
|
||
"strings"
|
||
)
|
||
|
||
// Task tools are nomos-LOCAL, not MCP tools. They are session-scoped, and the
|
||
// shared MCP server (api:8090/mcp) has no session id — so these are handled
|
||
// in-process by nomos, which knows the session/task and holds the store.
|
||
// buildTools appends these to the model's tool list; the agent loop routes a
|
||
// call whose name isTaskTool to handleTaskTool instead of the MCP client.
|
||
//
|
||
// Phase 3 ships complete_task; set_goal / propose_plan / update_plan_step /
|
||
// ask_operator land in later phases through the same mechanism.
|
||
|
||
func taskToolDefs() []toolDef {
|
||
return []toolDef{
|
||
{
|
||
Name: "set_goal",
|
||
Description: "State the goal of this task in one sentence, as early as you " +
|
||
"can. This is what the task is trying to achieve (e.g. 'Deploy TypeType " +
|
||
"as an LXC on strong'); it heads the task on the board and the context " +
|
||
"panel. Call it once you understand what the operator wants.",
|
||
InputSchema: map[string]any{
|
||
"type": "object",
|
||
"properties": map[string]any{
|
||
"goal": map[string]any{"type": "string", "description": "The task's goal, one sentence."},
|
||
},
|
||
"required": []string{"goal"},
|
||
},
|
||
},
|
||
{
|
||
Name: "propose_plan",
|
||
Description: "Lay out ALL the ordered steps you'll take to reach the goal, in ONE " +
|
||
"call, listing every step end-to-end — not just the next one. The operator " +
|
||
"sees the full list in the context panel and watches it progress; a plan " +
|
||
"with only 1 step looks broken to them even if you intend to add more later. " +
|
||
"Your FIRST step should be research (prior knowledge, relations, blast radius " +
|
||
"— not just this target's status) and your LAST step should be writing back " +
|
||
"what you learned (update_entity_attributes / create_relationship / " +
|
||
"upsert_knowledge) BEFORE complete_task — this is what keeps the knowledge " +
|
||
"graph from drifting out of date. " +
|
||
"Call this ONCE, before you start executing (after gathering what you need). " +
|
||
"As you work, call update_plan_step (not propose_plan again) to advance each " +
|
||
"step. Only re-call propose_plan if the plan itself has fundamentally changed " +
|
||
"(e.g. a new approach is needed) — in that case new steps are appended after " +
|
||
"whatever already ran, never erasing completed work.",
|
||
InputSchema: map[string]any{
|
||
"type": "object",
|
||
"properties": map[string]any{
|
||
"steps": map[string]any{
|
||
"type": "array",
|
||
"description": "Ordered steps, first to last.",
|
||
"items": map[string]any{
|
||
"type": "object",
|
||
"properties": map[string]any{
|
||
"title": map[string]any{"type": "string", "description": "Short imperative step title (e.g. 'Create the LXC')."},
|
||
"detail": map[string]any{"type": "string", "description": "Optional one-line detail."},
|
||
"target_slug": map[string]any{"type": "string", "description": "Optional entity slug this step acts on (e.g. lxc:typetype)."},
|
||
},
|
||
"required": []string{"title"},
|
||
},
|
||
},
|
||
},
|
||
"required": []string{"steps"},
|
||
},
|
||
},
|
||
{
|
||
Name: "update_plan_step",
|
||
Description: "Advance a plan step as you work it. Set status to 'running' when " +
|
||
"you start it (pass execution_id if the step queued a gated action, so " +
|
||
"the board can auto-close it when that finishes), then 'done' / 'failed' " +
|
||
"/ 'skipped' / 'blocked' when it resolves. Keeps the operator's progress " +
|
||
"view honest.",
|
||
InputSchema: map[string]any{
|
||
"type": "object",
|
||
"properties": map[string]any{
|
||
"seq": map[string]any{"type": "integer", "description": "1-based step number from propose_plan."},
|
||
"status": map[string]any{"type": "string", "enum": []string{"running", "done", "failed", "skipped", "blocked"}, "description": "New status for the step."},
|
||
"execution_id": map[string]any{"type": "string", "description": "Optional execution UUID this step is running, so it auto-closes on completion."},
|
||
},
|
||
"required": []string{"seq", "status"},
|
||
},
|
||
},
|
||
{
|
||
Name: "ask_operator",
|
||
Description: "Ask the operator a question when you hit a real decision only " +
|
||
"they can make — an ambiguous target, a trade-off, missing information, " +
|
||
"or a destructive choice not already approved. This pins a structured " +
|
||
"question card in the context panel (with your options and the entities " +
|
||
"involved) and PAUSES the task until they answer; their answer resumes " +
|
||
"you automatically. Do NOT use it for things you can determine yourself " +
|
||
"with tools — only for genuine decisions.",
|
||
InputSchema: map[string]any{
|
||
"type": "object",
|
||
"properties": map[string]any{
|
||
"prompt": map[string]any{"type": "string", "description": "The question, stated plainly."},
|
||
"why": map[string]any{"type": "string", "description": "Why you're asking / what's at stake."},
|
||
"options": map[string]any{
|
||
"type": "array", "items": map[string]any{"type": "string"},
|
||
"description": "The choices, if it's a pick-one decision.",
|
||
},
|
||
"context_entities": map[string]any{
|
||
"type": "array", "items": map[string]any{"type": "string"},
|
||
"description": "Entity slugs relevant to the decision (shown as chips).",
|
||
},
|
||
},
|
||
"required": []string{"prompt"},
|
||
},
|
||
},
|
||
{
|
||
Name: "complete_task",
|
||
Description: "Mark the current task finished. Call this once the goal is " +
|
||
"verified done — or when you've genuinely failed or only partially " +
|
||
"succeeded. Sets the task's outcome and a one-line summary shown on the " +
|
||
"task board. Record what you learned with upsert_knowledge BEFORE " +
|
||
"completing, so future tasks on the same entities benefit.",
|
||
InputSchema: map[string]any{
|
||
"type": "object",
|
||
"properties": map[string]any{
|
||
"outcome": map[string]any{
|
||
"type": "string",
|
||
"enum": []string{"success", "failure", "partial"},
|
||
"description": "Did the task achieve its goal?",
|
||
},
|
||
"summary": map[string]any{
|
||
"type": "string",
|
||
"description": "One line describing the result (shown on the task card).",
|
||
},
|
||
},
|
||
"required": []string{"outcome", "summary"},
|
||
},
|
||
},
|
||
}
|
||
}
|
||
|
||
// toInt coerces a JSON tool-arg number (float64 after unmarshal) to int.
|
||
func toInt(v any) int {
|
||
switch n := v.(type) {
|
||
case float64:
|
||
return int(n)
|
||
case int:
|
||
return n
|
||
default:
|
||
return 0
|
||
}
|
||
}
|
||
|
||
// toStringSlice coerces a JSON tool-arg array to a non-empty []string.
|
||
func toStringSlice(v any) []string {
|
||
arr, ok := v.([]any)
|
||
if !ok {
|
||
return nil
|
||
}
|
||
out := make([]string, 0, len(arr))
|
||
for _, e := range arr {
|
||
if s, ok := e.(string); ok && strings.TrimSpace(s) != "" {
|
||
out = append(out, s)
|
||
}
|
||
}
|
||
return out
|
||
}
|
||
|
||
// handleTaskTool executes a nomos-local task tool. Returns (result, true) if it
|
||
// handled the call, or (nil, false) if name is not a local task tool (so the
|
||
// caller forwards it to the MCP client).
|
||
func (a *agent) handleTaskTool(ctx context.Context, sessionID, name string, args map[string]any) (any, bool) {
|
||
switch name {
|
||
case "set_goal":
|
||
goal, _ := args["goal"].(string)
|
||
if strings.TrimSpace(goal) == "" {
|
||
return "error: set_goal needs a goal", true
|
||
}
|
||
if err := a.store.setGoal(ctx, sessionID, goal); err != nil {
|
||
return fmt.Sprintf("error setting goal: %v", err), true
|
||
}
|
||
return "Goal set: " + goal, true
|
||
|
||
case "propose_plan":
|
||
raw, _ := args["steps"].([]any)
|
||
var steps []planStepInput
|
||
for _, r := range raw {
|
||
m, ok := r.(map[string]any)
|
||
if !ok {
|
||
continue
|
||
}
|
||
title, _ := m["title"].(string)
|
||
if strings.TrimSpace(title) == "" {
|
||
continue
|
||
}
|
||
detail, _ := m["detail"].(string)
|
||
target, _ := m["target_slug"].(string)
|
||
steps = append(steps, planStepInput{Title: title, Detail: detail, TargetSlug: target})
|
||
}
|
||
if len(steps) == 0 {
|
||
return "error: propose_plan needs at least one step with a title", true
|
||
}
|
||
persisted, err := a.store.proposePlan(ctx, sessionID, steps)
|
||
if err != nil {
|
||
return fmt.Sprintf("error proposing plan: %v", err), true
|
||
}
|
||
// Nudge: if the last step doesn't mention entity writeback tools,
|
||
// the graph will keep drifting — discovered facts won't be persisted.
|
||
lastStep := steps[len(steps)-1]
|
||
hasWriteback := strings.Contains(lastStep.Title+lastStep.Detail, "update_entity_attributes") ||
|
||
strings.Contains(lastStep.Title+lastStep.Detail, "create_relationship")
|
||
result := fmt.Sprintf("Plan set: %d step(s). Execute them now, marking each with update_plan_step as you go.", len(persisted))
|
||
if !hasWriteback {
|
||
result += "\n\n⚠️ The final step doesn't mention update_entity_attributes or create_relationship. Without those, any facts you discovered about entities (IPs, versions, hosts, states) will be LOST — the next session starts from scratch. Consider revising the last step to include entity writeback BEFORE completing the task."
|
||
}
|
||
return result, true
|
||
|
||
case "update_plan_step":
|
||
seq := toInt(args["seq"])
|
||
status, _ := args["status"].(string)
|
||
execID, _ := args["execution_id"].(string)
|
||
if seq <= 0 || status == "" {
|
||
return "error: update_plan_step needs seq (>=1) and status", true
|
||
}
|
||
if err := a.store.updatePlanStep(ctx, sessionID, seq, status, execID); err != nil {
|
||
return fmt.Sprintf("error updating step %d: %v", seq, err), true
|
||
}
|
||
return fmt.Sprintf("Step %d → %s", seq, status), true
|
||
|
||
case "ask_operator":
|
||
prompt, _ := args["prompt"].(string)
|
||
if strings.TrimSpace(prompt) == "" {
|
||
return "error: ask_operator needs a prompt", true
|
||
}
|
||
qctx := map[string]any{}
|
||
if why, _ := args["why"].(string); strings.TrimSpace(why) != "" {
|
||
qctx["why"] = why
|
||
}
|
||
if opts := toStringSlice(args["options"]); len(opts) > 0 {
|
||
qctx["options"] = opts
|
||
}
|
||
if ents := toStringSlice(args["context_entities"]); len(ents) > 0 {
|
||
qctx["entities"] = ents
|
||
}
|
||
if _, err := a.store.askOperator(ctx, sessionID, prompt, qctx); err != nil {
|
||
return fmt.Sprintf("error posting question: %v", err), true
|
||
}
|
||
return "Question posted to the operator; the task is paused until they answer. " +
|
||
"Do not continue or call more tools — end your turn now and wait for their answer.", true
|
||
|
||
case "complete_task":
|
||
outcome, _ := args["outcome"].(string)
|
||
summary, _ := args["summary"].(string)
|
||
switch outcome {
|
||
case "":
|
||
outcome = "success" // no outcome given at all — assume success, the common case
|
||
case "success", "failure", "partial":
|
||
// valid, use as-is
|
||
default:
|
||
// The tool schema declares an enum, but a weaker model (or a
|
||
// typo) can still send anything — an unrecognized value used to
|
||
// persist as-is, silently, with only "failure" special-cased
|
||
// (store.completeTask derives status='failed' from it; anything
|
||
// else became status='done' regardless of what the value
|
||
// actually said). Default to "partial" rather than silently
|
||
// treating an unrecognized value as "success" — safer to
|
||
// under-claim than over-claim a task's outcome.
|
||
slog.Warn("nomos: complete_task got an unrecognized outcome, defaulting to partial",
|
||
"session", sessionID, "outcome", outcome)
|
||
outcome = "partial"
|
||
}
|
||
if err := a.store.completeTask(ctx, sessionID, outcome, summary); err != nil {
|
||
return fmt.Sprintf("error completing task: %v", err), true
|
||
}
|
||
result := fmt.Sprintf("Task marked %s: %s", outcome, summary)
|
||
if !a.store.hadEntityWriteback(ctx, sessionID) {
|
||
result += "\n\n⚠️ No entity attributes or relationships were updated in this session. Call update_entity_attributes and create_relationship to persist what you learned about entities before the next session starts from scratch."
|
||
}
|
||
return result, true
|
||
default:
|
||
return nil, false
|
||
}
|
||
}
|
||
|
||
// autoCompleteTrivialTask is the case-1 fix from
|
||
// plans/2026-07-11-task-completion-safety-net.md: a session that never
|
||
// called set_goal never framed itself as a structured task, so a turn that
|
||
// ends with a plain-text answer and no further tool calls IS the task
|
||
// ending — but the model consistently skips complete_task for exactly this
|
||
// case (confirmed live: 43/50 production sessions were a single trivial
|
||
// Q&A exchange, none of which ever reached a terminal status). Rather than
|
||
// leave agent_sessions.status stuck at its creation-time default forever,
|
||
// close it out mechanically here: no judgment call needed, since SOUL.md
|
||
// already treats a one-shot answered question as done by definition.
|
||
func (a *agent) autoCompleteTrivialTask(ctx context.Context, sessionID, responseText string) {
|
||
summary := strings.TrimSpace(responseText)
|
||
summary = strings.SplitN(summary, "\n", 2)[0] // first line only — the board shows one line
|
||
const maxLen = 120
|
||
if len(summary) > maxLen {
|
||
summary = summary[:maxLen] + "…"
|
||
}
|
||
if summary == "" {
|
||
summary = "Answered without further action needed."
|
||
}
|
||
if err := a.store.completeTask(ctx, sessionID, "success", summary); err != nil {
|
||
slog.Error("nomos: auto-complete trivial task failed", "session", sessionID, "error", err)
|
||
}
|
||
}
|