Merge PR #3496: fix(tui): guard streaming Esc cancellation (@roboomp)

This commit is contained in:
can1357
2026-06-27 01:39:35 +02:00
6 changed files with 248 additions and 4 deletions
+4
View File
@@ -100,6 +100,10 @@
- Fixed a background-task spawn slot leaking from the `task.maxConcurrency` limiter when progress reporting threw between acquiring the slot and entering the guarded run: `markRunning`/`reportProgress` now run inside the try whose `finally` releases the semaphore, so a failed progress report can no longer permanently shrink subagent concurrency. ([#3464](https://github.com/can1357/oh-my-pi/issues/3464))
- Fixed active goal runs that successfully call `yield` and then receive a trailing empty assistant `stop` skipping threshold compaction; post-yield empty-stop suppression now still anchors active-goal compaction on the yield-bearing assistant turn, so long-running tasks continue after maintenance instead of settling early.
### Fixed
- Fixed streaming Esc handling so the first Esc arms a 2s cancel hint and only a second Esc in that window aborts the active response ([#3493](https://github.com/can1357/oh-my-pi/issues/3493)).
## [16.1.19] - 2026-06-25
### Fixed
@@ -125,6 +125,7 @@ const TINY_TITLE_PROGRESS_REVEAL_DELAY_MS = 1_000;
// deliberate human double-tap is always tens of milliseconds apart.
const LEFT_DOUBLE_TAP_MIN_GAP_MS = 40;
const LEFT_DOUBLE_TAP_MAX_GAP_MS = 500;
const STREAMING_ESCAPE_CANCEL_WINDOW_MS = 2_000;
export class InputController {
constructor(
@@ -149,6 +150,16 @@ export class InputController {
// (>= LEFT_DOUBLE_TAP_MAX_GAP_MS) starts a fresh sequence. See
// #detectLeftDoubleTap.
#leftTapCount = 0;
// Streaming turns use a two-step Esc: first press arms this token, second press
// within the window aborts the same live assistant turn. The token is a per-turn
// sentinel minted lazily on demand and reset on every `agent_start`/`agent_end`
// (see setupKeyHandlers), so it survives `message_start`/`message_update`
// transitions inside a single turn but cannot leak across turn boundaries.
#streamingEscapeTurnSentinel: object | undefined;
#streamingEscapeArmedToken: object | undefined;
#streamingEscapeArmedUntil = 0;
#streamingEscapeTimer: NodeJS.Timeout | undefined;
#streamingEscapeSessionSubscribed = false;
// Sequential index for `local://attachment-N` references created by large-paste and
// pasted-file attachments. Seeded from 0 and bumped past existing attachment files.
#attachmentCounter = 0;
@@ -198,8 +209,50 @@ export class InputController {
const unsubscribe = tinyTitleClient.onProgress(update);
}
#clearStreamingEscapeArm(): void {
this.#streamingEscapeArmedToken = undefined;
this.#streamingEscapeArmedUntil = 0;
if (this.#streamingEscapeTimer) {
clearTimeout(this.#streamingEscapeTimer);
this.#streamingEscapeTimer = undefined;
}
}
#handleStreamingEscape(): void {
if (!this.#streamingEscapeTurnSentinel) {
this.#streamingEscapeTurnSentinel = {};
}
const token = this.#streamingEscapeTurnSentinel;
const now = Date.now();
if (this.#streamingEscapeArmedToken === token && now <= this.#streamingEscapeArmedUntil) {
this.#clearStreamingEscapeArm();
void this.ctx.session.abort({ reason: USER_INTERRUPT_LABEL });
return;
}
this.#clearStreamingEscapeArm();
this.#streamingEscapeArmedToken = token;
this.#streamingEscapeArmedUntil = now + STREAMING_ESCAPE_CANCEL_WINDOW_MS;
this.#streamingEscapeTimer = setTimeout(() => {
if (this.#streamingEscapeArmedToken === token && Date.now() >= this.#streamingEscapeArmedUntil) {
this.#clearStreamingEscapeArm();
}
}, STREAMING_ESCAPE_CANCEL_WINDOW_MS);
this.#streamingEscapeTimer.unref?.();
this.ctx.showStatus("Press Esc again within 2s to cancel streaming.");
}
setupKeyHandlers(): void {
this.ctx.editor.setActionKeys("app.interrupt", this.ctx.keybindings.getKeys("app.interrupt"));
if (!this.#streamingEscapeSessionSubscribed && typeof this.ctx.session.subscribe === "function") {
this.#streamingEscapeSessionSubscribed = true;
this.ctx.session.subscribe(event => {
if (event.type === "agent_start" || event.type === "agent_end") {
this.#streamingEscapeTurnSentinel = undefined;
this.#clearStreamingEscapeArm();
}
});
}
if (!this.#focusedLeftTapListenerInstalled) {
this.#focusedLeftTapListenerInstalled = true;
this.ctx.ui.addInputListener(data => {
@@ -275,7 +328,7 @@ export class InputController {
if (this.ctx.loopModeEnabled) {
this.ctx.pauseLoop();
if (this.ctx.session.isStreaming) {
void this.ctx.session.abort({ reason: USER_INTERRUPT_LABEL });
this.#handleStreamingEscape();
} else {
this.ctx.cancelPendingSubmission();
}
@@ -326,12 +379,13 @@ export class InputController {
this.ctx.isPythonMode = false;
this.ctx.updateEditorBorderColor();
} else if (this.ctx.session.isStreaming) {
void this.ctx.session.abort({ reason: USER_INTERRUPT_LABEL });
this.#handleStreamingEscape();
} else if (this.ctx.editor.getText().trim()) {
// Esc with typed text clears the draft instead of (or before) any double-Esc action
this.ctx.editor.setText("");
this.ctx.ui.requestRender();
this.ctx.lastEscapeTime = 0;
this.#clearStreamingEscapeArm();
} else {
// Double-interrupt with empty editor triggers /tree, /branch, or nothing based on setting
const action = settings.get("doubleEscapeAction");
@@ -76,10 +76,12 @@ function createContext(): {
requestRender: Spy;
resetDisplay: Spy;
shutdown: Spy;
showStatus: Spy;
startPendingSubmission: StartPendingSubmissionSpy;
updatePendingMessagesDisplay: Spy;
};
inputListeners: Array<(data: string) => { consume?: boolean; data?: string } | undefined>;
sessionListeners: Array<(event: { type: string }) => void>;
} {
let editorText = "";
const abort = vi.fn();
@@ -93,7 +95,9 @@ function createContext(): {
const onInputCallback = vi.fn();
const requestRender = vi.fn();
const resetDisplay = vi.fn();
const showStatus = vi.fn();
const inputListeners: Array<(data: string) => { consume?: boolean; data?: string } | undefined> = [];
const sessionListeners: Array<(event: { type: string }) => void> = [];
const handleBtwCommand = vi.fn(async () => {});
const handleBtwEscape = vi.fn(() => true);
const hasActiveBtw = vi.fn(() => false);
@@ -158,6 +162,13 @@ function createContext(): {
clearQueue,
getQueuedMessages,
prompt,
subscribe: vi.fn((listener: (event: { type: string }) => void) => {
sessionListeners.push(listener);
return () => {
const index = sessionListeners.indexOf(listener);
if (index >= 0) sessionListeners.splice(index, 1);
};
}),
} as unknown as InteractiveModeContext["session"],
viewSession: {
isCompacting: false,
@@ -205,6 +216,7 @@ function createContext(): {
showSessionSelector: vi.fn(),
shutdown: vi.fn(async () => {}),
clearEditor: vi.fn(),
showStatus,
} as unknown as InteractiveModeContext;
return {
@@ -231,11 +243,13 @@ function createContext(): {
prompt,
requestRender,
resetDisplay,
showStatus,
shutdown: ctx.shutdown as Spy,
startPendingSubmission,
updatePendingMessagesDisplay,
},
inputListeners,
sessionListeners,
};
}
beforeEach(async () => {
@@ -407,7 +421,9 @@ describe("InputController escape behavior", () => {
expect(spies.abort).not.toHaveBeenCalled();
});
it("aborts streaming even when the working loader is no longer present", () => {
it("requires a second Esc within two seconds to abort streaming", () => {
const now = vi.spyOn(Date, "now");
now.mockReturnValue(1_000);
const { ctx, editor, spies } = createContext();
(ctx.session as { isStreaming: boolean }).isStreaming = true;
const controller = new InputController(ctx);
@@ -417,7 +433,99 @@ describe("InputController escape behavior", () => {
expect(spies.cancelPendingSubmission).not.toHaveBeenCalled();
expect(spies.clearQueue).not.toHaveBeenCalled();
expect(spies.abort).not.toHaveBeenCalled();
expect(spies.showStatus).toHaveBeenCalledWith("Press Esc again within 2s to cancel streaming.");
now.mockReturnValue(2_500);
editor.onEscape?.();
expect(spies.abort).toHaveBeenCalledTimes(1);
expect(spies.abort).toHaveBeenCalledWith({ reason: USER_INTERRUPT_LABEL });
});
it("expires the streaming Esc arm instead of aborting on a late second press", () => {
const now = vi.spyOn(Date, "now");
now.mockReturnValue(1_000);
const { ctx, editor, spies } = createContext();
(ctx.session as { isStreaming: boolean }).isStreaming = true;
const controller = new InputController(ctx);
controller.setupKeyHandlers();
editor.onEscape?.();
now.mockReturnValue(3_001);
editor.onEscape?.();
expect(spies.abort).not.toHaveBeenCalled();
expect(spies.showStatus).toHaveBeenCalledTimes(2);
});
it("preserves the streaming Esc arm when streamingComponent appears between presses", () => {
// Pre-`message_start`: first Esc arms on the per-turn sentinel. `message_start`
// then publishes `ctx.streamingComponent`; the second Esc must still abort the
// same live turn instead of re-arming on the new component reference.
const now = vi.spyOn(Date, "now");
now.mockReturnValue(1_000);
const { ctx, editor, spies } = createContext();
(ctx.session as { isStreaming: boolean }).isStreaming = true;
const controller = new InputController(ctx);
controller.setupKeyHandlers();
editor.onEscape?.();
(ctx as unknown as { streamingComponent: object }).streamingComponent = {};
now.mockReturnValue(1_500);
editor.onEscape?.();
expect(spies.abort).toHaveBeenCalledTimes(1);
expect(spies.abort).toHaveBeenCalledWith({ reason: USER_INTERRUPT_LABEL });
});
it("aborts on the second Esc even when ctx.streamingMessage was replaced by a delta in between", () => {
// `EventController` replaces `ctx.streamingMessage` with a fresh immutable
// snapshot on every `message_update`; the per-turn sentinel is unaffected so
// swapping the message must not invalidate the armed token.
const now = vi.spyOn(Date, "now");
now.mockReturnValue(1_000);
const { ctx, editor, spies } = createContext();
(ctx.session as { isStreaming: boolean }).isStreaming = true;
(ctx as unknown as { streamingComponent: object }).streamingComponent = {};
(ctx as unknown as { streamingMessage: object }).streamingMessage = { content: [] };
const controller = new InputController(ctx);
controller.setupKeyHandlers();
editor.onEscape?.();
(ctx as unknown as { streamingMessage: object }).streamingMessage = { content: ["delta"] };
now.mockReturnValue(1_500);
editor.onEscape?.();
expect(spies.abort).toHaveBeenCalledTimes(1);
expect(spies.abort).toHaveBeenCalledWith({ reason: USER_INTERRUPT_LABEL });
});
it("clears the streaming Esc arm when the current turn ends", () => {
const now = vi.spyOn(Date, "now");
now.mockReturnValue(1_000);
const { ctx, editor, spies, sessionListeners } = createContext();
(ctx.session as { isStreaming: boolean }).isStreaming = true;
const controller = new InputController(ctx);
controller.setupKeyHandlers();
// Fallback arm (no streamingMessage/streamingComponent yet — pre-message_start).
editor.onEscape?.();
expect(sessionListeners).toHaveLength(1);
// Turn 1 ends; a new turn starts. session.subscribe receives both transitions,
// either of which must invalidate the still-armed fallback token so it cannot
// fast-abort the new turn's first Esc.
for (const listener of sessionListeners) {
listener({ type: "agent_end" });
listener({ type: "agent_start" });
}
now.mockReturnValue(1_500);
editor.onEscape?.();
expect(spies.abort).not.toHaveBeenCalled();
expect(spies.showStatus).toHaveBeenCalledTimes(2);
});
it("returns focused subagent view to main on Esc instead of aborting", () => {
+4
View File
@@ -26,6 +26,10 @@
- Recognized Warp (`TERM_PROGRAM=WarpTerminal`) as a first-class terminal. Inline images now negotiate the Kitty graphics protocol on macOS/Linux (direct placement — Warp has no Unicode-placeholder support); the protocol is dropped on Windows, where Warp ships without Kitty support and the APC sequences would render as visible garbage. True color is enabled. OSC 8 hyperlinks stay off by default because Warp's renderer prints the escape as literal text rather than a clickable link (opt in with `PI_FORCE_HYPERLINKS=1` once Warp lands real support), and synchronized output remains gated on the runtime DECRQM probe ([#3471](https://github.com/can1357/oh-my-pi/issues/3471)).
### Fixed
- Fixed ordinary render scheduling to yield behind already-queued terminal input, preventing delayed Esc delivery during heavy streaming paints ([#3493](https://github.com/can1357/oh-my-pi/issues/3493)).
## [16.1.19] - 2026-06-25
### Fixed
+1 -1
View File
@@ -110,7 +110,7 @@ export interface TUIStartOptions {
const DEFAULT_RENDER_SCHEDULER: RenderScheduler = {
now: () => performance.now(),
scheduleImmediate: callback => {
process.nextTick(callback);
setImmediate(callback);
},
scheduleRender: (callback, delayMs) => {
const timer = setTimeout(callback, delayMs);
@@ -0,0 +1,74 @@
import { describe, expect, it } from "bun:test";
import { type Component, type RenderTimer, TUI } from "@oh-my-pi/pi-tui";
import { VirtualTerminal } from "./virtual-terminal";
class InputProbe implements Component {
constructor(private readonly events: string[]) {}
invalidate(): void {}
render(_width: number): readonly string[] {
this.events.push("render");
return ["probe"];
}
handleInput(_data: string): void {
this.events.push("input");
}
}
class DeferredRenderScheduler {
nowMs = 0;
readonly immediates: Array<() => void> = [];
readonly timers: Array<{ callback: () => void; canceled: boolean }> = [];
now(): number {
return this.nowMs;
}
scheduleImmediate(callback: () => void): void {
this.immediates.push(callback);
}
scheduleRender(callback: () => void, _delayMs: number): RenderTimer {
const timer = { callback, canceled: false };
this.timers.push(timer);
return {
cancel: () => {
timer.canceled = true;
},
};
}
}
describe("TUI input/render scheduling", () => {
it("can process terminal input before a deferred ordinary repaint", () => {
const term = new VirtualTerminal(20, 4);
const scheduler = new DeferredRenderScheduler();
const events: string[] = [];
const probe = new InputProbe(events);
const tui = new TUI(term, undefined, { renderScheduler: scheduler });
tui.addChild(probe);
tui.setFocus(probe);
try {
tui.start();
scheduler.immediates.shift()?.();
const initialTimer = scheduler.timers.shift();
if (initialTimer && !initialTimer.canceled) initialTimer.callback();
events.length = 0;
scheduler.nowMs = 100;
tui.requestRender();
term.sendInput("x");
scheduler.immediates.shift()?.();
const repaintTimer = scheduler.timers.shift();
if (repaintTimer && !repaintTimer.canceled) repaintTimer.callback();
expect(events[0]).toBe("input");
expect(events).toContain("render");
} finally {
tui.stop();
}
});
});