diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index 26a33507c..27e90c97c 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -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 diff --git a/packages/coding-agent/src/modes/controllers/input-controller.ts b/packages/coding-agent/src/modes/controllers/input-controller.ts index a84221b24..55c70adaa 100644 --- a/packages/coding-agent/src/modes/controllers/input-controller.ts +++ b/packages/coding-agent/src/modes/controllers/input-controller.ts @@ -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"); diff --git a/packages/coding-agent/test/input-controller-escape.test.ts b/packages/coding-agent/test/input-controller-escape.test.ts index ac29fdb2b..6783058f1 100644 --- a/packages/coding-agent/test/input-controller-escape.test.ts +++ b/packages/coding-agent/test/input-controller-escape.test.ts @@ -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", () => { diff --git a/packages/tui/CHANGELOG.md b/packages/tui/CHANGELOG.md index 5b8d1f73c..e72840df0 100644 --- a/packages/tui/CHANGELOG.md +++ b/packages/tui/CHANGELOG.md @@ -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 diff --git a/packages/tui/src/tui.ts b/packages/tui/src/tui.ts index 643bb2bf3..706056e7b 100644 --- a/packages/tui/src/tui.ts +++ b/packages/tui/src/tui.ts @@ -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); diff --git a/packages/tui/test/input-render-scheduling.test.ts b/packages/tui/test/input-render-scheduling.test.ts new file mode 100644 index 000000000..430535089 --- /dev/null +++ b/packages/tui/test/input-render-scheduling.test.ts @@ -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(); + } + }); +});