refactor(server): move agent/task lifecycle state out of the ws handler

This commit is contained in:
程序员阿江(Relakkes)
2026-08-04 01:03:36 +08:00
parent 1c1a649658
commit eabdd12c30
3 changed files with 236 additions and 159 deletions
@@ -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', () => {
+185
View File
@@ -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
View File
File diff suppressed because it is too large Load Diff