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.
This commit is contained in:
@@ -1001,7 +1001,14 @@ fn split_into_tokens_with_ansi(line: &[u16]) -> SmallVec<[Vec<u16>; 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<String> = 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 =
|
||||
|
||||
@@ -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", () => {
|
||||
|
||||
@@ -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.
|
||||
@@ -152,6 +152,7 @@ export const TAB_GROUPS: Record<SettingTab, readonly string[]> = {
|
||||
"Output Limits",
|
||||
"Execution",
|
||||
"Discovery & MCP",
|
||||
"Extensions",
|
||||
"Developer",
|
||||
],
|
||||
tasks: ["Modes", "Subagents", "Isolation", "Commands & Skills"],
|
||||
|
||||
@@ -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://<id>`.");
|
||||
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))}`
|
||||
: "";
|
||||
|
||||
@@ -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 <id>` clears (#8634).
|
||||
*/
|
||||
live: boolean;
|
||||
}
|
||||
|
||||
/** Result details for messaging and job ops; fields are disjoint per op. */
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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();
|
||||
|
||||
@@ -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 () => {};
|
||||
|
||||
@@ -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");
|
||||
});
|
||||
});
|
||||
});
|
||||
|
||||
@@ -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 () => {
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -56,6 +56,7 @@ function createMockSession(
|
||||
abort: async () => {},
|
||||
dispose: async () => {},
|
||||
setIrcWakeTurnObserver: () => {},
|
||||
subscribeRunState: () => () => {},
|
||||
} as unknown as AgentSession;
|
||||
}
|
||||
|
||||
|
||||
@@ -161,6 +161,7 @@ function createAsyncSession(
|
||||
},
|
||||
dispose: options.dispose ?? (async () => {}),
|
||||
setIrcWakeTurnObserver: () => {},
|
||||
subscribeRunState: () => () => {},
|
||||
};
|
||||
harness.session = session as unknown as AgentSession;
|
||||
return harness;
|
||||
|
||||
@@ -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();
|
||||
|
||||
@@ -50,6 +50,7 @@ function createMockSession(onPrompt: (params: { emit: (event: AgentSessionEvent)
|
||||
abort: async () => {},
|
||||
dispose: async () => {},
|
||||
setIrcWakeTurnObserver: () => {},
|
||||
subscribeRunState: () => () => {},
|
||||
};
|
||||
return session as unknown as AgentSession;
|
||||
}
|
||||
|
||||
@@ -80,6 +80,7 @@ function yieldEmittingSession(
|
||||
abort: async () => {},
|
||||
dispose: async () => {},
|
||||
setIrcWakeTurnObserver: () => {},
|
||||
subscribeRunState: () => () => {},
|
||||
};
|
||||
return session as unknown as AgentSession;
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -93,6 +93,7 @@ function createMockSession(
|
||||
setIrcWakeTurnObserver: observer => {
|
||||
ircWakeTurnObserver = observer;
|
||||
},
|
||||
subscribeRunState: () => () => {},
|
||||
deliverIrcMessage: async msg => {
|
||||
const record: CustomMessage = {
|
||||
role: "custom",
|
||||
|
||||
@@ -80,6 +80,7 @@ function createMockSession(
|
||||
abort: async () => {},
|
||||
dispose: async () => {},
|
||||
setIrcWakeTurnObserver: () => {},
|
||||
subscribeRunState: () => () => {},
|
||||
};
|
||||
|
||||
return session as unknown as AgentSession;
|
||||
|
||||
@@ -32,6 +32,7 @@ function createHangingSession(): HangingSessionHandle {
|
||||
const { promise: hang, resolve: releaseHang } = Promise.withResolvers<void>();
|
||||
const session: Partial<AgentSession> = {
|
||||
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<AgentSession> = {
|
||||
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<AgentSession> = {
|
||||
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<AgentSession> = {
|
||||
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<AgentSession> = {
|
||||
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<AgentSession> = {
|
||||
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<AgentSession> = {
|
||||
setIrcWakeTurnObserver: () => {},
|
||||
subscribeRunState: () => () => {},
|
||||
state: { messages: [] } as never,
|
||||
agent: { state: { systemPrompt: ["test"] } } as never,
|
||||
extensionRunner: undefined as never,
|
||||
|
||||
@@ -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 };
|
||||
|
||||
@@ -86,6 +86,7 @@ function createYieldingSession(): AgentSession {
|
||||
abort: async () => {},
|
||||
dispose: async () => {},
|
||||
setIrcWakeTurnObserver: () => {},
|
||||
subscribeRunState: () => () => {},
|
||||
} as unknown as AgentSession;
|
||||
}
|
||||
|
||||
|
||||
@@ -109,6 +109,7 @@ function createFakeSession(config: FakeSessionConfig = {}): FakeSessionHandle {
|
||||
},
|
||||
dispose: async () => {},
|
||||
setIrcWakeTurnObserver: () => {},
|
||||
subscribeRunState: () => () => {},
|
||||
};
|
||||
return {
|
||||
session: session as AgentSession,
|
||||
|
||||
@@ -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({
|
||||
|
||||
@@ -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");
|
||||
});
|
||||
});
|
||||
|
||||
@@ -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 <underline>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
|
||||
|
||||
|
||||
@@ -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<void> {
|
||||
const nextTick = Promise.withResolvers<void>();
|
||||
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<void> {
|
||||
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<string, string | undefined>): void {
|
||||
describe("streaming scrollback — visual record", () => {
|
||||
let savedTerminalEnv: Record<string, string | undefined> = {};
|
||||
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 {
|
||||
|
||||
@@ -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<number, PendingTimer>();
|
||||
|
||||
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<void> {
|
||||
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<void> {
|
||||
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<void> {
|
||||
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;
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user