diff --git a/.claude/skills/env-reference/SKILL.md b/.claude/skills/env-reference/SKILL.md index 1d499e015c..57b158079e 100644 --- a/.claude/skills/env-reference/SKILL.md +++ b/.claude/skills/env-reference/SKILL.md @@ -761,7 +761,7 @@ Generated deployments validate and pass these values through cloud-init to newly - `ACP_PONG_TIMEOUT` — WebSocket pong deadline after ping (default: 10s) - `ACP_PROMPT_TIMEOUT` — Max ACP prompt runtime for workspace sessions; 0 = no timeout (default: 0) - `ACP_TASK_PROMPT_TIMEOUT` — Max ACP prompt runtime for task-driven sessions (default: 8h) -- `ACP_PROMPT_CANCEL_GRACE_PERIOD` — Grace wait after cancel before force-stop (default: 5s) +- `ACP_PROMPT_CANCEL_GRACE_PERIOD` — Grace wait for the cancelled prompt to settle before it is finished as `cancelled` and the agent is restarted (default: 5s). Bound to the cancelled prompt; never affects a later prompt - `ACP_PROMPT_RETRY_MAX_RETRIES` — Max transient provider prompt retries after the initial attempt (default: 2) - `ACP_PROMPT_RETRY_INITIAL_BACKOFF` — Initial backoff before retrying transient provider prompt errors (default: 15s) - `ACP_PROMPT_RETRY_MAX_BACKOFF` — Max exponential backoff for transient provider prompt retries (default: 2m) diff --git a/packages/vm-agent/internal/acp/gateway.go b/packages/vm-agent/internal/acp/gateway.go index 044e24c0a8..54d28f06ff 100644 --- a/packages/vm-agent/internal/acp/gateway.go +++ b/packages/vm-agent/internal/acp/gateway.go @@ -190,7 +190,8 @@ type GatewayConfig struct { PongTimeout time.Duration // PromptTimeout bounds how long a prompt can run before force-stop fallback. PromptTimeout time.Duration - // PromptCancelGracePeriod waits after cancel before force-stopping unresponsive prompt. + // PromptCancelGracePeriod waits after cancel for the cancelled prompt to + // settle before finishing it as cancelled and restarting the agent. PromptCancelGracePeriod time.Duration // PromptRetryMaxRetries bounds transient provider prompt retries after the initial attempt. PromptRetryMaxRetries int diff --git a/packages/vm-agent/internal/acp/session_host.go b/packages/vm-agent/internal/acp/session_host.go index 9c93cd6e61..e4c26486ca 100644 --- a/packages/vm-agent/internal/acp/session_host.go +++ b/packages/vm-agent/internal/acp/session_host.go @@ -1,7 +1,6 @@ package acp import ( - "bufio" "context" "encoding/json" "fmt" @@ -10,7 +9,6 @@ import ( "strings" "sync" "sync/atomic" - "syscall" "time" acpsdk "github.com/coder/acp-go-sdk" @@ -31,8 +29,9 @@ const ( ) const ( - // DefaultPromptCancelGracePeriod is how long we wait after cancel before - // force-stopping an unresponsive agent process. + // DefaultPromptCancelGracePeriod is how long a cancelled prompt may take to + // settle before it is finished as "cancelled" and the agent is restarted. + // The watchdog is bound to that one prompt attempt and disarms when it ends. DefaultPromptCancelGracePeriod = 5 * time.Second // DefaultPromptRetryInitialDelay is the first delay before retrying a @@ -62,10 +61,6 @@ const ( defaultControlPlaneHTTPTimeout = 30 * time.Second ) -// DefaultStderrBufferBytes is the default maximum agent stderr captured for -// crash reports. Override via ACP_STDERR_BUFFER_BYTES. -const DefaultStderrBufferBytes = 4096 - // DefaultMessageBufferSize is the default maximum number of messages buffered // per session for late-join replay. Override via ACP_MESSAGE_BUFFER_SIZE. const DefaultMessageBufferSize = 5000 @@ -263,6 +258,10 @@ type SessionHost struct { // promptActivityCancel stops the periodic prompting re-report loop. // Protected by promptCancelMu. promptActivityCancel context.CancelFunc + // cancelGraceTimer replaces the cancel-grace timer in tests so they can own + // the ordering between a cancel, the next prompt, and the deadline. Set + // only before the host is used; nil means a real time.Timer. + cancelGraceTimer func(time.Duration) (<-chan time.Time, func()) // Harness-owned background work is normalized from optional ACP extension // notifications. It is isolated from the prompt lifecycle because it may @@ -314,133 +313,6 @@ type SessionHost struct { cancel context.CancelFunc } -// PromptTerminalObserver receives the terminal state of one accepted prompt. -// It is used by the VM HTTP delivery protocol to durably complete a receipt. -// Implementations must return quickly; notification runs asynchronously. -type PromptTerminalObserver func(stopReason string, promptErr error) - -type promptAttempt struct { - id uint64 - startedAt time.Time - cancel context.CancelFunc - done chan struct{} - rpcDone chan struct{} - terminalMu sync.Mutex - terminal bool - checkpointOwned bool - observer PromptTerminalObserver - checkpointRequested atomic.Bool -} - -type checkpointRolloverEpisode struct { - sessionID string - attempt *promptAttempt - result chan CheckpointRolloverResult - operationCtx context.Context - forced bool - decisionMu sync.Mutex - outcome sync.Once - terminal atomic.Bool - finalResult CheckpointRolloverResult -} - -// complete is the checkpoint episode's linearization point. A terminal caller -// that wins permanently suppresses process-exit restart for this episode; a -// successful strict resume that wins cannot be overturned by a later deadline. -func (e *checkpointRolloverEpisode) complete(result CheckpointRolloverResult, terminal bool) bool { - e.decisionMu.Lock() - defer e.decisionMu.Unlock() - return e.completeLocked(result, terminal) -} - -func (e *checkpointRolloverEpisode) completeLocked(result CheckpointRolloverResult, terminal bool) bool { - won := false - e.outcome.Do(func() { - won = true - e.finalResult = result - if terminal { - e.terminal.Store(true) - } - e.result <- result - }) - return won -} - -// completeStrictResume atomically orders the checkpoint outcome after the -// prompt attempt's terminal arbiter. Natural completion or user cancellation -// therefore wins as superseded; an already-terminal episode can never publish -// a later successful resume. -func (e *checkpointRolloverEpisode) completeStrictResume(h *SessionHost, result CheckpointRolloverResult) (CheckpointRolloverResult, bool) { - e.decisionMu.Lock() - defer e.decisionMu.Unlock() - if e.terminal.Load() { - return e.finalResult, false - } - if !e.attempt.completeCheckpoint(h, checkpointPreemptedStopReason, nil) { - superseded := CheckpointRolloverResult{State: "superseded", ACPSessionID: e.sessionID} - e.completeLocked(superseded, true) - return superseded, false - } - e.completeLocked(result, false) - return result, true -} - -func (a *promptAttempt) complete(h *SessionHost, stopReason string, promptErr error) bool { - return a.completeWith(h, stopReason, promptErr, nil) -} - -func (a *promptAttempt) completeWith(h *SessionHost, stopReason string, promptErr error, finalize func()) bool { - a.terminalMu.Lock() - if a.terminal || a.checkpointOwned { - a.terminalMu.Unlock() - return false - } - a.terminal = true - a.terminalMu.Unlock() - a.publishCompletion(h, stopReason, promptErr, finalize) - return true -} - -func (a *promptAttempt) claimCheckpointTerminal() bool { - a.terminalMu.Lock() - defer a.terminalMu.Unlock() - if a.terminal || a.checkpointOwned { - return false - } - a.checkpointOwned = true - return true -} - -func (a *promptAttempt) completeCheckpoint(h *SessionHost, stopReason string, promptErr error) bool { - a.terminalMu.Lock() - if a.terminal || !a.checkpointOwned { - a.terminalMu.Unlock() - return false - } - a.checkpointOwned = false - a.terminal = true - a.terminalMu.Unlock() - a.publishCompletion(h, stopReason, promptErr, nil) - return true -} - -func (a *promptAttempt) publishCompletion(h *SessionHost, stopReason string, promptErr error, finalize func()) { - // Release the admission gate before publishing the terminal status. The - // public AcceptPrompt path also requires HostReady, so the intermediate - // prompting/starting/error status cannot admit a new prompt. - h.releasePrompt(a) - if finalize != nil { - finalize() - } - close(a.done) - if cb := h.config.OnPromptComplete; cb != nil { - go cb(stopReason, promptErr) - } - if a.observer != nil { - go a.observer(stopReason, promptErr) - } -} - func (h *SessionHost) now() time.Time { if h.config.Now != nil { return h.config.Now() @@ -652,154 +524,6 @@ func (h *SessionHost) autoSuspend() { }) } -// CancelPrompt cancels the currently running Prompt() call, if any. -// This is safe to call from any goroutine. If no prompt is in flight, -// it's a no-op. The cancel function is guarded by promptCancelMu -// (separate from promptMu) so we never deadlock with HandlePrompt. -func (h *SessionHost) CancelPrompt() { - h.cancelPrompt(true) -} - -// CancelPromptFromControlPlane mirrors the viewer WebSocket session/cancel path -// for HTTP control-plane cancellation requests. -func (h *SessionHost) CancelPromptFromControlPlane() { - if h.AgentType() == "opencode" { - h.cancelPrompt(false) - h.StopProcessForPromptCancel() - return - } - - h.CancelPrompt() - cancelMessage, err := h.cancelNotification() - if err != nil { - slog.Warn("CancelPromptFromControlPlane: could not build session/cancel notification", "error", err) - } else { - h.ForwardToAgent(cancelMessage) - } - h.StopProcessForPromptCancel() -} - -func (h *SessionHost) cancelNotification() ([]byte, error) { - sessionID := h.currentSessionIDForCancel() - if sessionID == "" { - return nil, fmt.Errorf("missing session ID") - } - - return json.Marshal(struct { - JSONRPC string `json:"jsonrpc"` - Method string `json:"method"` - Params struct { - SessionID string `json:"sessionId"` - } `json:"params"` - }{ - JSONRPC: "2.0", - Method: "session/cancel", - Params: struct { - SessionID string `json:"sessionId"` - }{SessionID: sessionID}, - }) -} - -func (h *SessionHost) currentSessionIDForCancel() string { - h.mu.RLock() - acpSessionID := string(h.sessionID) - h.mu.RUnlock() - if acpSessionID != "" { - return acpSessionID - } - return h.config.SessionID -} - -func (h *SessionHost) cancelPrompt(startGraceTimer bool) { - h.promptCancelMu.Lock() - cancelFn := h.promptCancel - promptID := h.activePromptID - if cancelFn != nil { - h.promptCancelRequested = true - } - h.promptCancelMu.Unlock() - - if cancelFn == nil { - slog.Info("CancelPrompt: no prompt in flight") - return - } - - slog.Info("CancelPrompt: cancelling in-flight prompt") - h.reportLifecycle("info", "Prompt cancel requested", nil) - cancelFn() - - if !startGraceTimer { - return - } - - grace := h.promptCancelGracePeriod() - if grace <= 0 { - return - } - - go func(id uint64, wait time.Duration) { - timer := time.NewTimer(wait) - defer timer.Stop() - <-timer.C - h.triggerPromptForceStopIfStuck(id, fmt.Sprintf("Prompt cancel grace elapsed after %s", wait)) - }(promptID, grace) -} - -// ForwardToAgent sends a raw message to the agent's stdin. -func (h *SessionHost) ForwardToAgent(message []byte) { - h.mu.RLock() - process := h.process - h.mu.RUnlock() - - if process == nil { - slog.Warn("No agent process running, dropping message") - return - } - - data := append(message, '\n') - if _, err := process.Stdin().Write(data); err != nil { - slog.Error("Failed to write to agent stdin", "error", err) - } -} - -// SignalProcess sends a signal to the agent process. This is used for agents -// that don't implement session/cancel (e.g., opencode) — SIGTERM is sent -// directly to the process instead of forwarding the cancel RPC. -func (h *SessionHost) SignalProcess(sig syscall.Signal) { - h.mu.RLock() - process := h.process - h.mu.RUnlock() - - if process == nil { - slog.Warn("SignalProcess: no agent process running") - return - } - - process.KillContainerProcesses(sig) - slog.Info("SignalProcess: sent signal to agent process", "signal", sig, "agentType", h.AgentType()) -} - -// StopProcessForPromptCancel terminates the current agent process for a user -// prompt cancel without marking the host stopped. The process monitor will -// restart the agent and return the host to ready for follow-up prompts. -func (h *SessionHost) StopProcessForPromptCancel() { - h.mu.Lock() - process := h.process - if process != nil { - h.intentionalPromptCancelProcessStop = true - } - h.mu.Unlock() - - if process == nil { - slog.Warn("StopProcessForPromptCancel: no agent process running") - return - } - - if err := process.Stop(); err != nil { - slog.Warn("StopProcessForPromptCancel: failed to stop agent process", "error", err) - } -} - // Stop kills the agent process, disconnects all viewers, and marks the session // as stopped. This is the only way to terminate the agent — browser disconnects // do NOT call this. @@ -873,116 +597,6 @@ func (h *SessionHost) Stop() { h.viewerMu.Unlock() } -// phaseTimeout returns a per-phase timeout duration. If phaseMs is > 0, it is -// used; otherwise the fallback timeout is returned. -func phaseTimeout(phaseMs int, fallback time.Duration) time.Duration { - if phaseMs > 0 { - return time.Duration(phaseMs) * time.Millisecond - } - return fallback -} - -// applySessionSettings calls SetSessionConfigOption for the model and -// SetSessionMode on the ACP connection. A rejected explicit Codex model is fatal: -// continuing would silently run a different model. Other adapter settings remain -// best-effort for backward compatibility. -func (h *SessionHost) applySessionSettings(ctx context.Context, settings *agentSettingsPayload) error { - if settings == nil || h.acpConn == nil || h.sessionID == "" { - return nil - } - - // These RPCs run while h.mu is held for write, and the ACP SDK blocks each - // response on its notification worker catching up (waitNotificationsUpTo). - // The caller's ctx is the WebSocket connection lifetime — effectively - // unbounded — so a stalled notification worker would hang the handshake (and - // every h.mu reader) forever. Bound it explicitly; Codex model selection - // failures are returned while the remaining settings stay best-effort. - settingsCtx, cancel := context.WithTimeout(ctx, h.sessionSettingsTimeout()) - defer cancel() - ctx = settingsCtx - - if settings.Model != "" { - if err := h.applySessionModelConfigOption(ctx, settings.Model); err != nil && h.agentType == "openai-codex" { - return fmt.Errorf("cannot apply requested Codex model %q: %w", settings.Model, err) - } - } - - if settings.PermissionMode != "" && settings.PermissionMode != "default" { - // Codex (openai-codex) does not support SetSessionMode — skip to avoid - // a guaranteed error on every session start. - if h.agentType == "openai-codex" { - slog.Info("ACP: skipping SetSessionMode for openai-codex (unsupported)", "mode", settings.PermissionMode) - } else { - slog.Info("ACP: setting session mode", "mode", settings.PermissionMode) - if _, err := h.acpConn.SetSessionMode(ctx, acpsdk.SetSessionModeRequest{ - SessionId: h.sessionID, - ModeId: acpsdk.SessionModeId(settings.PermissionMode), - }); err != nil { - slog.Warn("ACP SetSessionMode failed (non-fatal)", "mode", settings.PermissionMode, "error", err) - h.reportLifecycle("warn", "ACP SetSessionMode failed", map[string]interface{}{ - "mode": settings.PermissionMode, - "error": err.Error(), - }) - } else { - slog.Info("ACP: session mode set", "mode", settings.PermissionMode) - h.reportLifecycle("info", "ACP session mode applied", map[string]interface{}{ - "mode": settings.PermissionMode, - }) - } - } - } - return nil -} - -func (h *SessionHost) applySessionModelConfigOption(ctx context.Context, model string) error { - modelConfigID, ok := findModelConfigOptionID(h.configOptions) - if !ok { - slog.Warn("ACP session model config option unavailable", "model", model) - h.reportLifecycle("warn", "ACP session model config option unavailable", map[string]interface{}{ - "model": model, - }) - return fmt.Errorf("model config option unavailable") - } - - slog.Info("ACP: setting session model config option", "model", model, "configId", string(modelConfigID)) - resp, err := h.acpConn.SetSessionConfigOption(ctx, acpsdk.SetSessionConfigOptionRequest{ - ValueId: &acpsdk.SetSessionConfigOptionValueId{ - SessionId: h.sessionID, - ConfigId: modelConfigID, - Value: acpsdk.SessionConfigValueId(model), - }, - }) - if err != nil { - slog.Warn("ACP SetSessionConfigOption failed", "model", model, "configId", string(modelConfigID), "error", err) - h.reportLifecycle("warn", "ACP session model config option failed", map[string]interface{}{ - "model": model, - "configId": string(modelConfigID), - "error": err.Error(), - }) - return fmt.Errorf("set session model config option: %w", err) - } - - h.configOptions = resp.ConfigOptions - slog.Info("ACP: session model config option set", "model", model, "configId", string(modelConfigID)) - h.reportLifecycle("info", "ACP session model applied", map[string]interface{}{ - "model": model, - "configId": string(modelConfigID), - }) - return nil -} - -func findModelConfigOptionID(options []acpsdk.SessionConfigOption) (acpsdk.SessionConfigId, bool) { - for _, option := range options { - if option.Select == nil || option.Select.Category == nil { - continue - } - if *option.Select.Category == acpsdk.SessionConfigOptionCategoryModel { - return option.Select.Id, true - } - } - return "", false -} - // ensureAgentInstalled checks if the ACP adapter binary exists and installs it // on-demand if missing. func (h *SessionHost) ensureAgentInstalled(ctx context.Context, info agentCommandInfo) error { @@ -1009,89 +623,6 @@ func (h *SessionHost) ensureAgentInstalled(ctx context.Context, info agentComman return installAgentBinary(ctx, containerID, info) } -// monitorStderr reads the agent's stderr and collects it for error reporting. -func (h *SessionHost) monitorStderr(process agentProcess) { - scanner := bufio.NewScanner(process.Stderr()) - for scanner.Scan() { - line := redactAgentDiagnosticText(scanner.Text()) - slog.Warn("Agent stderr", "line", line) - h.stderrMu.Lock() - if h.stderrBuf.Len() < h.config.StderrBufferBytes { - if h.stderrBuf.Len() > 0 { - h.stderrBuf.WriteByte('\n') - } - h.stderrBuf.WriteString(line) - } - h.stderrMu.Unlock() - } -} - -func (h *SessionHost) getAndClearStderr() string { - h.stderrMu.Lock() - defer h.stderrMu.Unlock() - s := h.stderrBuf.String() - h.stderrBuf.Reset() - return s -} - -func (h *SessionHost) peekStderr() string { - h.stderrMu.Lock() - defer h.stderrMu.Unlock() - return h.stderrBuf.String() -} - -// silentErrorPatterns are stderr substrings that indicate an API-level error -// the agent may have swallowed (returning a normal end_turn instead of an error). -var silentErrorPatterns = []string{ - "AI_APICallError", - "Unauthorized", - "401", - "403", - "invalid_api_key", - "authentication_error", -} - -// checkStderrForSilentErrors peeks at the accumulated stderr buffer for known -// API error patterns. Some agents (notably OpenCode with Scaleway) silently -// swallow API errors and return {stopReason: "end_turn"} instead of an error. -// When detected, we log a warning and report a lifecycle event so the UI can -// surface the issue. The stderr buffer is NOT cleared — it remains available -// for crash reporting in monitorProcessExit. -func (h *SessionHost) checkStderrForSilentErrors(stopReason acpsdk.StopReason) { - h.stderrMu.Lock() - stderr := h.stderrBuf.String() - h.stderrMu.Unlock() - - if stderr == "" { - return - } - - for _, pattern := range silentErrorPatterns { - if strings.Contains(stderr, pattern) { - slog.Warn("ACP: possible silent API error detected in stderr after prompt completion", - "stopReason", string(stopReason), - "pattern", pattern, - "stderrSnippet", truncateString(stderr, 512), - "agentType", h.AgentType(), - ) - h.reportLifecycle("warn", "Possible silent API error — check agent credentials", map[string]interface{}{ - "stopReason": string(stopReason), - "errorPattern": pattern, - "stderrSnippet": truncateString(stderr, 256), - }) - return // report once per prompt, not per pattern - } - } -} - -// truncateString returns s truncated to maxLen with "..." appended if needed. -func truncateString(s string, maxLen int) string { - if len(s) <= maxLen { - return s - } - return s[:maxLen] + "..." -} - // stopCurrentAgentLocked stops the current agent process. Must hold h.mu. func (h *SessionHost) stopCurrentAgentLocked() { // Stop the process-scoped harness heartbeat before clearing the ACP @@ -1211,17 +742,3 @@ func (h *SessionHost) loadMirroredStatus() SessionHostStatus { } return HostIdle } - -// DefaultSessionSettingsTimeout bounds the post-handshake SetSessionMode / -// SetSessionConfigOption RPCs. Override via the existing NewSessionTimeoutMs -// gateway setting (ACP_NEW_SESSION_TIMEOUT_MS). -const DefaultSessionSettingsTimeout = 30 * time.Second - -// sessionSettingsTimeout resolves the bound for applySessionSettings' RPCs, -// reusing the configured NewSession timeout when one is set. -func (h *SessionHost) sessionSettingsTimeout() time.Duration { - if h.config.NewSessionTimeoutMs > 0 { - return time.Duration(h.config.NewSessionTimeoutMs) * time.Millisecond - } - return DefaultSessionSettingsTimeout -} diff --git a/packages/vm-agent/internal/acp/session_host_activity_test.go b/packages/vm-agent/internal/acp/session_host_activity_test.go index 27d2659d84..61a0e15c61 100644 --- a/packages/vm-agent/internal/acp/session_host_activity_test.go +++ b/packages/vm-agent/internal/acp/session_host_activity_test.go @@ -43,7 +43,7 @@ func TestPromptActivityRereportStopsBeforeIdle(t *testing.T) { }, }) - host.markPromptStarted(acpsdk.SessionId("sdk-1"), 1, "viewer-1") + host.markPromptStarted(nil, acpsdk.SessionId("sdk-1"), 1, "viewer-1") waitFor(t, 250*time.Millisecond, func() bool { return countActivity(&mu, &activities, "prompting") >= 2 }) diff --git a/packages/vm-agent/internal/acp/session_host_attempt.go b/packages/vm-agent/internal/acp/session_host_attempt.go new file mode 100644 index 0000000000..da639dfb5f --- /dev/null +++ b/packages/vm-agent/internal/acp/session_host_attempt.go @@ -0,0 +1,153 @@ +package acp + +import ( + "context" + "sync" + "sync/atomic" + "time" +) + +// PromptTerminalObserver receives the terminal state of one accepted prompt. +// It is used by the VM HTTP delivery protocol to durably complete a receipt. +// Implementations must return quickly; notification runs asynchronously. +type PromptTerminalObserver func(stopReason string, promptErr error) + +type promptAttempt struct { + id uint64 + // deliveryID is the control-plane prompt delivery that created this + // attempt, empty for viewer prompts. Immutable after beginPrompt. + deliveryID string + startedAt time.Time + cancel context.CancelFunc + done chan struct{} + rpcDone chan struct{} + terminalMu sync.Mutex + terminal bool + checkpointOwned bool + observer PromptTerminalObserver + checkpointRequested atomic.Bool +} + +type checkpointRolloverEpisode struct { + sessionID string + attempt *promptAttempt + result chan CheckpointRolloverResult + operationCtx context.Context + forced bool + decisionMu sync.Mutex + outcome sync.Once + terminal atomic.Bool + finalResult CheckpointRolloverResult +} + +// complete is the checkpoint episode's linearization point. A terminal caller +// that wins permanently suppresses process-exit restart for this episode; a +// successful strict resume that wins cannot be overturned by a later deadline. +func (e *checkpointRolloverEpisode) complete(result CheckpointRolloverResult, terminal bool) bool { + e.decisionMu.Lock() + defer e.decisionMu.Unlock() + return e.completeLocked(result, terminal) +} + +func (e *checkpointRolloverEpisode) completeLocked(result CheckpointRolloverResult, terminal bool) bool { + won := false + e.outcome.Do(func() { + won = true + e.finalResult = result + if terminal { + e.terminal.Store(true) + } + e.result <- result + }) + return won +} + +// completeStrictResume atomically orders the checkpoint outcome after the +// prompt attempt's terminal arbiter. Natural completion or user cancellation +// therefore wins as superseded; an already-terminal episode can never publish +// a later successful resume. +func (e *checkpointRolloverEpisode) completeStrictResume(h *SessionHost, result CheckpointRolloverResult) (CheckpointRolloverResult, bool) { + e.decisionMu.Lock() + defer e.decisionMu.Unlock() + if e.terminal.Load() { + return e.finalResult, false + } + if !e.attempt.completeCheckpoint(h, checkpointPreemptedStopReason, nil) { + superseded := CheckpointRolloverResult{State: "superseded", ACPSessionID: e.sessionID} + e.completeLocked(superseded, true) + return superseded, false + } + e.completeLocked(result, false) + return result, true +} + +func (a *promptAttempt) complete(h *SessionHost, stopReason string, promptErr error) bool { + return a.completeWith(h, stopReason, promptErr, nil) +} + +func (a *promptAttempt) completeWith(h *SessionHost, stopReason string, promptErr error, finalize func()) bool { + a.terminalMu.Lock() + if a.terminal || a.checkpointOwned { + a.terminalMu.Unlock() + return false + } + a.terminal = true + a.terminalMu.Unlock() + a.publishCompletion(h, stopReason, promptErr, finalize) + return true +} + +func (a *promptAttempt) isTerminal() bool { + a.terminalMu.Lock() + defer a.terminalMu.Unlock() + return a.terminal +} + +// logFields adds the attempt's delivery identity to lifecycle log fields. Safe +// on a nil attempt so callers can log a cancel whose attempt already settled. +func (a *promptAttempt) logFields(fields map[string]interface{}) map[string]interface{} { + if a != nil && a.deliveryID != "" { + fields["deliveryId"] = a.deliveryID + } + return fields +} + +func (a *promptAttempt) claimCheckpointTerminal() bool { + a.terminalMu.Lock() + defer a.terminalMu.Unlock() + if a.terminal || a.checkpointOwned { + return false + } + a.checkpointOwned = true + return true +} + +func (a *promptAttempt) completeCheckpoint(h *SessionHost, stopReason string, promptErr error) bool { + a.terminalMu.Lock() + if a.terminal || !a.checkpointOwned { + a.terminalMu.Unlock() + return false + } + a.checkpointOwned = false + a.terminal = true + a.terminalMu.Unlock() + a.publishCompletion(h, stopReason, promptErr, nil) + return true +} + +func (a *promptAttempt) publishCompletion(h *SessionHost, stopReason string, promptErr error, finalize func()) { + // Release the admission gate before publishing the terminal status. The + // public AcceptPrompt path also requires HostReady, so the intermediate + // prompting/starting/error status cannot admit a new prompt. + h.releasePrompt(a) + if finalize != nil { + finalize() + } + close(a.done) + if cb := h.config.OnPromptComplete; cb != nil { + go cb(stopReason, promptErr) + } + if a.observer != nil { + go a.observer(stopReason, promptErr) + } +} diff --git a/packages/vm-agent/internal/acp/session_host_cancel.go b/packages/vm-agent/internal/acp/session_host_cancel.go new file mode 100644 index 0000000000..cac6e08a85 --- /dev/null +++ b/packages/vm-agent/internal/acp/session_host_cancel.go @@ -0,0 +1,213 @@ +package acp + +import ( + "context" + "encoding/json" + "fmt" + "log/slog" + "syscall" + "time" +) + +// CancelPrompt cancels the currently running Prompt() call, if any. +// This is safe to call from any goroutine. If no prompt is in flight, +// it's a no-op. The cancel function is guarded by promptCancelMu +// (separate from promptMu) so we never deadlock with HandlePrompt. +func (h *SessionHost) CancelPrompt() { + h.cancelPrompt(true) +} + +// CancelPromptFromControlPlane mirrors the viewer WebSocket session/cancel path +// for HTTP control-plane cancellation requests. +func (h *SessionHost) CancelPromptFromControlPlane() { + if h.AgentType() == "opencode" { + h.cancelPrompt(false) + h.StopProcessForPromptCancel() + return + } + + h.CancelPrompt() + cancelMessage, err := h.cancelNotification() + if err != nil { + slog.Warn("CancelPromptFromControlPlane: could not build session/cancel notification", "error", err) + } else { + h.ForwardToAgent(cancelMessage) + } + h.StopProcessForPromptCancel() +} + +func (h *SessionHost) cancelNotification() ([]byte, error) { + sessionID := h.currentSessionIDForCancel() + if sessionID == "" { + return nil, fmt.Errorf("missing session ID") + } + + return json.Marshal(struct { + JSONRPC string `json:"jsonrpc"` + Method string `json:"method"` + Params struct { + SessionID string `json:"sessionId"` + } `json:"params"` + }{ + JSONRPC: "2.0", + Method: "session/cancel", + Params: struct { + SessionID string `json:"sessionId"` + }{SessionID: sessionID}, + }) +} + +func (h *SessionHost) currentSessionIDForCancel() string { + h.mu.RLock() + acpSessionID := string(h.sessionID) + h.mu.RUnlock() + if acpSessionID != "" { + return acpSessionID + } + return h.config.SessionID +} + +func (h *SessionHost) cancelPrompt(startGraceTimer bool) { + h.promptCancelMu.Lock() + cancelFn := h.promptCancel + promptID := h.activePromptID + if cancelFn != nil { + h.promptCancelRequested = true + } + h.promptCancelMu.Unlock() + + if cancelFn == nil { + slog.Info("CancelPrompt: no prompt in flight") + return + } + + // Resolve the exact attempt being cancelled. promptMu is taken only after + // promptCancelMu is released (beginPrompt nests them the other way round). + // If the attempt already settled and a newer one began in between, the + // lookup misses and no watchdog is armed: this cancel's attempt is done. + attempt := h.promptAttemptByID(promptID) + slog.Info("CancelPrompt: cancelling in-flight prompt", "promptId", promptID) + h.reportLifecycle("info", "Prompt cancel requested", attempt.logFields(map[string]interface{}{ + "promptId": promptID, + })) + cancelFn() + + if !startGraceTimer || attempt == nil { + return + } + + grace := h.promptCancelGracePeriod() + if grace <= 0 { + return + } + + go h.watchPromptCancelGrace(attempt, grace) +} + +// watchPromptCancelGrace is bound to the one attempt a cancel targeted. It is +// disarmed as soon as that attempt reaches a terminal state, so a stale timer +// can never act on the next prompt (the 2026-09 regression where a follow-up +// accepted inside the grace window was force-stopped and its task failed). +func (h *SessionHost) watchPromptCancelGrace(attempt *promptAttempt, grace time.Duration) { + fired, stop := h.startCancelGraceTimer(grace) + defer stop() + select { + case <-attempt.done: + return + case <-h.lifecycleContext().Done(): + return + case <-fired: + } + h.triggerPromptForceStopIfStuck(attempt, fmt.Sprintf("Prompt cancel grace elapsed after %s", grace)) +} + +func (h *SessionHost) startCancelGraceTimer(d time.Duration) (<-chan time.Time, func()) { + if h.cancelGraceTimer != nil { + return h.cancelGraceTimer(d) + } + timer := time.NewTimer(d) + return timer.C, func() { timer.Stop() } +} + +// settleStuckPromptCancel handles a requested cancel whose attempt did not +// finish within the grace period. A stop is never a task failure: the attempt +// finishes "cancelled" and the agent process is restarted through the same +// intentional prompt-cancel path the control-plane Stop uses, which returns +// the host to ready for follow-up prompts. +func (h *SessionHost) settleStuckPromptCancel(attempt *promptAttempt, reason string, fields map[string]interface{}) { + attempt.completeWith(h, "cancelled", context.Canceled, func() { + h.stopPromptActivityRereport() + h.broadcastControl(MsgSessionPromptDone, nil) + + h.mu.RLock() + hasProcess := h.process != nil + h.mu.RUnlock() + if !hasProcess { + // Whoever cleared the process owns the host transition: Stop() + // marks it stopped, and monitorProcessExit is already restarting it + // (including when a crash-recovery restart skipped failing this + // attempt). Starting a second recovery here would race that owner. + slog.Warn("ACP prompt cancel did not settle; agent restart already in progress", "promptId", attempt.id, "reason", reason) + h.reportLifecycle("warn", "ACP prompt cancel did not settle; agent restart already in progress", fields) + return + } + slog.Warn("ACP prompt cancel did not settle; restarting agent", "promptId", attempt.id, "reason", reason) + h.reportLifecycle("warn", "ACP prompt cancel did not settle; restarting agent", fields) + h.StopProcessForPromptCancel() + }) +} + +// ForwardToAgent sends a raw message to the agent's stdin. +func (h *SessionHost) ForwardToAgent(message []byte) { + h.mu.RLock() + process := h.process + h.mu.RUnlock() + + if process == nil { + slog.Warn("No agent process running, dropping message") + return + } + + data := append(message, '\n') + if _, err := process.Stdin().Write(data); err != nil { + slog.Error("Failed to write to agent stdin", "error", err) + } +} + +// SignalProcess sends a signal to the agent process. This is used for agents +// that don't implement session/cancel (e.g., opencode) — SIGTERM is sent +// directly to the process instead of forwarding the cancel RPC. +func (h *SessionHost) SignalProcess(sig syscall.Signal) { + h.mu.RLock() + process := h.process + h.mu.RUnlock() + + if process == nil { + slog.Warn("SignalProcess: no agent process running") + return + } + + process.KillContainerProcesses(sig) + slog.Info("SignalProcess: sent signal to agent process", "signal", sig, "agentType", h.AgentType()) +} + +// StopProcessForPromptCancel terminates the current agent process for a user +// prompt cancel without marking the host stopped. The process monitor will +// restart the agent and return the host to ready for follow-up prompts. +func (h *SessionHost) StopProcessForPromptCancel() { + h.mu.Lock() + process := h.process + if process != nil { + h.intentionalPromptCancelProcessStop = true + } + h.mu.Unlock() + + if process == nil { + slog.Warn("StopProcessForPromptCancel: no agent process running") + return + } + + if err := process.Stop(); err != nil { + slog.Warn("StopProcessForPromptCancel: failed to stop agent process", "error", err) + } +} diff --git a/packages/vm-agent/internal/acp/session_host_cancel_grace_test.go b/packages/vm-agent/internal/acp/session_host_cancel_grace_test.go new file mode 100644 index 0000000000..6c053cf0c8 --- /dev/null +++ b/packages/vm-agent/internal/acp/session_host_cancel_grace_test.go @@ -0,0 +1,531 @@ +package acp + +import ( + "bufio" + "context" + "encoding/json" + "io" + "sync" + "testing" + "time" + + acpsdk "github.com/coder/acp-go-sdk" +) + +// cancelGraceFakeAgent is an ACP agent on the far side of the host's real +// stdin/stdout pipes. Every session/prompt request is parked until the test +// answers it, and session/cancel notifications are recorded. +type cancelGraceFakeAgent struct { + t *testing.T + reader *bufio.Reader + writer *io.PipeWriter + prompts chan json.RawMessage + cancels chan struct{} + // stopReading makes the agent stop draining stdin after the next prompt, + // modelling a wedged agent: the SDK's post-cancel session/cancel write then + // blocks and the cancelled prompt cannot settle on its own. + stopReading chan struct{} +} + +func (a *cancelGraceFakeAgent) serve() { + for { + line, err := a.reader.ReadBytes('\n') + if err != nil { + return + } + var msg struct { + ID json.RawMessage `json:"id"` + Method string `json:"method"` + } + if json.Unmarshal(line, &msg) != nil { + continue + } + switch msg.Method { + case "session/prompt": + a.prompts <- append(json.RawMessage(nil), msg.ID...) + select { + case <-a.stopReading: + return + default: + } + case "session/cancel": + a.cancels <- struct{}{} + } + } +} + +func (a *cancelGraceFakeAgent) respond(id json.RawMessage, stopReason string) { + data, _ := json.Marshal(map[string]any{ + "jsonrpc": "2.0", + "id": id, + "result": map[string]any{"stopReason": stopReason}, + }) + _, _ = a.writer.Write(append(data, '\n')) +} + +func (a *cancelGraceFakeAgent) nextPrompt(t *testing.T) json.RawMessage { + t.Helper() + select { + case id := <-a.prompts: + return id + case <-time.After(2 * time.Second): + t.Fatal("fake agent did not receive a session/prompt") + return nil + } +} + +// lifecycleRecorder is a goroutine-safe ErrorReporter: lifecycle reports are +// emitted from prompt, cancel, and watchdog goroutines. +type lifecycleRecorder struct { + mu sync.Mutex + reports []lifecycleReport +} + +func (r *lifecycleRecorder) record(level, message string, ctx map[string]interface{}) { + r.mu.Lock() + defer r.mu.Unlock() + r.reports = append(r.reports, lifecycleReport{level: level, message: message, context: ctx}) +} + +func (r *lifecycleRecorder) ReportError(err error, _ string, _ string, ctx map[string]interface{}) { + r.record("error", err.Error(), ctx) +} +func (r *lifecycleRecorder) ReportInfo(message, _ string, _ string, ctx map[string]interface{}) { + r.record("info", message, ctx) +} +func (r *lifecycleRecorder) ReportWarn(message, _ string, _ string, ctx map[string]interface{}) { + r.record("warn", message, ctx) +} + +func (r *lifecycleRecorder) Count(message string) int { + return len(r.find(message)) +} + +func (r *lifecycleRecorder) find(message string) []lifecycleReport { + r.mu.Lock() + defer r.mu.Unlock() + var out []lifecycleReport + for _, report := range r.reports { + if report.message == message { + out = append(out, report) + } + } + return out +} + +type promptCompletion struct { + stopReason string + err error +} + +// gatedGraceTimer lets a test own the load-bearing midpoint: the cancel-grace +// deadline fires only when the test releases it. +type gatedGraceTimer struct { + armed chan chan time.Time +} + +func (g *gatedGraceTimer) start(time.Duration) (<-chan time.Time, func()) { + fire := make(chan time.Time, 1) + g.armed <- fire + return fire, func() {} +} + +func (g *gatedGraceTimer) awaitArmed(t *testing.T) chan time.Time { + t.Helper() + select { + case fire := <-g.armed: + return fire + case <-time.After(2 * time.Second): + t.Fatal("cancel-grace watchdog was never armed") + return nil + } +} + +type cancelGraceHarness struct { + host *SessionHost + agent *cancelGraceFakeAgent + process *fakeAgentProcess + timer *gatedGraceTimer + completions chan promptCompletion + events *lifecycleRecorder +} + +func newCancelGraceHarness(t *testing.T) *cancelGraceHarness { + t.Helper() + return newCancelGraceHarnessWithProcess(t, false) +} + +// newCancelGraceHarnessWithProcess lets a test choose whether Stop() makes the +// fake agent process exit, which is what lets a real process monitor restart it. +func newCancelGraceHarnessWithProcess(t *testing.T, exitOnStop bool) *cancelGraceHarness { + t.Helper() + completions := make(chan promptCompletion, 8) + events := &lifecycleRecorder{} + host := NewSessionHost(SessionHostConfig{ + GatewayConfig: GatewayConfig{ + SessionID: "test-session", + WorkspaceID: "test-workspace", + ErrorReporter: events, + MessageReporter: &mockMessageReporter{}, + PromptCancelGracePeriod: 5 * time.Second, + OnPromptComplete: func(stopReason string, err error) { + completions <- promptCompletion{stopReason: stopReason, err: err} + }, + }, + MessageBufferSize: 100, + ViewerSendBuffer: 32, + }) + timer := &gatedGraceTimer{armed: make(chan chan time.Time, 4)} + host.cancelGraceTimer = timer.start + + process, agentReader, agentWriter := newFakeAgentProcess(time.Now().Add(-time.Minute), exitOnStop) + agent := &cancelGraceFakeAgent{ + t: t, + reader: bufio.NewReader(agentReader), + writer: agentWriter, + prompts: make(chan json.RawMessage, 4), + cancels: make(chan struct{}, 4), + stopReading: make(chan struct{}), + } + go agent.serve() + t.Cleanup(func() { + host.mu.Lock() + host.process = nil + host.mu.Unlock() + host.Stop() + _ = agentReader.Close() + _ = agentWriter.Close() + _ = process.stdin.Close() + _ = process.stdout.Close() + }) + + acpConn := acpsdk.NewClientSideConnection(&sessionHostClient{host: host}, process.stdin, process.stdout) + host.mu.Lock() + host.agentType = "claude-code" + host.process = process + host.acpConn = acpConn + host.setSessionIDLocked("acp-session-cancel") + host.setStatusLocked(HostReady) + host.mu.Unlock() + + return &cancelGraceHarness{host: host, agent: agent, process: process, timer: timer, completions: completions, events: events} +} + +func (h *cancelGraceHarness) nextCompletion(t *testing.T) promptCompletion { + t.Helper() + select { + case c := <-h.completions: + return c + case <-time.After(2 * time.Second): + t.Fatal("prompt completion callback did not fire") + return promptCompletion{} + } +} + +func (h *cancelGraceHarness) assertNoCompletion(t *testing.T) { + t.Helper() + select { + case c := <-h.completions: + t.Fatalf("unexpected prompt completion: %+v", c) + case <-time.After(100 * time.Millisecond): + } +} + +func cancelGracePromptParams(text string) json.RawMessage { + params, _ := json.Marshal(map[string]any{ + "prompt": []map[string]string{{"type": "text", "text": text}}, + }) + return params +} + +func (h *cancelGraceHarness) waitForReady(t *testing.T) { + t.Helper() + waitFor(t, 2*time.Second, func() bool { return h.host.Status() == HostReady }) +} + +// cancelTransport drives one production entry point end to end: how a prompt +// arrives and how the user's Stop reaches the host. +type cancelTransport struct { + name string + prompt func(t *testing.T, h *cancelGraceHarness, text string) + cancel func(h *cancelGraceHarness) +} + +var cancelTransports = []cancelTransport{ + { + // Browser viewer: prompt and Stop both arrive as WebSocket JSON-RPC frames. + name: "ws session/cancel", + prompt: func(t *testing.T, h *cancelGraceHarness, text string) { + g := &Gateway{host: h.host, viewerID: "viewer-1"} + frame, _ := json.Marshal(map[string]any{ + "jsonrpc": "2.0", "id": text, "method": "session/prompt", + "params": cancelGracePromptParams(text), + }) + g.handleMessage(context.Background(), frame) + }, + cancel: func(h *cancelGraceHarness) { + g := &Gateway{host: h.host, viewerID: "viewer-1"} + g.handleMessage(context.Background(), []byte(`{"jsonrpc":"2.0","method":"session/cancel","params":{}}`)) + }, + }, + { + // Control plane: a durable prompt delivery, and the HTTP cancel used by + // the chat Stop button and urgent stop-and-deliver. + name: "http control-plane cancel", + prompt: func(t *testing.T, h *cancelGraceHarness, text string) { + accepted, ok := h.host.AcceptPrompt(context.Background(), json.RawMessage(`"`+text+`"`), + cancelGracePromptParams(text), "control-plane", false, "delivery-"+text, nil) + if !ok { + t.Fatalf("control-plane prompt %q was not accepted", text) + } + go accepted.Run() + }, + cancel: func(h *cancelGraceHarness) { h.host.CancelPromptFromControlPlane() }, + }, +} + +// TestCancelGraceWatchdogNeverTouchesTheNextPrompt reproduces the 2026-09-28 +// incident ordering: prompt A is cancelled, A settles, prompt B is accepted +// inside the grace window, and only then does A's grace deadline fire. +func TestCancelGraceWatchdogNeverTouchesTheNextPrompt(t *testing.T) { + t.Parallel() + for _, transport := range cancelTransports { + transport := transport + t.Run(transport.name, func(t *testing.T) { + t.Parallel() + h := newCancelGraceHarness(t) + + transport.prompt(t, h, "prompt-a") + h.agent.nextPrompt(t) + waitFor(t, 2*time.Second, func() bool { return h.host.Status() == HostPrompting }) + + transport.cancel(h) + fireA := h.timer.awaitArmed(t) + + a := h.nextCompletion(t) + if a.stopReason != "cancelled" { + t.Fatalf("prompt A stopReason = %q, want cancelled", a.stopReason) + } + h.waitForReady(t) + + // Prompt B is accepted before A's grace deadline. + transport.prompt(t, h, "prompt-b") + bID := h.agent.nextPrompt(t) + waitFor(t, 2*time.Second, func() bool { return h.host.Status() == HostPrompting }) + stopsBefore := h.process.stopCount.Load() + + // Now A's deadline fires while B is running. + fireA <- time.Now() + h.assertNoCompletion(t) + if got := h.host.Status(); got != HostPrompting { + t.Fatalf("status after stale grace deadline = %s, want %s", got, HostPrompting) + } + if got := h.process.stopCount.Load(); got != stopsBefore { + t.Fatalf("agent process stopped %d times by a stale watchdog", got-stopsBefore) + } + if got := h.events.Count("ACP prompt force-stopped"); got != 0 { + t.Fatalf("force-stop events = %d, want 0", got) + } + + // Every cancel/start log names the prompt it belongs to. + requested := h.events.find("Prompt cancel requested") + if len(requested) != 1 || requested[0].context["promptId"] == nil { + t.Fatalf("cancel requested reports = %+v, want one with promptId", requested) + } + started := h.events.find("ACP Prompt started") + if len(started) != 2 || started[0].context["promptId"] == started[1].context["promptId"] { + t.Fatalf("prompt started reports = %+v, want two distinct promptIds", started) + } + if transport.name == "http control-plane cancel" && requested[0].context["deliveryId"] != "delivery-prompt-a" { + t.Fatalf("cancel requested deliveryId = %v, want delivery-prompt-a", requested[0].context["deliveryId"]) + } + + // B completes normally, exactly once. + h.agent.respond(bID, "end_turn") + b := h.nextCompletion(t) + if b.stopReason != "end_turn" || b.err != nil { + t.Fatalf("prompt B completion = %+v, want end_turn", b) + } + h.waitForReady(t) + h.assertNoCompletion(t) + }) + } +} + +// TestCancelGraceWatchdogSettlesAGenuinelyStuckCancel is the convergence +// control: when the cancelled prompt itself cannot settle, the watchdog still +// fires, reports "cancelled" (never a fatal task failure), and restarts the +// agent through the intentional prompt-cancel stop. That the intentional stop +// then restarts the agent back to ready is proven separately by +// TestSessionHost_MonitorIntentionalPromptCancelReportsIdleAfterSuccessfulRestart. +func TestCancelGraceWatchdogSettlesAGenuinelyStuckCancel(t *testing.T) { + t.Parallel() + for _, transport := range cancelTransports { + transport := transport + t.Run(transport.name, func(t *testing.T) { + t.Parallel() + h := newCancelGraceHarness(t) + close(h.agent.stopReading) + + transport.prompt(t, h, "prompt-a") + h.agent.nextPrompt(t) + waitFor(t, 2*time.Second, func() bool { return h.host.Status() == HostPrompting }) + + // The wedged agent no longer drains stdin, so the transport's own + // session/cancel write can block; production calls it from a + // request goroutine. + go transport.cancel(h) + fireA := h.timer.awaitArmed(t) + h.assertNoCompletion(t) + stopsBefore := h.process.stopCount.Load() + + fireA <- time.Now() + a := h.nextCompletion(t) + if a.stopReason != "cancelled" { + t.Fatalf("stuck cancel stopReason = %q, want cancelled", a.stopReason) + } + waitFor(t, 2*time.Second, func() bool { return h.process.stopCount.Load() > stopsBefore }) + if got := h.host.Status(); got == HostError { + t.Fatalf("host status = %s after stuck cancel, want a restart rather than error", got) + } + h.host.mu.RLock() + intentional := h.host.intentionalPromptCancelProcessStop + h.host.mu.RUnlock() + if !intentional { + t.Fatal("stuck cancel must stop the agent as an intentional prompt cancel so the monitor restarts it") + } + if got := h.events.Count("ACP prompt cancel did not settle; restarting agent"); got != 1 { + t.Fatalf("stuck-cancel lifecycle events = %d, want 1", got) + } + h.assertNoCompletion(t) + }) + } +} + +// TestCancelGraceWatchdogDisarmsWhenAttemptSettles proves the watchdog exits on +// the attempt's own completion instead of blocking on its timer. +func TestCancelGraceWatchdogDisarmsWhenAttemptSettles(t *testing.T) { + t.Parallel() + h := newCancelGraceHarness(t) + var stopped sync.WaitGroup + stopped.Add(1) + h.host.cancelGraceTimer = func(time.Duration) (<-chan time.Time, func()) { + fire := make(chan time.Time) + h.timer.armed <- fire + return fire, stopped.Done + } + + cancelTransports[0].prompt(t, h, "prompt-a") + h.agent.nextPrompt(t) + waitFor(t, 2*time.Second, func() bool { return h.host.Status() == HostPrompting }) + cancelTransports[0].cancel(h) + h.timer.awaitArmed(t) + if a := h.nextCompletion(t); a.stopReason != "cancelled" { + t.Fatalf("prompt A stopReason = %q, want cancelled", a.stopReason) + } + + done := make(chan struct{}) + go func() { stopped.Wait(); close(done) }() + select { + case <-done: + case <-time.After(2 * time.Second): + t.Fatal("cancel-grace watchdog stayed armed after its attempt settled") + } +} + +// TestCancelGraceStuckViewerCancelRestartsAgentBackToReady closes the loop for +// the viewer WebSocket transport, where the settled stuck cancel is the only +// thing that stops the wedged agent: the real process monitor must restart it +// back to ready, and a follow-up prompt must then run on the new agent. +func TestCancelGraceStuckViewerCancelRestartsAgentBackToReady(t *testing.T) { + t.Setenv("HOME", t.TempDir()) + t.Setenv("CODEX_HOME", "") + h := newCancelGraceHarnessWithProcess(t, true) + h.host.config.InitializeTimeoutMs = 500 + h.host.config.LoadSessionTimeoutMs = 500 + h.host.config.ContainerResolver = func() (string, error) { return "", nil } + var restarts sync.WaitGroup + restarts.Add(1) + h.host.config.StartProcess = func(*agentStartup) (agentProcess, error) { + defer restarts.Done() + proc, reader, writer := newFakeAgentProcess(time.Now(), false) + serveRecoveryACP(t, reader, writer) + return proc, nil + } + h.host.mu.Lock() + h.host.agentSupportsLoadSession = true + h.host.mu.Unlock() + go h.host.monitorProcessExit(h.process, "claude-code", &agentCredential{credentialKind: "api-key"}, nil) + + close(h.agent.stopReading) + ws := cancelTransports[0] + ws.prompt(t, h, "prompt-a") + h.agent.nextPrompt(t) + waitFor(t, 2*time.Second, func() bool { return h.host.Status() == HostPrompting }) + + go ws.cancel(h) + fireA := h.timer.awaitArmed(t) + if got := h.process.stopCount.Load(); got != 0 { + t.Fatalf("viewer cancel stopped the agent %d times before the grace deadline", got) + } + + fireA <- time.Now() + if a := h.nextCompletion(t); a.stopReason != "cancelled" { + t.Fatalf("stuck cancel stopReason = %q, want cancelled", a.stopReason) + } + waitFor(t, 5*time.Second, func() bool { return h.host.Status() == HostReady }) + restarts.Wait() + + ws.prompt(t, h, "prompt-b") + b := h.nextCompletion(t) + if b.err != nil || b.stopReason == fatalErrorStopReason { + t.Fatalf("follow-up prompt after restart = %+v, want a clean completion", b) + } + if got := h.host.Status(); got != HostReady { + t.Fatalf("status after follow-up = %s, want %s", got, HostReady) + } +} + +// TestCancelGraceStuckCancelDefersToAnInFlightRestart covers a stuck cancel +// whose agent process was already cleared by a concurrent restart (for example +// a crash-recovery restart that does not fail the active prompt). The attempt +// still settles "cancelled", and the watchdog does not claim a second restart. +func TestCancelGraceStuckCancelDefersToAnInFlightRestart(t *testing.T) { + t.Parallel() + h := newCancelGraceHarness(t) + close(h.agent.stopReading) + ws := cancelTransports[0] + ws.prompt(t, h, "prompt-a") + h.agent.nextPrompt(t) + waitFor(t, 2*time.Second, func() bool { return h.host.Status() == HostPrompting }) + + go ws.cancel(h) + fireA := h.timer.awaitArmed(t) + + // Model monitorProcessExit's restart branch having already taken the process. + h.host.mu.Lock() + h.host.process = nil + h.host.setStatusLocked(HostStarting) + h.host.mu.Unlock() + + fireA <- time.Now() + if a := h.nextCompletion(t); a.stopReason != "cancelled" { + t.Fatalf("stuck cancel stopReason = %q, want cancelled", a.stopReason) + } + waitFor(t, 2*time.Second, func() bool { + return h.events.Count("ACP prompt cancel did not settle; agent restart already in progress") == 1 + }) + if got := h.events.Count("ACP prompt cancel did not settle; restarting agent"); got != 0 { + t.Fatalf("restart claims = %d, want 0 when another owner is restarting", got) + } + if got := h.host.Status(); got != HostStarting { + t.Fatalf("status = %s, want the in-flight restart's %s left untouched", got, HostStarting) + } + h.host.mu.RLock() + intentional := h.host.intentionalPromptCancelProcessStop + h.host.mu.RUnlock() + if intentional { + t.Fatal("no intentional stop may be recorded without a process to stop") + } + h.assertNoCompletion(t) +} diff --git a/packages/vm-agent/internal/acp/session_host_execution_test.go b/packages/vm-agent/internal/acp/session_host_execution_test.go index 8ba7cd25bd..0396629b8c 100644 --- a/packages/vm-agent/internal/acp/session_host_execution_test.go +++ b/packages/vm-agent/internal/acp/session_host_execution_test.go @@ -169,8 +169,8 @@ func TestPromptTimeoutConvergesToErrorAndSingleFatalCallback(t *testing.T) { t.Fatal("prompt was not accepted") } host.setStatus(HostPrompting, "") - host.triggerPromptForceStopIfStuck(attempt.id, "hard deadline") - host.triggerPromptForceStopIfStuck(attempt.id, "late process exit") + host.triggerPromptForceStopIfStuck(attempt, "hard deadline") + host.triggerPromptForceStopIfStuck(attempt, "late process exit") if host.Status() != HostError { t.Fatalf("status = %s, want error", host.Status()) } diff --git a/packages/vm-agent/internal/acp/session_host_harness_work_test.go b/packages/vm-agent/internal/acp/session_host_harness_work_test.go index e0e53e82b1..faa1da816d 100644 --- a/packages/vm-agent/internal/acp/session_host_harness_work_test.go +++ b/packages/vm-agent/internal/acp/session_host_harness_work_test.go @@ -689,7 +689,7 @@ func TestPromptReportsUpdateHarnessCoalescerSuccessfulSnapshot(t *testing.T) { }}) t.Cleanup(host.Stop) - host.markPromptStarted(acpsdk.SessionId("acp-session"), 1, "viewer-1") + host.markPromptStarted(nil, acpsdk.SessionId("acp-session"), 1, "viewer-1") waitForActivitySnapshot(t, host, "prompting") host.nudgeHarnessActivityReport() time.Sleep(3 * debounce) diff --git a/packages/vm-agent/internal/acp/session_host_prompt.go b/packages/vm-agent/internal/acp/session_host_prompt.go index 62e339c0af..4e5948f562 100644 --- a/packages/vm-agent/internal/acp/session_host_prompt.go +++ b/packages/vm-agent/internal/acp/session_host_prompt.go @@ -23,7 +23,7 @@ import ( // untrusted browser prompt must not be able to mark its own content as // origin=system (which would hide it from search, dedup, topic, and attention). func (h *SessionHost) HandlePrompt(ctx context.Context, reqID json.RawMessage, params json.RawMessage, viewerID string, trustedSource bool) { - accepted, ok := h.AcceptPrompt(ctx, reqID, params, viewerID, trustedSource, nil) + accepted, ok := h.AcceptPrompt(ctx, reqID, params, viewerID, trustedSource, "", nil) if !ok { return } @@ -54,6 +54,7 @@ func (h *SessionHost) AcceptPrompt( params json.RawMessage, viewerID string, trustedSource bool, + deliveryID string, observer PromptTerminalObserver, ) (*AcceptedPrompt, bool) { if h.Status() != HostReady { @@ -66,7 +67,7 @@ func (h *SessionHost) AcceptPrompt( } promptCtx, promptCancel, promptTimeout := h.newPromptContext(ctx) - attempt, ok := h.beginPrompt(promptCancel, observer) + attempt, ok := h.beginPromptForDelivery(promptCancel, deliveryID, observer) if !ok { promptCancel() h.sendJSONRPCErrorToViewer(viewerID, reqID, -32603, "Prompt already in progress") @@ -107,7 +108,7 @@ func (p *AcceptedPrompt) Abort(err error) { // recovery takes ownership of the terminal outcome. func (p *AcceptedPrompt) Run() { h := p.host - promptDone := h.startPromptWatchdog(p.attempt.id, p.ctx, p.viewerID, p.reqID, p.timeout) + promptDone := h.startPromptWatchdog(p.attempt, p.ctx, p.viewerID, p.reqID, p.timeout) defer close(promptDone) defer p.cancel() defer close(p.attempt.rpcDone) @@ -118,7 +119,7 @@ func (p *AcceptedPrompt) Run() { // fields before the blocked Prompt returns "peer disconnected"; the captured // snapshot lets finishPromptWithError still begin LoadSession recovery. recovery := h.captureCrashRecoveryPrerequisites() - h.markPromptStarted(p.request.sessionID, len(p.request.blocks), p.viewerID) + h.markPromptStarted(p.attempt, p.request.sessionID, len(p.request.blocks), p.viewerID) resp, err := h.promptWithTransientRetry(p.ctx, p.request, p.attempt.startedAt) if !h.isPromptActive(p.attempt.id) { @@ -485,7 +486,7 @@ func (h *SessionHost) newPromptContext(ctx context.Context) (context.Context, co } func (h *SessionHost) startPromptWatchdog( - promptID uint64, + attempt *promptAttempt, promptCtx context.Context, viewerID string, reqID json.RawMessage, @@ -493,23 +494,27 @@ func (h *SessionHost) startPromptWatchdog( ) chan struct{} { promptDone := make(chan struct{}) if promptTimeout > 0 { - go h.watchPromptTimeout(promptID, promptCtx, promptDone, viewerID, reqID, promptTimeout) + go h.watchPromptTimeout(attempt, promptCtx, promptDone, viewerID, reqID, promptTimeout) } return promptDone } -func (h *SessionHost) markPromptStarted(sessionID acpsdk.SessionId, blockCount int, viewerID string) { +func (h *SessionHost) markPromptStarted(attempt *promptAttempt, sessionID acpsdk.SessionId, blockCount int, viewerID string) { h.setStatus(HostPrompting, "") h.broadcastControl(MsgSessionPrompting, nil) h.reportActivity("prompting") h.startPromptActivityRereport() - slog.Info("ACP: sending Prompt", "sessionID", string(sessionID), "blockCount", blockCount) - h.reportLifecycle("info", "ACP Prompt started", map[string]interface{}{ + fields := attempt.logFields(map[string]interface{}{ "acpSessionId": string(sessionID), "blockCount": blockCount, "viewerId": viewerID, }) + if attempt != nil { + fields["promptId"] = attempt.id + } + slog.Info("ACP: sending Prompt", "sessionID", string(sessionID), "blockCount", blockCount, "promptId", fields["promptId"]) + h.reportLifecycle("info", "ACP Prompt started", fields) } // markPromptDone is the single turn-end hook shared by normal completion, @@ -583,11 +588,12 @@ func (h *SessionHost) finishPrompt( return } - slog.Info("ACP: Prompt completed", "stopReason", string(resp.StopReason)) - h.reportLifecycle("info", "ACP Prompt completed", map[string]interface{}{ + slog.Info("ACP: Prompt completed", "stopReason", string(resp.StopReason), "promptId", attempt.id) + h.reportLifecycle("info", "ACP Prompt completed", attempt.logFields(map[string]interface{}{ "stopReason": string(resp.StopReason), "duration": time.Since(info.startedAt).String(), - }) + "promptId": attempt.id, + })) h.checkStderrForSilentErrors(resp.StopReason) h.broadcastPromptResponse(reqID, resp) } @@ -596,10 +602,11 @@ func (h *SessionHost) finishPromptCancelled(attempt *promptAttempt, reqID json.R if !attempt.completeWith(h, "cancelled", context.Canceled, h.markPromptDone) { return } - slog.Info("ACP: Prompt cancelled") - h.reportLifecycle("info", "ACP Prompt cancelled", map[string]interface{}{ + slog.Info("ACP: Prompt cancelled", "promptId", attempt.id) + h.reportLifecycle("info", "ACP Prompt cancelled", attempt.logFields(map[string]interface{}{ "duration": time.Since(info.startedAt).String(), - }) + "promptId": attempt.id, + })) h.broadcastMessage(h.marshalJSONRPCError(reqID, -32800, "Prompt cancelled")) } diff --git a/packages/vm-agent/internal/acp/session_host_prompt_state.go b/packages/vm-agent/internal/acp/session_host_prompt_state.go index 6d08fe4ee1..0f11d64fad 100644 --- a/packages/vm-agent/internal/acp/session_host_prompt_state.go +++ b/packages/vm-agent/internal/acp/session_host_prompt_state.go @@ -5,6 +5,7 @@ import ( "encoding/json" "errors" "fmt" + "log/slog" "sync/atomic" "time" ) @@ -30,6 +31,10 @@ func (h *SessionHost) promptCancelGracePeriod() time.Duration { } func (h *SessionHost) beginPrompt(cancel context.CancelFunc, observer PromptTerminalObserver) (*promptAttempt, bool) { + return h.beginPromptForDelivery(cancel, "", observer) +} + +func (h *SessionHost) beginPromptForDelivery(cancel context.CancelFunc, deliveryID string, observer PromptTerminalObserver) (*promptAttempt, bool) { h.promptMu.Lock() defer h.promptMu.Unlock() if h.promptInFlight { @@ -38,12 +43,13 @@ func (h *SessionHost) beginPrompt(cancel context.CancelFunc, observer PromptTerm h.promptInFlight = true promptID := atomic.AddUint64(&h.promptSeq, 1) attempt := &promptAttempt{ - id: promptID, - startedAt: h.now(), - cancel: cancel, - done: make(chan struct{}), - rpcDone: make(chan struct{}), - observer: observer, + id: promptID, + startedAt: h.now(), + cancel: cancel, + deliveryID: deliveryID, + done: make(chan struct{}), + rpcDone: make(chan struct{}), + observer: observer, } h.promptAttempt = attempt @@ -71,24 +77,15 @@ func (h *SessionHost) releasePrompt(attempt *promptAttempt) { h.promptCancelMu.Unlock() } -func (h *SessionHost) promptAttemptForID(promptID uint64) *promptAttempt { +// promptAttemptByID returns the current attempt only when its identity +// matches. It never fabricates an attempt: in-flight prompt state does not +// survive a vm-agent restart, so an unknown ID is always a settled prompt. +func (h *SessionHost) promptAttemptByID(promptID uint64) *promptAttempt { h.promptMu.Lock() defer h.promptMu.Unlock() if h.promptAttempt != nil && h.promptAttempt.id == promptID { return h.promptAttempt } - // Preserve the focused-test and upgrade seam for hosts whose prompt state - // was established before the per-attempt arbiter existed. - if h.promptInFlight { - attempt := &promptAttempt{ - id: promptID, - startedAt: h.now(), - done: make(chan struct{}), - rpcDone: make(chan struct{}), - } - h.promptAttempt = attempt - return attempt - } return nil } @@ -114,7 +111,7 @@ func (h *SessionHost) isPromptCancelRequested(promptID uint64) bool { } func (h *SessionHost) watchPromptTimeout( - promptID uint64, + attempt *promptAttempt, promptCtx context.Context, done <-chan struct{}, viewerID string, @@ -130,22 +127,42 @@ func (h *SessionHost) watchPromptTimeout( } msg := fmt.Sprintf("Prompt timed out after %s", timeout) h.sendJSONRPCErrorToViewer(viewerID, reqID, -32603, msg) - h.triggerPromptForceStopIfStuck(promptID, msg) + h.triggerPromptForceStopIfStuck(attempt, msg) } } -func (h *SessionHost) triggerPromptForceStopIfStuck(promptID uint64, reason string) { - attempt := h.promptAttemptForID(promptID) - stopReason := fatalErrorStopReason - terminalErr := error(fmt.Errorf("%s", reason)) - if h.isPromptCancelRequested(promptID) { - stopReason = "cancelled" - terminalErr = context.Canceled - } +// triggerPromptForceStopIfStuck acts only on the exact attempt a watchdog was +// armed for, and only while that attempt is still the current, non-terminal +// one. A watchdog outliving its attempt is a no-op: it must never touch a +// later prompt (see idea 01M31M9G3T4SEWT9ZW1BM4QKZ3). +func (h *SessionHost) triggerPromptForceStopIfStuck(attempt *promptAttempt, reason string) { if attempt == nil { return } - attempt.completeWith(h, stopReason, terminalErr, func() { + h.promptMu.Lock() + current := h.promptAttempt + h.promptMu.Unlock() + var currentID uint64 + if current != nil { + currentID = current.id + } + cancelRequested := h.isPromptCancelRequested(attempt.id) + fields := attempt.logFields(map[string]interface{}{ + "reason": reason, + "promptId": attempt.id, + "currentPromptId": currentID, + "cancelRequested": cancelRequested, + }) + if current != attempt || attempt.isTerminal() { + slog.Info("ACP prompt force-stop skipped: attempt already settled", + "promptId", attempt.id, "currentPromptId", currentID, "reason", reason) + return + } + if cancelRequested { + h.settleStuckPromptCancel(attempt, reason, fields) + return + } + attempt.completeWith(h, fatalErrorStopReason, errors.New(reason), func() { h.mu.Lock() agentType := h.agentType if h.status == HostPrompting { @@ -155,9 +172,7 @@ func (h *SessionHost) triggerPromptForceStopIfStuck(promptID uint64, reason stri h.stopCurrentAgentLocked() h.mu.Unlock() - h.reportLifecycle("error", "ACP prompt force-stopped", map[string]interface{}{ - "reason": reason, - }) + h.reportLifecycle("error", "ACP prompt force-stopped", fields) h.broadcastControl(MsgSessionPromptDone, nil) h.broadcastAgentStatus(StatusError, agentType, reason) // A hard deadline is a terminal error, never an idle transition. diff --git a/packages/vm-agent/internal/acp/session_host_settings.go b/packages/vm-agent/internal/acp/session_host_settings.go new file mode 100644 index 0000000000..6c12aeae18 --- /dev/null +++ b/packages/vm-agent/internal/acp/session_host_settings.go @@ -0,0 +1,134 @@ +package acp + +import ( + "context" + "fmt" + "log/slog" + "time" + + acpsdk "github.com/coder/acp-go-sdk" +) + +// phaseTimeout returns a per-phase timeout duration. If phaseMs is > 0, it is +// used; otherwise the fallback timeout is returned. +func phaseTimeout(phaseMs int, fallback time.Duration) time.Duration { + if phaseMs > 0 { + return time.Duration(phaseMs) * time.Millisecond + } + return fallback +} + +// applySessionSettings calls SetSessionConfigOption for the model and +// SetSessionMode on the ACP connection. A rejected explicit Codex model is fatal: +// continuing would silently run a different model. Other adapter settings remain +// best-effort for backward compatibility. +func (h *SessionHost) applySessionSettings(ctx context.Context, settings *agentSettingsPayload) error { + if settings == nil || h.acpConn == nil || h.sessionID == "" { + return nil + } + + // These RPCs run while h.mu is held for write, and the ACP SDK blocks each + // response on its notification worker catching up (waitNotificationsUpTo). + // The caller's ctx is the WebSocket connection lifetime — effectively + // unbounded — so a stalled notification worker would hang the handshake (and + // every h.mu reader) forever. Bound it explicitly; Codex model selection + // failures are returned while the remaining settings stay best-effort. + settingsCtx, cancel := context.WithTimeout(ctx, h.sessionSettingsTimeout()) + defer cancel() + ctx = settingsCtx + + if settings.Model != "" { + if err := h.applySessionModelConfigOption(ctx, settings.Model); err != nil && h.agentType == "openai-codex" { + return fmt.Errorf("cannot apply requested Codex model %q: %w", settings.Model, err) + } + } + + if settings.PermissionMode != "" && settings.PermissionMode != "default" { + // Codex (openai-codex) does not support SetSessionMode — skip to avoid + // a guaranteed error on every session start. + if h.agentType == "openai-codex" { + slog.Info("ACP: skipping SetSessionMode for openai-codex (unsupported)", "mode", settings.PermissionMode) + } else { + slog.Info("ACP: setting session mode", "mode", settings.PermissionMode) + if _, err := h.acpConn.SetSessionMode(ctx, acpsdk.SetSessionModeRequest{ + SessionId: h.sessionID, + ModeId: acpsdk.SessionModeId(settings.PermissionMode), + }); err != nil { + slog.Warn("ACP SetSessionMode failed (non-fatal)", "mode", settings.PermissionMode, "error", err) + h.reportLifecycle("warn", "ACP SetSessionMode failed", map[string]interface{}{ + "mode": settings.PermissionMode, + "error": err.Error(), + }) + } else { + slog.Info("ACP: session mode set", "mode", settings.PermissionMode) + h.reportLifecycle("info", "ACP session mode applied", map[string]interface{}{ + "mode": settings.PermissionMode, + }) + } + } + } + return nil +} + +func (h *SessionHost) applySessionModelConfigOption(ctx context.Context, model string) error { + modelConfigID, ok := findModelConfigOptionID(h.configOptions) + if !ok { + slog.Warn("ACP session model config option unavailable", "model", model) + h.reportLifecycle("warn", "ACP session model config option unavailable", map[string]interface{}{ + "model": model, + }) + return fmt.Errorf("model config option unavailable") + } + + slog.Info("ACP: setting session model config option", "model", model, "configId", string(modelConfigID)) + resp, err := h.acpConn.SetSessionConfigOption(ctx, acpsdk.SetSessionConfigOptionRequest{ + ValueId: &acpsdk.SetSessionConfigOptionValueId{ + SessionId: h.sessionID, + ConfigId: modelConfigID, + Value: acpsdk.SessionConfigValueId(model), + }, + }) + if err != nil { + slog.Warn("ACP SetSessionConfigOption failed", "model", model, "configId", string(modelConfigID), "error", err) + h.reportLifecycle("warn", "ACP session model config option failed", map[string]interface{}{ + "model": model, + "configId": string(modelConfigID), + "error": err.Error(), + }) + return fmt.Errorf("set session model config option: %w", err) + } + + h.configOptions = resp.ConfigOptions + slog.Info("ACP: session model config option set", "model", model, "configId", string(modelConfigID)) + h.reportLifecycle("info", "ACP session model applied", map[string]interface{}{ + "model": model, + "configId": string(modelConfigID), + }) + return nil +} + +func findModelConfigOptionID(options []acpsdk.SessionConfigOption) (acpsdk.SessionConfigId, bool) { + for _, option := range options { + if option.Select == nil || option.Select.Category == nil { + continue + } + if *option.Select.Category == acpsdk.SessionConfigOptionCategoryModel { + return option.Select.Id, true + } + } + return "", false +} + +// DefaultSessionSettingsTimeout bounds the post-handshake SetSessionMode / +// SetSessionConfigOption RPCs. Override via the existing NewSessionTimeoutMs +// gateway setting (ACP_NEW_SESSION_TIMEOUT_MS). +const DefaultSessionSettingsTimeout = 30 * time.Second + +// sessionSettingsTimeout resolves the bound for applySessionSettings' RPCs, +// reusing the configured NewSession timeout when one is set. +func (h *SessionHost) sessionSettingsTimeout() time.Duration { + if h.config.NewSessionTimeoutMs > 0 { + return time.Duration(h.config.NewSessionTimeoutMs) * time.Millisecond + } + return DefaultSessionSettingsTimeout +} diff --git a/packages/vm-agent/internal/acp/session_host_stderr.go b/packages/vm-agent/internal/acp/session_host_stderr.go new file mode 100644 index 0000000000..0ba05b1212 --- /dev/null +++ b/packages/vm-agent/internal/acp/session_host_stderr.go @@ -0,0 +1,96 @@ +package acp + +import ( + "bufio" + "log/slog" + "strings" + + acpsdk "github.com/coder/acp-go-sdk" +) + +// DefaultStderrBufferBytes is the default maximum agent stderr captured for +// crash reports. Override via ACP_STDERR_BUFFER_BYTES. +const DefaultStderrBufferBytes = 4096 + +// monitorStderr reads the agent's stderr and collects it for error reporting. +func (h *SessionHost) monitorStderr(process agentProcess) { + scanner := bufio.NewScanner(process.Stderr()) + for scanner.Scan() { + line := redactAgentDiagnosticText(scanner.Text()) + slog.Warn("Agent stderr", "line", line) + h.stderrMu.Lock() + if h.stderrBuf.Len() < h.config.StderrBufferBytes { + if h.stderrBuf.Len() > 0 { + h.stderrBuf.WriteByte('\n') + } + h.stderrBuf.WriteString(line) + } + h.stderrMu.Unlock() + } +} + +func (h *SessionHost) getAndClearStderr() string { + h.stderrMu.Lock() + defer h.stderrMu.Unlock() + s := h.stderrBuf.String() + h.stderrBuf.Reset() + return s +} + +func (h *SessionHost) peekStderr() string { + h.stderrMu.Lock() + defer h.stderrMu.Unlock() + return h.stderrBuf.String() +} + +// silentErrorPatterns are stderr substrings that indicate an API-level error +// the agent may have swallowed (returning a normal end_turn instead of an error). +var silentErrorPatterns = []string{ + "AI_APICallError", + "Unauthorized", + "401", + "403", + "invalid_api_key", + "authentication_error", +} + +// checkStderrForSilentErrors peeks at the accumulated stderr buffer for known +// API error patterns. Some agents (notably OpenCode with Scaleway) silently +// swallow API errors and return {stopReason: "end_turn"} instead of an error. +// When detected, we log a warning and report a lifecycle event so the UI can +// surface the issue. The stderr buffer is NOT cleared — it remains available +// for crash reporting in monitorProcessExit. +func (h *SessionHost) checkStderrForSilentErrors(stopReason acpsdk.StopReason) { + h.stderrMu.Lock() + stderr := h.stderrBuf.String() + h.stderrMu.Unlock() + + if stderr == "" { + return + } + + for _, pattern := range silentErrorPatterns { + if strings.Contains(stderr, pattern) { + slog.Warn("ACP: possible silent API error detected in stderr after prompt completion", + "stopReason", string(stopReason), + "pattern", pattern, + "stderrSnippet", truncateString(stderr, 512), + "agentType", h.AgentType(), + ) + h.reportLifecycle("warn", "Possible silent API error — check agent credentials", map[string]interface{}{ + "stopReason": string(stopReason), + "errorPattern": pattern, + "stderrSnippet": truncateString(stderr, 256), + }) + return // report once per prompt, not per pattern + } + } +} + +// truncateString returns s truncated to maxLen with "..." appended if needed. +func truncateString(s string, maxLen int) string { + if len(s) <= maxLen { + return s + } + return s[:maxLen] + "..." +} diff --git a/packages/vm-agent/internal/acp/session_host_test.go b/packages/vm-agent/internal/acp/session_host_test.go index 0393d2dc68..03821915bd 100644 --- a/packages/vm-agent/internal/acp/session_host_test.go +++ b/packages/vm-agent/internal/acp/session_host_test.go @@ -1532,23 +1532,19 @@ func TestSessionHost_ForceStoppedPromptReportsFatalCompletionExactlyOnce(t *test completed <- completion{stopReason: stopReason, err: err} } - const promptID = uint64(42) const timeoutReason = "Prompt timed out after 6h0m0s" - host.promptCancelMu.Lock() - host.activePromptID = promptID - host.promptCancelMu.Unlock() - host.promptMu.Lock() - host.promptInFlight = true - host.promptMu.Unlock() - if host.promptAttemptForID(promptID) == nil { - t.Fatal("promptAttemptForID returned nil for in-flight prompt") + _, cancel := context.WithCancel(context.Background()) + defer cancel() + attempt, ok := host.beginPrompt(cancel, nil) + if !ok { + t.Fatal("prompt was not accepted") } host.mu.Lock() host.setStatusLocked(HostPrompting) host.agentType = "openai-codex" host.mu.Unlock() - host.triggerPromptForceStopIfStuck(promptID, timeoutReason) + host.triggerPromptForceStopIfStuck(attempt, timeoutReason) select { case got := <-completed: @@ -1563,7 +1559,7 @@ func TestSessionHost_ForceStoppedPromptReportsFatalCompletionExactlyOnce(t *test } // A late watchdog/retry for the same prompt must not terminalize twice. - host.triggerPromptForceStopIfStuck(promptID, timeoutReason) + host.triggerPromptForceStopIfStuck(attempt, timeoutReason) select { case got := <-completed: t.Fatalf("received duplicate force-stop completion: %+v", got) @@ -1586,17 +1582,12 @@ func TestSessionHost_CompetingPromptCompletionPathsClaimExactlyOnce(t *testing.T completed <- completion{stopReason: stopReason, err: err} } - const promptID = uint64(43) const timeoutReason = "Prompt timed out after 6h0m0s" - host.promptCancelMu.Lock() - host.activePromptID = promptID - host.promptCancelMu.Unlock() - host.promptMu.Lock() - host.promptInFlight = true - host.promptMu.Unlock() - attempt := host.promptAttemptForID(promptID) - if attempt == nil { - t.Fatal("promptAttemptForID returned nil for in-flight prompt") + _, cancel := context.WithCancel(context.Background()) + defer cancel() + attempt, ok := host.beginPrompt(cancel, nil) + if !ok { + t.Fatal("prompt was not accepted") } host.mu.Lock() host.setStatusLocked(HostPrompting) @@ -1614,7 +1605,7 @@ func TestSessionHost_CompetingPromptCompletionPathsClaimExactlyOnce(t *testing.T go func() { defer contenders.Done() <-start - host.triggerPromptForceStopIfStuck(promptID, timeoutReason) + host.triggerPromptForceStopIfStuck(attempt, timeoutReason) }() close(start) contenders.Wait() @@ -2588,88 +2579,6 @@ func TestSessionUpdate_EmptyUpdate_NoEnqueue(t *testing.T) { } } -func TestSessionHost_CancelPrompt_ForceStopsAfterGracePeriod(t *testing.T) { - t.Parallel() - - host := NewSessionHost(SessionHostConfig{ - GatewayConfig: GatewayConfig{ - SessionID: "test-session", - WorkspaceID: "test-workspace", - PromptCancelGracePeriod: 10 * time.Millisecond, - }, - MessageBufferSize: 100, - ViewerSendBuffer: 32, - }) - defer host.Stop() - - host.mu.Lock() - host.setStatusLocked(HostPrompting) - host.agentType = "claude-code" - host.mu.Unlock() - - host.promptMu.Lock() - host.promptInFlight = true - host.promptMu.Unlock() - - ctx, cancel := context.WithCancel(context.Background()) - host.promptCancelMu.Lock() - host.promptCancel = cancel - host.activePromptID = 42 - host.promptCancelMu.Unlock() - - host.CancelPrompt() - - select { - case <-ctx.Done(): - case <-time.After(500 * time.Millisecond): - t.Fatal("expected prompt context to be cancelled") - } - - deadline := time.Now().Add(1 * time.Second) - for host.Status() != HostError { - if time.Now().After(deadline) { - t.Fatal("expected host to transition to error after cancel grace elapsed") - } - time.Sleep(10 * time.Millisecond) - } - - host.mu.RLock() - statusErr := host.statusErr - host.mu.RUnlock() - if !strings.Contains(statusErr, "Prompt cancel grace elapsed") { - t.Fatalf("statusErr = %q, expected cancel grace reason", statusErr) - } - - host.promptCancelMu.Lock() - if host.activePromptID != 0 { - t.Fatalf("activePromptID = %d, want 0", host.activePromptID) - } - if host.promptCancel != nil { - t.Fatal("promptCancel should be cleared after force-stop") - } - host.promptCancelMu.Unlock() - - host.promptMu.Lock() - if host.promptInFlight { - t.Fatal("promptInFlight should be false after force-stop") - } - host.promptMu.Unlock() - - bufferDeadline := time.Now().Add(500 * time.Millisecond) - for { - host.bufMu.RLock() - buffered := len(host.messageBuf) - host.bufMu.RUnlock() - if buffered >= 2 { - break - } - if time.Now().After(bufferDeadline) { - t.Fatalf("expected prompt_done + error status messages, buffered=%d", buffered) - } - time.Sleep(10 * time.Millisecond) - } -} - func TestHandlePrompt_InjectsSyntheticUserMessage(t *testing.T) { t.Parallel() diff --git a/packages/vm-agent/internal/config/config.go b/packages/vm-agent/internal/config/config.go index 546e90ec5b..e534e5051a 100644 --- a/packages/vm-agent/internal/config/config.go +++ b/packages/vm-agent/internal/config/config.go @@ -266,7 +266,7 @@ type Config struct { ACPPongTimeout time.Duration // WebSocket pong deadline after ping (default: 10s) ACPPromptTimeout time.Duration // Max prompt runtime; 0 = no timeout (default: 0). Used for workspace sessions; task sessions use ACPTaskPromptTimeout via effectivePromptTimeout(). ACPTaskPromptTimeout time.Duration // Max prompt runtime for task-driven sessions; 0 = no timeout (default: 8h) - ACPPromptCancelGrace time.Duration // Wait after cancel before force-stop fallback (default: 5s) + ACPPromptCancelGrace time.Duration // Wait for a cancelled prompt to settle before finishing it cancelled and restarting the agent (default: 5s) ACPPromptRetryMaxRetries int // Retryable transient provider prompt errors after initial attempt (default: 2) ACPPromptRetryInitial time.Duration // Initial backoff for transient provider prompt retries (default: 15s) ACPPromptRetryMax time.Duration // Max backoff for transient provider prompt retries (default: 2m) diff --git a/packages/vm-agent/internal/server/workspaces.go b/packages/vm-agent/internal/server/workspaces.go index 3137cc7bc9..69c28886f7 100644 --- a/packages/vm-agent/internal/server/workspaces.go +++ b/packages/vm-agent/internal/server/workspaces.go @@ -1528,7 +1528,7 @@ func (s *Server) startAgentWithPromptObserved(host *acp.SessionHost, workspaceID host.HandlePrompt(ctx, syntheticReqID, promptParams, "server", true) return } - accepted, ok := host.AcceptPrompt(ctx, syntheticReqID, promptParams, "server", true, observer) + accepted, ok := host.AcceptPrompt(ctx, syntheticReqID, promptParams, "server", true, "", observer) if !ok { promptErr := errors.New("initial prompt was not accepted by the session host") observer("error", promptErr) @@ -1695,7 +1695,7 @@ func (s *Server) handleVersionedPromptDelivery( observer := s.promptReceiptObserver(workspaceID, sessionID, deliveryID) accepted, ok := host.AcceptPrompt(context.Background(), reqID, promptParams, - "control-plane", false, observer) + "control-plane", false, deliveryID, observer) if !ok { receipt.RuntimeIdentity = s.executionRuntimeID writeVersionedPromptResponse(w, http.StatusConflict, "not_ready", sessionID, receipt) diff --git a/tasks/archive/2026-09-29-bind-prompt-cancel-watchdog-to-attempt.md b/tasks/archive/2026-09-29-bind-prompt-cancel-watchdog-to-attempt.md new file mode 100644 index 0000000000..c14168cd23 --- /dev/null +++ b/tasks/archive/2026-09-29-bind-prompt-cancel-watchdog-to-attempt.md @@ -0,0 +1,124 @@ +# Bind the prompt-cancel grace watchdog to the cancelled attempt + +## Problem + +The VM agent's 5 s "cancel grace" watchdog (`cancelPrompt(true)`) is keyed by a +numeric prompt ID and is never disarmed. If the next prompt starts inside the +window, the stale timer force-stops the **new** prompt: `promptAttemptForID` +fabricates an attempt for the stale ID and overwrites `h.promptAttempt`, the stop is +classified `fatalErrorStopReason` because `activePromptID` moved on, the new +prompt's agent is killed, and the task goes `failed` with +`Prompt cancel grace elapsed after 5s`. + +Production 2026-08-30 → 2026-09-29: 91 cancels, 4 had the next prompt start within +5 s, all 4 were killed, 0 true positives. Regression from PR #1785. Full evidence: +idea `01M31M9G3T4SEWT9ZW1BM4QKZ3`. This task is Fix Plan section A (VM agent). The +control-plane amplifier (section B: per-session delivery single-flight, one stop per +urgent delivery) is a separate follow-up PR. + +## Research Findings + +- `session_host.go` `cancelPrompt` spawns `go func(){ <-timer.C; triggerPromptForceStopIfStuck(id) }`. +- `session_host_prompt_state.go` `promptAttemptForID` fabricates an attempt whenever + `promptInFlight` and the ID does not match (the "focused-test and upgrade seam"). + In-flight prompt state never survives a vm-agent restart, so there is no upgrade case. +- `triggerPromptForceStopIfStuck` finalizer sets `HostError`, calls + `stopCurrentAgentLocked()` (no restart), and reports `fatalErrorStopReason` + unless cancel was requested for that ID. +- Both cancel transports arm the watchdog: HTTP `CancelPromptFromControlPlane` + (`server/workspaces.go` cancel handler) and WS `session/cancel` (`gateway.go`). +- `StopProcessForPromptCancel` sets `intentionalPromptCancelProcessStop`, and + `monitorProcessExit` then restarts the agent back to Ready with LoadSession of the + previous ACP session. That is the correct recovery for a genuinely stuck cancel. +- `watchPromptTimeout` (hard prompt timeout) also calls the force-stop; it is already + bound to `promptDone`, but should pass the attempt too. +- acp-go-sdk `Prompt` returns on ctx cancel, but first writes `session/cancel` with + `context.Background()`; a blocked agent stdin can therefore hold `Run` — the real + "stuck cancel" case the watchdog still needs to cover. +- Tests at `session_host_test.go` (`TestSessionHost_CancelPrompt_ForceStopsAfterGracePeriod`, + `TestSessionHost_ForceStoppedPromptReportsFatalCompletionExactlyOnce`, + `TestSessionHost_CompetingPromptCompletionPathsClaimExactlyOnce`) hand-build + `promptInFlight=true` and depend on the fabricate seam. +- `finishPromptWithError` test seam only creates an attempt when none exists; it + cannot overwrite a live attempt, so it is not part of this bug. +- Log lines `ACP Prompt started/cancelled/completed`, `Prompt cancel requested`, + `ACP prompt force-stopped` carry no prompt identity. +- Rule 18: `session_host.go` is 1295 lines with no exception header. + +## Implementation Checklist + +- [x] Commit 1 (pure move): split `session_host.go` below 800 lines — cancel block → + `session_host_cancel.go`; promptAttempt/checkpoint episode → + `session_host_attempt.go`; session settings → `session_host_settings.go`; + stderr helpers → `session_host_stderr.go`; MCP server builders → + `session_host_mcp.go` +- [x] Arm the cancel watchdog with the exact `*promptAttempt`; select on + `attempt.done`, `h.ctx.Done()`, and an injectable grace timer +- [x] Force-stop is attempt-bound: no-op (with log) unless `h.promptAttempt == attempt` + and the attempt is non-terminal; `watchPromptTimeout` passes the attempt +- [x] Delete `promptAttemptForID` and its fabricate branch +- [x] A stuck *requested* cancel finishes `cancelled` and restarts the agent via the + intentional prompt-cancel process stop (never `HostError`/fatal) +- [x] Observability: `promptId` (+ `deliveryId` for control-plane prompts) on + `ACP Prompt started/cancelled/completed`, `Prompt cancel requested`; force-stop + logs `{promptId, currentPromptId, cancelRequested}` +- [x] Tests: real prompts via fake ACP agent + gated timer, both HTTP and WS cancel + paths, next prompt accepted before deadline, then release timer → B untouched, + no HostError, one completion per prompt +- [x] Convergence control: fake agent that blocks cancel → watchdog fires, outcome + `cancelled`, agent restart requested, host not in error +- [x] Rewrite hand-built-state tests to drive real accepted attempts +- [x] Discrimination: revert to ID lookup + fabricate seam → new test goes red +- [ ] Update idea 01M31M9G3T4SEWT9ZW1BM4QKZ3 with PR evidence + +## Acceptance Criteria + +- [x] A stale cancel-grace timer never affects a later prompt (Go test, both transports) +- [x] A genuinely stuck cancel reports `cancelled`, not `failed`, and restarts the agent +- [x] Hard prompt timeout still reports fatal exactly once +- [x] `go test ./...` and `go vet` pass for `packages/vm-agent` +- [x] Staging: VM provisioned, heartbeat, prompt → Stop → immediate follow-up + completes without task failure + +## Implementation Notes + +- Stale-watchdog discrimination (2026-09-29): removing the `attempt.done` disarm and + re-pointing the force-stop at the current attempt (pre-fix semantics) made + `TestCancelGraceWatchdogNeverTouchesTheNextPrompt/{ws,http}` fail with prompt B + completing `fatal_error`, the incident signature, and + `TestCancelGraceWatchdogDisarmsWhenAttemptSettles` fail. Restored → green. +- Convergence discrimination: disabling the stuck-cancel settle branch made + `TestCancelGraceWatchdogSettlesAGenuinelyStuckCancel/{ws,http}` report + `fatal_error` instead of `cancelled`. Restored → green. +- Stuck cancel is modelled realistically: the fake agent stops draining stdin, so + the ACP SDK's post-cancel `session/cancel` write blocks and `Run` cannot settle. +- `finishPromptWithError` seam only creates an attempt when none exists, so it + cannot overwrite a live attempt; left as is. +- The HTTP transport test calls `CancelPromptFromControlPlane` directly; the HTTP + handler is a thin `IsPrompting()` guard in front of it. + +## Staging Evidence (2026-09-29) + +- Branch deployed to staging (run 36530902316, success). Fresh node + `01M3NYCH7N2JK1QAACABAKF4NJ` reported `agent_version` `4ef45be85` (this branch) and + healthy heartbeats; Claude Code VM task `01M3NYC8SHR8N7JFR0HY137SV1`. +- Three HTTP control-plane cancels of live turns, each followed immediately by a + durable follow-up. Cycle 2 hit the incident window: prompt 10 was cancelled at + 06:52:44.101 and follow-up prompt 11 started at 06:52:48.976 (4.875 s later, inside + the 5 s grace). The stale timer would have fired about 125 ms into prompt 11; + instead prompt 11 ran 17 s until the next deliberate cancel. +- Workspace logs: zero warn/error, zero `ACP prompt force-stopped`, zero grace or + "did not settle" events. Task ended `in_progress` / `awaiting_followup`, and the + final follow-up answered `PONG`. +- New log identity is live: `Prompt cancel requested`, `ACP Prompt started` and + `ACP Prompt cancelled` carry `promptId` and `deliveryId`. +- Cleanup: session stopped, node deleted and confirmed gone. +- Found: Claude Code blocks a bare `sleep N` as a standalone command; use a + `python3 -c "import time; time.sleep(N)"` busy-wait for staging timing tests. + +## References + +- Idea `01M31M9G3T4SEWT9ZW1BM4QKZ3` +- `.claude/rules/62-tests-must-observe-the-real-trigger.md` +- `packages/vm-agent/.claude/rules/54-vm-agent-rollout-compatibility.md` +- `.claude/rules/18-file-size-limits.md` diff --git a/tasks/backlog/2026-09-29-flaky-harness-activity-coalesce-test.md b/tasks/backlog/2026-09-29-flaky-harness-activity-coalesce-test.md new file mode 100644 index 0000000000..c057aca4b8 --- /dev/null +++ b/tasks/backlog/2026-09-29-flaky-harness-activity-coalesce-test.md @@ -0,0 +1,28 @@ +# Flaky: TestHarnessActivityReportCoalescesACPToolCallBursts under -race load + +## Problem + +`packages/vm-agent/internal/acp/session_host_harness_work_test.go` +`TestHarnessActivityReportCoalescesACPToolCallBursts` failed once in 10 runs of +`go test -race -count=10 ./internal/acp/` with +`tool-call burst was not coalesced: reports=2`. + +The test sends 13 ACP tool-call notifications and expects one debounced activity +report with a 20 ms debounce window. Under full-package `-race` load the burst can +take longer than 20 ms to deliver, so the debounce fires mid-burst and a second +report is sent. The code under test is likely fine; the test's timing is too tight. +Other debounce/timing tests in the package showed the same pattern +(`TestACPToolCallSettlingLeaseStopsRereporting`). + +## Context + +Discovered 2026-09-29 while re-verifying the prompt-cancel watchdog fix +(`tasks/active/2026-09-29-bind-prompt-cancel-watchdog-to-attempt.md`). That change +does not touch the harness activity reporter. + +## Acceptance Criteria + +- [ ] The coalescing test owns its timing (fake clock or a debounce gate) instead of + relying on a 20 ms wall-clock window +- [ ] `go test -race -count=50 -run 'Harness|ToolCallSettling' ./internal/acp/` passes +- [ ] The test still fails if coalescing is removed (discrimination check)