{"version":3,"file":"parallel-utils.d.ts","sourceRoot":"","sources":["../../../../src/runs/shared/parallel-utils.ts"],"names":[],"mappings":"AAAA,MAAM,MAAM,oBAAoB,GAAG,OAAO,uBAAuB,EAAE,iBAAiB,CAAC;AAErF,MAAM,WAAW,kBAAkB;IAClC,oFAAoF;IACpF,eAAe,CAAC,EAAE,MAAM,CAAC;IACzB,4DAA4D;IAC5D,eAAe,CAAC,EAAE,OAAO,kBAAkB,EAAE,eAAe,CAAC;IAC7D,KAAK,EAAE,MAAM,CAAC;IACd,IAAI,EAAE,MAAM,CAAC;IACb,MAAM,CAAC,EAAE,oBAAoB,CAAC;IAC9B,8CAA8C;IAC9C,OAAO,CAAC,EAAE,OAAO,GAAG,MAAM,CAAC;IAC3B,eAAe,CAAC,EAAE;QACjB,KAAK,EAAE,MAAM,CAAC;QACd,QAAQ,EAAE,MAAM,CAAC;QACjB,UAAU,EAAE,MAAM,CAAC;QACnB,KAAK,EAAE,MAAM,CAAC;KACd,CAAC;IACF,KAAK,CAAC,EAAE,MAAM,CAAC;IACf,KAAK,CAAC,EAAE,MAAM,CAAC;IACf,UAAU,CAAC,EAAE,MAAM,CAAC;IACpB,UAAU,CAAC,EAAE,OAAO,CAAC;IACrB,GAAG,CAAC,EAAE,MAAM,CAAC;IACb,KAAK,CAAC,EAAE,MAAM,CAAC;IACf,QAAQ,CAAC,EAAE,MAAM,CAAC;IAClB,eAAe,CAAC,EAAE,MAAM,EAAE,CAAC;IAC3B,KAAK,CAAC,EAAE,MAAM,EAAE,CAAC;IACjB,UAAU,CAAC,EAAE,MAAM,EAAE,CAAC;IACtB,sBAAsB,CAAC,EAAE,MAAM,EAAE,CAAC;IAClC,cAAc,CAAC,EAAE,MAAM,EAAE,CAAC;IAC1B,eAAe,CAAC,EAAE,OAAO,CAAC;IAC1B,YAAY,CAAC,EAAE,MAAM,GAAG,IAAI,CAAC;IAC7B,gBAAgB,CAAC,EAAE,QAAQ,GAAG,SAAS,CAAC;IACxC,qBAAqB,EAAE,OAAO,CAAC;IAC/B,aAAa,EAAE,OAAO,CAAC;IACvB,MAAM,CAAC,EAAE,MAAM,EAAE,CAAC;IAClB,UAAU,CAAC,EAAE,MAAM,CAAC;IACpB,8FAA8F;IAC9F,mBAAmB,CAAC,EAAE,OAAO,CAAC;IAC9B,UAAU,CAAC,EAAE,QAAQ,GAAG,WAAW,CAAC;IACpC,WAAW,CAAC,EAAE,MAAM,CAAC;IACrB,gBAAgB,CAAC,EAAE,MAAM,CAAC;IAC1B,eAAe,CAAC,EAAE,OAAO,CAAC;IAC1B,gBAAgB,CAAC,EAAE;QAClB,MAAM,EAAE,OAAO,uBAAuB,EAAE,gBAAgB,CAAC;QACzD,UAAU,EAAE,MAAM,CAAC;QACnB,UAAU,EAAE,MAAM,CAAC;KACnB,CAAC;IACF,sBAAsB,CAAC,EAAE,OAAO,uBAAuB,EAAE,gBAAgB,CAAC;IAC1E,aAAa,CAAC,EAAE,OAAO,uBAAuB,EAAE,aAAa,CAAC;IAC9D,gBAAgB,CAAC,EAAE,MAAM,CAAC;IAC1B,iBAAiB,CAAC,EAAE,MAAM,CAAC;IAC3B,oBAAoB,CAAC,EAAE,MAAM,CAAC;IAC9B,wBAAwB,CAAC,EAAE,OAAO,uBAAuB,EAAE,+BAA+B,CAAC;IAC3F,6BAA6B,CAAC,EAAE,OAAO,uBAAuB,EAAE,oCAAoC,CAAC;IACrG,mBAAmB,CAAC,EAAE,OAAO,uBAAuB,EAAE,wBAAwB,CAAC;IAC/E,eAAe,CAAC,EAAE,OAAO,uBAAuB,EAAE,eAAe,CAAC;IAClE,cAAc,CAAC,EAAE,OAAO,uBAAuB,EAAE,cAAc,CAAC;IAChE,MAAM,CAAC,EAAE,OAAO,uBAAuB,EAAE,cAAc,CAAC;IACxD,UAAU,CAAC,EAAE,OAAO,uBAAuB,EAAE,kBAAkB,CAAC;IAChE,iBAAiB,CAAC,EAAE,OAAO,yBAAyB,EAAE,iCAAiC,CAAC;IACxF,eAAe,CAAC,EAAE,OAAO,yBAAyB,EAAE,uBAAuB,CAAC;CAC5E;AAED,MAAM,WAAW,oBAAoB;IACpC,UAAU,EAAE,MAAM,CAAC;IACnB,OAAO,CAAC,EAAE,MAAM,CAAC;IACjB,KAAK,CAAC,EAAE,MAAM,CAAC;IACf,KAAK,CAAC,EAAE,MAAM,CAAC;CACf;AAED,MAAM,WAAW,iBAAiB;IACjC,QAAQ,EAAE,kBAAkB,EAAE,CAAC;IAC/B,WAAW,CAAC,EAAE,MAAM,CAAC;IACrB,QAAQ,CAAC,EAAE,OAAO,CAAC;IACnB,QAAQ,CAAC,EAAE,OAAO,CAAC;CACnB;AAED,MAAM,WAAW,kBAAkB;IAClC,MAAM,EAAE,OAAO,0BAA0B,EAAE,iBAAiB,CAAC;IAC7D,QAAQ,EAAE,kBAAkB,CAAC;IAC7B,OAAO,EAAE,OAAO,0BAA0B,EAAE,kBAAkB,CAAC;IAC/D,WAAW,CAAC,EAAE,MAAM,CAAC;IACrB,QAAQ,CAAC,EAAE,OAAO,CAAC;IACnB,KAAK,CAAC,EAAE,MAAM,CAAC;IACf,KAAK,CAAC,EAAE,MAAM,CAAC;IACf,YAAY,CAAC,EAAE,CAAC,MAAM,GAAG,SAAS,CAAC,EAAE,CAAC;IACtC,iBAAiB,CAAC,EAAE,CAAC,MAAM,GAAG,SAAS,CAAC,EAAE,CAAC;IAC3C,mBAAmB,CAAC,EAAE,OAAO,uBAAuB,EAAE,wBAAwB,CAAC;IAC/E,eAAe,CAAC,EAAE,OAAO,uBAAuB,EAAE,eAAe,CAAC;IAClE,cAAc,CAAC,EAAE,OAAO,uBAAuB,EAAE,cAAc,CAAC;IAChE,aAAa,CAAC,EAAE,OAAO,uBAAuB,EAAE,aAAa,CAAC;IAC9D,MAAM,CAAC,EAAE,OAAO,uBAAuB,EAAE,cAAc,CAAC;IACxD,iBAAiB,CAAC,EAAE,OAAO,yBAAyB,EAAE,iCAAiC,CAAC;IACxF,eAAe,CAAC,EAAE,OAAO,yBAAyB,EAAE,uBAAuB,CAAC;CAC5E;AAED,MAAM,MAAM,UAAU,GAAG,kBAAkB,GAAG,iBAAiB,GAAG,kBAAkB,GAAG,oBAAoB,CAAC;AAE5G,wBAAgB,sBAAsB,CAAC,IAAI,EAAE,UAAU,GAAG,IAAI,IAAI,oBAAoB,CAErF;AAED,wBAAgB,eAAe,CAAC,IAAI,EAAE,UAAU,GAAG,IAAI,IAAI,iBAAiB,CAE3E;AAED,wBAAgB,oBAAoB,CAAC,IAAI,EAAE,UAAU,GAAG,IAAI,IAAI,kBAAkB,CAOjF;AAED,wBAAgB,YAAY,CAAC,KAAK,EAAE,UAAU,EAAE,GAAG,kBAAkB,EAAE,CAYtE;AAED,eAAO,MAAM,gCAAgC,KAAK,CAAC;AAEnD;;;;;GAKG;AACH,qBAAa,SAAS;IACrB,OAAO,CAAC,SAAS,CAAS;IAC1B,OAAO,CAAC,QAAQ,CAAC,KAAK,CAAyB;IAE/C,YAAY,KAAK,EAAE,MAAM,EAExB;IAED,OAAO,IAAI,OAAO,CAAC,IAAI,CAAC,CAQvB;IAED,OAAO,IAAI,IAAI,CAOd;CACD;AAED,wBAAsB,aAAa,CAAC,CAAC,EAAE,CAAC,EACvC,KAAK,EAAE,CAAC,EAAE,EACV,KAAK,EAAE,MAAM,EACb,EAAE,EAAE,CAAC,IAAI,EAAE,CAAC,EAAE,CAAC,EAAE,MAAM,KAAK,OAAO,CAAC,CAAC,CAAC,EACtC,eAAe,CAAC,EAAE,SAAS;AAC3B,8FAA8F;AAC9F,mBAAmB,CAAC,EAAE,MAAM,IAAI,GAC9B,OAAO,CAAC,CAAC,EAAE,CAAC,CAuCd;AAED,MAAM,WAAW,kBAAkB;IAClC,KAAK,EAAE,MAAM,CAAC;IACd,SAAS,CAAC,EAAE,MAAM,CAAC;IACnB,MAAM,EAAE,MAAM,CAAC;IACf,QAAQ,EAAE,MAAM,GAAG,IAAI,CAAC;IACxB,KAAK,CAAC,EAAE,MAAM,CAAC;IACf,QAAQ,CAAC,EAAE,OAAO,CAAC;IACnB,KAAK,CAAC,EAAE,MAAM,CAAC;IACf,eAAe,CAAC,EAAE,MAAM,EAAE,CAAC;IAC3B,gBAAgB,CAAC,EAAE,MAAM,CAAC;IAC1B,kBAAkB,CAAC,EAAE,OAAO,CAAC;CAC7B;AAED,wBAAgB,wBAAwB,CACvC,OAAO,EAAE,kBAAkB,EAAE,EAC7B,YAAY,GAAE,CAAC,KAAK,EAAE,MAAM,EAAE,KAAK,EAAE,MAAM,KAAK,MAAkE,GAChH,MAAM,CAsBR;AAED,eAAO,MAAM,wBAAwB,IAAI,CAAC","sourcesContent":["export type ResolvedRunnerConfig = import(\"../../shared/types.ts\").AgentRunnerConfig;\n\nexport interface RunnerSubagentStep {\n\t/** Session id of the direct parent session for permission-system ask forwarding. */\n\tparentSessionId?: string;\n\t/** Resolved opt-in rules for native Pi child tool calls. */\n\tpermissionRules?: import(\"./permissions.ts\").PermissionRules;\n\tagent: string;\n\ttask: string;\n\trunner?: ResolvedRunnerConfig;\n\t/** Resolved launch context for this child. */\n\tcontext?: \"fresh\" | \"fork\";\n\timportAsyncRoot?: {\n\t\trunId: string;\n\t\tasyncDir: string;\n\t\tresultPath: string;\n\t\tindex: number;\n\t};\n\tphase?: string;\n\tlabel?: string;\n\toutputName?: string;\n\tstructured?: boolean;\n\tcwd?: string;\n\tmodel?: string;\n\tthinking?: string;\n\tmodelCandidates?: string[];\n\ttools?: string[];\n\textensions?: string[];\n\tsubagentOnlyExtensions?: string[];\n\tmcpDirectTools?: string[];\n\tcompletionGuard?: boolean;\n\tsystemPrompt?: string | null;\n\tsystemPromptMode?: \"append\" | \"replace\";\n\tinheritProjectContext: boolean;\n\tinheritSkills: boolean;\n\tskills?: string[];\n\toutputPath?: string;\n\t/** Defer the authoritative output instruction until a dynamic fanout item is materialized. */\n\tnamespaceOutputPath?: boolean;\n\toutputMode?: \"inline\" | \"file-only\";\n\tsessionFile?: string;\n\tmaxSubagentDepth?: number;\n\twaitToolEnabled?: boolean;\n\tstructuredOutput?: {\n\t\tschema: import(\"../../shared/types.ts\").JsonSchemaObject;\n\t\tschemaPath: string;\n\t\toutputPath: string;\n\t};\n\tstructuredOutputSchema?: import(\"../../shared/types.ts\").JsonSchemaObject;\n\tagentContract?: import(\"../../shared/types.ts\").AgentContract;\n\tdefinitionDigest?: string;\n\tlaunchBindingTask?: string;\n\tlaunchContractDigest?: string;\n\tlaunchResolvedExtensions?: import(\"../../shared/types.ts\").LaunchResolvedChildExtensionsV1;\n\truntimeAcknowledgedExtensions?: import(\"../../shared/types.ts\").RuntimeAcknowledgedChildExtensionsV1;\n\teffectiveAcceptance?: import(\"../../shared/types.ts\").ResolvedAcceptanceConfig;\n\tacceptanceInput?: import(\"../../shared/types.ts\").AcceptanceInput;\n\tacceptanceRole?: import(\"../../shared/types.ts\").AcceptanceRole;\n\tgateOn?: import(\"../../shared/types.ts\").ChainGateLayer;\n\ttoolBudget?: import(\"../../shared/types.ts\").ResolvedToolBudget;\n\tcapabilityCeiling?: import(\"./capability-ceiling.ts\").ResolvedSubagentCapabilityCeiling;\n\tcapabilityAudit?: import(\"./capability-ceiling.ts\").SubagentCapabilityAudit;\n}\n\nexport interface RunnerCheckpointStep {\n\tcheckpoint: string;\n\tmessage?: string;\n\tphase?: string;\n\tlabel?: string;\n}\n\nexport interface ParallelStepGroup {\n\tparallel: RunnerSubagentStep[];\n\tconcurrency?: number;\n\tfailFast?: boolean;\n\tworktree?: boolean;\n}\n\nexport interface DynamicRunnerGroup {\n\texpand: import(\"../../shared/settings.ts\").DynamicExpandSpec;\n\tparallel: RunnerSubagentStep;\n\tcollect: import(\"../../shared/settings.ts\").DynamicCollectSpec;\n\tconcurrency?: number;\n\tfailFast?: boolean;\n\tphase?: string;\n\tlabel?: string;\n\tsessionFiles?: (string | undefined)[];\n\tthinkingOverrides?: (string | undefined)[];\n\teffectiveAcceptance?: import(\"../../shared/types.ts\").ResolvedAcceptanceConfig;\n\tacceptanceInput?: import(\"../../shared/types.ts\").AcceptanceInput;\n\tacceptanceRole?: import(\"../../shared/types.ts\").AcceptanceRole;\n\tagentContract?: import(\"../../shared/types.ts\").AgentContract;\n\tgateOn?: import(\"../../shared/types.ts\").ChainGateLayer;\n\tcapabilityCeiling?: import(\"./capability-ceiling.ts\").ResolvedSubagentCapabilityCeiling;\n\tcapabilityAudit?: import(\"./capability-ceiling.ts\").SubagentCapabilityAudit;\n}\n\nexport type RunnerStep = RunnerSubagentStep | ParallelStepGroup | DynamicRunnerGroup | RunnerCheckpointStep;\n\nexport function isCheckpointRunnerStep(step: RunnerStep): step is RunnerCheckpointStep {\n\treturn \"checkpoint\" in step;\n}\n\nexport function isParallelGroup(step: RunnerStep): step is ParallelStepGroup {\n\treturn \"parallel\" in step && Array.isArray(step.parallel);\n}\n\nexport function isDynamicRunnerGroup(step: RunnerStep): step is DynamicRunnerGroup {\n\treturn (\n\t\t\"expand\" in step &&\n\t\t\"collect\" in step &&\n\t\t\"parallel\" in step &&\n\t\t!Array.isArray((step as { parallel?: unknown }).parallel)\n\t);\n}\n\nexport function flattenSteps(steps: RunnerStep[]): RunnerSubagentStep[] {\n\tconst flat: RunnerSubagentStep[] = [];\n\tfor (const step of steps) {\n\t\tif (isCheckpointRunnerStep(step)) {\n\t\t} else if (isParallelGroup(step)) {\n\t\t\tfor (const task of step.parallel) flat.push(task);\n\t\t} else if (isDynamicRunnerGroup(step)) {\n\t\t} else {\n\t\t\tflat.push(step);\n\t\t}\n\t}\n\treturn flat;\n}\n\nexport const DEFAULT_GLOBAL_CONCURRENCY_LIMIT = 20;\n\n/**\n * A promise-based semaphore for limiting concurrent access across multiple\n * mapConcurrent calls within a single run. Enforces a global cap on the total\n * number of subagent tasks executing simultaneously, regardless of each step's\n * per-step concurrency limit.\n */\nexport class Semaphore {\n\tprivate available: number;\n\tprivate readonly queue: Array<() => void> = [];\n\n\tconstructor(limit: number) {\n\t\tthis.available = Math.max(1, Math.floor(limit) || 1);\n\t}\n\n\tacquire(): Promise<void> {\n\t\tif (this.available > 0) {\n\t\t\tthis.available--;\n\t\t\treturn Promise.resolve();\n\t\t}\n\t\treturn new Promise<void>((resolve) => {\n\t\t\tthis.queue.push(resolve);\n\t\t});\n\t}\n\n\trelease(): void {\n\t\tconst next = this.queue.shift();\n\t\tif (next) {\n\t\t\tnext();\n\t\t} else {\n\t\t\tthis.available++;\n\t\t}\n\t}\n}\n\nexport async function mapConcurrent<T, R>(\n\titems: T[],\n\tlimit: number,\n\tfn: (item: T, i: number) => Promise<R>,\n\tglobalSemaphore?: Semaphore,\n\t/** Invoked after every worker has stopped, including workers outliving an early rejection. */\n\tonSchedulingSettled?: () => void,\n): Promise<R[]> {\n\tconst safeLimit = Math.max(1, Math.floor(limit) || 1);\n\tconst results: R[] = new Array(items.length);\n\tlet next = 0;\n\tlet liveWorkers = Math.min(safeLimit, items.length);\n\tconst notifySchedulingSettled = () => {\n\t\ttry {\n\t\t\tonSchedulingSettled?.();\n\t\t} catch {\n\t\t\t// Scheduling lifecycle cleanup must not replace the worker result.\n\t\t}\n\t};\n\n\tasync function worker(_workerIndex: number): Promise<void> {\n\t\ttry {\n\t\t\twhile (next < items.length) {\n\t\t\t\tconst i = next++;\n\t\t\t\tif (!(i in items)) throw new Error(`Missing parallel item at index ${i}`);\n\t\t\t\tconst item = items[i] as T;\n\t\t\t\tif (globalSemaphore) {\n\t\t\t\t\tawait globalSemaphore.acquire();\n\t\t\t\t\ttry {\n\t\t\t\t\t\tresults[i] = await fn(item, i);\n\t\t\t\t\t} finally {\n\t\t\t\t\t\tglobalSemaphore.release();\n\t\t\t\t\t}\n\t\t\t\t} else {\n\t\t\t\t\tresults[i] = await fn(item, i);\n\t\t\t\t}\n\t\t\t}\n\t\t} finally {\n\t\t\tliveWorkers--;\n\t\t\tif (liveWorkers === 0) notifySchedulingSettled();\n\t\t}\n\t}\n\n\tif (liveWorkers === 0) notifySchedulingSettled();\n\tawait Promise.all(Array.from({ length: liveWorkers }, (_, wi) => worker(wi)));\n\treturn results;\n}\n\nexport interface ParallelTaskResult {\n\tagent: string;\n\ttaskIndex?: number;\n\toutput: string;\n\texitCode: number | null;\n\terror?: string;\n\ttimedOut?: boolean;\n\tmodel?: string;\n\tattemptedModels?: string[];\n\toutputTargetPath?: string;\n\toutputTargetExists?: boolean;\n}\n\nexport function aggregateParallelOutputs(\n\tresults: ParallelTaskResult[],\n\theaderFormat: (index: number, agent: string) => string = (i, agent) => `=== Parallel Task ${i + 1} (${agent}) ===`,\n): string {\n\treturn results\n\t\t.map((r, i) => {\n\t\t\tconst header = headerFormat(r.taskIndex ?? i, r.agent);\n\t\t\tconst hasOutput = Boolean(r.output?.trim());\n\t\t\tconst status = r.timedOut\n\t\t\t\t? `TIMED OUT${r.error ? `: ${r.error}` : \"\"}`\n\t\t\t\t: r.exitCode === -1\n\t\t\t\t\t? \"SKIPPED\"\n\t\t\t\t\t: r.exitCode !== 0 && r.exitCode !== null\n\t\t\t\t\t\t? `FAILED (exit code ${r.exitCode})${r.error ? `: ${r.error}` : \"\"}`\n\t\t\t\t\t\t: r.error\n\t\t\t\t\t\t\t? `WARNING: ${r.error}`\n\t\t\t\t\t\t\t: !hasOutput && r.outputTargetPath && r.outputTargetExists === false\n\t\t\t\t\t\t\t\t? `EMPTY OUTPUT (expected output file missing: ${r.outputTargetPath})`\n\t\t\t\t\t\t\t\t: !hasOutput && !r.outputTargetPath\n\t\t\t\t\t\t\t\t\t? \"EMPTY OUTPUT (no textual response returned)\"\n\t\t\t\t\t\t\t\t\t: \"\";\n\t\t\tconst body = status ? (hasOutput ? `${status}\\n${r.output}` : status) : r.output;\n\t\t\treturn `${header}\\n${body}`;\n\t\t})\n\t\t.join(\"\\n\\n\");\n}\n\nexport const MAX_PARALLEL_CONCURRENCY = 4;\n"]}