diff --git a/apps/backend/internal/agent/runtime/lifecycle/manager_execution.go b/apps/backend/internal/agent/runtime/lifecycle/manager_execution.go index 4b06a84848..d2c0be4952 100644 --- a/apps/backend/internal/agent/runtime/lifecycle/manager_execution.go +++ b/apps/backend/internal/agent/runtime/lifecycle/manager_execution.go @@ -348,6 +348,24 @@ func (m *Manager) GetExecutionIDForSession(_ context.Context, sessionID string) return "", fmt.Errorf("%w: %s", ErrNoExecutionForSession, sessionID) } +// GetACPSessionIDForSession returns the ACP conversation currently owned by a +// live execution. The orchestrator uses this optional accessor after a context +// reset to persist the new conversation immediately, instead of depending on +// an asynchronous session-created event arriving before a backend restart. +func (m *Manager) GetACPSessionIDForSession(sessionID string) (string, bool) { + execution, exists := m.executionStore.GetBySessionID(sessionID) + if !exists || execution == nil { + return "", false + } + var acpSessionID string + if err := m.executionStore.WithRLock(execution.ID, func(exec *AgentExecution) { + acpSessionID = exec.ACPSessionID + }); err != nil || acpSessionID == "" { + return "", false + } + return acpSessionID, true +} + // IsAgentCommandConfigured reports whether an execution has been promoted from // workspace-only infrastructure to an agent execution ready to start. func (m *Manager) IsAgentCommandConfigured(executionID string) bool { diff --git a/apps/backend/internal/orchestrator/event_handlers.go b/apps/backend/internal/orchestrator/event_handlers.go index 2589e890f6..f4b48c5530 100644 --- a/apps/backend/internal/orchestrator/event_handlers.go +++ b/apps/backend/internal/orchestrator/event_handlers.go @@ -61,6 +61,22 @@ func (s *Service) handleACPSessionCreated(ctx context.Context, data watcher.ACPS // session/load vs session/new in session.go — agents without native resume (e.g., // Claude Code) use the token for their own --resume CLI flag instead. func (s *Service) storeResumeToken(ctx context.Context, taskID, sessionID, expectedExecID, acpSessionID, lastMessageUUID string) { + // The lifecycle manager updates its in-memory ACP session ID before it + // publishes reset/start events. Events from the previous ACP session can + // still be queued after that point, so reject those events before the + // execution-level CAS. The execution ID alone does not identify an ACP + // session generation because context resets keep the same execution. + if currentACPSessionID := s.currentACPSessionID(sessionID); currentACPSessionID != "" && + acpSessionID != "" && acpSessionID != currentACPSessionID { + s.logger.Info("dropping resume token from stale ACP session generation", + zap.String("task_id", taskID), + zap.String("session_id", sessionID), + zap.String("expected_exec_id", expectedExecID), + zap.String("resume_token", acpSessionID), + zap.String("current_resume_token", currentACPSessionID)) + return + } + err := s.repo.UpdateResumeToken(ctx, sessionID, expectedExecID, acpSessionID, lastMessageUUID) switch { case err == nil: @@ -104,6 +120,23 @@ func (s *Service) storeResumeToken(ctx context.Context, taskID, sessionID, expec } } +// currentACPSessionID returns the lifecycle manager's current ACP session ID +// when the concrete manager exposes it. The optional seam keeps the generic +// AgentManagerClient contract unchanged for remote clients and test doubles. +func (s *Service) currentACPSessionID(sessionID string) string { + provider, ok := s.agentManager.(interface { + GetACPSessionIDForSession(string) (string, bool) + }) + if !ok { + return "" + } + acpSessionID, ok := provider.GetACPSessionIDForSession(sessionID) + if !ok { + return "" + } + return acpSessionID +} + // persistACPSessionID mirrors the agent's ACP session id into the session's // "acp" metadata map. Best-effort: resume correctness never depends on this // copy — it exists so the id survives executors_running cleanup for consumers diff --git a/apps/backend/internal/orchestrator/event_handlers_agent.go b/apps/backend/internal/orchestrator/event_handlers_agent.go index 652dc7735f..82bf213a5e 100644 --- a/apps/backend/internal/orchestrator/event_handlers_agent.go +++ b/apps/backend/internal/orchestrator/event_handlers_agent.go @@ -1466,20 +1466,24 @@ func (s *Service) wasResumeAttempt(ctx context.Context, sessionID string) bool { } // clearResumeToken removes the resume token from the executor running record so -// the next agent start won't use --resume. It is reserved for explicit -// user-initiated fresh-start recovery; ordinary ACP startup failures retain the -// token so the session can be retried. +// the next agent start won't use --resume. Callers use this for explicit fresh +// starts and after a successful context reset; ordinary ACP startup failures +// retain the token so the session can be retried. // // Unconditional clear: passes expectedExecID="" so the narrow update is not // CAS-guarded — clearing a token is always intentional regardless of which // execution is currently registered. -func (s *Service) clearResumeToken(ctx context.Context, sessionID string) { +func (s *Service) clearResumeToken(ctx context.Context, sessionID string) error { err := s.repo.UpdateResumeToken(ctx, sessionID, "", "", "") - if err != nil && !errors.Is(err, models.ErrExecutorRunningNotFound) { + if errors.Is(err, models.ErrExecutorRunningNotFound) { + return nil + } + if err != nil { s.logger.Error("failed to clear resume token", zap.String("session_id", sessionID), zap.Error(err)) } + return err } // handleRecoverableFailure handles agent failures by keeping the session recoverable. diff --git a/apps/backend/internal/orchestrator/event_handlers_test.go b/apps/backend/internal/orchestrator/event_handlers_test.go index 7994ca66b9..1ad8538c90 100644 --- a/apps/backend/internal/orchestrator/event_handlers_test.go +++ b/apps/backend/internal/orchestrator/event_handlers_test.go @@ -322,6 +322,8 @@ type mockAgentManager struct { repoForExecutionLookup interface { GetExecutorRunningBySessionID(ctx context.Context, sessionID string) (*models.ExecutorRunning, error) } + // Optional current ACP session lookup used by reset-token generation tests. + getACPSessionIDForSessionFunc func(string) (string, bool) // CancelAgent tracking. cancelAgentCalls counts every invocation. If // cancelAgentBlock is non-nil, CancelAgent blocks on it before returning; @@ -664,6 +666,14 @@ func (m *mockAgentManager) GetExecutionIDForSession(ctx context.Context, session } return "", fmt.Errorf("no execution found") } + +func (m *mockAgentManager) GetACPSessionIDForSession(sessionID string) (string, bool) { + if m.getACPSessionIDForSessionFunc == nil { + return "", false + } + return m.getACPSessionIDForSessionFunc(sessionID) +} + func (m *mockAgentManager) GetGitLog(ctx context.Context, sessionID, baseCommit string, limit int, targetBranch string) (*client.GitLogResult, error) { if m.getGitLogFunc != nil { return m.getGitLogFunc(ctx, sessionID, baseCommit, limit, targetBranch) diff --git a/apps/backend/internal/orchestrator/event_handlers_workflow.go b/apps/backend/internal/orchestrator/event_handlers_workflow.go index e052296171..9c678f2d37 100644 --- a/apps/backend/internal/orchestrator/event_handlers_workflow.go +++ b/apps/backend/internal/orchestrator/event_handlers_workflow.go @@ -2947,10 +2947,34 @@ func (s *Service) resetAgentContext(ctx context.Context, taskID string, session return true } + releaseLifecycleLock := s.acquireSessionLifecycleLock(sessionID) + defer releaseLifecycleLock() + s.setSessionResetInProgress(sessionID, true) + defer s.setSessionResetInProgress(sessionID, false) + executionID, err := s.agentManager.GetExecutionIDForSession(ctx, sessionID) if err != nil || executionID == "" { - s.logger.Debug("no agent execution for context reset, skipping", - zap.String("session_id", sessionID)) + // No in-memory execution exists yet — most commonly a lazily-resumed + // session whose process has not been relaunched since the last run. + // The resume path (applyRunningRecordToResumeRequest) reads the ACP + // resume token straight from the executors_running row, bypassing any + // in-memory execution lookup entirely, so leaving that token in place + // here would let the next lazy launch reconnect to the pre-reset + // conversation and silently skip the reset. Clear the same persisted + // state the live-execution path clears below so the reset survives + // until the agent's first turn regardless of when the process starts. + s.logger.Debug("no live agent execution for context reset, clearing persisted resume state", + zap.String("session_id", sessionID), + zap.String("step_name", stepName)) + if err := s.clearResumeToken(ctx, sessionID); err != nil { + s.logger.Error("failed to clear lazy resume token before context reset", + zap.String("task_id", taskID), + zap.String("session_id", sessionID), + zap.String("step_name", stepName), + zap.Error(err)) + return false + } + s.clearPersistedResetState(ctx, sessionID, session) return true } @@ -2960,9 +2984,6 @@ func (s *Service) resetAgentContext(ctx context.Context, taskID string, session zap.String("step_name", stepName), zap.String("agent_execution_id", executionID)) - s.setSessionResetInProgress(sessionID, true) - defer s.setSessionResetInProgress(sessionID, false) - if err := s.agentManager.ResetAgentContext(ctx, executionID); err != nil { s.logger.Error("failed to reset agent context", zap.String("task_id", taskID), @@ -2972,6 +2993,33 @@ func (s *Service) resetAgentContext(ctx context.Context, taskID string, session return false } + // Clear the old resume token only after the provider reset succeeds. This + // keeps a valid recovery token when the runtime reset fails. A fresh ACP + // session event can race this clear, so persist the lifecycle manager's + // current session ID again below after the clear. + if err := s.clearResumeToken(ctx, sessionID); err != nil { + s.logger.Error("failed to clear resume token after context reset", + zap.String("task_id", taskID), + zap.String("session_id", sessionID), + zap.String("step_name", stepName), + zap.Error(err)) + return false + } + if acpSessionID := s.currentACPSessionID(sessionID); acpSessionID != "" { + s.storeResumeToken(ctx, taskID, sessionID, executionID, acpSessionID, "") + } + + // Clear the remaining persisted state (ACP session metadata, context window) + // after the provider reset succeeds. The token is handled explicitly above. + s.clearPersistedResetState(ctx, sessionID, session) + return true +} + +// clearPersistedResetState clears the durable, DB-backed session state that a +// later lazy resume would otherwise pick back up: the stored ACP session ID +// in session metadata and the persisted context window. The resume token is +// cleared explicitly by resetAgentContext so a reset failure can retain it. +func (s *Service) clearPersistedResetState(ctx context.Context, sessionID string, session *models.TaskSession) { // Clear the stored ACP session ID using json_set to avoid clobbering other keys. if updateErr := s.repo.SetSessionMetadataKey(ctx, sessionID, "acp_session_id", ""); updateErr != nil { s.logger.Warn("failed to clear ACP session ID from session metadata", @@ -2988,7 +3036,6 @@ func (s *Service) resetAgentContext(ctx context.Context, taskID string, session // its cache after the provider reset succeeds and must not receive stale data // back through the final processOnEnter state event. clearInMemoryContextWindow(session) - return true } // resolveSessionMCPSupport checks if the agent for a session supports MCP. diff --git a/apps/backend/internal/orchestrator/event_handlers_workflow_lazy_resume_test.go b/apps/backend/internal/orchestrator/event_handlers_workflow_lazy_resume_test.go new file mode 100644 index 0000000000..8643819bcb --- /dev/null +++ b/apps/backend/internal/orchestrator/event_handlers_workflow_lazy_resume_test.go @@ -0,0 +1,283 @@ +package orchestrator + +import ( + "context" + "errors" + "testing" + + "github.com/kandev/kandev/internal/task/models" + wfmodels "github.com/kandev/kandev/internal/workflow/models" +) + +type failResumeTokenUpdateRepo struct { + sessionExecutorStore + err error +} + +func (r *failResumeTokenUpdateRepo) UpdateResumeToken(context.Context, string, string, string, string) error { + return r.err +} + +// TestProcessOnEnterResetAgentContext_ClearsLazyResumeTokenWithoutLiveExecution +// is the regression test for the lazy-resume reset: when no in-memory execution +// exists, resetAgentContext must erase the stale resume token so the next lazy +// launch does not reconnect to the pre-reset ACP conversation. +func TestProcessOnEnterResetAgentContext_ClearsLazyResumeTokenWithoutLiveExecution(t *testing.T) { + ctx := context.Background() + repo := setupTestRepo(t) + seedSession(t, repo, "task-lazy-resume", "session-lazy-resume", "step-work") + if err := repo.UpsertExecutorRunning(ctx, &models.ExecutorRunning{ + ID: "session-lazy-resume", + SessionID: "session-lazy-resume", + TaskID: "task-lazy-resume", + ResumeToken: "old-acp-session", + Status: "stopped", + }); err != nil { + t.Fatalf("seed resumable execution: %v", err) + } + + agentManager := &mockAgentManager{repoForExecutionLookup: repo} + svc := createTestServiceWithAgent(repo, newMockStepGetter(), newMockTaskRepo(), agentManager) + step := &wfmodels.WorkflowStep{ + ID: "step-review", WorkflowID: "workflow-1", Name: "Review", + Events: wfmodels.StepEvents{OnEnter: []wfmodels.OnEnterAction{ + {Type: wfmodels.OnEnterResetAgentContext}, + }}, + } + session, err := repo.GetTaskSession(ctx, "session-lazy-resume") + if err != nil { + t.Fatalf("load session: %v", err) + } + + svc.processOnEnter(ctx, "task-lazy-resume", session, step, "review task") + + if len(agentManager.restartProcessCalls) != 0 { + t.Fatalf("expected no reset against a missing live execution, got %d calls", len(agentManager.restartProcessCalls)) + } + running, err := repo.GetExecutorRunningBySessionID(ctx, "session-lazy-resume") + if err != nil { + t.Fatalf("load resumable execution: %v", err) + } + if running.ResumeToken != "" { + t.Fatalf("reset must clear lazy resume before the next agent turn, got %q", running.ResumeToken) + } +} + +func TestProcessOnEnterResetAgentContext_ReportsLazyResumeTokenClearFailure(t *testing.T) { + ctx := context.Background() + repo := setupTestRepo(t) + seedSession(t, repo, "task-lazy-resume-error", "session-lazy-resume-error", "step-work") + if err := repo.UpsertExecutorRunning(ctx, &models.ExecutorRunning{ + ID: "session-lazy-resume-error", + SessionID: "session-lazy-resume-error", + TaskID: "task-lazy-resume-error", + ResumeToken: "old-acp-session", + Status: "stopped", + }); err != nil { + t.Fatalf("seed resumable execution: %v", err) + } + + agentManager := &mockAgentManager{repoForExecutionLookup: repo} + svc := createTestServiceWithAgent(repo, newMockStepGetter(), newMockTaskRepo(), agentManager) + svc.repo = &failResumeTokenUpdateRepo{ + sessionExecutorStore: repo, + err: errors.New("resume token store unavailable"), + } + session, err := repo.GetTaskSession(ctx, "session-lazy-resume-error") + if err != nil { + t.Fatalf("load session: %v", err) + } + + if svc.resetAgentContext(ctx, "task-lazy-resume-error", session, "review") { + t.Fatal("expected reset to fail when lazy resume token cannot be cleared") + } + if len(agentManager.restartProcessCalls) != 0 { + t.Fatalf("expected no provider reset without a live execution, got %d calls", len(agentManager.restartProcessCalls)) + } +} + +func TestResetAgentContext_FailedProviderResetRetainsResumeToken(t *testing.T) { + ctx := context.Background() + repo := setupTestRepo(t) + seedSession(t, repo, "t-reset-token", "s-reset-token", "step1") + seedExecutorRunning(t, repo, "s-reset-token", "t-reset-token", "exec-reset-token") + if err := repo.UpdateResumeToken(ctx, "s-reset-token", "exec-reset-token", "old-acp-session", ""); err != nil { + t.Fatalf("seed stale resume token: %v", err) + } + + svc := createTestServiceWithAgent( + repo, + newMockStepGetter(), + newMockTaskRepo(), + &mockAgentManager{ + repoForExecutionLookup: repo, + restartProcessErr: errors.New("provider reset failed"), + }, + ) + session, err := repo.GetTaskSession(ctx, "s-reset-token") + if err != nil { + t.Fatalf("load session: %v", err) + } + if svc.resetAgentContext(ctx, "t-reset-token", session, "review") { + t.Fatal("expected provider reset to fail") + } + running, err := repo.GetExecutorRunningBySessionID(ctx, "s-reset-token") + if err != nil { + t.Fatalf("load executor row: %v", err) + } + if running.ResumeToken != "old-acp-session" { + t.Fatalf("failed reset must retain recovery token, got %q", running.ResumeToken) + } +} + +// TestResetAgentContext_InterleavingB_ClearBeforeStore proves that a successful +// reset persists the fresh token before an async session-created event arrives. +func TestResetAgentContext_InterleavingB_ClearBeforeStore(t *testing.T) { + ctx := context.Background() + repo := setupTestRepo(t) + seedSession(t, repo, "t1", "s1", "step1") + + const execID = "exec-1" + const staleToken = "old-acp-session" + const freshToken = "fresh-acp-session" + seedExecutorRunning(t, repo, "s1", "t1", execID) + if err := repo.UpdateResumeToken(ctx, "s1", execID, staleToken, ""); err != nil { + t.Fatalf("seed stale resume token: %v", err) + } + + agentManager := &mockAgentManager{ + repoForExecutionLookup: repo, + getACPSessionIDForSessionFunc: func(string) (string, bool) { + return freshToken, true + }, + } + svc := createTestServiceWithAgent(repo, newMockStepGetter(), newMockTaskRepo(), agentManager) + session, err := repo.GetTaskSession(ctx, "s1") + if err != nil { + t.Fatalf("load session: %v", err) + } + + // 1. resetAgentContext resets the provider, clears stale state, and + // persists the current ACP session synchronously. + if !svc.resetAgentContext(ctx, "t1", session, "test") { + t.Fatal("expected reset to succeed") + } + + // 2. Verify the old token was replaced without waiting for an async event. + running, err := repo.GetExecutorRunningBySessionID(ctx, "s1") + if err != nil { + t.Fatalf("load executor row after reset: %v", err) + } + if running.ResumeToken != freshToken { + t.Fatalf("expected fresh resume_token after reset, got %q", running.ResumeToken) + } + + // 3. An async session-created event with the same current ID remains + // idempotent. + svc.storeResumeToken(ctx, "t1", "s1", execID, freshToken, "") + + running, err = repo.GetExecutorRunningBySessionID(ctx, "s1") + if err != nil { + t.Fatalf("load executor row after storeResumeToken: %v", err) + } + if running.ResumeToken != freshToken { + t.Fatalf("expected fresh token %q after storeResumeToken, got %q", freshToken, running.ResumeToken) + } +} + +// TestResetAgentContext_InterleavingA_StoreBeforeClear proves that metadata +// cleanup does not clear a token that was already persisted by the fresh ACP +// session event. +func TestResetAgentContext_InterleavingA_StoreBeforeClear(t *testing.T) { + ctx := context.Background() + repo := setupTestRepo(t) + seedSession(t, repo, "t1", "s1", "step1") + + const execID = "exec-1" + const staleToken = "old-acp-session" + seedExecutorRunning(t, repo, "s1", "t1", execID) + if err := repo.UpdateResumeToken(ctx, "s1", execID, staleToken, ""); err != nil { + t.Fatalf("seed stale resume token: %v", err) + } + + agentManager := &mockAgentManager{repoForExecutionLookup: repo} + svc := createTestServiceWithAgent(repo, newMockStepGetter(), newMockTaskRepo(), agentManager) + session, err := repo.GetTaskSession(ctx, "s1") + if err != nil { + t.Fatalf("load session: %v", err) + } + + // Simulate the async ACP session.created event arriving before the remaining + // reset metadata is cleared. The token is handled separately and survives. + const freshToken = "fresh-acp-session" + svc.storeResumeToken(ctx, "t1", "s1", execID, freshToken, "") + svc.clearPersistedResetState(ctx, "s1", session) + + running, err := repo.GetExecutorRunningBySessionID(ctx, "s1") + if err != nil { + t.Fatalf("load executor row: %v", err) + } + if running.ResumeToken == "" { + t.Fatal("RACE: clearResumeToken erased the fresh token written by storeResumeToken") + } + if running.ResumeToken != freshToken { + t.Fatalf("expected fresh token %q to survive, got %q", freshToken, running.ResumeToken) + } +} + +// TestResetAgentContext_InterleavingC_StaleOldEventOverwritesFresh proves that +// the current ACP session ID filters a delayed event from the old session even +// though both sessions share one lifecycle execution ID. +func TestResetAgentContext_InterleavingC_StaleOldEventDoesNotOverwriteFresh(t *testing.T) { + ctx := context.Background() + repo := setupTestRepo(t) + seedSession(t, repo, "t1", "s1", "step1") + + const execID = "exec-1" + const staleToken = "old-acp-session" + seedExecutorRunning(t, repo, "s1", "t1", execID) + if err := repo.UpdateResumeToken(ctx, "s1", execID, staleToken, ""); err != nil { + t.Fatalf("seed stale resume token: %v", err) + } + + const freshToken = "fresh-acp-session" + agentManager := &mockAgentManager{ + repoForExecutionLookup: repo, + getACPSessionIDForSessionFunc: func(string) (string, bool) { + return freshToken, true + }, + } + svc := createTestServiceWithAgent(repo, newMockStepGetter(), newMockTaskRepo(), agentManager) + session, err := repo.GetTaskSession(ctx, "s1") + if err != nil { + t.Fatalf("load session: %v", err) + } + + // 1. Reset the provider and persist the fresh ACP session. + if !svc.resetAgentContext(ctx, "t1", session, "test") { + t.Fatal("expected reset to succeed") + } + + // 2. New session event arrives — writes the fresh token. + svc.storeResumeToken(ctx, "t1", "s1", execID, freshToken, "") + + running, err := repo.GetExecutorRunningBySessionID(ctx, "s1") + if err != nil { + t.Fatalf("load executor row: %v", err) + } + if running.ResumeToken != freshToken { + t.Fatalf("expected fresh token %q, got %q", freshToken, running.ResumeToken) + } + + // 3. A stale old-session event arrives with the old ACP session ID. The + // execution ID is unchanged, but the ACP generation check rejects it. + svc.storeResumeToken(ctx, "t1", "s1", execID, staleToken, "") + + running, err = repo.GetExecutorRunningBySessionID(ctx, "s1") + if err != nil { + t.Fatalf("load executor row after stale event: %v", err) + } + if running.ResumeToken != freshToken { + t.Fatalf("stale old-session event overwrote fresh token: got %q, want %q", running.ResumeToken, freshToken) + } +} diff --git a/apps/backend/internal/orchestrator/service.go b/apps/backend/internal/orchestrator/service.go index 7849a13c3c..622fc59089 100644 --- a/apps/backend/internal/orchestrator/service.go +++ b/apps/backend/internal/orchestrator/service.go @@ -691,6 +691,13 @@ type Service struct { // Session reset flags: sessionID -> true while resetAgentContext is restarting process. // Used to suppress stale ready events and avoid draining queued prompts mid-reset. resetInProgressSessions sync.Map + // sessionLifecycleLocks serializes context resets with every path that can + // resume a session from executors_running. The lock is intentionally shared + // by resetAgentContext and ensureSessionRunning so a lazy resume cannot read + // the old token while a workflow reset is clearing it. + // Entries are not deleted: deleting a lock can let a new caller create a + // second mutex while an existing waiter still owns the old one. + sessionLifecycleLocks sync.Map // map[sessionID]*sync.Mutex // contextWindowGuards serializes live context usage writes and gives each // session a generation boundary that reset_agent_context can invalidate. contextWindowGuards sync.Map @@ -1778,6 +1785,20 @@ func (s *Service) setSessionResetInProgress(sessionID string, inProgress bool) { s.resetInProgressSessions.Delete(sessionID) } +// acquireSessionLifecycleLock serializes operations that change or consume a +// session's persisted conversation identity. Keep the critical section around +// the complete reset/resume operation, including the bounded runtime wait, so +// no caller can launch with a token that a concurrent reset is invalidating. +func (s *Service) acquireSessionLifecycleLock(sessionID string) func() { + if sessionID == "" { + return func() {} + } + value, _ := s.sessionLifecycleLocks.LoadOrStore(sessionID, &sync.Mutex{}) + lock := value.(*sync.Mutex) + lock.Lock() + return lock.Unlock +} + func (s *Service) isSessionResetInProgress(sessionID string) bool { if sessionID == "" { return false diff --git a/apps/backend/internal/orchestrator/session_launch.go b/apps/backend/internal/orchestrator/session_launch.go index ac6334733c..6afa83c936 100644 --- a/apps/backend/internal/orchestrator/session_launch.go +++ b/apps/backend/internal/orchestrator/session_launch.go @@ -378,7 +378,9 @@ func (s *Service) RecoverSession(ctx context.Context, taskID, sessionID, action } switch action { case "fresh_start": - s.clearResumeToken(ctx, sessionID) + if err := s.clearResumeToken(ctx, sessionID); err != nil { + return nil, fmt.Errorf("failed to clear resume token for fresh start: %w", err) + } case "resume": // no-op — relaunch with existing resume token default: diff --git a/apps/backend/internal/orchestrator/task_operations.go b/apps/backend/internal/orchestrator/task_operations.go index 27cec8807a..bf37ea5054 100644 --- a/apps/backend/internal/orchestrator/task_operations.go +++ b/apps/backend/internal/orchestrator/task_operations.go @@ -1625,6 +1625,8 @@ func (s *Service) ResumeTaskSession(ctx context.Context, taskID, sessionID strin s.logger.Debug("resuming task session", zap.String("task_id", taskID), zap.String("session_id", sessionID)) + releaseLifecycleLock := s.acquireSessionLifecycleLock(sessionID) + defer releaseLifecycleLock() session, err := s.repo.GetTaskSession(ctx, sessionID) if err != nil { @@ -1934,6 +1936,9 @@ func (s *Service) advanceTaskWorkflowStep(ctx context.Context, task *models.Task // After lazy recovery, a session may be in WAITING_FOR_INPUT with no agent process; // this function detects that case and triggers a resume. func (s *Service) ensureSessionRunning(ctx context.Context, sessionID string, session *models.TaskSession) error { + releaseLifecycleLock := s.acquireSessionLifecycleLock(sessionID) + defer releaseLifecycleLock() + isOfficeTask, err := s.lookupOfficeTask(ctx, session.TaskID) if err != nil { return fmt.Errorf("failed to determine office task status: %w", err)