{"version":3,"file":"result-watcher.d.ts","sourceRoot":"","sources":["../../../../src/runs/background/result-watcher.ts"],"names":[],"mappings":"AAAA,OAAO,KAAK,EAAE,MAAM,SAAS,CAAC;AAU9B,OAAO,EACN,KAAK,gBAAgB,EAMrB,KAAK,aAAa,EAClB,MAAM,uBAAuB,CAAC;AAI/B,OAAO,KAAK,EAAE,sBAAsB,EAAE,kBAAkB,EAAE,MAAM,aAAa,CAAC;AAM9E,KAAK,eAAe,GAAG,IAAI,CAC1B,OAAO,EAAE,EACT,YAAY,GAAG,cAAc,GAAG,YAAY,GAAG,aAAa,GAAG,WAAW,GAAG,cAAc,GAAG,OAAO,CACrG,CAAC;AAEF,KAAK,mBAAmB,GAAG;IAC1B,UAAU,EAAE,OAAO,UAAU,CAAC;IAC9B,YAAY,EAAE,OAAO,YAAY,CAAC;IAClC,WAAW,EAAE,OAAO,WAAW,CAAC;IAChC,aAAa,EAAE,OAAO,aAAa,CAAC;CACpC,CAAC;AAEF,KAAK,iBAAiB,GAAG;IACxB,EAAE,CAAC,EAAE,eAAe,CAAC;IACrB,MAAM,CAAC,EAAE,mBAAmB,CAAC;IAC7B,QAAQ,CAAC,EAAE,IAAI,CAAC,kBAAkB,EAAE,SAAS,CAAC,CAAC;IAC/C,+EAA+E;IAC/E,iBAAiB,CAAC,EAAE,CAAC,MAAM,EAAE,sBAAsB,GAAG;QAAE,KAAK,EAAE,MAAM,CAAA;KAAE,KAAK,IAAI,CAAC;IACjF,oGAAoG;IACpG,sBAAsB,CAAC,EAAE,OAAO,CAAC;CACjC,CAAC;AAqEF;;;;GAIG;AACH,wBAAgB,mBAAmB,CAClC,EAAE,EAAE;IAAE,MAAM,EAAE,gBAAgB,CAAA;CAAE,EAChC,KAAK,EAAE,aAAa,EACpB,UAAU,EAAE,MAAM,EAClB,eAAe,EAAE,MAAM,EACvB,IAAI,GAAE,iBAAsB,GAC1B;IACF,kBAAkB,EAAE,MAAM,IAAI,CAAC;IAC/B,oBAAoB,EAAE,CAAC,OAAO,CAAC,EAAE;QAAE,WAAW,CAAC,EAAE,OAAO,CAAA;KAAE,KAAK,IAAI,CAAC;IACpE,iBAAiB,EAAE,MAAM,IAAI,CAAC;CAC9B,CAiVA","sourcesContent":["import * as fs from \"node:fs\";\nimport * as path from \"node:path\";\nimport {\n\tattachNestedChildrenToResultChildren,\n\tbuildSubagentResultIntercomPayload,\n\tcompactNestedResultChildren,\n\tdeliverSubagentResultIntercomEvent,\n\tresolveSubagentResultStatus,\n} from \"../../intercom/result-intercom.ts\";\nimport { createFileCoalescer } from \"../../shared/file-coalescer.ts\";\nimport {\n\ttype IntercomEventBus,\n\ttype NestedRunSummary,\n\ttype ParallelHandoffReference,\n\tSUBAGENT_ASYNC_COMPLETE_EVENT,\n\ttype SubagentOutputState,\n\ttype SubagentResultIntercomChild,\n\ttype SubagentState,\n} from \"../../shared/types.ts\";\nimport { resolveWatchPath } from \"../../shared/utils.ts\";\nimport { projectNestedRegistryForRoot, sanitizeSummary } from \"../shared/nested-events.ts\";\nimport { buildCompletionKey, markSeenWithTtl } from \"./completion-dedupe.ts\";\nimport type { CompletionNotification, CompletionNotifier } from \"./notify.ts\";\n\nconst WATCHER_RESTART_DELAY_MS = 3000;\nconst POLL_INTERVAL_MS = 3000;\nconst RETRY_DELAY_MS = 100;\n\ntype ResultWatcherFs = Pick<\n\ttypeof fs,\n\t\"existsSync\" | \"readFileSync\" | \"unlinkSync\" | \"readdirSync\" | \"mkdirSync\" | \"realpathSync\" | \"watch\"\n>;\n\ntype ResultWatcherTimers = {\n\tsetTimeout: typeof setTimeout;\n\tclearTimeout: typeof clearTimeout;\n\tsetInterval: typeof setInterval;\n\tclearInterval: typeof clearInterval;\n};\n\ntype ResultWatcherDeps = {\n\tfs?: ResultWatcherFs;\n\ttimers?: ResultWatcherTimers;\n\tnotifier?: Pick<CompletionNotifier, \"deliver\">;\n\t/** Receives persisted completions before active-session delivery filtering. */\n\tobserveCompletion?: (result: CompletionNotification & { runId: string }) => void;\n\t/** External grouped-result transport. Disable when native completion notifications own delivery. */\n\tdeliverIntercomResults?: boolean;\n};\n\ntype ResultFileChild = {\n\tagent?: string;\n\toutput?: string;\n\tstructuredOutput?: unknown;\n\toutputState?: SubagentOutputState;\n\terror?: string;\n\tsuccess?: boolean;\n\tstate?: string;\n\tinterrupted?: boolean;\n\ttimedOut?: boolean;\n\tstopped?: boolean;\n\tturnBudgetExceeded?: boolean;\n\tprocessSignal?: string | null;\n\tsessionFile?: string;\n\tartifactPaths?: { outputPath?: string };\n\tintercomTarget?: string;\n\tchildren?: unknown;\n};\n\ntype ResultFileData = CompletionNotification & {\n\trunId?: string;\n\tmode?: string;\n\tresults?: ResultFileChild[];\n\tnestedChildren?: unknown;\n\tasyncDir?: string;\n\tintercomTarget?: string;\n\tparallelHandoff?: ParallelHandoffReference;\n};\n\nfunction sanitizeNestedResultChildren(\n\tvalue: unknown,\n\tresultPath: string,\n\tlabel: string,\n): NestedRunSummary[] | undefined {\n\tif (value === undefined) return undefined;\n\tif (!Array.isArray(value)) {\n\t\tconsole.error(\n\t\t\t`Ignoring invalid nested children in subagent result file '${resultPath}' at ${label}: expected an array.`,\n\t\t);\n\t\treturn undefined;\n\t}\n\tconst children = value\n\t\t.map((child) => sanitizeSummary(child))\n\t\t.filter((child): child is NestedRunSummary => Boolean(child));\n\tif (children.length !== value.length) {\n\t\tconsole.error(\n\t\t\t`Ignoring ${value.length - children.length} invalid nested child record(s) in subagent result file '${resultPath}' at ${label}.`,\n\t\t);\n\t}\n\treturn children.length ? children : undefined;\n}\n\nfunction errorCode(error: unknown): string | undefined {\n\treturn typeof error === \"object\" && error !== null && \"code\" in error\n\t\t? (error as NodeJS.ErrnoException).code\n\t\t: undefined;\n}\n\nfunction isNotFound(error: unknown): boolean {\n\treturn errorCode(error) === \"ENOENT\";\n}\n\nfunction shouldPoll(error: unknown): boolean {\n\tconst code = errorCode(error);\n\treturn code === \"EMFILE\" || code === \"ENOSPC\";\n}\n\n/**\n * Watches persisted async results for the session currently owned by this\n * runtime. `stopResultWatcher()` revokes ownership before closing resources,\n * so old callbacks can never emit or delete after reload/session replacement.\n */\nexport function createResultWatcher(\n\tpi: { events: IntercomEventBus },\n\tstate: SubagentState,\n\tresultsDir: string,\n\tcompletionTtlMs: number,\n\tdeps: ResultWatcherDeps = {},\n): {\n\tstartResultWatcher: () => void;\n\tprimeExistingResults: (options?: { triggerTurn?: boolean }) => void;\n\tstopResultWatcher: () => void;\n} {\n\tconst fsApi = deps.fs ?? fs;\n\tconst timers = deps.timers ?? { setTimeout, clearTimeout, setInterval, clearInterval };\n\tconst notifier = deps.notifier ?? { deliver: async () => true };\n\tconst deliverIntercomResults = deps.deliverIntercomResults !== false;\n\tconst pendingTriggerTurn = new Map<string, boolean>();\n\tconst processing = new Set<string>();\n\tlet deliveryActive = true;\n\tlet deliveryEpoch = 0;\n\t// The sole in-memory ownership lease. It is acquired for one active session\n\t// and revoked before the watcher, queues, or callbacks are torn down.\n\tlet activeSessionId: string | null = null;\n\n\tconst ownsSession = (sessionId: string, epoch: number) => {\n\t\tif (!deliveryActive || epoch !== deliveryEpoch) return false;\n\t\tif (!activeSessionId && state.currentSessionId) activeSessionId = state.currentSessionId;\n\t\treturn activeSessionId === sessionId && state.currentSessionId === sessionId;\n\t};\n\n\tconst scheduleResult = (file: string, triggerTurn: boolean, delayMs = 0) => {\n\t\tconst pendingMode = pendingTriggerTurn.get(file);\n\t\tpendingTriggerTurn.set(file, !(pendingMode === false || !triggerTurn));\n\t\tstate.resultFileCoalescer.schedule(file, delayMs);\n\t};\n\n\tconst handleResult = async (file: string, triggerTurn: boolean) => {\n\t\tconst resultPath = path.join(resultsDir, file);\n\t\tif (processing.has(file) || !fsApi.existsSync(resultPath)) return;\n\t\tprocessing.add(file);\n\t\ttry {\n\t\t\tconst data = JSON.parse(fsApi.readFileSync(resultPath, \"utf-8\")) as ResultFileData;\n\t\t\tif (typeof data.sessionId !== \"string\" || !data.sessionId) return;\n\t\t\tconst runId = data.runId ?? data.id ?? file.replace(/\\.json$/i, \"\");\n\t\t\ttry {\n\t\t\t\tdeps.observeCompletion?.({ ...data, runId });\n\t\t\t} catch (error) {\n\t\t\t\tconsole.error(`Completion observer failed for '${resultPath}':`, error);\n\t\t\t}\n\t\t\tconst epoch = deliveryEpoch;\n\t\t\tif (!ownsSession(data.sessionId, epoch)) return;\n\t\t\tconst hasExplicitNestedChildren = data.nestedChildren !== undefined;\n\t\t\tlet nestedChildren = compactNestedResultChildren(\n\t\t\t\tsanitizeNestedResultChildren(data.nestedChildren, resultPath, \"nestedChildren\"),\n\t\t\t);\n\t\t\tif (!nestedChildren?.length && !hasExplicitNestedChildren) {\n\t\t\t\ttry {\n\t\t\t\t\tnestedChildren = compactNestedResultChildren(projectNestedRegistryForRoot(runId)?.children);\n\t\t\t\t} catch (error) {\n\t\t\t\t\tconsole.error(\n\t\t\t\t\t\t`Failed to enrich subagent result file '${resultPath}' with nested registry children; will retry later:`,\n\t\t\t\t\t\terror,\n\t\t\t\t\t);\n\t\t\t\t\tscheduleResult(file, triggerTurn, RETRY_DELAY_MS);\n\t\t\t\t\treturn;\n\t\t\t\t}\n\t\t\t}\n\n\t\t\tconst completionKey = buildCompletionKey(data, `result:${file}`);\n\t\t\tconst lastSeenAt = state.completionSeen.get(completionKey);\n\t\t\tif (lastSeenAt !== undefined && Date.now() - lastSeenAt > completionTtlMs) {\n\t\t\t\tstate.completionSeen.delete(completionKey);\n\t\t\t} else if (lastSeenAt !== undefined) {\n\t\t\t\tif (!ownsSession(data.sessionId, epoch) || !fsApi.existsSync(resultPath)) return;\n\t\t\t\ttry {\n\t\t\t\t\tfsApi.unlinkSync(resultPath);\n\t\t\t\t} catch (error) {\n\t\t\t\t\tif (!isNotFound(error)) {\n\t\t\t\t\t\tconsole.error(`Failed to remove delivered subagent result '${resultPath}'; will retry:`, error);\n\t\t\t\t\t\tscheduleResult(file, triggerTurn, RETRY_DELAY_MS);\n\t\t\t\t\t}\n\t\t\t\t}\n\t\t\t\treturn;\n\t\t\t}\n\n\t\t\tconst hasResultChildren = Array.isArray(data.results) && data.results.length > 0;\n\t\t\tconst resultChildren: ResultFileChild[] = hasResultChildren\n\t\t\t\t? data.results!\n\t\t\t\t: [{ agent: data.agent ?? undefined, output: data.summary, outputState: \"unknown\", success: data.success }];\n\t\t\tconst normalizedChildren = attachNestedChildrenToResultChildren(\n\t\t\t\trunId,\n\t\t\t\tresultChildren.map((result = {}, index): SubagentResultIntercomChild => {\n\t\t\t\t\tconst baseOutput = hasResultChildren ? result.output : (result.output ?? data.summary);\n\t\t\t\t\tconst hasRealOutput = typeof baseOutput === \"string\" && baseOutput.trim().length > 0;\n\t\t\t\t\tconst structuredPreview =\n\t\t\t\t\t\tresult.structuredOutput === undefined\n\t\t\t\t\t\t\t? undefined\n\t\t\t\t\t\t\t: JSON.stringify(result.structuredOutput, null, 2).slice(0, 4_000);\n\t\t\t\t\tconst output = hasRealOutput\n\t\t\t\t\t\t? baseOutput\n\t\t\t\t\t\t: structuredPreview\n\t\t\t\t\t\t\t? `Structured output:\\n${structuredPreview}`\n\t\t\t\t\t\t\t: \"(no output)\";\n\t\t\t\t\tconst summary =\n\t\t\t\t\t\tresult.success === false && result.error\n\t\t\t\t\t\t\t? `${result.error}${hasRealOutput ? `\\n\\nOutput:\\n${baseOutput}` : \"\"}`\n\t\t\t\t\t\t\t: output;\n\t\t\t\t\tconst sessionPath = result.sessionFile ?? (resultChildren.length === 1 ? data.sessionFile : undefined);\n\t\t\t\t\tconst childNestedChildren = sanitizeNestedResultChildren(\n\t\t\t\t\t\tresult.children,\n\t\t\t\t\t\tresultPath,\n\t\t\t\t\t\t`results[${index}].children`,\n\t\t\t\t\t);\n\t\t\t\t\tconst childState =\n\t\t\t\t\t\tresult.state === \"paused\" || result.state === \"stopped\"\n\t\t\t\t\t\t\t? result.state\n\t\t\t\t\t\t\t: result.stopped === true\n\t\t\t\t\t\t\t\t? \"stopped\"\n\t\t\t\t\t\t\t\t: data.state === \"paused\" ||\n\t\t\t\t\t\t\t\t\t\t(!hasResultChildren && (data.state === \"stopped\" || typeof result.success !== \"boolean\"))\n\t\t\t\t\t\t\t\t\t? data.state\n\t\t\t\t\t\t\t\t\t: undefined;\n\t\t\t\t\treturn {\n\t\t\t\t\t\tagent: result.agent ?? data.agent ?? `step-${index + 1}`,\n\t\t\t\t\t\tstatus: resolveSubagentResultStatus({\n\t\t\t\t\t\t\tsuccess: result.success,\n\t\t\t\t\t\t\tstate: childState,\n\t\t\t\t\t\t\tinterrupted: result.interrupted,\n\t\t\t\t\t\t\ttimedOut: result.timedOut,\n\t\t\t\t\t\t\tstopped: result.stopped,\n\t\t\t\t\t\t\tturnBudgetExceeded: result.turnBudgetExceeded,\n\t\t\t\t\t\t\tprocessSignal: result.processSignal,\n\t\t\t\t\t\t}),\n\t\t\t\t\t\toutputState:\n\t\t\t\t\t\t\tresult.outputState === \"present\" ||\n\t\t\t\t\t\t\tresult.outputState === \"absent\" ||\n\t\t\t\t\t\t\tresult.outputState === \"unknown\"\n\t\t\t\t\t\t\t\t? result.outputState\n\t\t\t\t\t\t\t\t: \"unknown\",\n\t\t\t\t\t\tsummary,\n\t\t\t\t\t\tindex,\n\t\t\t\t\t\tartifactPath: result.artifactPaths?.outputPath,\n\t\t\t\t\t\t...(typeof sessionPath === \"string\" && fsApi.existsSync(sessionPath) ? { sessionPath } : {}),\n\t\t\t\t\t\t...(result.intercomTarget ? { intercomTarget: result.intercomTarget } : {}),\n\t\t\t\t\t\t...(childNestedChildren ? { children: childNestedChildren } : {}),\n\t\t\t\t\t};\n\t\t\t\t}),\n\t\t\t\tnestedChildren,\n\t\t\t);\n\n\t\t\tconst intercomTarget = data.intercomTarget?.trim();\n\t\t\tlet intercomDelivered = false;\n\t\t\tif (deliverIntercomResults && intercomTarget && triggerTurn) {\n\t\t\t\tconst mode =\n\t\t\t\t\tdata.mode === \"single\" || data.mode === \"parallel\" || data.mode === \"chain\" || data.mode === \"workflow\"\n\t\t\t\t\t\t? data.mode\n\t\t\t\t\t\t: resultChildren.length > 1\n\t\t\t\t\t\t\t? \"chain\"\n\t\t\t\t\t\t\t: \"single\";\n\t\t\t\tintercomDelivered = await deliverSubagentResultIntercomEvent(\n\t\t\t\t\tpi.events,\n\t\t\t\t\tbuildSubagentResultIntercomPayload({\n\t\t\t\t\t\tto: intercomTarget,\n\t\t\t\t\t\trunId,\n\t\t\t\t\t\tmode,\n\t\t\t\t\t\tsource: \"async\",\n\t\t\t\t\t\tchildren: normalizedChildren,\n\t\t\t\t\t\tasyncId: data.id ?? undefined,\n\t\t\t\t\t\tasyncDir: data.asyncDir,\n\t\t\t\t\t\t...(data.parallelHandoff ? { parallelHandoff: data.parallelHandoff } : {}),\n\t\t\t\t\t}),\n\t\t\t\t);\n\t\t\t\tif (!ownsSession(data.sessionId, epoch)) return;\n\t\t\t\tif (!intercomDelivered)\n\t\t\t\t\tconsole.error(\n\t\t\t\t\t\t`Subagent async grouped result intercom delivery was not acknowledged for '${resultPath}'.`,\n\t\t\t\t\t);\n\t\t\t}\n\n\t\t\tconst accepted = await notifier.deliver({\n\t\t\t\t...data,\n\t\t\t\tid: data.id ?? runId,\n\t\t\t\trunId,\n\t\t\t\ttriggerTurn,\n\t\t\t\tintercomDelivered,\n\t\t\t\t...(nestedChildren?.length ? { nestedChildren } : {}),\n\t\t\t\t...(Array.isArray(data.results)\n\t\t\t\t\t? {\n\t\t\t\t\t\t\tresults: hasResultChildren\n\t\t\t\t\t\t\t\t? normalizedChildren.map((child, index) => ({\n\t\t\t\t\t\t\t\t\t\t...data.results![index],\n\t\t\t\t\t\t\t\t\t\tagent: child.agent,\n\t\t\t\t\t\t\t\t\t\tstatus: child.status,\n\t\t\t\t\t\t\t\t\t\tsummary: child.summary,\n\t\t\t\t\t\t\t\t\t\tindex: child.index,\n\t\t\t\t\t\t\t\t\t\tartifactPath: child.artifactPath,\n\t\t\t\t\t\t\t\t\t\tsessionPath: child.sessionPath,\n\t\t\t\t\t\t\t\t\t\tchildren: child.children,\n\t\t\t\t\t\t\t\t\t}))\n\t\t\t\t\t\t\t\t: [],\n\t\t\t\t\t\t}\n\t\t\t\t\t: {}),\n\t\t\t});\n\t\t\tif (!ownsSession(data.sessionId, epoch)) return;\n\t\t\tif (!accepted) {\n\t\t\t\tscheduleResult(file, triggerTurn, RETRY_DELAY_MS);\n\t\t\t\treturn;\n\t\t\t}\n\t\t\tmarkSeenWithTtl(state.completionSeen, completionKey, Date.now(), completionTtlMs);\n\t\t\ttry {\n\t\t\t\tpi.events.emit(SUBAGENT_ASYNC_COMPLETE_EVENT, {\n\t\t\t\t\t...data,\n\t\t\t\t\trunId,\n\t\t\t\t\ttriggerTurn,\n\t\t\t\t\tintercomDelivered,\n\t\t\t\t\t...(nestedChildren?.length ? { nestedChildren } : {}),\n\t\t\t\t\t...(Array.isArray(data.results)\n\t\t\t\t\t\t? {\n\t\t\t\t\t\t\t\tresults: hasResultChildren\n\t\t\t\t\t\t\t\t\t? normalizedChildren.map((child, index) => ({\n\t\t\t\t\t\t\t\t\t\t\t...data.results![index],\n\t\t\t\t\t\t\t\t\t\t\tagent: child.agent,\n\t\t\t\t\t\t\t\t\t\t\tstatus: child.status,\n\t\t\t\t\t\t\t\t\t\t\tsummary: child.summary,\n\t\t\t\t\t\t\t\t\t\t\tindex: child.index,\n\t\t\t\t\t\t\t\t\t\t\tartifactPath: child.artifactPath,\n\t\t\t\t\t\t\t\t\t\t\tsessionPath: child.sessionPath,\n\t\t\t\t\t\t\t\t\t\t\tchildren: child.children,\n\t\t\t\t\t\t\t\t\t\t}))\n\t\t\t\t\t\t\t\t\t: [],\n\t\t\t\t\t\t\t}\n\t\t\t\t\t\t: {}),\n\t\t\t\t});\n\t\t\t} catch (error) {\n\t\t\t\tconsole.error(`Completion observer failed for '${resultPath}':`, error);\n\t\t\t}\n\t\t\tif (!ownsSession(data.sessionId, epoch) || !fsApi.existsSync(resultPath)) return;\n\t\t\ttry {\n\t\t\t\tfsApi.unlinkSync(resultPath);\n\t\t\t} catch (error) {\n\t\t\t\tif (!isNotFound(error)) {\n\t\t\t\t\tconsole.error(`Failed to remove delivered subagent result '${resultPath}'; will retry:`, error);\n\t\t\t\t\tscheduleResult(file, triggerTurn, RETRY_DELAY_MS);\n\t\t\t\t}\n\t\t\t}\n\t\t} catch (error) {\n\t\t\tif (!isNotFound(error)) console.error(`Failed to process subagent result file '${resultPath}':`, error);\n\t\t} finally {\n\t\t\tprocessing.delete(file);\n\t\t}\n\t};\n\n\tstate.resultFileCoalescer = createFileCoalescer((file) => {\n\t\tconst triggerTurn = pendingTriggerTurn.get(file) !== false;\n\t\tpendingTriggerTurn.delete(file);\n\t\tvoid handleResult(file, triggerTurn);\n\t}, 50);\n\n\tconst primeExistingResults = (options: { triggerTurn?: boolean } = {}) => {\n\t\ttry {\n\t\t\tconst triggerTurn = options.triggerTurn !== false;\n\t\t\tfsApi\n\t\t\t\t.readdirSync(resultsDir)\n\t\t\t\t.filter((f) => f.endsWith(\".json\"))\n\t\t\t\t.forEach((file) => scheduleResult(file, triggerTurn));\n\t\t} catch (error) {\n\t\t\tif (!isNotFound(error)) console.error(`Failed to scan subagent result directory '${resultsDir}':`, error);\n\t\t}\n\t};\n\n\tconst startPolling = (reason: unknown) => {\n\t\tstate.watcher?.close();\n\t\tstate.watcher = null;\n\t\tif (state.watcherRestartTimer) return;\n\t\tconsole.error(\n\t\t\t`Subagent result watcher for '${resultsDir}' fell back to polling because native fs.watch is unavailable (${errorCode(reason) ?? \"unknown error\"}).`,\n\t\t);\n\t\tprimeExistingResults();\n\t\tstate.watcherRestartTimer = timers.setInterval(primeExistingResults, POLL_INTERVAL_MS);\n\t\tstate.watcherRestartTimer.unref?.();\n\t};\n\n\tconst scheduleRestart = () => {\n\t\tif (state.watcherRestartTimer) return;\n\t\tstate.watcherRestartTimer = timers.setTimeout(() => {\n\t\t\tstate.watcherRestartTimer = null;\n\t\t\ttry {\n\t\t\t\tfsApi.mkdirSync(resultsDir, { recursive: true });\n\t\t\t\tstartResultWatcher();\n\t\t\t} catch (error) {\n\t\t\t\tif (shouldPoll(error)) return startPolling(error);\n\t\t\t\tconsole.error(`Failed to restart subagent result watcher for '${resultsDir}':`, error);\n\t\t\t\tscheduleRestart();\n\t\t\t}\n\t\t}, WATCHER_RESTART_DELAY_MS);\n\t\tstate.watcherRestartTimer.unref?.();\n\t};\n\n\tconst startResultWatcher = () => {\n\t\tif (state.watcher) return;\n\t\tactiveSessionId = state.currentSessionId;\n\t\tdeliveryActive = true;\n\t\tdeliveryEpoch += 1;\n\t\tif (state.watcherRestartTimer) {\n\t\t\ttimers.clearTimeout(state.watcherRestartTimer);\n\t\t\ttimers.clearInterval(state.watcherRestartTimer);\n\t\t\tstate.watcherRestartTimer = null;\n\t\t}\n\t\ttry {\n\t\t\tconst watchDir = resolveWatchPath(resultsDir, fsApi.realpathSync.native);\n\t\t\tstate.watcher = fsApi.watch(watchDir, (event, file) => {\n\t\t\t\tif (event !== \"rename\" || !file) return;\n\t\t\t\tconst fileName = file.toString();\n\t\t\t\tif (fileName.endsWith(\".json\")) scheduleResult(fileName, true);\n\t\t\t});\n\t\t\tstate.watcher.on(\"error\", (error) => {\n\t\t\t\tif (shouldPoll(error)) return startPolling(error);\n\t\t\t\tconsole.error(`Subagent result watcher failed for '${resultsDir}':`, error);\n\t\t\t\tstate.watcher?.close();\n\t\t\t\tstate.watcher = null;\n\t\t\t\tscheduleRestart();\n\t\t\t});\n\t\t\tstate.watcher.unref?.();\n\t\t} catch (error) {\n\t\t\tif (shouldPoll(error)) return startPolling(error);\n\t\t\tconsole.error(`Failed to start subagent result watcher for '${resultsDir}':`, error);\n\t\t\tstate.watcher = null;\n\t\t\tscheduleRestart();\n\t\t}\n\t};\n\n\tconst stopResultWatcher = () => {\n\t\tdeliveryActive = false;\n\t\tactiveSessionId = null;\n\t\tdeliveryEpoch += 1;\n\t\tstate.watcher?.close();\n\t\tstate.watcher = null;\n\t\tif (state.watcherRestartTimer) {\n\t\t\ttimers.clearTimeout(state.watcherRestartTimer);\n\t\t\ttimers.clearInterval(state.watcherRestartTimer);\n\t\t}\n\t\tstate.watcherRestartTimer = null;\n\t\tstate.resultFileCoalescer.clear();\n\t\tpendingTriggerTurn.clear();\n\t\tprocessing.clear();\n\t};\n\n\treturn { startResultWatcher, primeExistingResults, stopResultWatcher };\n}\n"]}