diff --git a/cli/src/agent/loopBase.ts b/cli/src/agent/loopBase.ts index f7cd349b..6e2c35bd 100644 --- a/cli/src/agent/loopBase.ts +++ b/cli/src/agent/loopBase.ts @@ -3,6 +3,27 @@ import type { AgentSessionBase } from './sessionBase'; export type LoopLauncher = (session: TSession) => Promise<'switch' | 'exit'>; +export async function runLocalRemoteSession>(opts: { + session: TSession; + startingMode?: 'local' | 'remote'; + logTag: string; + runLocal: LoopLauncher; + runRemote: LoopLauncher; + onSessionReady?: (session: TSession) => void; +}): Promise { + if (opts.onSessionReady) { + opts.onSessionReady(opts.session); + } + + await runLocalRemoteLoop({ + session: opts.session, + startingMode: opts.startingMode, + logTag: opts.logTag, + runLocal: opts.runLocal, + runRemote: opts.runRemote + }); +} + export async function runLocalRemoteLoop>(opts: { session: TSession; startingMode?: 'local' | 'remote'; diff --git a/cli/src/agent/runnerLifecycle.ts b/cli/src/agent/runnerLifecycle.ts new file mode 100644 index 00000000..f21ec30a --- /dev/null +++ b/cli/src/agent/runnerLifecycle.ts @@ -0,0 +1,141 @@ +import type { ApiSessionClient } from '@/api/apiSession' +import { logger } from '@/ui/logger' +import { restoreTerminalState } from '@/ui/terminalState' + +type RunnerLifecycleOptions = { + session: ApiSessionClient + logTag: string + stopKeepAlive?: () => void + onBeforeClose?: () => Promise | void + onAfterClose?: () => Promise | void +} + +export type RunnerLifecycle = { + setExitCode: (code: number) => void + setArchiveReason: (reason: string) => void + markCrash: (error: unknown) => void + cleanup: () => Promise + cleanupAndExit: (codeOverride?: number) => Promise + registerProcessHandlers: () => void +} + +export function createRunnerLifecycle(options: RunnerLifecycleOptions): RunnerLifecycle { + let exitCode = 0 + let archiveReason = 'User terminated' + let cleanupStarted = false + let cleanupPromise: Promise | null = null + + const logPrefix = `[${options.logTag}]` + + const archiveAndClose = async () => { + options.session.updateMetadata((currentMetadata) => ({ + ...currentMetadata, + lifecycleState: 'archived', + lifecycleStateSince: Date.now(), + archivedBy: 'cli', + archiveReason + })) + + options.session.sendSessionDeath() + await options.session.flush() + await options.session.close() + } + + const cleanup = async () => { + if (cleanupPromise) { + return cleanupPromise + } + + cleanupStarted = true + cleanupPromise = (async () => { + logger.debug(`${logPrefix} Cleanup start`) + restoreTerminalState() + + try { + options.stopKeepAlive?.() + await options.onBeforeClose?.() + await archiveAndClose() + logger.debug(`${logPrefix} Cleanup complete`) + } finally { + try { + await options.onAfterClose?.() + } catch (error) { + logger.debug(`${logPrefix} Error during post-cleanup:`, error) + } + } + })() + + return cleanupPromise + } + + const cleanupAndExit = async (codeOverride?: number) => { + if (codeOverride !== undefined) { + exitCode = codeOverride + } + + try { + await cleanup() + process.exit(exitCode) + } catch (error) { + logger.debug(`${logPrefix} Error during cleanup:`, error) + process.exit(1) + } + } + + const setExitCode = (code: number) => { + exitCode = code + } + + const setArchiveReason = (reason: string) => { + archiveReason = reason + } + + const markCrash = (error: unknown) => { + logger.debug(`${logPrefix} Unhandled error:`, error) + exitCode = 1 + archiveReason = 'Session crashed' + } + + const registerProcessHandlers = () => { + process.on('SIGTERM', () => { + void cleanupAndExit() + }) + + process.on('SIGINT', () => { + void cleanupAndExit() + }) + + process.on('uncaughtException', (error) => { + markCrash(error) + void cleanupAndExit(1) + }) + + process.on('unhandledRejection', (reason) => { + markCrash(reason) + void cleanupAndExit(1) + }) + } + + return { + setExitCode, + setArchiveReason, + markCrash, + cleanup, + cleanupAndExit, + registerProcessHandlers + } +} + +export function setControlledByUser(session: ApiSessionClient, mode: 'local' | 'remote'): void { + session.updateAgentState((currentState) => ({ + ...currentState, + controlledByUser: mode === 'local' + })) +} + +export function createModeChangeHandler(session: ApiSessionClient): (mode: 'local' | 'remote') => void { + return (mode) => { + session.sendSessionEvent({ type: 'switch', mode }) + setControlledByUser(session, mode) + } +} diff --git a/cli/src/claude/loop.ts b/cli/src/claude/loop.ts index 06648cb0..d7400ace 100644 --- a/cli/src/claude/loop.ts +++ b/cli/src/claude/loop.ts @@ -1,7 +1,7 @@ import { ApiSessionClient } from "@/api/apiSession" import { MessageQueue2 } from "@/utils/MessageQueue2" import { logger } from "@/ui/logger" -import { runLocalRemoteLoop } from "@/agent/loopBase" +import { runLocalRemoteSession } from "@/agent/loopBase" import { Session } from "./session" import { claudeLocalLauncher } from "./claudeLocalLauncher" import { claudeRemoteLauncher } from "./claudeRemoteLauncher" @@ -47,7 +47,7 @@ export async function loop(opts: LoopOptions) { const modelMode: SessionModelMode = opts.model === 'sonnet' || opts.model === 'opus' ? opts.model : 'default'; - let session = new Session({ + const session = new Session({ api: opts.api, client: opts.session, path: opts.path, @@ -67,16 +67,12 @@ export async function loop(opts: LoopOptions) { modelMode }); - // Notify that session is ready - if (opts.onSessionReady) { - opts.onSessionReady(session); - } - - await runLocalRemoteLoop({ + await runLocalRemoteSession({ session, startingMode: opts.startingMode, logTag: 'loop', runLocal: claudeLocalLauncher, - runRemote: claudeRemoteLauncher + runRemote: claudeRemoteLauncher, + onSessionReady: opts.onSessionReady }); } diff --git a/cli/src/claude/runClaude.ts b/cli/src/claude/runClaude.ts index a1d26651..fc1e2dad 100644 --- a/cli/src/claude/runClaude.ts +++ b/cli/src/claude/runClaude.ts @@ -1,5 +1,4 @@ import { logger } from '@/ui/logger'; -import { restoreTerminalState } from '@/ui/terminalState'; import { loop } from '@/claude/loop'; import { AgentState, SessionModelMode } from '@/api/types'; import { EnhancedMode, PermissionMode } from './loop'; @@ -14,6 +13,7 @@ import { generateHookSettingsFile, cleanupHookSettingsFile } from '@/claude/util import { registerKillSessionHandler } from './registerKillSessionHandler'; import type { Session } from './session'; import { bootstrapSession } from '@/agent/sessionFactory'; +import { createModeChangeHandler, createRunnerLifecycle, setControlledByUser } from '@/agent/runnerLifecycle'; export interface StartOptions { model?: string @@ -72,8 +72,6 @@ export async function runClaude(options: StartOptions = {}): Promise { // Variable to track current session instance (updated via onSessionReady callback) const currentSessionRef: { current: Session | null } = { current: null }; - let exitCode = 0; - let archiveReason: string | undefined; const formatFailureReason = (message: string): string => { const maxLength = 200; @@ -108,12 +106,23 @@ export async function runClaude(options: StartOptions = {}): Promise { logger.infoDeveloper(`Session: ${sessionInfo.id}`); logger.infoDeveloper(`Logs: ${logPath}`); + const lifecycle = createRunnerLifecycle({ + session, + logTag: 'claude', + stopKeepAlive: () => currentSessionRef.current?.stopKeepAlive(), + onAfterClose: () => { + happyServer.stop(); + hookServer.stop(); + cleanupHookSettingsFile(hookSettingsPath); + } + }); + + lifecycle.registerProcessHandlers(); + registerKillSessionHandler(session.rpcHandlerManager, lifecycle.cleanupAndExit); + // Set initial agent state const startingMode = options.startingMode ?? (startedBy === 'daemon' ? 'remote' : 'local'); - session.updateAgentState((currentState) => ({ - ...currentState, - controlledByUser: startingMode !== 'remote' - })); + setControlledByUser(session, startingMode); // Import MessageQueue2 and create message queue const messageQueue = new MessageQueue2(mode => hashObject({ @@ -258,65 +267,6 @@ export async function runClaude(options: StartOptions = {}): Promise { logger.debugLargeJson('User message pushed to queue:', message) }); - // Setup signal handlers for graceful shutdown - const cleanup = async () => { - logger.debug('[START] Received termination signal, cleaning up...'); - restoreTerminalState(); - - try { - // Update lifecycle state to archived before closing - if (session) { - const reason = archiveReason ?? 'User terminated'; - session.updateMetadata((currentMetadata) => ({ - ...currentMetadata, - lifecycleState: 'archived', - lifecycleStateSince: Date.now(), - archivedBy: 'cli', - archiveReason: reason - })); - - // Send session death message - session.sendSessionDeath(); - await session.flush(); - await session.close(); - } - - // Stop HAPI MCP server - happyServer.stop(); - - // Stop Hook server and cleanup settings file - hookServer.stop(); - cleanupHookSettingsFile(hookSettingsPath); - - logger.debug('[START] Cleanup complete, exiting'); - process.exit(exitCode); - } catch (error) { - logger.debug('[START] Error during cleanup:', error); - process.exit(1); - } - }; - - // Handle termination signals - process.on('SIGTERM', cleanup); - process.on('SIGINT', cleanup); - - // Handle uncaught exceptions and rejections - process.on('uncaughtException', (error) => { - logger.debug('[START] Uncaught exception:', error); - exitCode = 1; - archiveReason = 'Session crashed'; - cleanup(); - }); - - process.on('unhandledRejection', (reason) => { - logger.debug('[START] Unhandled rejection:', reason); - exitCode = 1; - archiveReason = 'Session crashed'; - cleanup(); - }); - - registerKillSessionHandler(session.rpcHandlerManager, cleanup); - session.rpcHandlerManager.registerHandler('set-session-config', async (payload: unknown) => { if (!payload || typeof payload !== 'object') { throw new Error('Invalid session config payload'); @@ -344,72 +294,50 @@ export async function runClaude(options: StartOptions = {}): Promise { return { applied: { permissionMode: currentPermissionMode, modelMode: currentModelMode } }; }); - // Create claude loop - await loop({ - path: workingDirectory, - model: options.model, - permissionMode: options.permissionMode, - startingMode, - messageQueue, - api, - allowedTools: happyServer.toolNames.map(toolName => `mcp__hapi__${toolName}`), - onModeChange: (newMode) => { - session.sendSessionEvent({ type: 'switch', mode: newMode }); - session.updateAgentState((currentState) => ({ - ...currentState, - controlledByUser: newMode === 'local' - })); - }, - onSessionReady: (sessionInstance) => { - currentSessionRef.current = sessionInstance; - syncSessionModes(); - }, - mcpServers: { - 'hapi': { - type: 'http' as const, - url: happyServer.url, - } - }, - session, - claudeEnvVars: options.claudeEnvVars, - claudeArgs: options.claudeArgs, - startedBy, - hookSettingsPath - }); + let loopError: unknown = null; + let loopFailed = false; + try { + await loop({ + path: workingDirectory, + model: options.model, + permissionMode: options.permissionMode, + startingMode, + messageQueue, + api, + allowedTools: happyServer.toolNames.map(toolName => `mcp__hapi__${toolName}`), + onModeChange: createModeChangeHandler(session), + onSessionReady: (sessionInstance) => { + currentSessionRef.current = sessionInstance; + syncSessionModes(); + }, + mcpServers: { + 'hapi': { + type: 'http' as const, + url: happyServer.url, + } + }, + session, + claudeEnvVars: options.claudeEnvVars, + claudeArgs: options.claudeArgs, + startedBy, + hookSettingsPath + }); + } catch (error) { + loopError = error; + loopFailed = true; + lifecycle.markCrash(error); + } const localFailure = currentSessionRef.current?.localLaunchFailure; if (localFailure?.exitReason === 'exit') { - exitCode = 1; - archiveReason = `Local launch failed: ${formatFailureReason(localFailure.message)}`; - session.updateMetadata((currentMetadata) => ({ - ...currentMetadata, - lifecycleState: 'archived', - lifecycleStateSince: Date.now(), - archivedBy: 'cli', - archiveReason - })); + lifecycle.setExitCode(1); + lifecycle.setArchiveReason(`Local launch failed: ${formatFailureReason(localFailure.message)}`); } - // Send session death message - session.sendSessionDeath(); + if (loopFailed) { + await lifecycle.cleanup(); + throw loopError; + } - // Wait for socket to flush - logger.debug('Waiting for socket to flush...'); - await session.flush(); - - // Close session - logger.debug('Closing session...'); - await session.close(); - - // Stop HAPI MCP server - happyServer.stop(); - logger.debug('Stopped HAPI MCP server'); - - // Stop Hook server and cleanup settings file - hookServer.stop(); - cleanupHookSettingsFile(hookSettingsPath); - logger.debug('Stopped Hook server and cleaned up settings file'); - - // Exit - process.exit(exitCode); + await lifecycle.cleanupAndExit(); } diff --git a/cli/src/codex/loop.ts b/cli/src/codex/loop.ts index cadd7c09..f217898d 100644 --- a/cli/src/codex/loop.ts +++ b/cli/src/codex/loop.ts @@ -1,6 +1,6 @@ import { MessageQueue2 } from '@/utils/MessageQueue2'; import { logger } from '@/ui/logger'; -import { runLocalRemoteLoop } from '@/agent/loopBase'; +import { runLocalRemoteSession } from '@/agent/loopBase'; import { CodexSession } from './session'; import { codexLocalLauncher } from './codexLocalLauncher'; import { codexRemoteLauncher } from './codexRemoteLauncher'; @@ -48,15 +48,12 @@ export async function loop(opts: LoopOptions): Promise { permissionMode: opts.permissionMode ?? 'default' }); - if (opts.onSessionReady) { - opts.onSessionReady(session); - } - - await runLocalRemoteLoop({ + await runLocalRemoteSession({ session, startingMode: opts.startingMode, logTag: 'codex-loop', runLocal: codexLocalLauncher, - runRemote: codexRemoteLauncher + runRemote: codexRemoteLauncher, + onSessionReady: opts.onSessionReady }); } diff --git a/cli/src/codex/runCodex.ts b/cli/src/codex/runCodex.ts index 16d7cb15..47db1c74 100644 --- a/cli/src/codex/runCodex.ts +++ b/cli/src/codex/runCodex.ts @@ -1,5 +1,4 @@ import { logger } from '@/ui/logger'; -import { restoreTerminalState } from '@/ui/terminalState'; import { loop, type EnhancedMode, type PermissionMode } from './loop'; import { MessageQueue2 } from '@/utils/MessageQueue2'; import { hashObject } from '@/utils/deterministicJson'; @@ -8,6 +7,7 @@ import type { AgentState } from '@/api/types'; import type { CodexSession } from './session'; import { parseCodexCliOverrides } from './utils/codexCliOverrides'; import { bootstrapSession } from '@/agent/sessionFactory'; +import { createModeChangeHandler, createRunnerLifecycle, setControlledByUser } from '@/agent/runnerLifecycle'; export { emitReadyIfIdle } from './utils/emitReadyIfIdle'; @@ -33,10 +33,7 @@ export async function runCodex(opts: { const startingMode: 'local' | 'remote' = startedBy === 'daemon' ? 'remote' : 'local'; - session.updateAgentState((currentState) => ({ - ...currentState, - controlledByUser: startingMode === 'local' - })); + setControlledByUser(session, startingMode); const messageQueue = new MessageQueue2((mode) => hashObject({ permissionMode: mode.permissionMode, @@ -48,6 +45,15 @@ export async function runCodex(opts: { let currentPermissionMode: PermissionMode = opts.permissionMode ?? 'default'; + const lifecycle = createRunnerLifecycle({ + session, + logTag: 'codex', + stopKeepAlive: () => sessionWrapperRef.current?.stopKeepAlive() + }); + + lifecycle.registerProcessHandlers(); + registerKillSessionHandler(session.rpcHandlerManager, lifecycle.cleanupAndExit); + const syncSessionMode = () => { const sessionInstance = sessionWrapperRef.current; if (!sessionInstance) { @@ -67,10 +73,6 @@ export async function runCodex(opts: { messageQueue.push(message.content.text, enhancedMode); }); - let cleanupStarted = false; - let exitCode = 0; - let archiveReason = 'User terminated'; - const formatFailureReason = (message: string): string => { const maxLength = 200; if (message.length <= maxLength) { @@ -79,58 +81,6 @@ export async function runCodex(opts: { return `${message.slice(0, maxLength)}...`; }; - const cleanup = async (code: number = exitCode) => { - if (cleanupStarted) { - return; - } - cleanupStarted = true; - logger.debug('[codex] Cleanup start'); - restoreTerminalState(); - try { - const sessionWrapper = sessionWrapperRef.current; - if (sessionWrapper) { - sessionWrapper.stopKeepAlive(); - } - - session.updateMetadata((currentMetadata) => ({ - ...currentMetadata, - lifecycleState: 'archived', - lifecycleStateSince: Date.now(), - archivedBy: 'cli', - archiveReason - })); - - session.sendSessionDeath(); - await session.flush(); - await session.close(); - - logger.debug('[codex] Cleanup complete, exiting'); - process.exit(code); - } catch (error) { - logger.debug('[codex] Error during cleanup:', error); - process.exit(1); - } - }; - - process.on('SIGTERM', () => cleanup(0)); - process.on('SIGINT', () => cleanup(0)); - - process.on('uncaughtException', (error) => { - logger.debug('[codex] Uncaught exception:', error); - exitCode = 1; - archiveReason = 'Session crashed'; - cleanup(1); - }); - - process.on('unhandledRejection', (reason) => { - logger.debug('[codex] Unhandled rejection:', reason); - exitCode = 1; - archiveReason = 'Session crashed'; - cleanup(1); - }); - - registerKillSessionHandler(session.rpcHandlerManager, cleanup); - session.rpcHandlerManager.registerHandler('set-session-config', async (payload: unknown) => { if (!payload || typeof payload !== 'object') { throw new Error('Invalid session config payload'); @@ -149,7 +99,6 @@ export async function runCodex(opts: { return { applied: { permissionMode: currentPermissionMode } }; }); - let loopError: unknown = null; try { await loop({ path: workingDirectory, @@ -161,29 +110,21 @@ export async function runCodex(opts: { codexCliOverrides, startedBy, permissionMode: currentPermissionMode, - onModeChange: (newMode) => { - session.sendSessionEvent({ type: 'switch', mode: newMode }); - session.updateAgentState((currentState) => ({ - ...currentState, - controlledByUser: newMode === 'local' - })); - }, + onModeChange: createModeChangeHandler(session), onSessionReady: (instance) => { sessionWrapperRef.current = instance; syncSessionMode(); } }); } catch (error) { - loopError = error; - exitCode = 1; - archiveReason = 'Session crashed'; + lifecycle.markCrash(error); logger.debug('[codex] Loop error:', error); } finally { const localFailure = sessionWrapperRef.current?.localLaunchFailure; if (localFailure?.exitReason === 'exit') { - exitCode = 1; - archiveReason = `Local launch failed: ${formatFailureReason(localFailure.message)}`; + lifecycle.setExitCode(1); + lifecycle.setArchiveReason(`Local launch failed: ${formatFailureReason(localFailure.message)}`); } - await cleanup(loopError ? 1 : exitCode); + await lifecycle.cleanupAndExit(); } }