diff --git a/internal/mcp/server.go b/internal/mcp/server.go index 6064e17..50919bb 100644 --- a/internal/mcp/server.go +++ b/internal/mcp/server.go @@ -302,6 +302,32 @@ func newServer(pool *db.Pool, agentID uuid.UUID) *mcp.Server { return textResult(fmt.Sprintf("target not found: %s", targetSlug)), nil } + // restart, pct_exec, and systemctl (outside enable/disable) route + // through the same classify→gate path as `run` instead of executing + // immediately over SSH with a hardcoded risk_class='reversible_low' + // that was never actually checked against anything. Found live + // 2026-07-10: a chat request to "restart caddy" — the fleet's + // reverse proxy — executed instantly with zero approval, because + // this action bypassed the classifier entirely. classifyAndGate + // applies the same read-only/config-mutation/destructive + // classification and approval flow the `run` tool already uses. + if action == "restart" || action == "pct_exec" || (action == "systemctl" && params != "enable" && params != "disable") { + svc := strings.TrimPrefix(targetSlug, "lxc:") + var cmd, purpose string + switch action { + case "restart": + cmd = fmt.Sprintf("systemctl restart %s; sleep 1; systemctl is-active %s", svc, svc) + purpose = "restart " + svc + case "pct_exec": + cmd = params + purpose = "pct_exec (legacy) on " + targetSlug + case "systemctl": + cmd = fmt.Sprintf("systemctl %s %s; sleep 1; systemctl is-active %s", params, svc, svc) + purpose = "systemctl " + params + " " + svc + } + return classifyAndGate(ctx, pool, agentID, targetID, targetSlug, cmd, purpose, ""), nil + } + // Deduplicate: if a pending execution already exists for the same // target+action, return the existing one instead of creating a // duplicate. Prevents the LLM from re-requesting the same gated @@ -338,67 +364,16 @@ func newServer(pool *db.Pool, agentID uuid.UUID) *mcp.Server { pool.Exec(ctx, `INSERT INTO executions (entity_id, target_entity_id, action, risk_class, status, correlation_id, agent_id) VALUES ($1, $2, $3, 'reversible_low', 'running', $4, $5) ON CONFLICT DO NOTHING`, id, targetID, action+":"+params, correlationID, agentID) - // Execute reversible actions immediately + // Execute reversible actions immediately. restart/pct_exec/systemctl + // (outside enable/disable) never reach here — they're routed through + // classifyAndGate above, before this dedup+insert block. switch action { - case "restart": - host, user, err := resolveHost(ctx, pool, targetSlug) - if err != nil { - return textResult(fmt.Sprintf("resolve: %v", err)), nil - } - svc := strings.TrimPrefix(targetSlug, "lxc:") - out, err := sshExec(ctx, host, user, fmt.Sprintf("systemctl restart %s 2>&1; sleep 1; systemctl is-active %s", svc, svc)) - result := fmt.Sprintf("restart %s: %s", svc, out) - if err != nil { - result = fmt.Sprintf("restart %s: ERROR %v", svc, err) - } - pool.Exec(ctx, `UPDATE executions SET status='completed', result=$2::jsonb WHERE entity_id=$1`, - id, jsonOut(out)) - return textResult(result), nil - case "systemctl": + // Only enable/disable reach this case now. svc := strings.TrimPrefix(targetSlug, "lxc:") - if params == "enable" || params == "disable" { - pool.Exec(ctx, `UPDATE executions SET status='pending_approval', risk_class='config_mutation' WHERE entity_id=$1`, id) - createApproval(ctx, pool, id, targetID, "systemctl", svc+":"+params, "config_mutation") - return textResult(fmt.Sprintf("systemctl %s on %s requires approval — execution %s queued", params, svc, id)), nil - } - host, user, err := resolveHost(ctx, pool, targetSlug) - if err != nil { - return textResult(fmt.Sprintf("resolve: %v", err)), nil - } - cmd := fmt.Sprintf("systemctl %s %s 2>&1; sleep 1; systemctl is-active %s", params, svc, svc) - out, err := sshExec(ctx, host, user, cmd) - result := fmt.Sprintf("systemctl %s %s: %s", params, svc, out) - if err != nil { - result = fmt.Sprintf("systemctl %s %s: ERROR %v", params, svc, err) - } - pool.Exec(ctx, `UPDATE executions SET status='completed', result=$2::jsonb WHERE entity_id=$1`, - id, jsonOut(out)) - return textResult(result), nil - - case "pct_exec": - var pveID string - if err := pool.QueryRow(ctx, "SELECT attributes->>'pve_id' FROM entities WHERE slug = $1", targetSlug).Scan(&pveID); err != nil || pveID == "" { - return textResult(fmt.Sprintf("LXC not found: %s", targetSlug)), nil - } - // Resolve Proxmox host - var hostSlug string - pool.QueryRow(ctx, "SELECT attributes->>'host' FROM entities WHERE slug = $1", targetSlug).Scan(&hostSlug) - if hostSlug == "" { - hostSlug = "host:hubris" // default - } - host, user, err := resolveHost(ctx, pool, hostSlug) - if err != nil { - return textResult(fmt.Sprintf("resolve Proxmox host: %v", err)), nil - } - out, err := sshExec(ctx, host, user, fmt.Sprintf("pct exec %s -- %s 2>&1", pveID, params)) - result := fmt.Sprintf("pct exec %s: %s", pveID, out) - if err != nil { - result = fmt.Sprintf("pct exec %s: ERROR %v", pveID, err) - } - pool.Exec(ctx, `UPDATE executions SET status='completed', result=$2::jsonb WHERE entity_id=$1`, - id, jsonOut(out)) - return textResult(result), nil + pool.Exec(ctx, `UPDATE executions SET status='pending_approval', risk_class='config_mutation' WHERE entity_id=$1`, id) + createApproval(ctx, pool, id, targetID, "systemctl", svc+":"+params, "config_mutation") + return textResult(fmt.Sprintf("systemctl %s on %s requires approval — execution %s queued", params, svc, id)), nil case "apt_upgrade": if params == "audit" { @@ -488,106 +463,7 @@ func newServer(pool *db.Pool, agentID uuid.UUID) *mcp.Server { return textResult(fmt.Sprintf("target not found: %s", targetSlug)), nil } - riskClass := policy.ClassifyCommand(command, declaredRisk) - runParams, _ := json.Marshal(map[string]string{"command": command, "purpose": purpose}) - actionCol := "run:" + string(runParams) - - // Dedup: an identical pending command (same target, command, and - // purpose) blocks a re-request — stops a tool-calling loop from - // queuing the same approval repeatedly. - var existingID string - derr := pool.QueryRow(ctx, ` - SELECT e.id::text FROM entities e - JOIN executions ex ON ex.entity_id = e.id - WHERE e.type = 'execution' AND ex.target_entity_id = $1 - AND ex.action = $2 AND ex.status = 'pending_approval' - ORDER BY e.created_at DESC LIMIT 1`, - targetID, actionCol).Scan(&existingID) - if derr == nil && existingID != "" { - return textResult(fmt.Sprintf("An identical command is already queued for approval on %s — execution %s. Wait for the operator, don't re-request.", targetSlug, existingID)), nil - } - - id, _ := uuid.NewV7() - correlationID := uuid.New().String() - // Full UUID, not a truncated prefix — see the matching comment on - // request_execution's exec slug generation above; the 8-char prefix - // collided for real under back-to-back requests. - execName := "run on " + targetSlug + " (" + id.String() + ")" - execSlug := "exec:" + targetSlug + ":" + id.String() - if _, err := pool.Exec(ctx, `INSERT INTO entities (id, slug, type, name, attributes) VALUES ($1, $2, 'execution', $3, '{}')`, - id, execSlug, execName); err != nil { - return textResult(fmt.Sprintf("error: failed to create execution: %v", err)), nil - } - pool.Exec(ctx, `INSERT INTO executions (entity_id, target_entity_id, action, risk_class, status, correlation_id, agent_id) VALUES ($1, $2, $3, $4, 'running', $5, $6) ON CONFLICT DO NOTHING`, - id, targetID, actionCol, riskClass, correlationID, agentID) - - if riskClass == policy.RiskReadOnly { - host, user, wrap, rerr := resolveExecTarget(ctx, pool, targetSlug) - if rerr != nil { - pool.Exec(ctx, `UPDATE executions SET status='failed', result=$2::jsonb WHERE entity_id=$1`, id, jsonErr("%s", rerr.Error())) - return textResult(fmt.Sprintf("resolve target: %v", rerr)), nil - } - out, xerr := sshExec(ctx, host, user, wrap(command)) - if xerr != nil { - pool.Exec(ctx, `UPDATE executions SET status='failed', result=$2::jsonb WHERE entity_id=$1`, id, jsonErr("%s: %s", xerr.Error(), out)) - return textResult(fmt.Sprintf("run on %s: ERROR %v\n%s", targetSlug, xerr, out)), nil - } - pool.Exec(ctx, `UPDATE executions SET status='completed', result=$2::jsonb WHERE entity_id=$1`, id, jsonOut(out)) - return textResult(fmt.Sprintf("run on %s (read_only, auto): %s", targetSlug, out)), nil - } - - // Assent window: if the operator recently approved a plan in this - // agent's chat session, config_mutation commands auto-run without - // re-approval. This is the "approve the plan, carry it out" path — - // the operator approved the overall direction; individual config - // steps within the window don't each need a separate yes. - // Destructive commands never auto-run, regardless of window. - if riskClass == policy.RiskConfigMutation && assentWindowActive(ctx, pool, agentID) { - host, user, wrap, rerr := resolveExecTarget(ctx, pool, targetSlug) - if rerr != nil { - pool.Exec(ctx, `UPDATE executions SET status='failed', result=$2::jsonb WHERE entity_id=$1`, id, jsonErr("%s", rerr.Error())) - return textResult(fmt.Sprintf("resolve target: %v", rerr)), nil - } - out, xerr := sshExec(ctx, host, user, wrap(command)) - if xerr != nil { - pool.Exec(ctx, `UPDATE executions SET status='failed', result=$2::jsonb WHERE entity_id=$1`, id, jsonErr("%s: %s", xerr.Error(), out)) - return textResult(fmt.Sprintf("run on %s: ERROR %v\n%s", targetSlug, xerr, out)), nil - } - pool.Exec(ctx, `UPDATE executions SET status='completed', result=$2::jsonb WHERE entity_id=$1`, id, jsonOut(out)) - slog.Info("mcp: run auto-executed via assent window", "target", targetSlug, "execution_id", id) - return textResult(fmt.Sprintf("run on %s (config_mutation, auto via assent window): %s", targetSlug, out)), nil - } - - // Destructive window: a narrow, TARGET-scoped grant opened only after - // an operator's explicit typed confirmation ("I confirm") on this - // same target — never by loose assent. Exists for multi-step - // destructive recovery (e.g. a failed destroy needing stop, then - // destroy) so the operator isn't asked to re-type "I confirm" for - // every single command against the thing they just confirmed. - if riskClass == policy.RiskDestructive && destructiveWindowActive(ctx, pool, agentID, targetSlug) { - host, user, wrap, rerr := resolveExecTarget(ctx, pool, targetSlug) - if rerr != nil { - pool.Exec(ctx, `UPDATE executions SET status='failed', result=$2::jsonb WHERE entity_id=$1`, id, jsonErr("%s", rerr.Error())) - return textResult(fmt.Sprintf("resolve target: %v", rerr)), nil - } - out, xerr := sshExec(ctx, host, user, wrap(command)) - if xerr != nil { - pool.Exec(ctx, `UPDATE executions SET status='failed', result=$2::jsonb WHERE entity_id=$1`, id, jsonErr("%s: %s", xerr.Error(), out)) - return textResult(fmt.Sprintf("run on %s: ERROR %v\n%s", targetSlug, xerr, out)), nil - } - pool.Exec(ctx, `UPDATE executions SET status='completed', result=$2::jsonb WHERE entity_id=$1`, id, jsonOut(out)) - slog.Info("mcp: run auto-executed via destructive window", "target", targetSlug, "execution_id", id) - return textResult(fmt.Sprintf("run on %s (destructive, auto via confirmed-target window): %s", targetSlug, out)), nil - } - - pool.Exec(ctx, `UPDATE executions SET status='pending_approval', risk_class=$2 WHERE entity_id=$1`, id, riskClass) - createApproval(ctx, pool, id, targetID, "run", string(runParams), riskClass) - confirmNote := "" - if riskClass == policy.RiskDestructive { - confirmNote = " This is classified DESTRUCTIVE — flag that clearly to the operator; it needs explicit confirmation, not just a casual \"go ahead\"." - } - return textResult(fmt.Sprintf("run on %s requires approval (risk: %s) — execution %s queued.%s Present the command and purpose to the operator and wait; do not re-request.", - targetSlug, riskClass, id, confirmNote)), nil + return classifyAndGate(ctx, pool, agentID, targetID, targetSlug, command, purpose, declaredRisk), nil }) register(&mcp.Tool{Name: "http_get", Description: "Fetch a public web page or raw file (e.g. a GitHub README/raw URL) and return sanitized text. Use this to research how to deploy a service before provisioning. HTTP/HTTPS only; body is truncated to ~16KB.", @@ -1360,6 +1236,115 @@ func resolveExecTarget(ctx context.Context, pool *db.Pool, targetSlug string) (h return "", "", nil, fmt.Errorf("unsupported target %q: must be host: or lxc:", targetSlug) } +// classifyAndGate is the shared classify→execute-or-queue path for every +// mutating command, used by both the general `run` tool and +// request_execution's restart/systemctl/pct_exec actions. Those legacy +// actions used to execute immediately over SSH with a hardcoded +// risk_class='reversible_low' that was never actually evaluated against the +// command — found live 2026-07-10 when a chat request to restart caddy (the +// fleet's reverse proxy) executed instantly with no approval at all. Routing +// every mutating path through the same classifier + approval-queue logic +// closes that gap without special-casing each caller. +func classifyAndGate(ctx context.Context, pool *db.Pool, agentID, targetID uuid.UUID, targetSlug, command, purpose, declaredRisk string) *mcp.CallToolResult { + riskClass := policy.ClassifyCommand(command, declaredRisk) + runParams, _ := json.Marshal(map[string]string{"command": command, "purpose": purpose}) + actionCol := "run:" + string(runParams) + + // Dedup: an identical pending command (same target, command, and + // purpose) blocks a re-request — stops a tool-calling loop from queuing + // the same approval repeatedly. + var existingID string + derr := pool.QueryRow(ctx, ` + SELECT e.id::text FROM entities e + JOIN executions ex ON ex.entity_id = e.id + WHERE e.type = 'execution' AND ex.target_entity_id = $1 + AND ex.action = $2 AND ex.status = 'pending_approval' + ORDER BY e.created_at DESC LIMIT 1`, + targetID, actionCol).Scan(&existingID) + if derr == nil && existingID != "" { + return textResult(fmt.Sprintf("An identical command is already queued for approval on %s — execution %s. Wait for the operator, don't re-request.", targetSlug, existingID)) + } + + id, _ := uuid.NewV7() + correlationID := uuid.New().String() + execName := "run on " + targetSlug + " (" + id.String() + ")" + execSlug := "exec:" + targetSlug + ":" + id.String() + if _, err := pool.Exec(ctx, `INSERT INTO entities (id, slug, type, name, attributes) VALUES ($1, $2, 'execution', $3, '{}')`, + id, execSlug, execName); err != nil { + return textResult(fmt.Sprintf("error: failed to create execution: %v", err)) + } + pool.Exec(ctx, `INSERT INTO executions (entity_id, target_entity_id, action, risk_class, status, correlation_id, agent_id) VALUES ($1, $2, $3, $4, 'running', $5, $6) ON CONFLICT DO NOTHING`, + id, targetID, actionCol, riskClass, correlationID, agentID) + + if riskClass == policy.RiskReadOnly { + host, user, wrap, rerr := resolveExecTarget(ctx, pool, targetSlug) + if rerr != nil { + pool.Exec(ctx, `UPDATE executions SET status='failed', result=$2::jsonb WHERE entity_id=$1`, id, jsonErr("%s", rerr.Error())) + return textResult(fmt.Sprintf("resolve target: %v", rerr)) + } + out, xerr := sshExec(ctx, host, user, wrap(command)) + if xerr != nil { + pool.Exec(ctx, `UPDATE executions SET status='failed', result=$2::jsonb WHERE entity_id=$1`, id, jsonErr("%s: %s", xerr.Error(), out)) + return textResult(fmt.Sprintf("run on %s: ERROR %v\n%s", targetSlug, xerr, out)) + } + pool.Exec(ctx, `UPDATE executions SET status='completed', result=$2::jsonb WHERE entity_id=$1`, id, jsonOut(out)) + return textResult(fmt.Sprintf("run on %s (read_only, auto): %s", targetSlug, out)) + } + + // Assent window: if the operator recently approved a plan in this + // agent's chat session, config_mutation commands auto-run without + // re-approval. This is the "approve the plan, carry it out" path — the + // operator approved the overall direction; individual config steps + // within the window don't each need a separate yes. Destructive + // commands never auto-run, regardless of window. + if riskClass == policy.RiskConfigMutation && assentWindowActive(ctx, pool, agentID) { + host, user, wrap, rerr := resolveExecTarget(ctx, pool, targetSlug) + if rerr != nil { + pool.Exec(ctx, `UPDATE executions SET status='failed', result=$2::jsonb WHERE entity_id=$1`, id, jsonErr("%s", rerr.Error())) + return textResult(fmt.Sprintf("resolve target: %v", rerr)) + } + out, xerr := sshExec(ctx, host, user, wrap(command)) + if xerr != nil { + pool.Exec(ctx, `UPDATE executions SET status='failed', result=$2::jsonb WHERE entity_id=$1`, id, jsonErr("%s: %s", xerr.Error(), out)) + return textResult(fmt.Sprintf("run on %s: ERROR %v\n%s", targetSlug, xerr, out)) + } + pool.Exec(ctx, `UPDATE executions SET status='completed', result=$2::jsonb WHERE entity_id=$1`, id, jsonOut(out)) + slog.Info("mcp: run auto-executed via assent window", "target", targetSlug, "execution_id", id) + return textResult(fmt.Sprintf("run on %s (config_mutation, auto via assent window): %s", targetSlug, out)) + } + + // Destructive window: a narrow, TARGET-scoped grant opened only after an + // operator's explicit typed confirmation ("I confirm") on this same + // target — never by loose assent. Exists for multi-step destructive + // recovery (e.g. a failed destroy needing stop, then destroy) so the + // operator isn't asked to re-type "I confirm" for every single command + // against the thing they just confirmed. + if riskClass == policy.RiskDestructive && destructiveWindowActive(ctx, pool, agentID, targetSlug) { + host, user, wrap, rerr := resolveExecTarget(ctx, pool, targetSlug) + if rerr != nil { + pool.Exec(ctx, `UPDATE executions SET status='failed', result=$2::jsonb WHERE entity_id=$1`, id, jsonErr("%s", rerr.Error())) + return textResult(fmt.Sprintf("resolve target: %v", rerr)) + } + out, xerr := sshExec(ctx, host, user, wrap(command)) + if xerr != nil { + pool.Exec(ctx, `UPDATE executions SET status='failed', result=$2::jsonb WHERE entity_id=$1`, id, jsonErr("%s: %s", xerr.Error(), out)) + return textResult(fmt.Sprintf("run on %s: ERROR %v\n%s", targetSlug, xerr, out)) + } + pool.Exec(ctx, `UPDATE executions SET status='completed', result=$2::jsonb WHERE entity_id=$1`, id, jsonOut(out)) + slog.Info("mcp: run auto-executed via destructive window", "target", targetSlug, "execution_id", id) + return textResult(fmt.Sprintf("run on %s (destructive, auto via confirmed-target window): %s", targetSlug, out)) + } + + pool.Exec(ctx, `UPDATE executions SET status='pending_approval', risk_class=$2 WHERE entity_id=$1`, id, riskClass) + createApproval(ctx, pool, id, targetID, "run", string(runParams), riskClass) + confirmNote := "" + if riskClass == policy.RiskDestructive { + confirmNote = " This is classified DESTRUCTIVE — flag that clearly to the operator; it needs explicit confirmation, not just a casual \"go ahead\"." + } + return textResult(fmt.Sprintf("run on %s requires approval (risk: %s) — execution %s queued.%s Present the command and purpose to the operator and wait; do not re-request.", + targetSlug, riskClass, id, confirmNote)) +} + // autoApprove updates the approval + execution status in the DB to approved, // mirroring what DecideApproval does. Returns true on success. This is used // by the assent-window path to skip the operator-approval queue when the