From 02cd22dc9bb67d79654326d5d4b4c94dd8a306bf Mon Sep 17 00:00:00 2001 From: can1357 Date: Sun, 16 Aug 2026 10:18:56 +0300 Subject: [PATCH] feat: added live tracking and stale status warnings for agent activity snapshots - Added live tracking and stale status warnings for agent activity snapshots. - Fixed text wrapping with ANSI escape sequences to defer style open sequences after whitespace. - Added VirtualRenderScheduler for deterministic virtual-clock rendering tests. --- crates/pi-natives/src/text.rs | 25 +++- packages/ai/test/issue-931-repro.test.ts | 7 +- packages/coding-agent/CHANGELOG.md | 10 +- .../src/config/settings-schema.ts | 1 + packages/coding-agent/src/tools/hub/jobs.ts | 25 +++- packages/coding-agent/src/tools/hub/types.ts | 6 + packages/coding-agent/src/tools/think.ts | 18 ++- .../test/interactive-mode-vibe-toggle.test.ts | 3 + ...sue-2750-subagent-runtime-fallback.test.ts | 1 + .../test/job-renderer-preview.test.ts | 7 +- .../coding-agent/test/model-registry.test.ts | 4 +- packages/coding-agent/test/shake.test.ts | 8 +- .../test/task/autoload-skills.test.ts | 1 + .../task/executor-async-quiescence.test.ts | 1 + .../test/task/executor-launch-startup.test.ts | 1 + .../test/task/executor-pass-through.test.ts | 1 + .../test/task/executor-prewalk.test.ts | 1 + .../test/task/executor-recent-output.test.ts | 1 + .../test/task/executor-soft-budget.test.ts | 1 + .../task/executor-subagent-reminders.test.ts | 1 + .../test/task/executor-wall-clock.test.ts | 9 ++ .../test/task/persisted-revive.test.ts | 1 + .../test/task/subagent-lsp.test.ts | 1 + .../test/task/task-guards.test.ts | 1 + .../coding-agent/test/tools/hub-list.test.ts | 31 ++++- .../coding-agent/test/tools/hub-wait.test.ts | 7 +- packages/natives/CHANGELOG.md | 1 + .../test/streaming-scrollback-defer.test.ts | 69 +++++----- packages/tui/test/virtual-render-scheduler.ts | 128 ++++++++++++++++++ 29 files changed, 314 insertions(+), 57 deletions(-) create mode 100644 packages/tui/test/virtual-render-scheduler.ts diff --git a/crates/pi-natives/src/text.rs b/crates/pi-natives/src/text.rs index 22c9ca770..8d0f52e99 100644 --- a/crates/pi-natives/src/text.rs +++ b/crates/pi-natives/src/text.rs @@ -1001,7 +1001,14 @@ fn split_into_tokens_with_ansi(line: &[u16]) -> SmallVec<[Vec; 4]> { && let Some(seq_len) = ansi_seq_len_u16(line, i) { let seq = &line[i..i + seq_len]; - if current.is_empty() { + // A sequence that follows visible content closes it (color reset, + // underline off) and must ride along with that token so the closer + // cannot migrate into whitespace discarded at a soft wrap (#8582). + // A sequence after whitespace opens the *next* token's style, so it + // waits for it: gluing it to the whitespace token would keep that + // space alive past the wrap point and open the style on the line + // being broken instead of the one carrying the styled word. + if current.is_empty() || in_whitespace { pending_ansi.extend_from_slice(seq); } else { current.extend_from_slice(seq); @@ -2020,6 +2027,22 @@ mod tests { assert_eq!(actual, ["plain \x1b[33mcode\x1b[39m", "next"]); } + #[test] + fn test_wrap_text_with_ansi_defers_style_open_after_trailing_space() { + let data = to_u16("read this thread \x1b[4mhttps://example.com/very/long/path\x1b[24m"); + let lines = wrap_text_with_ansi_impl(&data, 40, DEFAULT_TAB_WIDTH); + let actual: Vec = lines + .iter() + .map(|line| String::from_utf16_lossy(line)) + .collect(); + + // The underline opens on the wrapped line that carries the URL: neither + // the discarded space nor a stray `4m`/`24m` pair stays on the head line. + assert_eq!(actual[0], "read this thread"); + assert!(actual[1].starts_with("\x1b[4m")); + assert!(actual[1].contains("https://")); + } + #[test] fn test_wrap_text_with_ansi_resets_strike_without_resetting_colors() { let data = diff --git a/packages/ai/test/issue-931-repro.test.ts b/packages/ai/test/issue-931-repro.test.ts index bfe8b620a..df50d9449 100644 --- a/packages/ai/test/issue-931-repro.test.ts +++ b/packages/ai/test/issue-931-repro.test.ts @@ -1,6 +1,7 @@ import { describe, expect, it } from "bun:test"; import { streamOpenAIResponses } from "@oh-my-pi/pi-ai/providers/openai-responses"; import type { Context, Model, OpenAICompat } from "@oh-my-pi/pi-ai/types"; +import { buildModel } from "@oh-my-pi/pi-catalog/build"; import { Effort } from "@oh-my-pi/pi-catalog/effort"; const testContext: Context = { @@ -29,7 +30,9 @@ function captureResponsesPayload( } function customResponsesModel(compat: OpenAICompat): Model<"openai-responses"> { - return { + // Resolve compat through the production constructor: sparse user overrides on + // top of host detection, exactly as a configured custom model is built. + return buildModel<"openai-responses">({ id: "deepseek-v4-flash:cloud", name: "deepseek-v4-flash:cloud", api: "openai-responses", @@ -46,7 +49,7 @@ function customResponsesModel(compat: OpenAICompat): Model<"openai-responses"> { effortMap: compat.reasoningEffortMap, }, compat, - } as unknown as Model<"openai-responses">; + }); } describe("issue #931 — openai-responses reasoning effort mapping", () => { diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index c5ae4422f..be45ecdf5 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -1,8 +1,9 @@ # Changelog ## [Unreleased] +### Added -- Fixed llama.cpp model discovery producing `baseUrl` without the `/v1` prefix for non-Qwen models, causing 404 errors on OpenAI-compatible endpoints (`/v1/responses`, `/v1/chat/completions`). The fix ensures all discovered llama.cpp models include `/v1` in their `baseUrl`, matching the behavior already applied to Qwen models. +- Added Extensions tab group to settings schema ### Changed @@ -12,6 +13,8 @@ ### Fixed +- Fixed `hub` job and wait lists hiding stale running subagent registrations that have no turn in flight, ensuring they remain visible so operators can cancel them +- Fixed external thinking scratchpads running alongside native reasoning on xAI Grok 4 and other reasoning-only Responses models that reject `reasoning.effort` - Fixed llama.cpp model discovery producing a baseUrl without the /v1 prefix for non-Qwen models, causing 404 errors on OpenAI-compatible endpoints. - Fixed prompt caching on open-weight providers (DeepSeek, Qwen, GLM, …) so tool schemas stay cached across directory changes and midnight rollovers. - Fixed omp --fork omitting the source session's artifact directory, so CLI-created forks now preserve artifact:// references like interactive /fork. @@ -60,6 +63,9 @@ - Fixed /vibe cancellation leaving an in-flight model turn unaware that Vibe mode and its tools were removed. - Fixed empty local-model stops lingering on the persisted active branch after retries, preventing them from resurfacing after reload or a mid-retry process kill. - Fixed the Biome linter client silently dropping every diagnostic due to an outdated JSON output schema; it now supports Biome 2.x's diagnostic format. +- Fixed `hub jobs` and empty `hub wait` snapshots hiding running subagents that have no live turn, which removed the only way to discover and `hub cancel` a stale registration; such agents are listed again and flagged as having no turn in flight. +- Fixed external thinking being offered on xAI reasoning-only Responses models (grok-4 family) that reject `reasoning.effort`, where the private scratchpad ran alongside native reasoning instead of replacing it. +- Fixed the extension tool-call handler timeout rendering outside a titled section in `/settings` by registering its Extensions group on the Tools tab. ## [17.3.4] - 2026-08-14 @@ -15134,4 +15140,4 @@ Initial public release. ## [0.7.6] - 2025-11-13 -Previous releases did not maintain a changelog. +Previous releases did not maintain a changelog. \ No newline at end of file diff --git a/packages/coding-agent/src/config/settings-schema.ts b/packages/coding-agent/src/config/settings-schema.ts index 5fc05954b..49426538d 100644 --- a/packages/coding-agent/src/config/settings-schema.ts +++ b/packages/coding-agent/src/config/settings-schema.ts @@ -152,6 +152,7 @@ export const TAB_GROUPS: Record = { "Output Limits", "Execution", "Discovery & MCP", + "Extensions", "Developer", ], tasks: ["Modes", "Subagents", "Isolation", "Commands & Skills"], diff --git a/packages/coding-agent/src/tools/hub/jobs.ts b/packages/coding-agent/src/tools/hub/jobs.ts index 7f7d6272c..d7402f6bd 100644 --- a/packages/coding-agent/src/tools/hub/jobs.ts +++ b/packages/coding-agent/src/tools/hub/jobs.ts @@ -87,6 +87,12 @@ export function visibleJobs(manager: AsyncJobManager, ids: string[], ownerId: st * that activity instead of implying the system is quiet. Existence is * already public via the peer roster, so listing ids here leaks nothing new; * job *control* stays owner-scoped. + * + * Reporting deliberately uses the claimed `status`, not the session-corroborated + * `registry.isRunning` used by the wait-sustaining gates: a ref that claims + * `running` with no live turn is exactly the stale entry an operator must see + * here to cancel it (#8634). Hiding it would match the badge count to nothing + * and remove the only discovery path for the id. */ export function runningAgentsOutsideJobs(session: ToolSession): AgentActivitySnapshot[] { const registry = session.agentRegistry; @@ -106,13 +112,14 @@ export function runningAgentsOutsideJobs(session: ToolSession): AgentActivitySna const now = Date.now(); const out: AgentActivitySnapshot[] = []; for (const ref of registry.list()) { - if (ref.kind !== "sub" || !registry.isRunning(ref)) continue; + if (ref.kind !== "sub" || ref.status !== "running") continue; if (ref.id === selfId || covered.has(ref.id)) continue; out.push({ id: ref.id, ...(ref.parentId ? { parentId: ref.parentId } : {}), ...(ref.activity ? { activity: ref.activity } : {}), ageMs: Math.max(0, now - ref.createdAt), + live: registry.isRunning(ref), }); } return out; @@ -124,9 +131,15 @@ function describeAgents(agents: AgentActivitySnapshot[]): string[] { for (const agent of agents) { const parent = agent.parentId ? ` (spawned by \`${agent.parentId}\`)` : ""; const activity = agent.activity ? ` — ${agent.activity}` : ""; - lines.push(`- \`${agent.id}\`${parent} — up ${formatDuration(agent.ageMs)}${activity}`); + const stale = agent.live ? "" : " — no turn in flight (stale registration?)"; + lines.push(`- \`${agent.id}\`${parent} — up ${formatDuration(agent.ageMs)}${activity}${stale}`); } lines.push("", "These agents have no job entry; message them via `hub` send, transcripts at `history://`."); + if (agents.some(agent => !agent.live)) { + lines.push( + "An agent with no turn in flight cannot answer a message and never satisfies a bare `wait`; clear it with `hub` cancel.", + ); + } return lines; } @@ -690,8 +703,12 @@ export function jobsRenderResult( maxCollapsed: COLLAPSED_LIST_LIMIT, itemType: "agent", renderItem: agent => { - const icon = formatStatusIcon("running", uiTheme, options.spinnerFrame); - const badge = formatBadge("agent", "accent", uiTheme); + const icon = agent.live + ? formatStatusIcon("running", uiTheme, options.spinnerFrame) + : formatStatusIcon("warning", uiTheme); + const badge = agent.live + ? formatBadge("agent", "accent", uiTheme) + : formatBadge("agent · no turn", "warning", uiTheme); const gist = agent.activity ? ` ${uiTheme.fg("toolOutput", truncateToWidth(replaceTabs(agent.activity), LABEL_MAX_WIDTH, Ellipsis.Unicode))}` : ""; diff --git a/packages/coding-agent/src/tools/hub/types.ts b/packages/coding-agent/src/tools/hub/types.ts index d3b7e0a10..5685e3563 100644 --- a/packages/coding-agent/src/tools/hub/types.ts +++ b/packages/coding-agent/src/tools/hub/types.ts @@ -74,6 +74,12 @@ export interface AgentActivitySnapshot { activity?: string; /** Time since the agent was registered. */ ageMs: number; + /** + * Whether an attached session corroborates the `running` claim. False marks + * a ref that says `running` with no turn in flight — either a spawn still + * wiring up or a stale registration that `hub cancel ` clears (#8634). + */ + live: boolean; } /** Result details for messaging and job ops; fields are disjoint per op. */ diff --git a/packages/coding-agent/src/tools/think.ts b/packages/coding-agent/src/tools/think.ts index a4868f71b..c272af687 100644 --- a/packages/coding-agent/src/tools/think.ts +++ b/packages/coding-agent/src/tools/think.ts @@ -8,14 +8,26 @@ import { getMarkdownTheme, type Theme } from "../modes/theme/theme"; /** Whether a model transport can suppress native reasoning while private scratchpad thoughts are active. */ export function supportsExternalThinking(model: Model | null | undefined): boolean { if (!model) return false; + const compat = model.compat; const requiresThinking = model.api === "anthropic-messages" && - model.compat !== undefined && - "requiresThinkingEnabled" in model.compat && - model.compat.requiresThinkingEnabled === true; + compat !== undefined && + "requiresThinkingEnabled" in compat && + compat.requiresThinkingEnabled === true; if (model.reasoning && (requiresThinking || (model.thinking?.requiresEffort && !model.thinking.suppressWhenOff))) { return false; } + // Transports that reject `reasoning.effort` (xAI Grok 4 and the other + // reasoning-only Responses models) cannot honour `forceReasoningOff`, so the + // scratchpad would run alongside native reasoning instead of replacing it. + if ( + model.reasoning && + compat !== undefined && + (("omitReasoningEffort" in compat && compat.omitReasoningEffort === true) || + ("supportsReasoningEffort" in compat && compat.supportsReasoningEffort === false)) + ) { + return false; + } if (model.api === "google-generative-ai" || model.api === "google-gemini-cli" || model.api === "google-vertex") { return !model.reasoning || model.thinking?.mode === "budget" || model.thinking?.suppressWhenOff === true; } diff --git a/packages/coding-agent/test/interactive-mode-vibe-toggle.test.ts b/packages/coding-agent/test/interactive-mode-vibe-toggle.test.ts index c95f65119..ed0612c81 100644 --- a/packages/coding-agent/test/interactive-mode-vibe-toggle.test.ts +++ b/packages/coding-agent/test/interactive-mode-vibe-toggle.test.ts @@ -101,6 +101,9 @@ describe("InteractiveMode vibe mode toggle", () => { await Settings.init({ inMemory: true, cwd: tempDir.path() }); const model = modelRegistry.find("anthropic", "claude-sonnet-4-5"); if (!model) throw new Error("Expected claude-sonnet-4-5 to exist in registry"); + // prompt() preflights credentials via modelRegistry.getApiKey; the + // in-memory auth storage has no anthropic key, so stub it. + vi.spyOn(modelRegistry, "getApiKey").mockResolvedValue("test-key"); const registryTools = [stubTool("read"), stubTool("todo")]; storage = new ExitFaultStorage(); diff --git a/packages/coding-agent/test/issue-2750-subagent-runtime-fallback.test.ts b/packages/coding-agent/test/issue-2750-subagent-runtime-fallback.test.ts index 2105882e8..dbf8cd186 100644 --- a/packages/coding-agent/test/issue-2750-subagent-runtime-fallback.test.ts +++ b/packages/coding-agent/test/issue-2750-subagent-runtime-fallback.test.ts @@ -49,6 +49,7 @@ function createYieldingSession(fallback: "served" | "unproven" = "served"): Agen getEnabledToolNames: () => ["yield"], setActiveToolsByName: async () => {}, setIrcWakeTurnObserver: () => {}, + subscribeRunState: () => () => {}, subscribe: (listener: (event: { type: string; [key: string]: unknown }) => void) => { listeners.push(listener); return () => {}; diff --git a/packages/coding-agent/test/job-renderer-preview.test.ts b/packages/coding-agent/test/job-renderer-preview.test.ts index 24bec5a56..499097787 100644 --- a/packages/coding-agent/test/job-renderer-preview.test.ts +++ b/packages/coding-agent/test/job-renderer-preview.test.ts @@ -240,7 +240,7 @@ describe("job renderer task-result preview", () => { details: { op: "jobs" as const, jobs: [], - agents: [{ id: "Worker", parentId: "Main", activity: "grepping the tree", ageMs: 65_000 }], + agents: [{ id: "Worker", parentId: "Main", activity: "grepping the tree", ageMs: 65_000, live: true }], }, }; const component = hubToolRenderer.renderResult( @@ -258,7 +258,7 @@ describe("job renderer task-result preview", () => { it("keeps a sealed bare-poll result visible when it carries an agent roster", () => { const result = { content: [{ type: "text" as const, text: "No running background jobs to wait for." }], - details: { op: "wait" as const, jobs: [], agents: [{ id: "Worker", ageMs: 1_000 }] }, + details: { op: "wait" as const, jobs: [], agents: [{ id: "Worker", ageMs: 1_000, live: false }] }, }; const component = hubToolRenderer.renderResult( result, @@ -268,6 +268,9 @@ describe("job renderer task-result preview", () => { ); const output = Bun.stripANSI((component.render(120) as readonly string[]).join("\n")); expect(output).toContain("Worker"); + // A ref claiming `running` with no turn in flight is flagged, not shown + // as live work. + expect(output).toContain("no turn"); }); }); }); diff --git a/packages/coding-agent/test/model-registry.test.ts b/packages/coding-agent/test/model-registry.test.ts index cb31133a1..ba4d5c121 100644 --- a/packages/coding-agent/test/model-registry.test.ts +++ b/packages/coding-agent/test/model-registry.test.ts @@ -1056,7 +1056,9 @@ describe("ModelRegistry", () => { const model = registry.find("custom-local", "gpt-5.4"); expect(model?.contextWindow).toBe(1_000_000); - expect(model?.baseUrl).toBe("http://127.0.0.1:8080"); + // llama.cpp discovery probes the bare root (`/models`, `/props`); chat + // traffic must go to the OpenAI-compatible `/v1` prefix. + expect(model?.baseUrl).toBe("http://127.0.0.1:8080/v1"); }); test("discoverable custom compat survives refresh", async () => { diff --git a/packages/coding-agent/test/shake.test.ts b/packages/coding-agent/test/shake.test.ts index 98be7bec2..906333c49 100644 --- a/packages/coding-agent/test/shake.test.ts +++ b/packages/coding-agent/test/shake.test.ts @@ -372,7 +372,7 @@ describe("AgentSession shake", () => { }; session.agent.emitExternalEvent({ type: "message_end", message: assistantMessage }); session.agent.emitExternalEvent({ type: "agent_end", messages: [assistantMessage] }); - await Bun.sleep(20); + await session.waitForIdle(); expect(shakeSpy).toHaveBeenCalledWith("elide", expect.anything()); const start = events.filter(e => e.type === "auto_compaction_start"); @@ -580,7 +580,7 @@ describe("AgentSession shake", () => { }; session.agent.emitExternalEvent({ type: "message_end", message: assistantMessage }); session.agent.emitExternalEvent({ type: "agent_end", messages: [assistantMessage] }); - await Bun.sleep(50); + await session.waitForIdle(); // Shake fires once. The pre-fix bug auto-continued, which would re-trigger shake // on the next agent_end. The fix replaces that loop with a one-shot fallback. @@ -634,7 +634,7 @@ describe("AgentSession shake", () => { }; session.agent.emitExternalEvent({ type: "message_end", message: assistantMessage }); session.agent.emitExternalEvent({ type: "agent_end", messages: [assistantMessage] }); - await Bun.sleep(50); + await session.waitForIdle(); expect(shakeSpy).toHaveBeenCalledTimes(1); @@ -705,7 +705,7 @@ describe("AgentSession shake", () => { session.agent.emitExternalEvent({ type: "message_end", message: assistantMessage }); session.agent.emitExternalEvent({ type: "agent_end", messages: [assistantMessage] }); - await Bun.sleep(50); + await session.waitForIdle(); expect(shakeSpy).toHaveBeenCalledTimes(1); const fullStart = events.find( diff --git a/packages/coding-agent/test/task/autoload-skills.test.ts b/packages/coding-agent/test/task/autoload-skills.test.ts index 2796182d8..0f3f1dba0 100644 --- a/packages/coding-agent/test/task/autoload-skills.test.ts +++ b/packages/coding-agent/test/task/autoload-skills.test.ts @@ -56,6 +56,7 @@ function createMockSession( abort: async () => {}, dispose: async () => {}, setIrcWakeTurnObserver: () => {}, + subscribeRunState: () => () => {}, } as unknown as AgentSession; } diff --git a/packages/coding-agent/test/task/executor-async-quiescence.test.ts b/packages/coding-agent/test/task/executor-async-quiescence.test.ts index d5b376ae8..0fe0fb101 100644 --- a/packages/coding-agent/test/task/executor-async-quiescence.test.ts +++ b/packages/coding-agent/test/task/executor-async-quiescence.test.ts @@ -161,6 +161,7 @@ function createAsyncSession( }, dispose: options.dispose ?? (async () => {}), setIrcWakeTurnObserver: () => {}, + subscribeRunState: () => () => {}, }; harness.session = session as unknown as AgentSession; return harness; diff --git a/packages/coding-agent/test/task/executor-launch-startup.test.ts b/packages/coding-agent/test/task/executor-launch-startup.test.ts index c78ae9ac2..85b771ff5 100644 --- a/packages/coding-agent/test/task/executor-launch-startup.test.ts +++ b/packages/coding-agent/test/task/executor-launch-startup.test.ts @@ -68,6 +68,7 @@ it("overlaps registry refresh with session-file opening and session setup", asyn abort: async () => {}, dispose: async () => {}, setIrcWakeTurnObserver: () => {}, + subscribeRunState: () => () => {}, } as unknown as AgentSession; vi.spyOn(sdkModule, "createAgentSession").mockImplementation(async () => { sessionCreationStarted.resolve(); diff --git a/packages/coding-agent/test/task/executor-pass-through.test.ts b/packages/coding-agent/test/task/executor-pass-through.test.ts index 51aba195c..1b783f6fb 100644 --- a/packages/coding-agent/test/task/executor-pass-through.test.ts +++ b/packages/coding-agent/test/task/executor-pass-through.test.ts @@ -50,6 +50,7 @@ function createMockSession(onPrompt: (params: { emit: (event: AgentSessionEvent) abort: async () => {}, dispose: async () => {}, setIrcWakeTurnObserver: () => {}, + subscribeRunState: () => () => {}, }; return session as unknown as AgentSession; } diff --git a/packages/coding-agent/test/task/executor-prewalk.test.ts b/packages/coding-agent/test/task/executor-prewalk.test.ts index 998c35935..5141641e3 100644 --- a/packages/coding-agent/test/task/executor-prewalk.test.ts +++ b/packages/coding-agent/test/task/executor-prewalk.test.ts @@ -80,6 +80,7 @@ function yieldEmittingSession( abort: async () => {}, dispose: async () => {}, setIrcWakeTurnObserver: () => {}, + subscribeRunState: () => () => {}, }; return session as unknown as AgentSession; } diff --git a/packages/coding-agent/test/task/executor-recent-output.test.ts b/packages/coding-agent/test/task/executor-recent-output.test.ts index e1986f26e..0762d6c59 100644 --- a/packages/coding-agent/test/task/executor-recent-output.test.ts +++ b/packages/coding-agent/test/task/executor-recent-output.test.ts @@ -193,6 +193,7 @@ function createScriptedSession( isAborted: () => aborted, dispose: async () => {}, setIrcWakeTurnObserver: () => {}, + subscribeRunState: () => () => {}, }; // AgentSession is a concrete class; the executor consumes only this // structural subset. Deliberate documented test-double escape hatch, diff --git a/packages/coding-agent/test/task/executor-soft-budget.test.ts b/packages/coding-agent/test/task/executor-soft-budget.test.ts index 7a2214823..f2c18da78 100644 --- a/packages/coding-agent/test/task/executor-soft-budget.test.ts +++ b/packages/coding-agent/test/task/executor-soft-budget.test.ts @@ -93,6 +93,7 @@ function createMockSession( setIrcWakeTurnObserver: observer => { ircWakeTurnObserver = observer; }, + subscribeRunState: () => () => {}, deliverIrcMessage: async msg => { const record: CustomMessage = { role: "custom", diff --git a/packages/coding-agent/test/task/executor-subagent-reminders.test.ts b/packages/coding-agent/test/task/executor-subagent-reminders.test.ts index bdfea201f..bc80846bb 100644 --- a/packages/coding-agent/test/task/executor-subagent-reminders.test.ts +++ b/packages/coding-agent/test/task/executor-subagent-reminders.test.ts @@ -80,6 +80,7 @@ function createMockSession( abort: async () => {}, dispose: async () => {}, setIrcWakeTurnObserver: () => {}, + subscribeRunState: () => () => {}, }; return session as unknown as AgentSession; diff --git a/packages/coding-agent/test/task/executor-wall-clock.test.ts b/packages/coding-agent/test/task/executor-wall-clock.test.ts index b03267093..d12b2bf06 100644 --- a/packages/coding-agent/test/task/executor-wall-clock.test.ts +++ b/packages/coding-agent/test/task/executor-wall-clock.test.ts @@ -32,6 +32,7 @@ function createHangingSession(): HangingSessionHandle { const { promise: hang, resolve: releaseHang } = Promise.withResolvers(); const session: Partial = { setIrcWakeTurnObserver: () => {}, + subscribeRunState: () => () => {}, state: { messages: [] } as never, agent: { state: { systemPrompt: ["test"] } } as never, extensionRunner: undefined as never, @@ -123,6 +124,7 @@ describe("runSubprocess wall clock (task.maxRuntimeMs)", () => { const settings = Settings.isolated({ "task.maxRuntimeMs": 0 }); const fastSession: Partial = { setIrcWakeTurnObserver: () => {}, + subscribeRunState: () => () => {}, state: { messages: [] } as never, agent: { state: { systemPrompt: ["test"] } } as never, extensionRunner: undefined as never, @@ -212,6 +214,7 @@ describe("runSubprocess wall clock (task.maxRuntimeMs)", () => { const lateSession = { dispose: async () => lateDisposed.resolve(), setIrcWakeTurnObserver: () => {}, + subscribeRunState: () => () => {}, } as unknown as AgentSession; let lateInstall = registry.get("late-generation"); vi.spyOn(sdkModule, "createAgentSession").mockImplementation(async (options = {}) => { @@ -251,6 +254,7 @@ describe("runSubprocess wall clock (task.maxRuntimeMs)", () => { const replacementSession = { dispose: async () => {}, setIrcWakeTurnObserver: () => {}, + subscribeRunState: () => () => {}, } as unknown as AgentSession; const replacement = registry.register({ id: "late-generation", @@ -279,6 +283,7 @@ describe("runSubprocess wall clock (task.maxRuntimeMs)", () => { let abortCount = 0; const session: Partial = { setIrcWakeTurnObserver: () => {}, + subscribeRunState: () => () => {}, state: { messages: [] } as never, agent: { state: { systemPrompt: ["test"] } } as never, extensionRunner: undefined as never, @@ -358,6 +363,7 @@ describe("runSubprocess wall clock (task.maxRuntimeMs)", () => { let abortCountBeforeYieldExecutionEnd: number | undefined; const session: Partial = { setIrcWakeTurnObserver: () => {}, + subscribeRunState: () => () => {}, state: { messages: [] } as never, agent: { state: { systemPrompt: ["test"] } } as never, extensionRunner: undefined as never, @@ -459,6 +465,7 @@ describe("runSubprocess wall clock (task.maxRuntimeMs)", () => { const promptCalls: Array<{ text: string; options?: PromptOptions }> = []; const session: Partial = { setIrcWakeTurnObserver: () => {}, + subscribeRunState: () => () => {}, state: { messages: [] } as never, agent: { state: { systemPrompt: ["test"] } } as never, extensionRunner: undefined as never, @@ -590,6 +597,7 @@ describe("runSubprocess wall clock (task.maxRuntimeMs)", () => { let abortCountAfterFollowingTurn: number | undefined; const session: Partial = { setIrcWakeTurnObserver: () => {}, + subscribeRunState: () => () => {}, state: { messages: [] } as never, agent: { state: { systemPrompt: ["test"] } } as never, extensionRunner: undefined as never, @@ -674,6 +682,7 @@ describe("runSubprocess wall clock (task.maxRuntimeMs)", () => { const settings = Settings.isolated({ "task.maxRuntimeMs": 0 }); const fastSession: Partial = { setIrcWakeTurnObserver: () => {}, + subscribeRunState: () => () => {}, state: { messages: [] } as never, agent: { state: { systemPrompt: ["test"] } } as never, extensionRunner: undefined as never, diff --git a/packages/coding-agent/test/task/persisted-revive.test.ts b/packages/coding-agent/test/task/persisted-revive.test.ts index c99fee83e..7ef12682a 100644 --- a/packages/coding-agent/test/task/persisted-revive.test.ts +++ b/packages/coding-agent/test/task/persisted-revive.test.ts @@ -57,6 +57,7 @@ function createRevivedSession(activeToolNames: string[][]): RevivedSessionHandle setIrcWakeTurnObserver: (next: IrcWakeObserver | undefined) => { observer = next; }, + subscribeRunState: () => () => {}, getLastAssistantMessage: () => undefined, } as unknown as AgentSession; return { session, observer: () => observer }; diff --git a/packages/coding-agent/test/task/subagent-lsp.test.ts b/packages/coding-agent/test/task/subagent-lsp.test.ts index 85961c2bc..0adcb61a9 100644 --- a/packages/coding-agent/test/task/subagent-lsp.test.ts +++ b/packages/coding-agent/test/task/subagent-lsp.test.ts @@ -86,6 +86,7 @@ function createYieldingSession(): AgentSession { abort: async () => {}, dispose: async () => {}, setIrcWakeTurnObserver: () => {}, + subscribeRunState: () => () => {}, } as unknown as AgentSession; } diff --git a/packages/coding-agent/test/task/task-guards.test.ts b/packages/coding-agent/test/task/task-guards.test.ts index b0726dcf3..6d93fcad5 100644 --- a/packages/coding-agent/test/task/task-guards.test.ts +++ b/packages/coding-agent/test/task/task-guards.test.ts @@ -109,6 +109,7 @@ function createFakeSession(config: FakeSessionConfig = {}): FakeSessionHandle { }, dispose: async () => {}, setIrcWakeTurnObserver: () => {}, + subscribeRunState: () => () => {}, }; return { session: session as AgentSession, diff --git a/packages/coding-agent/test/tools/hub-list.test.ts b/packages/coding-agent/test/tools/hub-list.test.ts index 4056190d3..2e2105938 100644 --- a/packages/coding-agent/test/tools/hub-list.test.ts +++ b/packages/coding-agent/test/tools/hub-list.test.ts @@ -1,16 +1,43 @@ import { describe, expect, it } from "bun:test"; import * as path from "node:path"; import { AgentRegistry, MAIN_AGENT_ID } from "@oh-my-pi/pi-coding-agent/registry/agent-registry"; +import { CURRENT_SESSION_VERSION } from "@oh-my-pi/pi-coding-agent/session/session-entries"; import { executeList } from "@oh-my-pi/pi-coding-agent/tools/hub/messaging"; import { TempDir } from "@oh-my-pi/pi-utils"; +function sessionHeader(id: string): string { + return JSON.stringify({ + type: "session", + version: CURRENT_SESSION_VERSION, + id, + timestamp: "2026-08-13T17:14:48.125Z", + cwd: "/tmp", + }); +} + describe("hub list", () => { it("restores persisted peers after the process registry is lost", async () => { using tempDir = TempDir.createSync("@omp-hub-list-persisted-"); const sessionFile = path.join(tempDir.path(), "main.jsonl"); const workerSessionFile = path.join(tempDir.path(), "main", "Worker.jsonl"); - await Bun.write(sessionFile, ""); - await Bun.write(workerSessionFile, ""); + await Bun.write(sessionFile, `${sessionHeader("main")}\n`); + // A settled subagent transcript: header + session_init. A header-only file + // is a mid-spawn stub and stays unregistered by design. + await Bun.write( + workerSessionFile, + `${[ + sessionHeader("worker"), + JSON.stringify({ + type: "session_init", + id: "si", + parentId: null, + timestamp: "2026-08-13T17:14:49.000Z", + systemPrompt: "review", + task: "review the diff", + tools: ["read"], + }), + ].join("\n")}\n`, + ); const registry = new AgentRegistry(); registry.register({ diff --git a/packages/coding-agent/test/tools/hub-wait.test.ts b/packages/coding-agent/test/tools/hub-wait.test.ts index 517384b84..8222bc09e 100644 --- a/packages/coding-agent/test/tools/hub-wait.test.ts +++ b/packages/coding-agent/test/tools/hub-wait.test.ts @@ -123,10 +123,15 @@ describe("hub unified wait", () => { }); const manager = new AsyncJobManager({ onJobComplete: () => {} }); + // `timeoutMs: 0` would block forever if the stale ref still opened the + // message-wait gate; the test times out instead of asserting. const result = await new HubTool(makeSession(manager)).execute("call_4", { op: "wait", timeoutMs: 0 }); const text = result.content[0]?.type === "text" ? result.content[0].text : ""; expect(text).toContain("No running background jobs to wait for."); - expect(result.useless).toBe(true); + // The stale ref is reported (not silently dropped): it is the only handle + // the caller has for clearing it with `hub cancel`. + expect(text).toContain("Zombie"); + expect(text).toContain("no turn in flight"); }); }); diff --git a/packages/natives/CHANGELOG.md b/packages/natives/CHANGELOG.md index 21db8937f..390e97df2 100644 --- a/packages/natives/CHANGELOG.md +++ b/packages/natives/CHANGELOG.md @@ -6,6 +6,7 @@ - Fixed the native `xargs` builtin panicking in `-I`/`-i` replace mode when stdin is empty; it now exits successfully without running the command, matching GNU behavior. - Fixed inline-code foreground color incorrectly carrying into plain text when a Markdown codespan ended exactly at a soft-wrap boundary. +- Fixed `wrapTextWithAnsi` leaving a trailing space plus a stray underline open/close pair on the line above a soft wrap when a style opened immediately after that space (e.g. `read this thread https://…`), a regression from the codespan color-bleed fix: only sequences that follow visible content now ride along with the current token, while sequences after whitespace still wait for the token they style. ## [17.3.4] - 2026-08-14 diff --git a/packages/tui/test/streaming-scrollback-defer.test.ts b/packages/tui/test/streaming-scrollback-defer.test.ts index 793e22e2e..c191c53e4 100644 --- a/packages/tui/test/streaming-scrollback-defer.test.ts +++ b/packages/tui/test/streaming-scrollback-defer.test.ts @@ -7,6 +7,7 @@ import { type NativeScrollbackLiveRegion, TUI, } from "@oh-my-pi/pi-tui"; +import { VirtualRenderScheduler } from "./virtual-render-scheduler"; import { VirtualTerminal } from "./virtual-terminal"; // Law-encoding suite for native-scrollback commits. @@ -107,21 +108,18 @@ class CommittedRowsWireProbe extends CommittedRowsProbe { } } +/** One scheduler per file; `beforeEach` rewinds it so each test starts at t=0 with an empty queue. */ +const scheduler = new VirtualRenderScheduler(); + async function settle(term: VirtualTerminal): Promise { - const nextTick = Promise.withResolvers(); - process.nextTick(nextTick.resolve); - await nextTick.promise; - await Bun.sleep(40); - await term.flush(); + await scheduler.settle(term); } // The non-multiplexer resize fast path paints the viewport at once and defers // the authoritative full replay (the ED3 scrollback rebuild) until the drag has -// been quiet for the resize settle window (120 ms). This is an integration test -// against the real render scheduler, so the window is driven with a real delay. +// been quiet for the resize settle window (120 ms). Open that window explicitly. async function settleResize(term: VirtualTerminal): Promise { - await Bun.sleep(160); - await settle(term); + await scheduler.advance(term, 160); } function capture(term: VirtualTerminal): string[] { @@ -197,6 +195,7 @@ function restoreTerminalEnv(saved: Record): void { describe("streaming scrollback — visual record", () => { let savedTerminalEnv: Record = {}; beforeEach(() => { + scheduler.reset(); savedTerminalEnv = saveTerminalEnv(); }); afterEach(() => { @@ -208,7 +207,7 @@ describe("streaming scrollback — visual record", () => { if (process.platform === "win32") return; const term = new VirtualTerminal(20, 4); overrideProbe(term, undefined); - const tui = new TUI(term); + const tui = new TUI(term, undefined, { renderScheduler: scheduler }); const sealed = new LineList(rows("prior-", 12)); const live = new SeamLineList([]); @@ -254,7 +253,7 @@ describe("streaming scrollback — visual record", () => { if (process.platform === "win32") return; const term = new VirtualTerminal(20, 4); overrideProbe(term, undefined); - const tui = new TUI(term); + const tui = new TUI(term, undefined, { renderScheduler: scheduler }); const live = new SeamLineList([]); try { @@ -286,7 +285,7 @@ describe("streaming scrollback — visual record", () => { if (process.platform === "win32") return; const term = new VirtualTerminal(60, 8, 1_000); overrideProbe(term, undefined); - const tui = new TUI(term); + const tui = new TUI(term, undefined, { renderScheduler: scheduler }); const wall = new PinnedSeamLineList([]); try { @@ -323,7 +322,7 @@ describe("streaming scrollback — visual record", () => { if (process.platform === "win32") return; const term = new VirtualTerminal(20, 4); overrideProbe(term, undefined); - const tui = new TUI(term); + const tui = new TUI(term, undefined, { renderScheduler: scheduler }); // An append-only streaming reply declares every rendered row final // (Infinity clamps to the rendered length): its scrolled-off head enters // the verified zone and never needs a finalize-time repair. @@ -357,7 +356,7 @@ describe("streaming scrollback — visual record", () => { if (process.platform === "win32") return; const term = new VirtualTerminal(24, 4); overrideProbe(term, undefined); - const tui = new TUI(term); + const tui = new TUI(term, undefined, { renderScheduler: scheduler }); const sealed = new LineList(rows("prior-", 12)); const live = new SeamLineList([]); @@ -407,7 +406,7 @@ describe("streaming scrollback — visual record", () => { if (process.platform === "win32") return; const term = new VirtualTerminal(24, 4); overrideProbe(term, undefined); - const tui = new TUI(term); + const tui = new TUI(term, undefined, { renderScheduler: scheduler }); const sealed = new LineList(rows("prior-", 12)); const live = new SeamLineList([]); // Status loader below the transcript: also reports a seam. Exactness is @@ -458,7 +457,7 @@ describe("streaming scrollback — visual record", () => { if (process.platform === "win32") return; const term = new VirtualTerminal(40, 10); overrideProbe(term, undefined); - const tui = new TUI(term); + const tui = new TUI(term, undefined, { renderScheduler: scheduler }); const component = new LineList([...rows("init-", 10), "prompt"]); try { @@ -501,7 +500,7 @@ describe("streaming scrollback — visual record", () => { if (process.platform === "win32") return; const term = new VirtualTerminal(40, 10); overrideProbe(term, undefined); - const tui = new TUI(term); + const tui = new TUI(term, undefined, { renderScheduler: scheduler }); const component = new LineList([...rows("init-", 10), "prompt"]); try { @@ -536,7 +535,7 @@ describe("streaming scrollback — visual record", () => { if (process.platform === "win32") return; const term = new VirtualTerminal(20, 4); overrideProbe(term, undefined); - const tui = new TUI(term); + const tui = new TUI(term, undefined, { renderScheduler: scheduler }); const sealed = new LineList(rows("prior-", 12)); const live = new SeamLineList([]); @@ -571,7 +570,7 @@ describe("streaming scrollback — visual record", () => { if (process.platform === "win32") return; const term = new VirtualTerminal(24, 4); overrideProbe(term, undefined); - const tui = new TUI(term); + const tui = new TUI(term, undefined, { renderScheduler: scheduler }); const sealed = new LineList(rows("base-", 12)); try { @@ -608,7 +607,7 @@ describe("streaming scrollback — visual record", () => { if (process.platform === "win32") return; const term = new VirtualTerminal(40, 10); overrideProbe(term, undefined); - const tui = new TUI(term); + const tui = new TUI(term, undefined, { renderScheduler: scheduler }); const component = new LineList([...rows("init-", 5), "prompt"]); try { @@ -649,7 +648,7 @@ describe("streaming scrollback — visual record", () => { if (process.platform === "win32") return; const term = new VirtualTerminal(20, 4); overrideProbe(term, undefined); - const tui = new TUI(term); + const tui = new TUI(term, undefined, { renderScheduler: scheduler }); const probe = new CommittedRowsProbe([]); try { @@ -677,7 +676,7 @@ describe("streaming scrollback — visual record", () => { if (process.platform === "win32") return; const term = new VirtualTerminal(40, 8); overrideProbe(term, undefined); - const tui = new TUI(term); + const tui = new TUI(term, undefined, { renderScheduler: scheduler }); const probe = new CommittedRowsWireProbe([]); try { @@ -726,7 +725,7 @@ describe("streaming scrollback — visual record", () => { if (process.platform === "win32") return; const term = new VirtualTerminal(40, 8); overrideProbe(term, undefined); - const tui = new TUI(term); + const tui = new TUI(term, undefined, { renderScheduler: scheduler }); // Content taller than the viewport before the first paint: the initial // frame takes the full-paint path, whose replay commits (frame - height) // rows in one shot on a separate exit from the ordinary update emit. @@ -753,7 +752,7 @@ describe("streaming scrollback — visual record", () => { if (process.platform === "win32") return; const term = new VirtualTerminal(40, 8); overrideProbe(term, undefined); - const tui = new TUI(term); + const tui = new TUI(term, undefined, { renderScheduler: scheduler }); // A short header above a tall overflowing body: the engine's committed // boundary sails past the header's 2-row extent. Both feeds are in the // child's own coordinates and must saturate at what the child actually @@ -801,7 +800,7 @@ describe("streaming scrollback — visual record", () => { if (process.platform === "win32") return; const term = new VirtualTerminal(20, 4); overrideProbe(term, undefined); - const tui = new TUI(term); + const tui = new TUI(term, undefined, { renderScheduler: scheduler }); // A block that rewrites an interior row every frame (a streaming table // re-aligning, a collapsing preview). Its scrolled rows are frozen // snapshots: drift never sprays re-anchors; the single strict scan at @@ -857,7 +856,7 @@ describe("streaming scrollback — visual record", () => { if (process.platform === "win32") return; const term = new VirtualTerminal(20, 4); overrideProbe(term, undefined); - const tui = new TUI(term); + const tui = new TUI(term, undefined, { renderScheduler: scheduler }); // The block declares its whole body final, commits, then violates the // contract by rewriting TWO committed rows (alignment breaks, so the // tail-sample tolerance cannot absorb it). The audit re-anchors, the @@ -908,7 +907,7 @@ describe("scrollback commit gap — live barriers", () => { if (process.platform === "win32") return; const term = new VirtualTerminal(20, 4); overrideProbe(term, undefined); - const tui = new TUI(term); + const tui = new TUI(term, undefined, { renderScheduler: scheduler }); const root = new SeamLineList([]); try { @@ -946,7 +945,7 @@ describe("scrollback commit gap — live barriers", () => { if (process.platform === "win32") return; const term = new VirtualTerminal(20, 4); overrideProbe(term, undefined); - const tui = new TUI(term); + const tui = new TUI(term, undefined, { renderScheduler: scheduler }); const root = new SeamLineList([]); try { @@ -983,7 +982,7 @@ describe("scrollback commit gap — live barriers", () => { if (process.platform === "win32") return; const term = new VirtualTerminal(20, 4); overrideProbe(term, undefined); - const tui = new TUI(term); + const tui = new TUI(term, undefined, { renderScheduler: scheduler }); const root = new SeamLineList([]); try { @@ -1021,7 +1020,7 @@ describe("scrollback commit gap — live barriers", () => { if (process.platform === "win32") return; const term = new VirtualTerminal(20, 5); overrideProbe(term, undefined); - const tui = new TUI(term); + const tui = new TUI(term, undefined, { renderScheduler: scheduler }); const barrier = new SeamLineList(["[tool pending]"]); const tail = new LineList(rows("out-", 10)); @@ -1058,7 +1057,7 @@ describe("scrollback commit gap — live barriers", () => { if (process.platform === "win32") return; const term = new VirtualTerminal(20, 5); overrideProbe(term, undefined); - const tui = new TUI(term); + const tui = new TUI(term, undefined, { renderScheduler: scheduler }); const root = new SeamLineList([]); try { @@ -1098,7 +1097,7 @@ describe("scrollback commit gap — live barriers", () => { // — so finalize needs NO repair and the tape never duplicates a byte. const term = new VirtualTerminal(20, 4); overrideProbe(term, undefined); - const tui = new TUI(term); + const tui = new TUI(term, undefined, { renderScheduler: scheduler }); const root = new SeamLineList([]); try { @@ -1137,7 +1136,7 @@ describe("scrollback commit gap — live barriers", () => { if (process.platform === "win32") return; const term = new VirtualTerminal(20, 4); overrideProbe(term, undefined); - const tui = new TUI(term); + const tui = new TUI(term, undefined, { renderScheduler: scheduler }); const root = new SeamLineList([]); try { @@ -1177,7 +1176,7 @@ describe("scrollback commit gap — live barriers", () => { if (process.platform === "win32") return; const term = new VirtualTerminal(20, 4); overrideProbe(term, undefined); - const tui = new TUI(term); + const tui = new TUI(term, undefined, { renderScheduler: scheduler }); const root = new SeamLineList([]); try { @@ -1228,7 +1227,7 @@ describe("scrollback divergence — multiplexer fallback", () => { if (process.platform === "win32") return; const term = new VirtualTerminal(20, 4); overrideProbe(term, undefined); - const tui = new TUI(term); + const tui = new TUI(term, undefined, { renderScheduler: scheduler }); const root = new SeamLineList([]); try { diff --git a/packages/tui/test/virtual-render-scheduler.ts b/packages/tui/test/virtual-render-scheduler.ts new file mode 100644 index 000000000..80ee85027 --- /dev/null +++ b/packages/tui/test/virtual-render-scheduler.ts @@ -0,0 +1,128 @@ +import type { RenderScheduler, RenderTimer } from "@oh-my-pi/pi-tui/tui"; +import type { VirtualTerminal } from "./virtual-terminal"; + +/** + * Longest delay a throttled frame can ask for: the 30 Hz cadence. Adaptive + * backpressure adds nothing under a virtual clock (a frame costs 0 virtual ms), + * and every deferred window the engine arms is longer — multiplexer resize + * debounce 50 ms, resize viewport settle 120 ms, ConPTY post-paint settle + * 150 ms — so this horizon separates "frames the pipeline owes us now" from + * "windows a test must open deliberately". + */ +const FRAME_HORIZON_MS = 40; + +/** Bound on drain rounds; a render loop that keeps re-arming is a bug, not a wait. */ +const MAX_ROUNDS = 500; + +interface PendingTimer { + at: number; + run: () => void; +} + +/** + * Deterministic {@link RenderScheduler} with a virtual clock. + * + * `TUI` reads every timestamp through `RenderScheduler.now()`, so owning the + * clock removes wall-clock races from render assertions: a loaded machine can + * neither lose a throttled frame (sleep too short) nor open a deferred settle + * window early (sleep too long). + * + * ```ts + * const scheduler = new VirtualRenderScheduler(); + * const tui = new TUI(term, undefined, { renderScheduler: scheduler }); + * tui.start(); + * await scheduler.settle(term); // every frame the engine owes, nothing more + * await scheduler.advance(term, 160); // open the resize settle window + * ``` + */ +export class VirtualRenderScheduler implements RenderScheduler { + #now = 0; + #nextTimerId = 0; + #immediates: (() => void)[] = []; + #timers = new Map(); + + now(): number { + return this.#now; + } + + scheduleImmediate(callback: () => void): void { + this.#immediates.push(callback); + } + + scheduleRender(callback: () => void, delayMs: number): RenderTimer { + const id = this.#nextTimerId; + this.#nextTimerId += 1; + this.#timers.set(id, { at: this.#now + Math.max(0, delayMs), run: callback }); + return { + cancel: () => { + this.#timers.delete(id); + }, + }; + } + + /** Drop every queued frame and rewind the clock, for reuse across tests. */ + reset(): void { + this.#now = 0; + this.#immediates = []; + this.#timers.clear(); + } + + /** + * Run the frames the engine owes — immediates plus timers landing within + * `horizonMs` of the clock, including the ones those frames schedule — until + * the pipeline is quiescent. Deferred windows stay armed. + */ + async settle(term: VirtualTerminal, horizonMs = FRAME_HORIZON_MS): Promise { + await this.#drain(term, () => this.#now + horizonMs); + } + + /** + * Open a deferred window: advance the clock by `ms`, running every frame that + * comes due on the way, then settle the frames that follow it. + */ + async advance(term: VirtualTerminal, ms: number): Promise { + const deadline = this.#now + ms; + await this.#drain(term, () => deadline); + this.#now = Math.max(this.#now, deadline); + await this.settle(term); + } + + async #drain(term: VirtualTerminal, deadline: () => number): Promise { + for (let round = 0; round < MAX_ROUNDS; round++) { + if (!this.#runImmediates() && !this.#runDue(deadline())) return; + // Terminal writes land synchronously; yielding lets promise + // continuations queue their renders before the next round. + await term.flush(); + } + throw new Error(`VirtualRenderScheduler did not settle after ${MAX_ROUNDS} rounds`); + } + + #runImmediates(): boolean { + if (this.#immediates.length === 0) return false; + const queued = this.#immediates; + this.#immediates = []; + for (const run of queued) run(); + return true; + } + + /** Fire the earliest timer batch due at or before `deadline`, advancing the clock to it. */ + #runDue(deadline: number): boolean { + let earliest: number | undefined; + for (const timer of this.#timers.values()) { + if (timer.at <= deadline && (earliest === undefined || timer.at < earliest)) earliest = timer.at; + } + if (earliest === undefined) return false; + this.#now = Math.max(this.#now, earliest); + const due: number[] = []; + for (const [id, timer] of this.#timers) { + if (timer.at === earliest) due.push(id); + } + for (const id of due) { + const timer = this.#timers.get(id); + if (!timer) continue; + this.#timers.delete(id); + timer.run(); + } + return true; + } +}