{"version":3,"file":"native-supervisor-channel.d.ts","sourceRoot":"","sources":["../../../src/intercom/native-supervisor-channel.ts"],"names":[],"mappings":"AAIA,OAAO,KAAK,EAAE,YAAY,EAAoC,MAAM,2BAA2B,CAAC;AAWhG,OAAO,EAIN,KAAK,aAAa,EAElB,MAAM,oBAAoB,CAAC;AAK5B,eAAO,MAAM,2BAA2B,wBAAwB,CAAC;AAOjE,KAAK,gBAAgB,GAAG,eAAe,GAAG,mBAAmB,GAAG,iBAAiB,CAAC;AAElF,UAAU,iBAAiB;IAC1B,IAAI,EAAE,6BAA6B,CAAC;IACpC,EAAE,EAAE,MAAM,CAAC;IACX,SAAS,EAAE,MAAM,CAAC;IAClB,SAAS,CAAC,EAAE,MAAM,CAAC;IACnB,MAAM,EAAE,gBAAgB,CAAC;IACzB,OAAO,EAAE,MAAM,CAAC;IAChB,YAAY,EAAE,OAAO,CAAC;IACtB,kBAAkB,CAAC,EAAE,MAAM,CAAC;IAC5B,qBAAqB,CAAC,EAAE,MAAM,CAAC;IAC/B,KAAK,EAAE,MAAM,CAAC;IACd,KAAK,EAAE,MAAM,CAAC;IACd,UAAU,EAAE,MAAM,CAAC;IACnB,WAAW,CAAC,EAAE,MAAM,CAAC;IACrB,SAAS,CAAC,EAAE,OAAO,CAAC;CACpB;AAED,UAAU,wBAAyB,SAAQ,iBAAiB;IAC3D,UAAU,EAAE,MAAM,CAAC;IACnB,WAAW,EAAE,MAAM,CAAC;CACpB;AAkDD,wBAAgB,2BAA2B,CAAC,KAAK,EAAE,MAAM,EAAE,KAAK,EAAE,MAAM,EAAE,UAAU,EAAE,MAAM,GAAG,MAAM,CAEpG;AAED,wBAAgB,0BAA0B,CAAC,UAAU,EAAE,MAAM,GAAG,IAAI,CAGnE;AAsND,wBAAgB,8BAA8B,CAC7C,EAAE,EAAE,YAAY,EAChB,OAAO,GAAE;IAAE,uBAAuB,CAAC,EAAE,OAAO,CAAA;CAAO,GACjD,IAAI,CAoDN;AA+WD,wBAAgB,6BAA6B,CAC5C,EAAE,EAAE,YAAY,EAChB,KAAK,EAAE,aAAa,GAClB;IAAE,KAAK,EAAE,MAAM,IAAI,CAAC;IAAC,OAAO,EAAE,MAAM,IAAI,CAAC;IAAC,OAAO,EAAE,GAAG,CAAC,MAAM,EAAE,wBAAwB,CAAC,CAAA;CAAE,CAyF5F","sourcesContent":["import { randomUUID } from \"node:crypto\";\nimport * as fs from \"node:fs\";\nimport * as path from \"node:path\";\nimport type { AgentToolResult } from \"@lpb-work/pi-agent-core\";\nimport type { ExtensionAPI, ExtensionContext, ToolDefinition } from \"@lpb-work/pi-coding-agent\";\nimport { Type } from \"typebox\";\nimport {\n\tSUBAGENT_CHILD_AGENT_ENV,\n\tSUBAGENT_CHILD_INDEX_ENV,\n\tSUBAGENT_ORCHESTRATOR_SESSION_ID_ENV,\n\tSUBAGENT_ORCHESTRATOR_TARGET_ENV,\n\tSUBAGENT_RUN_ID_ENV,\n\tSUBAGENT_SUPERVISOR_CHANNEL_DIR_ENV,\n} from \"../runs/shared/pi-args.ts\";\nimport { writeAtomicJson } from \"../shared/atomic-json.ts\";\nimport {\n\tINTERCOM_DETACH_REQUEST_EVENT,\n\ttype IntercomEventBus,\n\tPOLL_INTERVAL_MS,\n\ttype SubagentState,\n\tTEMP_ROOT_DIR,\n} from \"../shared/types.ts\";\n\nconst SUPERVISOR_CHANNEL_ROOT = path.join(TEMP_ROOT_DIR, \"supervisor-channels\");\nconst REQUESTS_DIR = \"requests\";\nconst REPLIES_DIR = \"replies\";\nexport const NATIVE_SUPERVISOR_TOOL_NAME = \"subagent_supervisor\";\nconst MAX_MESSAGE_BYTES = 64 * 1024;\nconst DEFAULT_ASK_TIMEOUT_MS = 10 * 60 * 1000;\nconst CHANNEL_POLL_MS = Math.min(POLL_INTERVAL_MS, 500);\nconst STALE_EMPTY_CHANNEL_AGE_MS = 60 * 1000;\nconst STALE_EMPTY_CHANNEL_CLEANUP_INTERVAL_MS = 60 * 1000;\n\ntype SupervisorReason = \"need_decision\" | \"interview_request\" | \"progress_update\";\n\ninterface SupervisorRequest {\n\ttype: \"subagent.supervisor.request\";\n\tid: string;\n\tcreatedAt: number;\n\texpiresAt?: number;\n\treason: SupervisorReason;\n\tmessage: string;\n\texpectsReply: boolean;\n\torchestratorTarget?: string;\n\torchestratorSessionId?: string;\n\trunId: string;\n\tagent: string;\n\tchildIndex: number;\n\tchildTarget?: string;\n\tinterview?: unknown;\n}\n\ninterface PendingSupervisorRequest extends SupervisorRequest {\n\tchannelDir: string;\n\trequestFile: string;\n}\n\ninterface SupervisorReply {\n\ttype: \"subagent.supervisor.reply\";\n\trequestId: string;\n\tcreatedAt: number;\n\tmessage: string;\n}\n\ninterface ContactSupervisorParams {\n\treason: SupervisorReason;\n\tmessage?: string;\n\tinterview?: unknown;\n}\n\ninterface IntercomParams {\n\taction: \"list\" | \"send\" | \"ask\" | \"reply\" | \"pending\" | \"status\";\n\tto?: string;\n\tmessage?: string;\n\treplyTo?: string;\n}\n\nconst ContactSupervisorParamsSchema = Type.Object(\n\t{\n\t\treason: Type.String({ enum: [\"need_decision\", \"interview_request\", \"progress_update\"] }),\n\t\tmessage: Type.Optional(Type.String()),\n\t\tinterview: Type.Optional(Type.Unsafe({ type: \"object\", additionalProperties: true })),\n\t},\n\t{ additionalProperties: false },\n);\n\nconst IntercomParamsSchema = Type.Object(\n\t{\n\t\taction: Type.String({ enum: [\"list\", \"send\", \"ask\", \"reply\", \"pending\", \"status\"] }),\n\t\tto: Type.Optional(Type.String()),\n\t\tmessage: Type.Optional(Type.String()),\n\t\treplyTo: Type.Optional(Type.String()),\n\t},\n\t{ additionalProperties: false },\n);\n\nfunction safeSegment(value: string): string {\n\treturn (\n\t\tvalue\n\t\t\t.trim()\n\t\t\t.replace(/[^A-Za-z0-9._-]+/g, \"-\")\n\t\t\t.replace(/^-+|-+$/g, \"\") || \"unknown\"\n\t);\n}\n\nexport function resolveSupervisorChannelDir(runId: string, agent: string, childIndex: number): string {\n\treturn path.join(SUPERVISOR_CHANNEL_ROOT, `${safeSegment(runId)}-${safeSegment(agent)}-${childIndex}`);\n}\n\nexport function ensureSupervisorChannelDir(channelDir: string): void {\n\tfs.mkdirSync(path.join(channelDir, REQUESTS_DIR), { recursive: true, mode: 0o700 });\n\tfs.mkdirSync(path.join(channelDir, REPLIES_DIR), { recursive: true, mode: 0o700 });\n}\n\nfunction requestPath(channelDir: string, requestId: string): string {\n\treturn path.join(channelDir, REQUESTS_DIR, `${safeSegment(requestId)}.json`);\n}\n\nfunction replyPath(channelDir: string, requestId: string): string {\n\treturn path.join(channelDir, REPLIES_DIR, `${safeSegment(requestId)}.json`);\n}\n\nfunction readTextEnv(name: string): string | undefined {\n\tconst value = process.env[name]?.trim();\n\treturn value ? value : undefined;\n}\n\nfunction readChildMetadata():\n\t| {\n\t\t\tchannelDir: string;\n\t\t\trunId: string;\n\t\t\tagent: string;\n\t\t\tchildIndex: number;\n\t\t\torchestratorTarget?: string;\n\t\t\torchestratorSessionId?: string;\n\t\t\tchildTarget?: string;\n\t  }\n\t| undefined {\n\tconst channelDir = readTextEnv(SUBAGENT_SUPERVISOR_CHANNEL_DIR_ENV);\n\tconst runId = readTextEnv(SUBAGENT_RUN_ID_ENV);\n\tconst agent = readTextEnv(SUBAGENT_CHILD_AGENT_ENV);\n\tconst rawIndex = readTextEnv(SUBAGENT_CHILD_INDEX_ENV);\n\tconst orchestratorSessionId = readTextEnv(SUBAGENT_ORCHESTRATOR_SESSION_ID_ENV);\n\tif (!channelDir || !runId || !agent || !orchestratorSessionId || rawIndex === undefined || !/^\\d+$/.test(rawIndex))\n\t\treturn undefined;\n\treturn {\n\t\tchannelDir,\n\t\trunId,\n\t\tagent,\n\t\tchildIndex: Number(rawIndex),\n\t\torchestratorTarget: readTextEnv(SUBAGENT_ORCHESTRATOR_TARGET_ENV),\n\t\torchestratorSessionId,\n\t\tchildTarget: readTextEnv(\"PI_SUBAGENT_INTERCOM_SESSION_NAME\"),\n\t};\n}\n\nfunction reasonHeading(reason: SupervisorReason): string {\n\tif (reason === \"interview_request\") return \"Subagent requests a structured supervisor interview.\";\n\tif (reason === \"progress_update\") return \"Subagent progress update.\";\n\treturn \"Subagent needs a supervisor decision.\";\n}\n\nfunction formatChildMessage(input: {\n\treason: SupervisorReason;\n\tmessage?: string;\n\tinterview?: unknown;\n\trunId: string;\n\tagent: string;\n\tchildIndex: number;\n\tchildTarget?: string;\n}): string {\n\tconst lines = [\n\t\treasonHeading(input.reason),\n\t\t`Run: ${input.runId}`,\n\t\t`Agent: ${input.agent}`,\n\t\t`Child index: ${input.childIndex}`,\n\t];\n\tif (input.childTarget) lines.push(`Child intercom target: ${input.childTarget}`);\n\tlines.push(\"\");\n\tif (input.message?.trim()) lines.push(input.message.trim());\n\tif (input.reason === \"interview_request\") {\n\t\tlines.push(\n\t\t\t\"\",\n\t\t\t\"Structured response requested. Reply with JSON, optionally fenced in ```json, matching the requested interview shape.\",\n\t\t);\n\t\tif (input.interview !== undefined) lines.push(JSON.stringify(input.interview, null, \"\\t\"));\n\t}\n\treturn lines.join(\"\\n\").trimEnd();\n}\n\nfunction parseStructuredReply(message: string): { value?: unknown; error?: string } {\n\tconst trimmed = message.trim();\n\tconst fenced = trimmed.match(/^```(?:json)?\\s*([\\s\\S]*?)\\s*```$/i)?.[1]?.trim();\n\ttry {\n\t\treturn { value: JSON.parse(fenced ?? trimmed) };\n\t} catch (error) {\n\t\treturn { error: error instanceof Error ? `${error.name}: ${error.message}` : String(error) };\n\t}\n}\n\nfunction askTimeoutMs(): number {\n\tconst parsed = Number(process.env.PI_INTERCOM_ASK_TIMEOUT_MS);\n\treturn Number.isFinite(parsed) && parsed > 0 ? parsed : DEFAULT_ASK_TIMEOUT_MS;\n}\n\nfunction delay(ms: number, signal?: AbortSignal): Promise<void> {\n\treturn new Promise((resolve, reject) => {\n\t\tif (signal?.aborted) {\n\t\t\treject(new Error(\"Supervisor request cancelled.\"));\n\t\t\treturn;\n\t\t}\n\t\tlet timer: ReturnType<typeof setTimeout> | undefined;\n\t\tconst cleanup = () => {\n\t\t\tif (timer) clearTimeout(timer);\n\t\t\tsignal?.removeEventListener(\"abort\", onAbort);\n\t\t};\n\t\tconst onAbort = () => {\n\t\t\tcleanup();\n\t\t\treject(new Error(\"Supervisor request cancelled.\"));\n\t\t};\n\t\ttimer = setTimeout(() => {\n\t\t\tcleanup();\n\t\t\tresolve();\n\t\t}, ms);\n\t\tsignal?.addEventListener(\"abort\", onAbort, { once: true });\n\t});\n}\n\nasync function waitForReply(\n\tchannelDir: string,\n\trequestId: string,\n\tdeadline: number,\n\tsignal?: AbortSignal,\n): Promise<SupervisorReply> {\n\tconst file = replyPath(channelDir, requestId);\n\twhile (Date.now() <= deadline) {\n\t\tif (signal?.aborted) throw new Error(\"Supervisor request cancelled.\");\n\t\tif (fs.existsSync(file)) {\n\t\t\tconst parsed = JSON.parse(fs.readFileSync(file, \"utf-8\")) as Partial<SupervisorReply>;\n\t\t\tif (\n\t\t\t\tparsed.type === \"subagent.supervisor.reply\" &&\n\t\t\t\tparsed.requestId === requestId &&\n\t\t\t\ttypeof parsed.message === \"string\"\n\t\t\t) {\n\t\t\t\treturn parsed as SupervisorReply;\n\t\t\t}\n\t\t}\n\t\tawait delay(250, signal);\n\t}\n\tthrow new Error(\"Timed out waiting for supervisor reply.\");\n}\n\nasync function sendSupervisorRequest(\n\tparams: ContactSupervisorParams,\n\tsignal?: AbortSignal,\n): Promise<AgentToolResult<Record<string, unknown>>> {\n\tconst metadata = readChildMetadata();\n\tif (!metadata) throw new Error(\"Native supervisor channel is not available for this subagent.\");\n\tif (params.reason !== \"progress_update\" && !params.message?.trim() && params.reason !== \"interview_request\") {\n\t\tthrow new Error(\"message is required for supervisor decisions.\");\n\t}\n\tensureSupervisorChannelDir(metadata.channelDir);\n\tconst requestId = randomUUID();\n\tconst expectsReply = params.reason !== \"progress_update\";\n\tconst createdAt = Date.now();\n\tconst replyDeadline = createdAt + askTimeoutMs();\n\tconst expiresAt = expectsReply ? replyDeadline : undefined;\n\tconst message = formatChildMessage({\n\t\t...metadata,\n\t\treason: params.reason,\n\t\tmessage: params.message,\n\t\tinterview: params.interview,\n\t});\n\tconst request: SupervisorRequest = {\n\t\ttype: \"subagent.supervisor.request\",\n\t\tid: requestId,\n\t\tcreatedAt,\n\t\t...(expiresAt !== undefined ? { expiresAt } : {}),\n\t\treason: params.reason,\n\t\tmessage,\n\t\texpectsReply,\n\t\t...(metadata.orchestratorTarget ? { orchestratorTarget: metadata.orchestratorTarget } : {}),\n\t\t...(metadata.orchestratorSessionId ? { orchestratorSessionId: metadata.orchestratorSessionId } : {}),\n\t\trunId: metadata.runId,\n\t\tagent: metadata.agent,\n\t\tchildIndex: metadata.childIndex,\n\t\t...(metadata.childTarget ? { childTarget: metadata.childTarget } : {}),\n\t\t...(params.interview !== undefined ? { interview: params.interview } : {}),\n\t};\n\tconst serialized = JSON.stringify(request, null, \"\\t\");\n\tif (Buffer.byteLength(serialized, \"utf-8\") > MAX_MESSAGE_BYTES) throw new Error(\"Supervisor request is too large.\");\n\twriteAtomicJson(requestPath(metadata.channelDir, requestId), request);\n\n\tif (!expectsReply) {\n\t\treturn {\n\t\t\tcontent: [{ type: \"text\", text: \"Supervisor progress update queued.\" }],\n\t\t\tdetails: { delivered: true, requestId, reason: params.reason },\n\t\t};\n\t}\n\n\ttry {\n\t\tconst reply = await waitForReply(metadata.channelDir, requestId, replyDeadline, signal);\n\t\tconst details: Record<string, unknown> = { requestId, reason: params.reason };\n\t\tif (params.reason === \"interview_request\") {\n\t\t\tconst structured = parseStructuredReply(reply.message);\n\t\t\tif (structured.error) details.structuredReplyParseError = structured.error;\n\t\t\telse details.structuredReply = structured.value;\n\t\t}\n\t\treturn {\n\t\t\tcontent: [{ type: \"text\", text: `**Reply from supervisor:**\\n${reply.message}` }],\n\t\t\tdetails,\n\t\t};\n\t} catch (error) {\n\t\tremoveRequestFile(requestPath(metadata.channelDir, requestId));\n\t\tthrow error;\n\t}\n}\n\nfunction hasTool(pi: ExtensionAPI, name: string): boolean {\n\ttry {\n\t\treturn pi.getAllTools?.().some((tool: { name?: unknown }) => tool.name === name) === true;\n\t} catch {\n\t\treturn false;\n\t}\n}\n\nexport function registerNativeSupervisorClient(\n\tpi: ExtensionAPI,\n\toptions: { includeIntercomFallback?: boolean } = {},\n): void {\n\tif (!readChildMetadata()) return;\n\tconst includeIntercomFallback = options.includeIntercomFallback !== false;\n\tif (!hasTool(pi, \"contact_supervisor\")) {\n\t\tconst tool: ToolDefinition<typeof ContactSupervisorParamsSchema, Record<string, unknown>> = {\n\t\t\tname: \"contact_supervisor\",\n\t\t\tlabel: \"Contact Supervisor\",\n\t\t\tdescription:\n\t\t\t\t\"Contact the parent/supervisor session for a blocking decision, structured interview, or progress update.\",\n\t\t\tparameters: ContactSupervisorParamsSchema,\n\t\t\texecute(_id, params, signal) {\n\t\t\t\treturn sendSupervisorRequest(params as ContactSupervisorParams, signal);\n\t\t\t},\n\t\t};\n\t\tpi.registerTool(tool);\n\t}\n\tif (includeIntercomFallback && !hasTool(pi, \"intercom\")) {\n\t\tconst tool: ToolDefinition<typeof IntercomParamsSchema, Record<string, unknown>> = {\n\t\t\tname: \"intercom\",\n\t\t\tlabel: \"Intercom\",\n\t\t\tdescription:\n\t\t\t\t\"Native supervisor-channel intercom fallback for subagents. Prefer contact_supervisor when available.\",\n\t\t\tparameters: IntercomParamsSchema,\n\t\t\tasync execute(_id, params, signal) {\n\t\t\t\tconst action = (params as IntercomParams).action;\n\t\t\t\tif (action === \"status\")\n\t\t\t\t\treturn {\n\t\t\t\t\t\tcontent: [{ type: \"text\", text: \"Native supervisor channel is active.\" }],\n\t\t\t\t\t\tdetails: { active: true },\n\t\t\t\t\t};\n\t\t\t\tif (action === \"list\")\n\t\t\t\t\treturn {\n\t\t\t\t\t\tcontent: [{ type: \"text\", text: \"Supervisor session available through contact_supervisor.\" }],\n\t\t\t\t\t\tdetails: { sessions: [] },\n\t\t\t\t\t};\n\t\t\t\tif (action === \"send\")\n\t\t\t\t\treturn sendSupervisorRequest(\n\t\t\t\t\t\t{ reason: \"progress_update\", message: (params as IntercomParams).message ?? \"\" },\n\t\t\t\t\t\tsignal,\n\t\t\t\t\t);\n\t\t\t\tif (action === \"ask\")\n\t\t\t\t\treturn sendSupervisorRequest(\n\t\t\t\t\t\t{ reason: \"need_decision\", message: (params as IntercomParams).message ?? \"\" },\n\t\t\t\t\t\tsignal,\n\t\t\t\t\t);\n\t\t\t\tthrow new Error(\n\t\t\t\t\t\"Native child intercom supports status, list, send, and ask. Use parent intercom reply from the supervisor session.\",\n\t\t\t\t);\n\t\t\t},\n\t\t};\n\t\tpi.registerTool(tool);\n\t}\n}\n\nfunction parseRequestFile(file: string, channelDir: string): PendingSupervisorRequest | undefined {\n\ttry {\n\t\tconst parsed = JSON.parse(fs.readFileSync(file, \"utf-8\")) as Partial<SupervisorRequest>;\n\t\tif (parsed.type !== \"subagent.supervisor.request\") return undefined;\n\t\tif (typeof parsed.id !== \"string\" || !parsed.id) return undefined;\n\t\tif (\n\t\t\tparsed.reason !== \"need_decision\" &&\n\t\t\tparsed.reason !== \"interview_request\" &&\n\t\t\tparsed.reason !== \"progress_update\"\n\t\t)\n\t\t\treturn undefined;\n\t\tif (typeof parsed.message !== \"string\" || !parsed.message) return undefined;\n\t\tif (typeof parsed.runId !== \"string\" || typeof parsed.agent !== \"string\" || typeof parsed.childIndex !== \"number\")\n\t\t\treturn undefined;\n\t\treturn { ...(parsed as SupervisorRequest), channelDir, requestFile: file };\n\t} catch {\n\t\treturn undefined;\n\t}\n}\n\nfunction listRequestFiles(): Array<{ channelDir: string; file: string }> {\n\tlet channelEntries: fs.Dirent[];\n\ttry {\n\t\tchannelEntries = fs.readdirSync(SUPERVISOR_CHANNEL_ROOT, { withFileTypes: true });\n\t} catch (error) {\n\t\tif ((error as NodeJS.ErrnoException).code === \"ENOENT\") return [];\n\t\tthrow error;\n\t}\n\tconst files: Array<{ channelDir: string; file: string }> = [];\n\tfor (const entry of channelEntries) {\n\t\tif (!entry.isDirectory()) continue;\n\t\tconst channelDir = path.join(SUPERVISOR_CHANNEL_ROOT, entry.name);\n\t\tconst requestsDir = path.join(channelDir, REQUESTS_DIR);\n\t\tlet requestEntries: fs.Dirent[];\n\t\ttry {\n\t\t\trequestEntries = fs.readdirSync(requestsDir, { withFileTypes: true });\n\t\t} catch {\n\t\t\tcontinue;\n\t\t}\n\t\tfor (const requestEntry of requestEntries) {\n\t\t\tif (requestEntry.isFile() && requestEntry.name.endsWith(\".json\"))\n\t\t\t\tfiles.push({ channelDir, file: path.join(requestsDir, requestEntry.name) });\n\t\t}\n\t}\n\treturn files;\n}\n\nfunction readDirectoryEntries(dir: string): fs.Dirent[] | undefined {\n\ttry {\n\t\treturn fs.readdirSync(dir, { withFileTypes: true });\n\t} catch (error) {\n\t\tif ((error as NodeJS.ErrnoException).code === \"ENOENT\") return [];\n\t\treturn undefined;\n\t}\n}\n\nfunction directoryMtimeMs(dir: string): number {\n\ttry {\n\t\treturn fs.statSync(dir).mtimeMs;\n\t} catch {\n\t\treturn 0;\n\t}\n}\n\nfunction removeEmptyDirectory(dir: string): boolean {\n\ttry {\n\t\tfs.rmdirSync(dir);\n\t\treturn true;\n\t} catch (error) {\n\t\tconst code = (error as NodeJS.ErrnoException).code;\n\t\tif (code === \"ENOENT\") return true;\n\t\tif (code === \"ENOTEMPTY\" || code === \"EEXIST\" || code === \"EPERM\" || code === \"EBUSY\") return false;\n\t\tthrow error;\n\t}\n}\n\nfunction removeStaleEmptySupervisorChannel(channelDir: string, nowMs: number): boolean {\n\tconst requestsDir = path.join(channelDir, REQUESTS_DIR);\n\tconst repliesDir = path.join(channelDir, REPLIES_DIR);\n\tconst newestKnownMtimeMs = Math.max(\n\t\tdirectoryMtimeMs(channelDir),\n\t\tdirectoryMtimeMs(requestsDir),\n\t\tdirectoryMtimeMs(repliesDir),\n\t);\n\tif (nowMs - newestKnownMtimeMs < STALE_EMPTY_CHANNEL_AGE_MS) return false;\n\n\tconst requestEntries = readDirectoryEntries(requestsDir);\n\tif (!requestEntries || requestEntries.length > 0) return false;\n\tconst replyEntries = readDirectoryEntries(repliesDir);\n\tif (!replyEntries || replyEntries.length > 0) return false;\n\n\tif (!removeEmptyDirectory(requestsDir)) return false;\n\tif (!removeEmptyDirectory(repliesDir)) return false;\n\tif (!removeEmptyDirectory(channelDir)) return false;\n\treturn true;\n}\n\nfunction cleanupStaleEmptySupervisorChannels(nowMs = Date.now()): number {\n\tlet channelEntries: fs.Dirent[];\n\ttry {\n\t\tchannelEntries = fs.readdirSync(SUPERVISOR_CHANNEL_ROOT, { withFileTypes: true });\n\t} catch (error) {\n\t\tif ((error as NodeJS.ErrnoException).code === \"ENOENT\") return 0;\n\t\tthrow error;\n\t}\n\n\tlet removed = 0;\n\tfor (const entry of channelEntries) {\n\t\tif (!entry.isDirectory()) continue;\n\t\ttry {\n\t\t\tif (removeStaleEmptySupervisorChannel(path.join(SUPERVISOR_CHANNEL_ROOT, entry.name), nowMs)) removed++;\n\t\t} catch {\n\t\t\t// Cleanup is opportunistic; active writers can race with us and will be picked up by a later pass.\n\t\t}\n\t}\n\treturn removed;\n}\n\nfunction currentContextSessionId(\n\tstate: Pick<SubagentState, \"currentSessionId\">,\n\tctx: ExtensionContext,\n): string | undefined {\n\ttry {\n\t\tconst sessionId = ctx.sessionManager.getSessionId();\n\t\tif (sessionId) return sessionId;\n\t} catch {\n\t\t// Fall through to the last known identity.\n\t}\n\treturn state.currentSessionId ?? undefined;\n}\n\nfunction requestMatchesContext(\n\trequest: SupervisorRequest,\n\tstate: Pick<SubagentState, \"currentSessionId\">,\n\tctx: ExtensionContext,\n): boolean {\n\tconst currentSessionId = currentContextSessionId(state, ctx);\n\treturn Boolean(currentSessionId && request.orchestratorSessionId === currentSessionId);\n}\n\nfunction rememberedForegroundChild(request: SupervisorRequest, state: SubagentState) {\n\tconst run = state.foregroundRuns?.get(request.runId);\n\tconst child =\n\t\trun?.children.find((candidate) => candidate.index === request.childIndex && candidate.agent === request.agent) ??\n\t\trun?.children[request.childIndex];\n\treturn run && child ? { run, child } : undefined;\n}\n\nfunction markForegroundSupervisorAttention(request: SupervisorRequest, state: SubagentState): void {\n\tconst remembered = rememberedForegroundChild(request, state);\n\tif (!remembered || remembered.child.status !== \"detached\") return;\n\tconst updatedAt = Date.now();\n\tremembered.run.updatedAt = updatedAt;\n\tremembered.child.activityState = \"needs_attention\";\n\tremembered.child.lastActivityAt = request.createdAt;\n\tremembered.child.currentTool = \"contact_supervisor\";\n\tremembered.child.currentToolStartedAt = request.createdAt;\n\tremembered.child.updatedAt = updatedAt;\n}\n\nfunction clearForegroundSupervisorAttention(\n\trequest: SupervisorRequest,\n\tpending: Map<string, PendingSupervisorRequest>,\n\tstate: SubagentState,\n): void {\n\tif (\n\t\t[...pending.values()].some(\n\t\t\t(candidate) =>\n\t\t\t\tcandidate.expectsReply &&\n\t\t\t\tcandidate.runId === request.runId &&\n\t\t\t\tcandidate.agent === request.agent &&\n\t\t\t\tcandidate.childIndex === request.childIndex,\n\t\t)\n\t)\n\t\treturn;\n\tconst remembered = rememberedForegroundChild(request, state);\n\tif (!remembered || remembered.child.status !== \"detached\" || remembered.child.currentTool !== \"contact_supervisor\")\n\t\treturn;\n\tconst updatedAt = Date.now();\n\tremembered.run.updatedAt = updatedAt;\n\tremembered.child.activityState = undefined;\n\tremembered.child.lastActivityAt = updatedAt;\n\tremembered.child.currentTool = undefined;\n\tremembered.child.currentToolStartedAt = undefined;\n\tremembered.child.updatedAt = updatedAt;\n}\n\nfunction removeRequestFile(file: string): void {\n\ttry {\n\t\tfs.rmSync(file, { force: true });\n\t} catch {\n\t\t// Request cleanup is best-effort; reply files and timeout errors remain authoritative.\n\t}\n}\n\ntype SupervisorRequestLifecycle = \"pending\" | \"resolved\" | \"expired\" | \"inactive\" | \"missing\" | \"wrong-session\";\n\nfunction requestExpiresAt(request: SupervisorRequest, now: number): number {\n\tconst expiresAt = (request as { expiresAt?: unknown }).expiresAt;\n\tif (typeof expiresAt === \"number\" && Number.isFinite(expiresAt)) return expiresAt;\n\treturn Number.isFinite(request.createdAt) ? request.createdAt + askTimeoutMs() : now;\n}\n\nfunction requestRunInactive(request: SupervisorRequest, state: SubagentState): boolean {\n\tif (state.foregroundControls.has(request.runId)) return false;\n\tconst foreground = rememberedForegroundChild(request, state);\n\tif (foreground) return foreground.child.status !== \"detached\";\n\n\tconst asyncJob = state.asyncJobs.get(request.runId);\n\tif (!asyncJob) return false;\n\tif (asyncJob.status === \"complete\" || asyncJob.status === \"failed\" || asyncJob.status === \"paused\") return true;\n\tconst stepStatus = asyncJob.steps?.[request.childIndex]?.status;\n\treturn stepStatus === \"complete\" || stepStatus === \"completed\" || stepStatus === \"failed\" || stepStatus === \"paused\";\n}\n\nfunction requestLifecycle(\n\trequest: PendingSupervisorRequest,\n\tstate: SubagentState,\n\tctx: ExtensionContext | undefined,\n\tnow: number,\n): SupervisorRequestLifecycle {\n\tif (ctx && !requestMatchesContext(request, state, ctx)) return \"wrong-session\";\n\tif (!fs.existsSync(request.requestFile)) return \"missing\";\n\tif (request.expectsReply && fs.existsSync(replyPath(request.channelDir, request.id))) return \"resolved\";\n\tif (request.expectsReply && now > requestExpiresAt(request, now)) return \"expired\";\n\tif (request.expectsReply && requestRunInactive(request, state)) return \"inactive\";\n\treturn \"pending\";\n}\n\nfunction cleanupRequestLifecycle(request: PendingSupervisorRequest, lifecycle: SupervisorRequestLifecycle): void {\n\tif (lifecycle === \"resolved\" || lifecycle === \"expired\" || lifecycle === \"inactive\")\n\t\tremoveRequestFile(request.requestFile);\n}\n\nfunction refreshPendingRequests(\n\tpending: Map<string, PendingSupervisorRequest>,\n\tstate: SubagentState,\n\tctx: ExtensionContext | undefined,\n): void {\n\tconst now = Date.now();\n\tfor (const request of pending.values()) {\n\t\tconst lifecycle = requestLifecycle(request, state, ctx, now);\n\t\tif (lifecycle === \"pending\") continue;\n\t\tpending.delete(request.id);\n\t\tcleanupRequestLifecycle(request, lifecycle);\n\t}\n}\n\nfunction formatPendingLine(request: PendingSupervisorRequest): string {\n\tconst replyHint = request.expectsReply\n\t\t? ` Reply: ${NATIVE_SUPERVISOR_TOOL_NAME}({ action: \"reply\", replyTo: \"${request.id}\", message: \"...\" })`\n\t\t: \"\";\n\treturn `- ${request.id}: ${request.agent} [${request.runId}#${request.childIndex}] ${request.reason}.${replyHint}`;\n}\n\nfunction requestVisibleText(request: PendingSupervisorRequest): string {\n\tconst lines = [request.message];\n\tif (request.expectsReply) {\n\t\tlines.push(\n\t\t\t\"\",\n\t\t\t`Reply with: ${NATIVE_SUPERVISOR_TOOL_NAME}({ action: \"reply\", replyTo: \"${request.id}\", message: \"...\" })`,\n\t\t);\n\t}\n\treturn lines.join(\"\\n\");\n}\n\nfunction writeReply(request: PendingSupervisorRequest, message: string): void {\n\tif (!message.trim()) throw new Error(\"message is required for supervisor replies.\");\n\tconst reply: SupervisorReply = {\n\t\ttype: \"subagent.supervisor.reply\",\n\t\trequestId: request.id,\n\t\tcreatedAt: Date.now(),\n\t\tmessage: message.trim(),\n\t};\n\twriteAtomicJson(replyPath(request.channelDir, request.id), reply);\n\tremoveRequestFile(request.requestFile);\n}\n\nfunction resolvePendingRequest(\n\tpending: Map<string, PendingSupervisorRequest>,\n\tparams: IntercomParams,\n): PendingSupervisorRequest {\n\tif (params.replyTo) {\n\t\tconst request = pending.get(params.replyTo);\n\t\tif (!request) throw new Error(`No pending supervisor request found for replyTo '${params.replyTo}'.`);\n\t\treturn request;\n\t}\n\tconst requests = [...pending.values()].filter((request) => request.expectsReply);\n\tif (params.to) {\n\t\tconst normalizedTo = params.to.toLowerCase();\n\t\tconst matches = requests.filter(\n\t\t\t(request) =>\n\t\t\t\trequest.id.toLowerCase().startsWith(normalizedTo) ||\n\t\t\t\trequest.agent.toLowerCase() === normalizedTo ||\n\t\t\t\trequest.childTarget?.toLowerCase() === normalizedTo,\n\t\t);\n\t\tif (matches.length === 1) return matches[0]!;\n\t\tif (matches.length > 1)\n\t\t\tthrow new Error(`Multiple pending supervisor requests match '${params.to}'. Use replyTo.`);\n\t}\n\tif (requests.length === 1) return requests[0]!;\n\tif (requests.length === 0) throw new Error(\"No pending supervisor requests need a reply.\");\n\tthrow new Error(\"Multiple pending supervisor requests need replies. Use replyTo.\");\n}\n\nfunction publicPendingRequests(pending: Map<string, PendingSupervisorRequest>): Array<Record<string, unknown>> {\n\treturn [...pending.values()].map((request) => ({\n\t\tid: request.id,\n\t\trunId: request.runId,\n\t\tagent: request.agent,\n\t\tchildIndex: request.childIndex,\n\t\treason: request.reason,\n\t\texpectsReply: request.expectsReply,\n\t}));\n}\n\nfunction buildParentIntercomTool(\n\tpending: Map<string, PendingSupervisorRequest>,\n\tstate: SubagentState,\n\tname = \"intercom\",\n): ToolDefinition<typeof IntercomParamsSchema, Record<string, unknown>> {\n\treturn {\n\t\tname,\n\t\tlabel: name === \"intercom\" ? \"Intercom\" : \"Subagent Supervisor\",\n\t\tdescription:\n\t\t\tname === \"intercom\"\n\t\t\t\t? \"Native pi-subagents supervisor channel. Use reply/pending/status to answer child subagent requests.\"\n\t\t\t\t: \"Native pi-subagents supervisor channel. Use reply/pending/status to answer child subagent requests without overriding pi-intercom.\",\n\t\tparameters: IntercomParamsSchema,\n\t\tasync execute(_id, params) {\n\t\t\trefreshPendingRequests(pending, state, state.lastUiContext ?? undefined);\n\t\t\tconst input = params as IntercomParams;\n\t\t\tif (input.action === \"status\") {\n\t\t\t\treturn {\n\t\t\t\t\tcontent: [{ type: \"text\", text: `Native supervisor channel active. Pending replies: ${pending.size}.` }],\n\t\t\t\t\tdetails: { active: true, pending: pending.size, root: SUPERVISOR_CHANNEL_ROOT },\n\t\t\t\t};\n\t\t\t}\n\t\t\tif (input.action === \"pending\" || input.action === \"list\") {\n\t\t\t\tconst lines = [...pending.values()].filter((request) => request.expectsReply).map(formatPendingLine);\n\t\t\t\treturn {\n\t\t\t\t\tcontent: [{ type: \"text\", text: lines.length ? lines.join(\"\\n\") : \"No pending supervisor requests.\" }],\n\t\t\t\t\tdetails: { pending: publicPendingRequests(pending) },\n\t\t\t\t};\n\t\t\t}\n\t\t\tif (input.action === \"reply\") {\n\t\t\t\tconst request = resolvePendingRequest(pending, input);\n\t\t\t\twriteReply(request, input.message ?? \"\");\n\t\t\t\tpending.delete(request.id);\n\t\t\t\tclearForegroundSupervisorAttention(request, pending, state);\n\t\t\t\treturn {\n\t\t\t\t\tcontent: [{ type: \"text\", text: `Replied to supervisor request ${request.id}.` }],\n\t\t\t\t\tdetails: { replyTo: request.id, runId: request.runId, agent: request.agent },\n\t\t\t\t};\n\t\t\t}\n\t\t\tif (input.action === \"send\" || input.action === \"ask\") {\n\t\t\t\tthrow new Error(\n\t\t\t\t\t\"Native pi-subagents intercom currently handles supervisor replies. Child agents initiate asks with contact_supervisor.\",\n\t\t\t\t);\n\t\t\t}\n\t\t\tthrow new Error(`Unsupported intercom action: ${input.action}`);\n\t\t},\n\t};\n}\n\nexport function createNativeSupervisorChannel(\n\tpi: ExtensionAPI,\n\tstate: SubagentState,\n): { start: () => void; dispose: () => void; pending: Map<string, PendingSupervisorRequest> } {\n\tconst pending = new Map<string, PendingSupervisorRequest>();\n\tconst seenFiles = new Set<string>();\n\tlet poller: ReturnType<typeof setInterval> | undefined;\n\tlet lastStaleCleanupAt = 0;\n\n\tconst registerParentTools = (): void => {\n\t\tif (!hasTool(pi, NATIVE_SUPERVISOR_TOOL_NAME))\n\t\t\tpi.registerTool(buildParentIntercomTool(pending, state, NATIVE_SUPERVISOR_TOOL_NAME));\n\t\tif (!hasTool(pi, \"intercom\")) pi.registerTool(buildParentIntercomTool(pending, state));\n\t};\n\n\tconst cleanupStaleChannelsIfDue = (): void => {\n\t\tconst nowMs = Date.now();\n\t\tif (nowMs - lastStaleCleanupAt < STALE_EMPTY_CHANNEL_CLEANUP_INTERVAL_MS) return;\n\t\tlastStaleCleanupAt = nowMs;\n\t\ttry {\n\t\t\tcleanupStaleEmptySupervisorChannels(nowMs);\n\t\t} catch {\n\t\t\t// Supervisor delivery must not fail because best-effort temp cleanup failed.\n\t\t}\n\t};\n\n\tconst poll = (): void => {\n\t\tcleanupStaleChannelsIfDue();\n\t\tconst ctx = state.lastUiContext;\n\t\tif (!ctx) return;\n\t\trefreshPendingRequests(pending, state, ctx);\n\t\tconst now = Date.now();\n\t\tfor (const { channelDir, file } of listRequestFiles()) {\n\t\t\tif (seenFiles.has(file)) continue;\n\t\t\tconst request = parseRequestFile(file, channelDir);\n\t\t\tif (!request || !requestMatchesContext(request, state, ctx)) continue;\n\t\t\tconst lifecycle = requestLifecycle(request, state, undefined, now);\n\t\t\tif (lifecycle !== \"pending\") {\n\t\t\t\tseenFiles.add(file);\n\t\t\t\tcleanupRequestLifecycle(request, lifecycle);\n\t\t\t\tcontinue;\n\t\t\t}\n\t\t\tseenFiles.add(file);\n\t\t\tif (request.expectsReply) {\n\t\t\t\tpending.set(request.id, request);\n\t\t\t\tmarkForegroundSupervisorAttention(request, state);\n\t\t\t} else {\n\t\t\t\tremoveRequestFile(request.requestFile);\n\t\t\t}\n\t\t\tpi.sendMessage(\n\t\t\t\t{\n\t\t\t\t\tcustomType: \"subagent_supervisor_request\",\n\t\t\t\t\tcontent: requestVisibleText(request),\n\t\t\t\t\tdisplay: true,\n\t\t\t\t\tdetails: {\n\t\t\t\t\t\tid: request.id,\n\t\t\t\t\t\treason: request.reason,\n\t\t\t\t\t\texpectsReply: request.expectsReply,\n\t\t\t\t\t\trunId: request.runId,\n\t\t\t\t\t\tagent: request.agent,\n\t\t\t\t\t\tchildIndex: request.childIndex,\n\t\t\t\t\t},\n\t\t\t\t},\n\t\t\t\t{ triggerTurn: true },\n\t\t\t);\n\t\t\tif (request.expectsReply) {\n\t\t\t\t(pi as { events?: IntercomEventBus }).events?.emit(INTERCOM_DETACH_REQUEST_EVENT, {\n\t\t\t\t\trequestId: request.id,\n\t\t\t\t\trunId: request.runId,\n\t\t\t\t\tagent: request.agent,\n\t\t\t\t\tchildIndex: request.childIndex,\n\t\t\t\t});\n\t\t\t}\n\t\t}\n\t};\n\n\treturn {\n\t\tstart: () => {\n\t\t\tif (poller) return;\n\t\t\tregisterParentTools();\n\t\t\tpoll();\n\t\t\tpoller = setInterval(poll, CHANNEL_POLL_MS);\n\t\t\tpoller.unref?.();\n\t\t},\n\t\tdispose: () => {\n\t\t\tif (poller) clearInterval(poller);\n\t\t\tpoller = undefined;\n\t\t\tpending.clear();\n\t\t\tseenFiles.clear();\n\t\t},\n\t\tpending,\n\t};\n}\n"]}