{"version":3,"file":"async-job-tracker.d.ts","sourceRoot":"","sources":["../../../../src/runs/background/async-job-tracker.ts"],"names":[],"mappings":"AAEA,OAAO,KAAK,EAAE,YAAY,EAAE,gBAAgB,EAAE,MAAM,2BAA2B,CAAC;AAChF,OAAO,EAUN,KAAK,aAAa,EAClB,MAAM,uBAAuB,CAAC;AAS/B,UAAU,sBAAsB;IAC/B,qBAAqB,CAAC,EAAE,MAAM,CAAC;IAC/B,cAAc,CAAC,EAAE,MAAM,CAAC;IACxB,UAAU,CAAC,EAAE,MAAM,CAAC;IACpB,aAAa,CAAC,EAAE,OAAO,CAAC;IACxB,IAAI,CAAC,EAAE,CAAC,GAAG,EAAE,MAAM,EAAE,MAAM,CAAC,EAAE,MAAM,CAAC,OAAO,GAAG,CAAC,KAAK,OAAO,CAAC;IAC7D,GAAG,CAAC,EAAE,MAAM,MAAM,CAAC;CACnB;AAsBD,wBAAgB,qBAAqB,CACpC,EAAE,EAAE,IAAI,CAAC,YAAY,EAAE,QAAQ,CAAC,EAChC,KAAK,EAAE,aAAa,EACpB,YAAY,EAAE,MAAM,EACpB,OAAO,GAAE,sBAA2B,GAClC;IACF,YAAY,EAAE,MAAM,IAAI,CAAC;IACzB,aAAa,EAAE,CAAC,GAAG,EAAE,gBAAgB,KAAK,IAAI,CAAC;IAC/C,aAAa,EAAE,CAAC,IAAI,EAAE,OAAO,KAAK,IAAI,CAAC;IACvC,cAAc,EAAE,CAAC,IAAI,EAAE,OAAO,KAAK,IAAI,CAAC;IACxC,SAAS,EAAE,CAAC,GAAG,CAAC,EAAE,gBAAgB,KAAK,IAAI,CAAC;IAC5C,iBAAiB,EAAE,CAAC,GAAG,CAAC,EAAE,gBAAgB,KAAK,IAAI,CAAC;CACpD,CA8hBA","sourcesContent":["import * as fs from \"node:fs\";\nimport * as path from \"node:path\";\nimport type { ExtensionAPI, ExtensionContext } from \"@lpb-work/pi-coding-agent\";\nimport {\n\ttype AsyncJobState,\n\ttype AsyncStartedEvent,\n\ttype ControlEvent,\n\tDIRS,\n\tPOLL_INTERVAL_MS,\n\ttype SteeringNotice,\n\tSUBAGENT_CONTROL_EVENT,\n\tSUBAGENT_CONTROL_INTERCOM_EVENT,\n\tSUBAGENT_STEERING_NOTICE_EVENT,\n\ttype SubagentState,\n} from \"../../shared/types.ts\";\nimport { readStatus } from \"../../shared/utils.ts\";\nimport { renderWidget, widgetRenderKey } from \"../../tui/render.ts\";\nimport { hasLiveNestedDescendants, updateAsyncJobNestedProjection } from \"../shared/nested-events.ts\";\nimport { formatControlNoticeMessage } from \"../shared/subagent-control.ts\";\nimport { type AsyncRunSummary, listAsyncRuns } from \"./async-status.ts\";\nimport { normalizeParallelGroups } from \"./parallel-groups.ts\";\nimport { reconcileAsyncRun, reconcileNestedAsyncDescendants } from \"./stale-run-reconciler.ts\";\n\ninterface AsyncJobTrackerOptions {\n\tcompletionRetentionMs?: number;\n\tpollIntervalMs?: number;\n\tresultsDir?: string;\n\twidgetEnabled?: boolean;\n\tkill?: (pid: number, signal?: NodeJS.Signals | 0) => boolean;\n\tnow?: () => number;\n}\n\nconst CONTROL_EVENT_READ_CHUNK_BYTES = 64 * 1024;\nconst MAX_CONTROL_EVENT_LINE_BYTES = 1024 * 1024;\nconst CONTROL_EVENT_SCAN_WINDOW_BYTES = 2 * 1024 * 1024;\nconst MAX_RECENT_FLEET_JOBS = 20;\n\nfunction rememberFleetJob(state: SubagentState, job: AsyncJobState): void {\n\tstate.fleetJobs ??= new Map();\n\tstate.fleetJobs.set(job.asyncId, job);\n\tconst terminal = [...state.fleetJobs.values()]\n\t\t.filter(\n\t\t\t(candidate) =>\n\t\t\t\tcandidate.status === \"complete\" ||\n\t\t\t\tcandidate.status === \"failed\" ||\n\t\t\t\tcandidate.status === \"paused\" ||\n\t\t\t\tcandidate.status === \"stopped\",\n\t\t)\n\t\t.sort((left, right) => (right.updatedAt ?? right.startedAt ?? 0) - (left.updatedAt ?? left.startedAt ?? 0));\n\tfor (const stale of terminal.slice(MAX_RECENT_FLEET_JOBS)) state.fleetJobs.delete(stale.asyncId);\n}\n\nexport function createAsyncJobTracker(\n\tpi: Pick<ExtensionAPI, \"events\">,\n\tstate: SubagentState,\n\tasyncDirRoot: string,\n\toptions: AsyncJobTrackerOptions = {},\n): {\n\tensurePoller: () => void;\n\trefreshWidget: (ctx: ExtensionContext) => void;\n\thandleStarted: (data: unknown) => void;\n\thandleComplete: (data: unknown) => void;\n\tresetJobs: (ctx?: ExtensionContext) => void;\n\trestoreActiveJobs: (ctx?: ExtensionContext) => void;\n} {\n\tconst completionRetentionMs = options.completionRetentionMs ?? 10000;\n\tconst pollIntervalMs = options.pollIntervalMs ?? POLL_INTERVAL_MS;\n\tconst resultsDir = options.resultsDir ?? DIRS.results;\n\tconst steeringNoticeSeen = new Map<string, number>();\n\tconst rerenderWidget = (ctx: ExtensionContext, jobs = Array.from(state.asyncJobs.values())) => {\n\t\trenderWidget(ctx, options.widgetEnabled === false ? [] : jobs);\n\t\t(ctx.ui as { requestRender?: () => void }).requestRender?.();\n\t};\n\tconst rerenderLastWidget = (jobs = Array.from(state.asyncJobs.values())) => {\n\t\tconst ctx = state.lastUiContext;\n\t\tif (!ctx) return;\n\t\ttry {\n\t\t\tif (ctx.hasUI) rerenderWidget(ctx, jobs);\n\t\t} catch (error) {\n\t\t\tif (error instanceof Error && error.message.includes(\"extension ctx is stale\")) {\n\t\t\t\tstate.lastUiContext = null;\n\t\t\t\treturn;\n\t\t\t}\n\t\t\tthrow error;\n\t\t}\n\t};\n\tconst refreshWidget = (ctx: ExtensionContext) => rerenderWidget(ctx);\n\tconst restoredControlEventCursor = (asyncDir: string) => {\n\t\ttry {\n\t\t\treturn fs.statSync(path.join(asyncDir, \"events.jsonl\")).size;\n\t\t} catch (error) {\n\t\t\tif ((error as NodeJS.ErrnoException).code === \"ENOENT\") return 0;\n\t\t\tthrow error;\n\t\t}\n\t};\n\tconst summaryToJob = (run: AsyncRunSummary): AsyncJobState => {\n\t\tconst groups = normalizeParallelGroups(\n\t\t\trun.parallelGroups,\n\t\t\trun.steps.length,\n\t\t\trun.chainStepCount ?? run.steps.length,\n\t\t);\n\t\tconst activeGroup =\n\t\t\trun.currentStep !== undefined\n\t\t\t\t? groups.find((group) => run.currentStep! >= group.start && run.currentStep! < group.start + group.count)\n\t\t\t\t: undefined;\n\t\tconst visibleSteps = activeGroup\n\t\t\t? run.steps\n\t\t\t\t\t.slice(activeGroup.start, activeGroup.start + activeGroup.count)\n\t\t\t\t\t.map((step, index) => ({ ...step, index: activeGroup.start + index }))\n\t\t\t: run.steps.map((step, index) => ({ ...step, index }));\n\t\treturn {\n\t\t\tasyncId: run.id,\n\t\t\tasyncDir: run.asyncDir,\n\t\t\tstatus: run.state,\n\t\t\tsessionId: run.sessionId,\n\t\t\tactivityState: run.activityState,\n\t\t\tlastActivityAt: run.lastActivityAt,\n\t\t\tcurrentTool: run.currentTool,\n\t\t\tcurrentToolStartedAt: run.currentToolStartedAt,\n\t\t\tcurrentPath: run.currentPath,\n\t\t\tturnCount: run.turnCount,\n\t\t\ttoolCount: run.toolCount,\n\t\t\tsteering: run.steering,\n\t\t\tmode: run.mode,\n\t\t\tcontext: run.context,\n\t\t\tcwd: run.cwd,\n\t\t\tagents: visibleSteps.map((step) => step.agent),\n\t\t\tcurrentStep: run.currentStep,\n\t\t\tchainStepCount: run.chainStepCount,\n\t\t\tparallelGroups: groups,\n\t\t\tsteps: visibleSteps,\n\t\t\tstepsTotal: visibleSteps.length,\n\t\t\trunningSteps: visibleSteps.filter((step) => step.status === \"running\").length,\n\t\t\tcompletedSteps: visibleSteps.filter((step) => step.status === \"complete\" || step.status === \"completed\")\n\t\t\t\t.length,\n\t\t\thasParallelGroups: groups.length > 0,\n\t\t\tactiveParallelGroup: Boolean(activeGroup),\n\t\t\tstartedAt: run.startedAt,\n\t\t\tupdatedAt: run.lastUpdate ?? run.startedAt,\n\t\t\ttimeoutMs: run.timeoutMs,\n\t\t\tdeadlineAt: run.deadlineAt,\n\t\t\ttimedOut: run.timedOut,\n\t\t\tstopped: run.stopped,\n\t\t\tturnBudget: run.turnBudget,\n\t\t\tturnBudgetExceeded: run.turnBudgetExceeded,\n\t\t\twrapUpRequested: run.wrapUpRequested,\n\t\t\tsessionDir: run.sessionDir,\n\t\t\toutputFile: run.outputFile,\n\t\t\ttotalTokens: run.totalTokens,\n\t\t\tsessionFile: run.sessionFile,\n\t\t\tcontrolEventCursor: restoredControlEventCursor(run.asyncDir),\n\t\t\tnestedChildren: run.nestedChildren,\n\t\t\tparentWorkflowRunId: run.parentWorkflowRunId,\n\t\t\tworkflowKey: run.workflowKey,\n\t\t\tworkflow: run.workflow,\n\t\t};\n\t};\n\tconst cancelCleanup = (asyncId: string) => {\n\t\tconst existingTimer = state.cleanupTimers.get(asyncId);\n\t\tif (!existingTimer) return;\n\t\tclearTimeout(existingTimer);\n\t\tstate.cleanupTimers.delete(asyncId);\n\t};\n\tconst scheduleCleanup = (asyncId: string) => {\n\t\tcancelCleanup(asyncId);\n\t\tconst timer = setTimeout(() => {\n\t\t\tstate.cleanupTimers.delete(asyncId);\n\t\t\tstate.asyncJobs.delete(asyncId);\n\t\t\trerenderLastWidget();\n\t\t}, completionRetentionMs);\n\t\tstate.cleanupTimers.set(asyncId, timer);\n\t};\n\tconst emitNewControlEvents = (job: AsyncJobState) => {\n\t\tconst eventsPath = path.join(job.asyncDir, \"events.jsonl\");\n\t\tlet fd: number;\n\t\ttry {\n\t\t\tfd = fs.openSync(eventsPath, \"r\");\n\t\t} catch (error) {\n\t\t\tif ((error as NodeJS.ErrnoException).code === \"ENOENT\") return;\n\t\t\tconsole.error(`Failed to open async control events for '${job.asyncDir}':`, error);\n\t\t\treturn;\n\t\t}\n\t\ttry {\n\t\t\tconst stat = fs.fstatSync(fd);\n\t\t\tconst savedCursor = job.controlEventCursor;\n\t\t\tlet cursor = stat.size < (savedCursor ?? 0) ? 0 : (savedCursor ?? 0);\n\t\t\tconst startedFromTail = savedCursor === undefined && stat.size > CONTROL_EVENT_SCAN_WINDOW_BYTES;\n\t\t\tif (startedFromTail) cursor = stat.size - CONTROL_EVENT_SCAN_WINDOW_BYTES;\n\t\t\tif (stat.size <= cursor) return;\n\t\t\tconst previousByte = Buffer.alloc(1);\n\t\t\tconst startsMidLine =\n\t\t\t\tcursor > 0 && fs.readSync(fd, previousByte, 0, 1, cursor - 1) === 1 && previousByte[0] !== 0x0a;\n\t\t\tconst scanEnd = Math.min(stat.size, cursor + CONTROL_EVENT_SCAN_WINDOW_BYTES);\n\t\t\tconst handleLine = (line: string) => {\n\t\t\t\tif (!line.trim()) return;\n\t\t\t\tlet parsed: unknown;\n\t\t\t\ttry {\n\t\t\t\t\tparsed = JSON.parse(line);\n\t\t\t\t} catch (error) {\n\t\t\t\t\tconsole.error(`Ignoring malformed async control event in '${eventsPath}':`, error);\n\t\t\t\t\treturn;\n\t\t\t\t}\n\t\t\t\tif (!parsed || typeof parsed !== \"object\") return;\n\t\t\t\tif ((parsed as { type?: unknown }).type === \"subagent.steering.notice\") {\n\t\t\t\t\tconst notice = parsed as Partial<SteeringNotice>;\n\t\t\t\t\tif (\n\t\t\t\t\t\ttypeof notice.requestId !== \"string\" ||\n\t\t\t\t\t\ttypeof notice.runId !== \"string\" ||\n\t\t\t\t\t\t(notice.state !== \"failed\" && notice.state !== \"partial\" && notice.state !== \"recovered\") ||\n\t\t\t\t\t\ttypeof notice.message !== \"string\"\n\t\t\t\t\t)\n\t\t\t\t\t\treturn;\n\t\t\t\t\tif (typeof state.currentSessionId === \"string\" && notice.currentSessionId !== state.currentSessionId)\n\t\t\t\t\t\treturn;\n\t\t\t\t\tconst key = `${notice.runId}:${notice.requestId}:${notice.state}`;\n\t\t\t\t\tif (steeringNoticeSeen.has(key)) return;\n\t\t\t\t\tconst now = Date.now();\n\t\t\t\t\tsteeringNoticeSeen.set(key, now);\n\t\t\t\t\tif (steeringNoticeSeen.size > 200) {\n\t\t\t\t\t\tfor (const [seenKey, seenAt] of steeringNoticeSeen) {\n\t\t\t\t\t\t\tif (now - seenAt > 10 * 60 * 1000 || steeringNoticeSeen.size > 200)\n\t\t\t\t\t\t\t\tsteeringNoticeSeen.delete(seenKey);\n\t\t\t\t\t\t}\n\t\t\t\t\t}\n\t\t\t\t\tpi.events.emit(SUBAGENT_STEERING_NOTICE_EVENT, {\n\t\t\t\t\t\t...notice,\n\t\t\t\t\t\tsource: \"async\",\n\t\t\t\t\t\tasyncDir: job.asyncDir,\n\t\t\t\t\t\tnoticeText: notice.message,\n\t\t\t\t\t});\n\t\t\t\t\treturn;\n\t\t\t\t}\n\t\t\t\tif ((parsed as { type?: unknown }).type !== \"subagent.control\") return;\n\t\t\t\tconst record = parsed as {\n\t\t\t\t\tevent?: ControlEvent;\n\t\t\t\t\tchannels?: string[];\n\t\t\t\t\tchildIntercomTarget?: string;\n\t\t\t\t\tnoticeText?: string;\n\t\t\t\t\tintercom?: { to?: string; message?: string };\n\t\t\t\t};\n\t\t\t\tif (!record.event || !Array.isArray(record.channels)) return;\n\t\t\t\tconst payload = {\n\t\t\t\t\tevent: record.event,\n\t\t\t\t\tsource: \"async\" as const,\n\t\t\t\t\tasyncDir: job.asyncDir,\n\t\t\t\t\tchildIntercomTarget: record.childIntercomTarget,\n\t\t\t\t\tnoticeText: record.noticeText ?? formatControlNoticeMessage(record.event, record.childIntercomTarget),\n\t\t\t\t};\n\t\t\t\tif (record.channels.includes(\"event\")) {\n\t\t\t\t\tpi.events.emit(SUBAGENT_CONTROL_EVENT, payload);\n\t\t\t\t}\n\t\t\t\tif (\n\t\t\t\t\trecord.event.type !== \"active_long_running\" &&\n\t\t\t\t\trecord.channels.includes(\"intercom\") &&\n\t\t\t\t\trecord.intercom?.to &&\n\t\t\t\t\trecord.intercom.message\n\t\t\t\t) {\n\t\t\t\t\tpi.events.emit(SUBAGENT_CONTROL_INTERCOM_EVENT, {\n\t\t\t\t\t\t...payload,\n\t\t\t\t\t\tto: record.intercom.to,\n\t\t\t\t\t\tmessage: record.intercom.message,\n\t\t\t\t\t});\n\t\t\t\t}\n\t\t\t};\n\t\t\tlet readCursor = cursor;\n\t\t\tlet lastCompleteCursor = cursor;\n\t\t\tlet lineParts: Buffer[] = [];\n\t\t\tlet lineBytes = 0;\n\t\t\tlet skippingOversizedLine = startedFromTail || startsMidLine;\n\t\t\tconst appendLineSegment = (segment: Buffer) => {\n\t\t\t\tif (segment.length === 0 || skippingOversizedLine) return;\n\t\t\t\tif (lineBytes + segment.length > MAX_CONTROL_EVENT_LINE_BYTES) {\n\t\t\t\t\tlineParts = [];\n\t\t\t\t\tlineBytes = 0;\n\t\t\t\t\tskippingOversizedLine = true;\n\t\t\t\t\treturn;\n\t\t\t\t}\n\t\t\t\tlineParts.push(segment);\n\t\t\t\tlineBytes += segment.length;\n\t\t\t};\n\t\t\twhile (readCursor < scanEnd) {\n\t\t\t\tconst toRead = Math.min(CONTROL_EVENT_READ_CHUNK_BYTES, scanEnd - readCursor);\n\t\t\t\tconst buffer = Buffer.alloc(toRead);\n\t\t\t\tconst bytesRead = fs.readSync(fd, buffer, 0, toRead, readCursor);\n\t\t\t\tif (bytesRead <= 0) break;\n\t\t\t\tconst chunk = bytesRead === buffer.length ? buffer : buffer.subarray(0, bytesRead);\n\t\t\t\tlet lineStart = 0;\n\t\t\t\tfor (let index = 0; index < chunk.length; index++) {\n\t\t\t\t\tif (chunk[index] !== 0x0a) continue;\n\t\t\t\t\tappendLineSegment(chunk.subarray(lineStart, index));\n\t\t\t\t\tif (!skippingOversizedLine && lineBytes > 0) {\n\t\t\t\t\t\thandleLine(Buffer.concat(lineParts, lineBytes).toString(\"utf-8\"));\n\t\t\t\t\t}\n\t\t\t\t\tlineParts = [];\n\t\t\t\t\tlineBytes = 0;\n\t\t\t\t\tskippingOversizedLine = false;\n\t\t\t\t\tlastCompleteCursor = readCursor + index + 1;\n\t\t\t\t\tlineStart = index + 1;\n\t\t\t\t}\n\t\t\t\tappendLineSegment(chunk.subarray(lineStart));\n\t\t\t\treadCursor += bytesRead;\n\t\t\t\tif (skippingOversizedLine) job.controlEventCursor = readCursor;\n\t\t\t}\n\t\t\tif (lastCompleteCursor > cursor) job.controlEventCursor = lastCompleteCursor;\n\t\t\telse if (scanEnd < stat.size || startedFromTail) job.controlEventCursor = scanEnd;\n\t\t} catch (error) {\n\t\t\tconsole.error(`Failed to read async control events for '${job.asyncDir}':`, error);\n\t\t} finally {\n\t\t\tfs.closeSync(fd);\n\t\t}\n\t};\n\n\tconst ensurePoller = () => {\n\t\tif (state.poller) return;\n\t\tstate.poller = setInterval(() => {\n\t\t\tif (state.asyncJobs.size === 0) {\n\t\t\t\trerenderLastWidget([]);\n\t\t\t\tif (state.poller) {\n\t\t\t\t\tclearInterval(state.poller);\n\t\t\t\t\tstate.poller = null;\n\t\t\t\t}\n\t\t\t\treturn;\n\t\t\t}\n\n\t\t\tlet widgetChanged = false;\n\t\t\tfor (const job of state.asyncJobs.values()) {\n\t\t\t\tconst widgetStateBefore = widgetRenderKey(job);\n\t\t\t\tlet nestedRefreshFailed = false;\n\t\t\t\tconst refreshNestedProjection = () => {\n\t\t\t\t\ttry {\n\t\t\t\t\t\tupdateAsyncJobNestedProjection(job);\n\t\t\t\t\t} catch (error) {\n\t\t\t\t\t\tnestedRefreshFailed = true;\n\t\t\t\t\t\tconsole.error(`Failed to refresh nested async descendants for '${job.asyncDir}':`, error);\n\t\t\t\t\t}\n\t\t\t\t};\n\t\t\t\tconst reconcileNestedDescendants = () => {\n\t\t\t\t\ttry {\n\t\t\t\t\t\tif (job.nestedRoute)\n\t\t\t\t\t\t\treconcileNestedAsyncDescendants(job.nestedRoute, {\n\t\t\t\t\t\t\t\tresultsDir,\n\t\t\t\t\t\t\t\tkill: options.kill,\n\t\t\t\t\t\t\t\tnow: options.now,\n\t\t\t\t\t\t\t});\n\t\t\t\t\t} catch (error) {\n\t\t\t\t\t\tnestedRefreshFailed = true;\n\t\t\t\t\t\tconsole.error(`Failed to refresh nested async descendants for '${job.asyncDir}':`, error);\n\t\t\t\t\t}\n\t\t\t\t\trefreshNestedProjection();\n\t\t\t\t};\n\t\t\t\ttry {\n\t\t\t\t\temitNewControlEvents(job);\n\t\t\t\t\treconcileNestedDescendants();\n\t\t\t\t\tconst reconciliation = reconcileAsyncRun(job.asyncDir, {\n\t\t\t\t\t\tresultsDir,\n\t\t\t\t\t\tkill: options.kill,\n\t\t\t\t\t\tnow: options.now,\n\t\t\t\t\t\tstartedRun: {\n\t\t\t\t\t\t\trunId: job.asyncId,\n\t\t\t\t\t\t\tpid: job.pid,\n\t\t\t\t\t\t\tsessionId: job.sessionId,\n\t\t\t\t\t\t\tmode: job.mode,\n\t\t\t\t\t\t\tagents: job.agents,\n\t\t\t\t\t\t\tchainStepCount: job.chainStepCount,\n\t\t\t\t\t\t\tparallelGroups: job.parallelGroups,\n\t\t\t\t\t\t\tstartedAt: job.startedAt,\n\t\t\t\t\t\t\tsessionFile: job.sessionFile,\n\t\t\t\t\t\t},\n\t\t\t\t\t});\n\t\t\t\t\tconst status = reconciliation.status ?? readStatus(job.asyncDir);\n\t\t\t\t\tif (status) {\n\t\t\t\t\t\tconst previousStatus = job.status;\n\t\t\t\t\t\tjob.status = status.state;\n\t\t\t\t\t\tif (\n\t\t\t\t\t\t\tjob.status !== \"complete\" &&\n\t\t\t\t\t\t\tjob.status !== \"failed\" &&\n\t\t\t\t\t\t\tjob.status !== \"paused\" &&\n\t\t\t\t\t\t\tjob.status !== \"stopped\"\n\t\t\t\t\t\t)\n\t\t\t\t\t\t\tcancelCleanup(job.asyncId);\n\t\t\t\t\t\tjob.sessionId = status.sessionId ?? job.sessionId;\n\t\t\t\t\t\tjob.activityState = status.activityState;\n\t\t\t\t\t\tjob.lastActivityAt = status.lastActivityAt ?? job.lastActivityAt;\n\t\t\t\t\t\tjob.currentTool = status.currentTool;\n\t\t\t\t\t\tjob.currentToolStartedAt = status.currentToolStartedAt;\n\t\t\t\t\t\tjob.currentPath = status.currentPath;\n\t\t\t\t\t\tjob.turnCount = status.turnCount ?? job.turnCount;\n\t\t\t\t\t\tjob.toolCount = status.toolCount ?? job.toolCount;\n\t\t\t\t\t\tjob.steering = status.steering ?? job.steering;\n\t\t\t\t\t\tjob.mode = status.mode;\n\t\t\t\t\t\tjob.parentWorkflowRunId = status.parentWorkflowRunId ?? job.parentWorkflowRunId;\n\t\t\t\t\t\tjob.workflowKey = status.workflowKey ?? job.workflowKey;\n\t\t\t\t\t\tjob.workflow = status.workflow ?? job.workflow;\n\t\t\t\t\t\tjob.currentStep = status.currentStep ?? job.currentStep;\n\t\t\t\t\t\tjob.chainStepCount = status.chainStepCount ?? job.chainStepCount;\n\t\t\t\t\t\tjob.startedAt = status.startedAt ?? job.startedAt;\n\t\t\t\t\t\tif (status.lastUpdate !== undefined) job.updatedAt = status.lastUpdate;\n\t\t\t\t\t\tif (status.steps?.length) {\n\t\t\t\t\t\t\tconst groups = normalizeParallelGroups(\n\t\t\t\t\t\t\t\tstatus.parallelGroups,\n\t\t\t\t\t\t\t\tstatus.steps.length,\n\t\t\t\t\t\t\t\tstatus.chainStepCount ?? status.steps.length,\n\t\t\t\t\t\t\t);\n\t\t\t\t\t\t\tjob.parallelGroups = groups.length ? groups : job.parallelGroups;\n\t\t\t\t\t\t\tjob.hasParallelGroups = groups.length > 0 || job.hasParallelGroups;\n\t\t\t\t\t\t\tconst activeGroup =\n\t\t\t\t\t\t\t\tstatus.currentStep !== undefined\n\t\t\t\t\t\t\t\t\t? groups.find(\n\t\t\t\t\t\t\t\t\t\t\t(group) =>\n\t\t\t\t\t\t\t\t\t\t\t\tstatus.currentStep! >= group.start &&\n\t\t\t\t\t\t\t\t\t\t\t\tstatus.currentStep! < group.start + group.count,\n\t\t\t\t\t\t\t\t\t\t)\n\t\t\t\t\t\t\t\t\t: undefined;\n\t\t\t\t\t\t\tconst visibleSteps = activeGroup\n\t\t\t\t\t\t\t\t? status.steps\n\t\t\t\t\t\t\t\t\t\t.slice(activeGroup.start, activeGroup.start + activeGroup.count)\n\t\t\t\t\t\t\t\t\t\t.map((step, index) => ({ ...step, index: activeGroup.start + index }))\n\t\t\t\t\t\t\t\t: status.steps.map((step, index) => ({ ...step, index }));\n\t\t\t\t\t\t\tjob.activeParallelGroup = Boolean(activeGroup);\n\t\t\t\t\t\t\tjob.agents = visibleSteps.map((step) => step.agent);\n\t\t\t\t\t\t\tjob.steps = visibleSteps;\n\t\t\t\t\t\t\trefreshNestedProjection();\n\t\t\t\t\t\t\tjob.stepsTotal = visibleSteps.length;\n\t\t\t\t\t\t\tjob.runningSteps = visibleSteps.filter((step) => step.status === \"running\").length;\n\t\t\t\t\t\t\tjob.completedSteps = visibleSteps.filter(\n\t\t\t\t\t\t\t\t(step) => step.status === \"complete\" || step.status === \"completed\",\n\t\t\t\t\t\t\t).length;\n\t\t\t\t\t\t\tif (status.state === \"complete\") job.completedSteps = visibleSteps.length;\n\t\t\t\t\t\t}\n\t\t\t\t\t\tjob.sessionDir = status.sessionDir ?? job.sessionDir;\n\t\t\t\t\t\tjob.outputFile = status.outputFile ?? job.outputFile;\n\t\t\t\t\t\tjob.totalTokens = status.totalTokens ?? job.totalTokens;\n\t\t\t\t\t\tjob.timeoutMs = status.timeoutMs ?? job.timeoutMs;\n\t\t\t\t\t\tjob.deadlineAt = status.deadlineAt ?? job.deadlineAt;\n\t\t\t\t\t\tjob.timedOut = status.timedOut ?? job.timedOut;\n\t\t\t\t\t\tjob.stopped = status.stopped ?? job.stopped;\n\t\t\t\t\t\tjob.turnBudget = status.turnBudget ?? job.turnBudget;\n\t\t\t\t\t\tjob.turnBudgetExceeded = status.turnBudgetExceeded ?? job.turnBudgetExceeded;\n\t\t\t\t\t\tjob.wrapUpRequested = status.wrapUpRequested ?? job.wrapUpRequested;\n\t\t\t\t\t\tjob.sessionFile = status.sessionFile ?? job.sessionFile;\n\t\t\t\t\t\tif (\n\t\t\t\t\t\t\tjob.status === \"complete\" ||\n\t\t\t\t\t\t\tjob.status === \"failed\" ||\n\t\t\t\t\t\t\tjob.status === \"paused\" ||\n\t\t\t\t\t\t\tjob.status === \"stopped\"\n\t\t\t\t\t\t) {\n\t\t\t\t\t\t\trememberFleetJob(state, job);\n\t\t\t\t\t\t\tif (\n\t\t\t\t\t\t\t\t!nestedRefreshFailed &&\n\t\t\t\t\t\t\t\t!hasLiveNestedDescendants(job.nestedChildren) &&\n\t\t\t\t\t\t\t\t(previousStatus !== job.status || !state.cleanupTimers.has(job.asyncId))\n\t\t\t\t\t\t\t) {\n\t\t\t\t\t\t\t\tscheduleCleanup(job.asyncId);\n\t\t\t\t\t\t\t}\n\t\t\t\t\t\t}\n\t\t\t\t\t\tif (widgetRenderKey(job) !== widgetStateBefore) widgetChanged = true;\n\t\t\t\t\t\tcontinue;\n\t\t\t\t\t}\n\t\t\t\t\tif (job.status === \"queued\") {\n\t\t\t\t\t\tjob.status = \"running\";\n\t\t\t\t\t\tjob.updatedAt = Date.now();\n\t\t\t\t\t}\n\t\t\t\t} catch (error) {\n\t\t\t\t\tif (job.status !== \"failed\") {\n\t\t\t\t\t\tconsole.error(`Failed to read async status for '${job.asyncDir}':`, error);\n\t\t\t\t\t\tjob.status = \"failed\";\n\t\t\t\t\t\tjob.updatedAt = Date.now();\n\t\t\t\t\t}\n\t\t\t\t\trememberFleetJob(state, job);\n\t\t\t\t\tif (!hasLiveNestedDescendants(job.nestedChildren) && !state.cleanupTimers.has(job.asyncId)) {\n\t\t\t\t\t\tscheduleCleanup(job.asyncId);\n\t\t\t\t\t}\n\t\t\t\t}\n\t\t\t\tif (widgetRenderKey(job) !== widgetStateBefore) widgetChanged = true;\n\t\t\t}\n\n\t\t\tif (widgetChanged) rerenderLastWidget();\n\t\t}, pollIntervalMs);\n\t\tstate.poller.unref?.();\n\t};\n\n\tconst handleStarted = (data: unknown) => {\n\t\tconst info = data as AsyncStartedEvent;\n\t\tif (!info.id) return;\n\t\tif (typeof state.currentSessionId === \"string\" && info.sessionId !== state.currentSessionId) return;\n\t\tconst now = Date.now();\n\t\tconst asyncDir = info.asyncDir ?? path.join(asyncDirRoot, info.id);\n\t\tconst rawAgents = info.agents?.length\n\t\t\t? info.agents\n\t\t\t: info.chain && info.chain.length > 0\n\t\t\t\t? info.chain\n\t\t\t\t: info.agent\n\t\t\t\t\t? [info.agent]\n\t\t\t\t\t: undefined;\n\t\tconst validParallelGroups = normalizeParallelGroups(\n\t\t\tinfo.parallelGroups,\n\t\t\tNumber.MAX_SAFE_INTEGER,\n\t\t\tinfo.chainStepCount ?? Number.MAX_SAFE_INTEGER,\n\t\t);\n\t\tconst firstGroup = validParallelGroups.find((group) => group.start === 0);\n\t\tconst firstGroupCount = firstGroup?.count;\n\t\tconst agents = firstGroupCount && firstGroupCount > 0 ? rawAgents?.slice(0, firstGroupCount) : rawAgents;\n\t\tstate.asyncJobs.set(info.id, {\n\t\t\tasyncId: info.id,\n\t\t\tasyncDir,\n\t\t\t...(typeof info.cwd === \"string\" ? { cwd: path.resolve(info.cwd) } : {}),\n\t\t\tstatus: \"queued\",\n\t\t\tpid: typeof info.pid === \"number\" ? info.pid : undefined,\n\t\t\t...(typeof info.sessionId === \"string\" ? { sessionId: info.sessionId } : {}),\n\t\t\tmode: info.mode ?? (info.chain ? \"chain\" : \"single\"),\n\t\t\tdescription: info.goal ?? info.task,\n\t\t\tagents,\n\t\t\tchainStepCount: info.chainStepCount,\n\t\t\tparallelGroups: validParallelGroups,\n\t\t\tnestedRoute: info.nestedRoute,\n\t\t\tstepsTotal: firstGroupCount ?? agents?.length,\n\t\t\thasParallelGroups: validParallelGroups.length > 0,\n\t\t\tactiveParallelGroup: Boolean(firstGroupCount && firstGroupCount > 0),\n\t\t\tstartedAt: now,\n\t\t\tupdatedAt: now,\n\t\t\ttimeoutMs: info.timeoutMs,\n\t\t\tdeadlineAt: info.deadlineAt,\n\t\t\tturnBudget: info.turnBudget,\n\t\t\tparentWorkflowRunId: info.parentWorkflowRunId,\n\t\t\tworkflowKey: info.workflowKey,\n\t\t\tcontrolEventCursor: 0,\n\t\t});\n\t\trememberFleetJob(state, state.asyncJobs.get(info.id)!);\n\t\tensurePoller();\n\t\trerenderLastWidget();\n\t};\n\n\tconst handleComplete = (data: unknown) => {\n\t\tconst result = data as {\n\t\t\tid?: string;\n\t\t\tsuccess?: boolean;\n\t\t\tstate?: AsyncJobState[\"status\"];\n\t\t\tasyncDir?: string;\n\t\t\tsessionId?: string;\n\t\t\tstopped?: boolean;\n\t\t};\n\t\tif (typeof state.currentSessionId === \"string\" && result.sessionId !== state.currentSessionId) return;\n\t\tconst asyncId = result.id;\n\t\tif (!asyncId) return;\n\t\tconst job = state.asyncJobs.get(asyncId);\n\t\tlet nestedRefreshFailed = false;\n\t\tif (job) {\n\t\t\tjob.status = result.state ?? (result.success ? \"complete\" : \"failed\");\n\t\t\tjob.stopped = result.stopped ?? job.stopped;\n\t\t\tjob.updatedAt = Date.now();\n\t\t\tif (result.asyncDir) job.asyncDir = result.asyncDir;\n\t\t\ttry {\n\t\t\t\tupdateAsyncJobNestedProjection(job);\n\t\t\t} catch (error) {\n\t\t\t\tnestedRefreshFailed = true;\n\t\t\t\tconsole.error(`Failed to refresh nested async descendants for '${job.asyncDir}':`, error);\n\t\t\t}\n\t\t}\n\t\tif (job) rememberFleetJob(state, job);\n\t\trerenderLastWidget();\n\t\tif (!nestedRefreshFailed && !hasLiveNestedDescendants(job?.nestedChildren)) scheduleCleanup(asyncId);\n\t};\n\n\tconst resetJobs = (ctx?: ExtensionContext) => {\n\t\tfor (const timer of state.cleanupTimers.values()) {\n\t\t\tclearTimeout(timer);\n\t\t}\n\t\tstate.cleanupTimers.clear();\n\t\tstate.asyncJobs.clear();\n\t\tstate.fleetJobs?.clear();\n\t\tstate.foregroundControls?.clear();\n\t\tstate.lastForegroundControlId = null;\n\t\tstate.resultFileCoalescer.clear();\n\t\tif (ctx?.hasUI) {\n\t\t\tstate.lastUiContext = ctx;\n\t\t\trerenderWidget(ctx, []);\n\t\t}\n\t};\n\n\tconst restoreActiveJobs = (ctx?: ExtensionContext) => {\n\t\tif (ctx?.hasUI) state.lastUiContext = ctx;\n\t\tif (!state.currentSessionId) return;\n\t\tlet runs: AsyncRunSummary[];\n\t\ttry {\n\t\t\truns = listAsyncRuns(asyncDirRoot, {\n\t\t\t\tstates: [\"queued\", \"running\"],\n\t\t\t\tsessionId: state.currentSessionId,\n\t\t\t\tresultsDir,\n\t\t\t\tkill: options.kill,\n\t\t\t\tnow: options.now,\n\t\t\t});\n\t\t} catch (error) {\n\t\t\tconsole.error(`Failed to restore active async jobs from '${asyncDirRoot}':`, error);\n\t\t\treturn;\n\t\t}\n\t\tfor (const run of runs) {\n\t\t\tconst job = summaryToJob(run);\n\t\t\tstate.asyncJobs.set(run.id, job);\n\t\t\trememberFleetJob(state, job);\n\t\t}\n\t\tif (runs.length === 0) return;\n\t\tensurePoller();\n\t\trerenderLastWidget();\n\t};\n\n\treturn { ensurePoller, refreshWidget, handleStarted, handleComplete, resetJobs, restoreActiveJobs };\n}\n"]}