/** * controller — concurrency controller for subagent batches. * * Architecture reference: AgentSwarm pattern. * * Two-phase scheduling: * Normal phase: ramp-up (5 initial, +1 every 700ms). * Rate-limit phase: capacity tracking with exponential backoff retries. * * Environment variables: * PI_SWARM_MAX_CONCURRENCY — cap on concurrent subagents (optional). */ import type { QueuedSubagentTask, SubagentResult, SubagentBatchOptions, SubagentBatchLauncher, SwarmHandle, CoordinatorOptions } from "./types.js"; /** * Marker property that identifies a user-cancellation error. * * #124: replaces the previous substring match on "user"/"cancel"/"interrupt" * which misclassified any error message containing those words (e.g. "user_id * is invalid", "User aborted the request" from an unrelated framework) as * user cancellation, hiding real failures behind status "aborted". * * #132: the marker is a `Symbol.for` value (not a string) so an unrelated * library that happens to attach a string property with the same name * cannot collide with this identity. The global symbol registry gives us * a stable marker across module instantiations. * * Identity is now determined by the marker property on the Error object. * Only errors created via `userCancellationReason()` carry the marker. */ export declare const USER_CANCELLATION_MARKER: unique symbol; /** * Check whether an abort reason indicates user cancellation * (as opposed to a programmatic abort). Uses a marker property on the * Error object so unrelated error messages are not misclassified (#124). */ export declare function isUserCancellation(reason?: unknown): boolean; /** * Create an Error that marks itself as a user cancellation via the * `USER_CANCELLATION_MARKER` property (#124). Use this anywhere the code * path needs to signal "the user explicitly cancelled this work". */ export declare function userCancellationReason(): Error; export declare class SubagentBatchController { private readonly launcher; private readonly states; private readonly pending; private readonly results; private readonly active; private readonly controller; private readonly batchSignals; private readonly maxConcurrency; private readonly maxRateLimitRetries; private batchAborted; private normalLaunchCount; private normalLaunchTimer; private rateLimitLaunchTimer; private rateLimitMode; private rateLimitCapacity; private lastRateLimitAt; private lastCapacityShrinkAt; private lastCapacityRecoveryAt; private globalRetryIntervalMs; private nextRateLimitLaunchAt; private resolve; private reject; private finished; private heartbeatInFlight; private started; private startedAt; private startedSuccessCount; private readonly onProgress?; private eventLog; private nextEventId; private completionTimesMs; private coordRunId?; private coordOnEvent?; private readonly coordAgentControllers; private coordResolve?; private coordReject?; private coordStarted; private lastHeartbeatUpdate; constructor(launcher: SubagentBatchLauncher, tasks: readonly QueuedSubagentTask[], options?: SubagentBatchOptions); /** * Run the batch. Returns a promise that resolves when all tasks * have a terminal result, or rejects on non-user cancellation. */ run(): Promise>>; /** * Run the batch in non-blocking coordinator mode. * Returns a SwarmHandle immediately that allows: * - Getting results so far via getResults() * - Stopping individual agents via stopAgent() * - Aborting the entire swarm via abort() * - Waiting for all agents via completion promise * * Per-agent events are delivered via onEvent callback. */ runAsync(runId: string, coordOpts?: CoordinatorOptions): SwarmHandle; private startInternal; private handleBatchAbort; private schedule; private scheduleNormalLaunch; private isAtConcurrencyLimit; private scheduleRateLimitLaunch; private startAttempt; private runAttempt; private handleAttemptOutcome; private handleAttemptError; private handleRateLimit; private failedAttemptOutcome; private enterRateLimitPhase; private shrinkRateLimitCapacity; private recoverRateLimitCapacity; private finishIfComplete; /** * Update the run manifest's heartbeat timestamp, throttled to 30s intervals. * * Business (#107): Without periodic heartbeat updates, recoverRuns falls * back to startedAt for staleness detection, causing long-running swarms * to be falsely marked abandoned after the 30-minute threshold. */ private updateHeartbeatIfNeeded; private emitProgress; /** Add an event to the event log, keeping it bounded. */ private addEvent; private finish; private finishWithUserCancellation; private fail; private clearNormalTimer; private clearRateLimitTimer; private scheduleRateLimitWakeup; private scheduleNextRateLimitWakeup; private markAttemptReady; private countReadyActive; private nextPendingReadyAt; private nextRateLimitCapacityRecoveryAt; private linkAttemptSignals; private getCompletedResults; private stopAgentById; private sendMessageToAgent; } /** * Resolve the optional swarm max concurrency from pi settings.json * or the environment variable. * * Priority: * 1. `.pi/settings.json` → `pi-swarm.maxConcurrency` (project-local) * 2. `~/.pi/agent/settings.json` → `pi-swarm.maxConcurrency` (global) * 3. `PI_SWARM_MAX_CONCURRENCY` env var * * Falls back to DEFAULT_MAX_CONCURRENCY (5) when unset. A present * value must be a positive integer; invalid input throws so a * misconfigured cap never silently reverts to uncapped. */ export declare function resolveSwarmMaxConcurrency(cwd?: string): number; /** * Resolve the optional small model from pi settings. * Used as the default model for simple subagent tasks. * * Priority: * 1. `.pi/settings.json` → `pi-swarm.smallModel` * 2. `~/.pi/agent/settings.json` → `pi-swarm.smallModel` */ export declare function resolveSwarmSmallModel(cwd?: string): string | undefined; //# sourceMappingURL=controller.d.ts.map