Operator-reported bug: on 'proceed with the rest' the agent re-proposed the
plan, duplicating it in the sidebar. Root cause was a three-bug chain, not
one bug:
1. Trigger — model empty-response on 'proceed' (approval vocabulary didn't
list 'proceed', so the agent wasn't sure it was approved and no-op'd).
2. Amplifier — chatWith emitted 'error' without 'done' on empty response
(agent.go:370). The frontend's onComplete saw !receivedDone and
misclassified the model failure as a network disconnect, calling
handleDisconnect -> resumeSession.
3. Divergence — the reconnect note was generic ('report your state'), so
the agent re-proposed + re-executed instead of advancing the plan.
Fixes (shipped, e2e-validated against the live agent on oikos-nomos-1):
- A.2: proposePlan refuses re-proposal once a step has started (returns
errPlanInFlight). Drops the append-mode safety net (commit 5384499) that
was the direct source of the sidebar duplication. The agent must advance
with update_plan_step + run; the tool result directs it.
- A.1: proposePlan sets the 'generation' column on INSERT (migration 020
added the column + frontend grouping, but the INSERT never wired it).
- A.3: propose_plan tool description restated as a crisp contract (ONCE,
STOP and wait, REFUSES once a step started, advance with update_plan_step).
- F.3: approval vocabulary expanded to approved/yes/go/proceed/continue/ok/
go ahead; propose_plan result string tightened to an imperative.
- B.1: chatWith emits 'done' after 'error' on every terminal path via a new
emitError helper. The frontend now treats model errors as ended (not
disconnected), so no auto-reconnect -> resumeSession fires.
- B.2: reconnect/resume note carries the operator's last message + an
explicit 'advance the plan, do NOT call propose_plan again' directive when
a plan is in flight. Wired into all 4 resume entry points (reconnect,
/resume, idle-sweep, question-answer) via enrichResumeNote.
- B.3: resumeSession escalates the recovery note across its 3 attempts (final
retry: 'pick the lowest-pending step, mark it running, call run — do that
now') instead of 3 identical notes -> 3 identical empties.
Verification: TestProposePlan_RefuseInFlight replaces TestProposePlan_
AppendVsReplace. e2e conversations against the rebuilt container:
conv2 ('proceed with the rest') -> 0 propose_plan calls, plan stayed at
3 steps (was 6+ before), update_plan_step x5 + run x2 + complete_task.
conv3 (full plan, 'go ahead') -> apt-get update on lxc:dns auto-ran under
the plan window, update_entity_attributes writeback, clean complete_task.
nomos logs show zero reconnect/resume entries for the plan-proposing
sessions (the three-bug chain is closed).
Remaining (not in this commit): D.1 refuse complete_task without writeback
(next blocker), C.1/C.2, F.1/F.2 SOUL.md consolidation, B.4-B.6, E.1/E.2.
See plans/2026-07-14-post-fix-session-remainders.md.
Also: re-audit 2026-07-10-general-gated-execution.md — request_execution enum
retirement (60effcb) closes item 9; only auto-act revival (item 10) remains.
Version 0.4.1 -> 0.5.0 (minor: new structural behavior, not a bugfix).
306 lines
14 KiB
Go
306 lines
14 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
"log/slog"
|
|
"regexp"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/dtoro/oikos/internal/safego"
|
|
"github.com/google/uuid"
|
|
)
|
|
|
|
// execIDRe matches "execution <uuid>" in a tool result — the phrasing shared
|
|
// by request_execution / run when they queue or start a gated execution.
|
|
// Only these async executions need continuation; the synchronous auto-run
|
|
// path returns its output inline and is already observed in-turn.
|
|
var execIDRe = regexp.MustCompile(`(?i)execution\s+([0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12})`)
|
|
|
|
func extractExecutionIDs(toolResult string) []uuid.UUID {
|
|
matches := execIDRe.FindAllStringSubmatch(toolResult, -1)
|
|
seen := map[uuid.UUID]bool{}
|
|
var out []uuid.UUID
|
|
for _, m := range matches {
|
|
if id, err := uuid.Parse(m[1]); err == nil && !seen[id] {
|
|
seen[id] = true
|
|
out = append(out, id)
|
|
}
|
|
}
|
|
return out
|
|
}
|
|
|
|
// idleTaskThreshold is how long a goal-bearing session can sit non-terminal
|
|
// with no activity before the idle sweep nudges it, per
|
|
// plans/2026-07-11-task-completion-safety-net.md. Arbitrary starting point,
|
|
// not measured against real task durations — long enough that it won't fire
|
|
// mid-turn, short enough the board doesn't lie for hours.
|
|
const idleTaskThreshold = 15 * time.Minute
|
|
|
|
// runIdleSweepWorker is the safety net for case 2 of
|
|
// plans/2026-07-11-task-completion-safety-net.md: sessions that called
|
|
// set_goal (so the inline safety net in agent.go correctly left them alone,
|
|
// since they framed themselves as a real task) but then stalled without
|
|
// ever calling complete_task. Coarser than runContinuationWorker's 4s tick
|
|
// since "gone idle" is a much slower signal than "an execution just
|
|
// finished." Blocks until ctx is cancelled.
|
|
func (a *agent) runIdleSweepWorker(ctx context.Context) {
|
|
if a.store == nil {
|
|
slog.Warn("nomos: idle sweep worker disabled (no store)")
|
|
return
|
|
}
|
|
slog.Info("nomos: idle sweep worker started")
|
|
ticker := time.NewTicker(2 * time.Minute)
|
|
defer ticker.Stop()
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-ticker.C:
|
|
a.processIdleSweep(ctx)
|
|
}
|
|
}
|
|
}
|
|
|
|
// processIdleSweep nudges a stalled goal-bearing session once; if it's still
|
|
// non-terminal on the NEXT sweep (meaning the nudge itself went unanswered,
|
|
// not just that the model is still working), auto-closes it with a
|
|
// visible "auto-closed" outcome instead of leaving it stuck forever — same
|
|
// reasoning resumeSession already applies below for a different failure
|
|
// mode (a resume that produces no response at all).
|
|
func (a *agent) processIdleSweep(ctx context.Context) {
|
|
stale := a.store.staleGoalSessions(ctx, idleTaskThreshold, 5)
|
|
for _, s := range stale {
|
|
s := s
|
|
if s.CompletionNudges == 0 {
|
|
safego.Go("nomos:idle-nudge:"+s.ID, func() {
|
|
if err := a.store.bumpCompletionNudge(ctx, s.ID); err != nil {
|
|
slog.Error("nomos: idle nudge bump failed", "session", s.ID, "error", err)
|
|
return
|
|
}
|
|
note := fmt.Sprintf("[System: this task ('%s') has been idle for %s with no complete_task call. "+
|
|
"If the goal is done (or can't be completed), call complete_task now with the outcome and a "+
|
|
"one-line summary. If you're still genuinely working through the plan, ignore this and continue.]",
|
|
s.Goal, idleTaskThreshold)
|
|
note = a.store.enrichResumeNote(ctx, s.ID, note)
|
|
a.resumeSession(ctx, s.ID, note)
|
|
})
|
|
continue
|
|
}
|
|
safego.Go("nomos:idle-autoclose:"+s.ID, func() {
|
|
summary := fmt.Sprintf("Auto-closed after %s idle with no response to a completion nudge.", idleTaskThreshold)
|
|
if err := a.store.completeTask(ctx, s.ID, "partial", summary); err != nil {
|
|
slog.Error("nomos: idle auto-close failed", "session", s.ID, "error", err)
|
|
}
|
|
})
|
|
}
|
|
}
|
|
|
|
// runContinuationWorker is the event loop that replaces the human typing
|
|
// "continue". It polls for gated executions that (a) were initiated by a chat
|
|
// session and (b) have just finished, and — while that agent has an open assent
|
|
// window (an approved plan is in flight) — feeds each result back into the
|
|
// agent so it proceeds to the next step or recovers from the failure, all
|
|
// without an operator tick. Blocks until ctx is cancelled.
|
|
func (a *agent) runContinuationWorker(ctx context.Context) {
|
|
if a.store == nil {
|
|
slog.Warn("nomos: continuation worker disabled (no store)")
|
|
return
|
|
}
|
|
slog.Info("nomos: continuation worker started")
|
|
ticker := time.NewTicker(4 * time.Second)
|
|
defer ticker.Stop()
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-ticker.C:
|
|
a.processContinuations(ctx)
|
|
}
|
|
}
|
|
}
|
|
|
|
// processContinuations dispatches each pending item as its OWN goroutine
|
|
// (safego.Go, so a panic deep in one task's resumed turn — JSON parsing of
|
|
// model output, an unexpected nil in a tool result — is recovered and logged
|
|
// instead of taking down this whole function, which used to run every
|
|
// item sequentially in the SAME goroutine as the ticker loop. Two problems
|
|
// that fixed: (1) throughput — task B's continuation no longer waits for
|
|
// task A's full (up to 10-minute) resumed turn to finish first, the exact
|
|
// per-task blocking this session's earlier concurrency work removed from the
|
|
// live-chat path but had left in place here; (2) survivability — since Go
|
|
// panics unwind the goroutine they occur in, an unrecovered one here used to
|
|
// mean this call (and every future tick, since the whole ticker loop runs in
|
|
// one goroutine) would simply stop — auto-continuation for every task would
|
|
// silently die until nomos restarted. Now a single bad item can only ever
|
|
// take down its own goroutine.
|
|
func (a *agent) processContinuations(ctx context.Context) {
|
|
pending := a.store.pendingContinuations(ctx, 5)
|
|
for _, p := range pending {
|
|
// Scope gate: only auto-continue while an approved plan is active FOR
|
|
// THIS SESSION. Checked per-item, not once for the whole batch — with
|
|
// multiple tasks in flight, one task's open window must never cover a
|
|
// pending continuation belonging to a different task.
|
|
if !a.store.assentWindowActive(ctx, a.agentID, p.SessionID) {
|
|
// Re-open the assent window if this session is genuinely
|
|
// executing (plan was approved, work is in progress) — the
|
|
// window may have expired while the execution ran. Don't
|
|
// penalize timing: the plan was approved, the work happened,
|
|
// the result should flow back.
|
|
sesh, seshErr := a.store.getSession(ctx, p.SessionID)
|
|
if seshErr == nil && sesh.Goal != "" && (sesh.Status == "executing" || sesh.Status == "planning") {
|
|
a.openAssentWindow(ctx, p.SessionID)
|
|
slog.Info("nomos: re-opened assent window for continuing session", "session", p.SessionID, "execution", p.ExecID)
|
|
} else {
|
|
// Genuinely no plan — inject a visible note so the
|
|
// operator knows WHY the agent didn't auto-continue.
|
|
note := fmt.Sprintf("[System: execution %s finished with status=%s, but the assent window for this session is not active. The agent will not auto-continue. Reply 'continue' or re-approve the plan to resume.]", p.ExecID, p.Status)
|
|
body, _ := json.Marshal(map[string]any{"role": "assistant", "text": note, "auto": true})
|
|
a.store.saveMessage(context.Background(), p.SessionID, "assistant", body)
|
|
a.store.markContinued(ctx, p.ExecID)
|
|
continue
|
|
}
|
|
}
|
|
a.store.markContinued(ctx, p.ExecID) // stamp first: a failure here must not cause a re-continue loop
|
|
safego.Go("nomos:continue-session:"+p.SessionID, func() { a.continueSession(ctx, p) })
|
|
}
|
|
}
|
|
|
|
// continueSession re-invokes the agent for one finished execution. Persists
|
|
// progress LIVE — a placeholder row immediately, updated in place as each
|
|
// tool call completes — instead of only saving once the whole continuation
|
|
// finishes. The frontend polls (see chat.ts startPolling); without
|
|
// incremental persistence here, a continuation that runs several tool calls
|
|
// before concluding would look like total silence in the UI for however long
|
|
// that takes, which is exactly the "I just wait while nothing happens"
|
|
// complaint this exists to fix — polling alone only helps if there's
|
|
// something new to poll for.
|
|
func (a *agent) continueSession(ctx context.Context, p pendingContinuation) {
|
|
slog.Info("nomos: auto-continuing session", "session", p.SessionID, "execution", p.ExecID, "status", p.Status)
|
|
a.resumeSession(ctx, p.SessionID, buildContinuationNote(p))
|
|
}
|
|
|
|
// resumeSession re-invokes the agent for a session with a system-injected note —
|
|
// a finished execution (continueSession) or an operator's answer to a question
|
|
// (handleAnswerQuestion) — persisting progress LIVE (a placeholder row updated
|
|
// in place as each tool call lands) so the frontend poller sees each step,
|
|
// instead of total silence until the whole resume concludes.
|
|
func (a *agent) resumeSession(ctx context.Context, sessionID, note string) {
|
|
placeholder, _ := json.Marshal(map[string]any{
|
|
"role": "assistant",
|
|
"text": "",
|
|
"auto": true,
|
|
})
|
|
msgID, err := a.store.insertMessageReturningID(ctx, sessionID, "assistant", placeholder)
|
|
if err != nil {
|
|
slog.Error("nomos: resume placeholder insert failed", "session", sessionID, "error", err)
|
|
}
|
|
|
|
var toolCalls []map[string]any
|
|
var finalText, errText string
|
|
|
|
persist := func() {
|
|
if msgID == uuid.Nil {
|
|
return
|
|
}
|
|
text := finalText
|
|
if text == "" && errText != "" {
|
|
text = fmt.Sprintf("(auto-continuation hit an internal error and did not respond: %s — the execution's own result is above; you may need to prompt the agent again)", errText)
|
|
}
|
|
body, _ := json.Marshal(map[string]any{
|
|
"role": "assistant",
|
|
"text": text,
|
|
"tool_calls": toolCalls,
|
|
"auto": true, // marks this as an autonomous continuation, not an operator turn
|
|
})
|
|
a.store.updateMessage(ctx, msgID, body)
|
|
}
|
|
|
|
// One retry if the LLM call itself produced nothing (transient flake /
|
|
// empty-response) — the whole point of this mechanism is "don't give up
|
|
// on the first error," which should apply to the continuation call
|
|
// itself, not just the homelab commands it's continuing. Found live: a
|
|
// destructive-recovery continuation hit an empty LLM response, its
|
|
// internal retry (chatWith's own maxLLMRetries=1) also came up empty, and
|
|
// without this outer retry the operator would see nothing at all.
|
|
cctx, cancel := context.WithTimeout(ctx, 10*time.Minute)
|
|
defer cancel()
|
|
for attempt := 0; attempt < 3; attempt++ {
|
|
toolCalls, finalText, errText = nil, "", ""
|
|
emit := func(ev agentEvent) {
|
|
if ev.Type == "tool_use" || ev.Type == "tool_result" {
|
|
if m, ok := ev.Data.(map[string]any); ok {
|
|
m["type"] = ev.Type
|
|
toolCalls = append(toolCalls, m)
|
|
}
|
|
persist() // live: a poller sees this step land within seconds
|
|
}
|
|
if ev.Type == "text" {
|
|
finalText, _ = ev.Data.(string)
|
|
}
|
|
if ev.Type == "error" {
|
|
errText, _ = ev.Data.(string)
|
|
}
|
|
}
|
|
a.chatWith(cctx, sessionID, "", note, emit)
|
|
if finalText != "" || len(toolCalls) > 0 {
|
|
break
|
|
}
|
|
if attempt == 0 {
|
|
slog.Warn("nomos: resume produced nothing, retrying once", "session", sessionID, "error", errText)
|
|
}
|
|
}
|
|
|
|
if errText != "" && finalText == "" {
|
|
slog.Error("nomos: resume produced no response after retry", "session", sessionID, "error", errText)
|
|
// Persist a visible system note in the transcript so the
|
|
// operator sees what happened, but do NOT auto-complete the
|
|
// task — leave it in 'executing' so a follow-up chat message
|
|
// can resume it. Before this fix, the task was marked 'failed'
|
|
// here, which ended it permanently and required starting over.
|
|
resumeFailedNote := fmt.Sprintf("[System: auto-resume failed after retrying: %s. The task is paused — send another message to continue.]", errText)
|
|
body, _ := json.Marshal(map[string]any{
|
|
"role": "assistant",
|
|
"text": resumeFailedNote,
|
|
"auto": true,
|
|
})
|
|
if msgID != uuid.Nil {
|
|
a.store.updateMessage(context.Background(), msgID, body)
|
|
} else {
|
|
// No placeholder was inserted (rare), save directly.
|
|
a.store.saveMessage(context.Background(), sessionID, "assistant", body)
|
|
}
|
|
return // do not call persist() again — already persisted above
|
|
}
|
|
persist() // final state — same row, updated one last time with the concluding text
|
|
}
|
|
|
|
// buildContinuationNote frames the finished execution for the model: what
|
|
// happened, and what to do about it. The persist-through-errors instruction
|
|
// lives here (and in SOUL) so the agent recovers instead of stopping.
|
|
func buildContinuationNote(p pendingContinuation) string {
|
|
action := p.Action
|
|
if i := strings.IndexByte(action, ':'); i > 0 && len(action) > 40 {
|
|
action = action[:i] // keep just the action verb for brevity; params are in the DB
|
|
}
|
|
result := p.Result
|
|
if len(result) > 3000 {
|
|
result = result[:3000] + "…[truncated]"
|
|
}
|
|
var b strings.Builder
|
|
fmt.Fprintf(&b, "[System: execution %s (%s) finished with status=%s.\nResult: %s\n\n",
|
|
p.ExecID, action, p.Status, result)
|
|
switch p.Status {
|
|
case "completed":
|
|
b.WriteString("It SUCCEEDED. Continue the approved plan: run the next step. If this was the final step, verify the end goal actually works (e.g. curl the service) and then report success to the operator. Do NOT stop and wait for the operator to say 'continue'.")
|
|
case "failed", "cancelled":
|
|
b.WriteString("It FAILED. Do NOT give up or hand back to the operator. Diagnose the cause from the result above (and by running read-only inspection commands if needed), form a hypothesis, fix it, and retry or take an alternative approach. You have an active assent window, so config_mutation steps run without re-approval. Only stop and ask the operator if you are genuinely blocked (need information only they have) or the fix would require a destructive action they haven't approved.")
|
|
default: // denied / revoked
|
|
b.WriteString("The operator denied or revoked this step. Stop executing this plan and briefly acknowledge.")
|
|
}
|
|
b.WriteString("]")
|
|
return b.String()
|
|
}
|