mirror of
https://github.com/NanmiCoder/claude-code-haha.git
synced 2026-10-10 20:03:13 +08:00
refactor(server): move agent/task lifecycle state out of the ws handler
This commit is contained in:
@@ -16,7 +16,16 @@ import { fileURLToPath } from 'node:url'
|
||||
* future reader can evaluate instead of re-deriving.
|
||||
*/
|
||||
|
||||
const HANDLER_PATH = fileURLToPath(new URL('../ws/handler.ts', import.meta.url))
|
||||
/**
|
||||
* Sources that together own per-session state. `handler.ts` is being split, so the
|
||||
* cleanup closure now spans files: `clearAgentRuntimeState` and the six agent/task
|
||||
* containers live in `agentTaskState.ts` while `cleanupSessionRuntimeState` still
|
||||
* calls it from `handler.ts`. Add a file here when a further cut moves state out.
|
||||
*/
|
||||
const SOURCE_PATHS = [
|
||||
fileURLToPath(new URL('../ws/handler.ts', import.meta.url)),
|
||||
fileURLToPath(new URL('../ws/agentTaskState.ts', import.meta.url)),
|
||||
]
|
||||
const CLEANUP_ENTRY = 'cleanupSessionRuntimeState'
|
||||
|
||||
type Classification =
|
||||
@@ -101,19 +110,28 @@ const CONTAINERS: Record<string, Classification> = {
|
||||
},
|
||||
}
|
||||
|
||||
const source = readFileSync(HANDLER_PATH, 'utf8')
|
||||
const sources = SOURCE_PATHS.map((path) => readFileSync(path, 'utf8'))
|
||||
const source = sources.join('\n')
|
||||
|
||||
function declaredContainers(): string[] {
|
||||
return [...source.matchAll(/^const ([a-zA-Z][a-zA-Z0-9]*) = new (?:Map|Set|WeakMap|WeakSet)\b/gm)]
|
||||
return sources
|
||||
.flatMap((text) => [
|
||||
...text.matchAll(/^(?:export )?const ([a-zA-Z][a-zA-Z0-9]*) = new (?:Map|Set|WeakMap|WeakSet)\b/gm),
|
||||
])
|
||||
.map((match) => match[1])
|
||||
.sort()
|
||||
}
|
||||
|
||||
function functionBody(name: string): string | null {
|
||||
const start = source.search(new RegExp(`^(?:export )?(?:async )?function ${name}\\b`, 'm'))
|
||||
if (start === -1) return null
|
||||
const end = source.indexOf('\n}', start)
|
||||
return end === -1 ? source.slice(start) : source.slice(start, end)
|
||||
// Searched per file: concatenating first would let a slice run past the end of one
|
||||
// file into the next, silently widening the closure.
|
||||
for (const text of sources) {
|
||||
const start = text.search(new RegExp(`^(?:export )?(?:async )?function ${name}\\b`, 'm'))
|
||||
if (start === -1) continue
|
||||
const end = text.indexOf('\n}', start)
|
||||
return end === -1 ? text.slice(start) : text.slice(start, end)
|
||||
}
|
||||
return null
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -123,7 +141,7 @@ function functionBody(name: string): string | null {
|
||||
*/
|
||||
function clearedByCleanupClosure(): Set<string> {
|
||||
const entry = functionBody(CLEANUP_ENTRY)
|
||||
if (!entry) throw new Error(`${CLEANUP_ENTRY} not found in handler.ts`)
|
||||
if (!entry) throw new Error(`${CLEANUP_ENTRY} not found in any registered source`)
|
||||
|
||||
const bodies = [entry]
|
||||
for (const call of new Set([...entry.matchAll(/^\s{2}([a-zA-Z][a-zA-Z0-9]*)\(/gm)].map((m) => m[1]))) {
|
||||
@@ -153,6 +171,9 @@ describe('handler session-state cleanup', () => {
|
||||
// A classification left behind after its container is gone is dead weight.
|
||||
expect(classified.filter((name) => !declared.includes(name))).toEqual([])
|
||||
expect(declared.length).toBeGreaterThan(30)
|
||||
// Guards the split itself: the agent/task containers must stay findable after
|
||||
// they moved out of handler.ts.
|
||||
expect(declared).toContain('activeAgentTasks')
|
||||
})
|
||||
|
||||
test('releases every container classified as cleared', () => {
|
||||
|
||||
@@ -0,0 +1,185 @@
|
||||
/**
|
||||
* Agent and background-task lifecycle state for the session WebSocket handler.
|
||||
*
|
||||
* Moved verbatim out of `handler.ts` as the first cut of that file. This slice was
|
||||
* chosen because it is provably closed: these functions read and write only the six
|
||||
* containers declared here and call nothing outside them, so relocating them cannot
|
||||
* change behavior. The functions that emit stop bookends stayed behind on purpose —
|
||||
* they reach into the broadcast layer and the socket registry, so moving them would
|
||||
* have required injecting dependencies rather than moving code.
|
||||
*
|
||||
* All six containers are per-session and released by `clearAgentRuntimeState`, which
|
||||
* `cleanupSessionRuntimeState` calls. `src/server/__tests__/sessionStateCleanup.test.ts`
|
||||
* follows that closure across this module boundary.
|
||||
*
|
||||
* The slice needs no imports: every helper it uses is declared inside it.
|
||||
*/
|
||||
|
||||
export type AgentTaskType = 'local_agent' | 'remote_agent'
|
||||
|
||||
export type ActiveNonAgentTaskState = {
|
||||
taskId: string
|
||||
taskType?: string
|
||||
toolUseId: string
|
||||
description?: string
|
||||
}
|
||||
|
||||
export type ActiveAgentTaskState = {
|
||||
taskId: string
|
||||
taskType: AgentTaskType
|
||||
toolUseId: string
|
||||
remoteSessionId?: string
|
||||
description?: string
|
||||
stopIntent: boolean
|
||||
stopRequested: boolean
|
||||
localStopConfirmed: boolean
|
||||
bookendPending: boolean
|
||||
finalizationRetryCount: number
|
||||
finalizationRetryTimer?: ReturnType<typeof setTimeout>
|
||||
finalization?: Promise<boolean>
|
||||
remoteArchive?: Promise<boolean>
|
||||
remoteArchiveError?: string
|
||||
stopFailureMessage?: string
|
||||
}
|
||||
|
||||
export const activeBackgroundTaskIds = new Map<string, Set<string>>()
|
||||
|
||||
export const activeAgentTasks = new Map<string, Map<string, ActiveAgentTaskState>>()
|
||||
|
||||
export const activeNonAgentTasks = new Map<string, Map<string, ActiveNonAgentTaskState>>()
|
||||
|
||||
export const authoritativeStoppedTaskIds = new Map<string, Set<string>>()
|
||||
|
||||
export const agentStopRequestedSessions = new Set<string>()
|
||||
|
||||
export const runtimeExitStoppedSessions = new Set<string>()
|
||||
|
||||
export type CliBackgroundTaskLifecycle = {
|
||||
taskId: string
|
||||
running: boolean
|
||||
taskType?: string
|
||||
toolUseId?: string
|
||||
remoteSessionId?: string
|
||||
description?: string
|
||||
status?: string
|
||||
suppressForward?: boolean
|
||||
}
|
||||
|
||||
export function getCliBackgroundTaskLifecycle(cliMsg: any): CliBackgroundTaskLifecycle | null {
|
||||
if (cliMsg?.type !== 'system') return null
|
||||
const taskId = typeof cliMsg.task_id === 'string' ? cliMsg.task_id.trim() : ''
|
||||
if (!taskId) return null
|
||||
const optionalString = (value: unknown) =>
|
||||
typeof value === 'string' && value.trim() ? value.trim() : undefined
|
||||
const taskType = typeof cliMsg.task_type === 'string' && cliMsg.task_type.trim()
|
||||
? cliMsg.task_type.trim()
|
||||
: undefined
|
||||
const toolUseId = optionalString(cliMsg.tool_use_id)
|
||||
const remoteSessionId = optionalString(cliMsg.remote_session_id)
|
||||
const description = optionalString(cliMsg.description) ??
|
||||
optionalString(cliMsg.message) ??
|
||||
optionalString(cliMsg.title)
|
||||
|
||||
if (cliMsg.subtype === 'task_started') {
|
||||
return { taskId, running: true, taskType, toolUseId, remoteSessionId, description }
|
||||
}
|
||||
|
||||
if (cliMsg.subtype === 'task_notification' && cliMsg.status === 'running') {
|
||||
return { taskId, running: true, taskType, toolUseId, remoteSessionId, description }
|
||||
}
|
||||
|
||||
if (
|
||||
cliMsg.subtype === 'task_notification' &&
|
||||
(cliMsg.status === 'completed' ||
|
||||
cliMsg.status === 'failed' ||
|
||||
cliMsg.status === 'stopped' ||
|
||||
cliMsg.status === 'killed')
|
||||
) {
|
||||
return { taskId, running: false, status: cliMsg.status }
|
||||
}
|
||||
|
||||
return null
|
||||
}
|
||||
|
||||
export function isAgentTaskType(taskType: string | undefined): taskType is AgentTaskType {
|
||||
return taskType === 'local_agent' || taskType === 'remote_agent'
|
||||
}
|
||||
|
||||
export function untrackCliBackgroundTask(sessionId: string, taskId: string): void {
|
||||
const taskIds = activeBackgroundTaskIds.get(sessionId)
|
||||
taskIds?.delete(taskId)
|
||||
if (taskIds?.size === 0) activeBackgroundTaskIds.delete(sessionId)
|
||||
|
||||
const sessionAgentTasks = activeAgentTasks.get(sessionId)
|
||||
const agentTask = sessionAgentTasks?.get(taskId)
|
||||
if (agentTask?.finalizationRetryTimer !== undefined) {
|
||||
clearTimeout(agentTask.finalizationRetryTimer)
|
||||
}
|
||||
sessionAgentTasks?.delete(taskId)
|
||||
if (sessionAgentTasks?.size === 0) activeAgentTasks.delete(sessionId)
|
||||
|
||||
const sessionNonAgentTasks = activeNonAgentTasks.get(sessionId)
|
||||
sessionNonAgentTasks?.delete(taskId)
|
||||
if (sessionNonAgentTasks?.size === 0) activeNonAgentTasks.delete(sessionId)
|
||||
}
|
||||
|
||||
export function clearAgentRuntimeState(
|
||||
sessionId: string,
|
||||
options?: { preserveRetryableStops?: boolean },
|
||||
): void {
|
||||
const retryableStops = options?.preserveRetryableStops
|
||||
? new Map(
|
||||
[...(activeAgentTasks.get(sessionId)?.entries() ?? [])].filter(([, task]) =>
|
||||
task.stopIntent && task.localStopConfirmed && Boolean(task.stopFailureMessage),
|
||||
),
|
||||
)
|
||||
: new Map<string, ActiveAgentTaskState>()
|
||||
|
||||
for (const task of activeAgentTasks.get(sessionId)?.values() ?? []) {
|
||||
clearAgentStopFinalizationRetry(task)
|
||||
}
|
||||
activeBackgroundTaskIds.delete(sessionId)
|
||||
activeAgentTasks.delete(sessionId)
|
||||
activeNonAgentTasks.delete(sessionId)
|
||||
authoritativeStoppedTaskIds.delete(sessionId)
|
||||
agentStopRequestedSessions.delete(sessionId)
|
||||
runtimeExitStoppedSessions.delete(sessionId)
|
||||
|
||||
if (retryableStops.size > 0) {
|
||||
activeAgentTasks.set(sessionId, retryableStops)
|
||||
activeBackgroundTaskIds.set(sessionId, new Set(retryableStops.keys()))
|
||||
agentStopRequestedSessions.add(sessionId)
|
||||
}
|
||||
}
|
||||
|
||||
export function markTaskAuthoritativelyStopped(sessionId: string, taskId: string): void {
|
||||
let taskIds = authoritativeStoppedTaskIds.get(sessionId)
|
||||
if (!taskIds) {
|
||||
taskIds = new Set()
|
||||
authoritativeStoppedTaskIds.set(sessionId, taskIds)
|
||||
}
|
||||
taskIds.add(taskId)
|
||||
}
|
||||
|
||||
export function hasActiveBackgroundTasks(sessionId: string): boolean {
|
||||
const taskIds = activeBackgroundTaskIds.get(sessionId)
|
||||
if (!taskIds || taskIds.size === 0) return false
|
||||
const sessionAgentTasks = activeAgentTasks.get(sessionId)
|
||||
return [...taskIds].some((taskId) => {
|
||||
const agentTask = sessionAgentTasks?.get(taskId)
|
||||
return !agentTask || !agentTask.localStopConfirmed || agentTask.bookendPending
|
||||
})
|
||||
}
|
||||
|
||||
export function clearAgentStopFinalizationRetry(task: ActiveAgentTaskState): void {
|
||||
if (task.finalizationRetryTimer === undefined) return
|
||||
clearTimeout(task.finalizationRetryTimer)
|
||||
task.finalizationRetryTimer = undefined
|
||||
}
|
||||
|
||||
export function markActiveAgentsStopping(sessionId: string): void {
|
||||
for (const task of activeAgentTasks.get(sessionId)?.values() ?? []) {
|
||||
task.stopIntent = true
|
||||
task.stopRequested = true
|
||||
}
|
||||
}
|
||||
+22
-151
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user