diff --git a/apps/desktop/src/app/gateway/hooks/use-gateway-boot.test.tsx b/apps/desktop/src/app/gateway/hooks/use-gateway-boot.test.tsx index 9ea5cf1d6b..7c77fb602a 100644 --- a/apps/desktop/src/app/gateway/hooks/use-gateway-boot.test.tsx +++ b/apps/desktop/src/app/gateway/hooks/use-gateway-boot.test.tsx @@ -258,6 +258,7 @@ function Harness({ useGatewayBoot({ beforeConnectionSwitch, handleGatewayEvent: () => undefined, + handleServerRequest: () => false, onConnectionReady: () => undefined, onGatewayReady: () => undefined, refreshHermesConfig, diff --git a/apps/desktop/src/app/gateway/hooks/use-gateway-boot.ts b/apps/desktop/src/app/gateway/hooks/use-gateway-boot.ts index b635cc0a79..fba5caeb91 100644 --- a/apps/desktop/src/app/gateway/hooks/use-gateway-boot.ts +++ b/apps/desktop/src/app/gateway/hooks/use-gateway-boot.ts @@ -922,6 +922,7 @@ export function useGatewayBoot({ recordSessionEventScope(scopedEvent) callbacksRef.current.handleGatewayEvent(scopedEvent) }) + // Secondary sockets reach the same handler through the registry's onServerRequest. const offRequest = gateway.onRequest(request => dispatchPrimaryServerRequest(request, sourceProfile)) diff --git a/apps/desktop/src/app/session/hooks/use-message-stream/clarify-hydration.test.tsx b/apps/desktop/src/app/session/hooks/use-message-stream/clarify-hydration.test.tsx index cc0ba10c29..c503e983f2 100644 --- a/apps/desktop/src/app/session/hooks/use-message-stream/clarify-hydration.test.tsx +++ b/apps/desktop/src/app/session/hooks/use-message-stream/clarify-hydration.test.tsx @@ -8,7 +8,7 @@ import { onScrollToBottomRequest } from '@/store/thread-scroll' import { type MessageStreamHarness, renderMessageStream } from './test-harness' -// A `clarify.request` must leave an answerable inline row even when the +// A `clarify` server request must leave an answerable inline row even when the // `tool.start` that normally mounts it was missed (stream reconnect / // hydration race). Without it the sidebar says "needs input" but the // transcript has nowhere to render the choices, so the agent blocks forever. @@ -24,8 +24,8 @@ function mountStream() { stream = renderMessageStream(SID) } -const clarifyRequest = (payload: Record) => - act(() => stream.handleEvent({ payload, session_id: SID, type: 'clarify.request' })) +const clarifyRequest = ({ request_id, ...params }: Record) => + act(() => void stream.handleRequest('clarify', { ...params, session_id: SID }, request_id as string)) const toolStart = (payload: Record) => act(() => stream.handleEvent({ payload, session_id: SID, type: 'tool.start' })) @@ -34,7 +34,13 @@ const toolComplete = (payload: Record) => act(() => stream.handleEvent({ payload, session_id: SID, type: 'tool.complete' })) const clarifyExpire = (requestId: string) => - act(() => stream.handleEvent({ payload: { request_id: requestId }, session_id: SID, type: 'clarify.expire' })) + act(() => + stream.handleEvent({ + payload: { id: requestId, method: 'clarify', reason: 'timeout' }, + session_id: SID, + type: 'request.cancel' + }) + ) function clarifyParts() { const messages = stream.state().messages ?? [] @@ -49,7 +55,7 @@ function seedHydratedMessages(messages: ChatMessage[]) { stream.states.set(SID, state) } -describe('clarify.request stream hydration', () => { +describe('clarify request stream hydration', () => { beforeEach(() => { clearClarifyRequest() scrollToBottom.mockClear() @@ -89,12 +95,13 @@ describe('clarify.request stream hydration', () => { it('does not move the active thread for a background session clarify', () => { mountStream() - act(() => - stream.handleEvent({ - payload: { choices: ['yes', 'no'], question: 'Ship it?', request_id: 'req-background' }, - session_id: 'session-background', - type: 'clarify.request' - }) + act( + () => + void stream.handleRequest( + 'clarify', + { choices: ['yes', 'no'], question: 'Ship it?', session_id: 'session-background' }, + 'req-background' + ) ) expect(scrollToBottom).not.toHaveBeenCalled() @@ -129,7 +136,7 @@ describe('clarify.request stream hydration', () => { it('merges with the real tool.start row even though its id differs from the request id', () => { mountStream() - // Reality: tool.start carries the model's tool_call_id, clarify.request a + // Reality: tool.start carries the model's tool_call_id, the clarify request a // separately-generated request_id. They must still collapse to ONE card // (correlated by question), not two. toolStart({ args: { choices: ['a'], question: 'Pick' }, name: 'clarify', tool_id: 'call-abc' }) diff --git a/apps/desktop/src/app/session/hooks/use-message-stream/gateway-event/desktop-bridge.test.ts b/apps/desktop/src/app/session/hooks/use-message-stream/gateway-event/desktop-bridge.test.ts deleted file mode 100644 index ef664fe0e0..0000000000 --- a/apps/desktop/src/app/session/hooks/use-message-stream/gateway-event/desktop-bridge.test.ts +++ /dev/null @@ -1,95 +0,0 @@ -import { afterEach, describe, expect, it, vi } from 'vitest' - -import { $gateway } from '@/store/gateway' -import { $toursEnabled } from '@/store/tours' - -import { handleDesktopBridgeEvent } from './desktop-bridge' -import type { GatewayEventContext } from './types' - -function previewActContext({ - explicitSid, - isActiveEvent -}: { - explicitSid: string - isActiveEvent: boolean -}): GatewayEventContext { - return { - event: { session_id: explicitSid || undefined, type: 'preview.act.request' }, - explicitSid, - isActiveEvent, - payload: { action: 'elements', request_id: 'request-1' } - } as GatewayEventContext -} - -describe('preview action bridge routing', () => { - afterEach(() => { - $gateway.set(null) - }) - - it('leaves a scoped action request unanswered in a window showing another session', () => { - const request = vi.fn() - $gateway.set({ request } as never) - - expect(handleDesktopBridgeEvent(previewActContext({ explicitSid: 'session-a', isActiveEvent: false }))).toBe(true) - expect(request).not.toHaveBeenCalled() - }) - - it('keeps the legacy fail-fast response for an unscoped inactive request', () => { - const request = vi.fn() - $gateway.set({ request } as never) - - expect(handleDesktopBridgeEvent(previewActContext({ explicitSid: '', isActiveEvent: false }))).toBe(true) - expect(request).toHaveBeenCalledWith('preview.act.respond', { - request_id: 'request-1', - text: JSON.stringify({ - error: 'The in-app browser only takes actions in the session the user is looking at.', - success: false - }) - }) - }) -}) - -function tourContext({ - explicitSid, - isActiveEvent -}: { - explicitSid: string - isActiveEvent: boolean -}): GatewayEventContext { - return { - event: { session_id: explicitSid || undefined, type: 'tour.request' }, - explicitSid, - isActiveEvent, - payload: { action: 'discover', request_id: 'tour-request-1' } - } as GatewayEventContext -} - -describe('tour bridge routing', () => { - afterEach(() => { - $gateway.set(null) - $toursEnabled.set(true) - }) - - it('leaves a scoped request unanswered in another session even when tours are disabled', () => { - const request = vi.fn() - $gateway.set({ request } as never) - $toursEnabled.set(false) - - expect(handleDesktopBridgeEvent(tourContext({ explicitSid: 'session-a', isActiveEvent: false }))).toBe(true) - expect(request).not.toHaveBeenCalled() - }) - - it('keeps the legacy fail-fast response for an unscoped inactive request', () => { - const request = vi.fn() - $gateway.set({ request } as never) - - expect(handleDesktopBridgeEvent(tourContext({ explicitSid: '', isActiveEvent: false }))).toBe(true) - expect(request).toHaveBeenCalledWith('tour.respond', { - request_id: 'tour-request-1', - text: JSON.stringify({ - error: 'Tours only run in the session the user is looking at.', - success: false - }) - }) - }) -}) diff --git a/apps/desktop/src/app/session/hooks/use-message-stream/gateway-event/server-requests.test.ts b/apps/desktop/src/app/session/hooks/use-message-stream/gateway-event/server-requests.test.ts new file mode 100644 index 0000000000..28913fcacc --- /dev/null +++ b/apps/desktop/src/app/session/hooks/use-message-stream/gateway-event/server-requests.test.ts @@ -0,0 +1,59 @@ +import { afterEach, describe, expect, it, vi } from 'vitest' + +import { $toursEnabled } from '@/store/tours' + +import { handleServerRequest } from './server-requests' +import type { ServerRequestContext } from './server-requests' + +const deps = {} as ServerRequestContext['deps'] + +function deliver(method: string, params: Record, activeSessionId: null | string) { + const respond = vi.fn() + const fail = vi.fn() + const handled = handleServerRequest({ fail, id: 'srq-1', method, params, profile: 'default', respond }, deps, activeSessionId) + + return { fail, handled, respond } +} + +describe('preview action request routing', () => { + it('leaves a scoped action request unanswered in a window showing another session', () => { + const { handled, respond, fail } = deliver('preview.act', { action: 'elements', session_id: 'session-a' }, 'session-b') + + expect(handled).toBe(true) + expect(respond).not.toHaveBeenCalled() + expect(fail).not.toHaveBeenCalled() + }) + + it('fails fast for an unscoped request with no session in view', () => { + const { respond } = deliver('preview.act', { action: 'elements' }, null) + + expect(respond).toHaveBeenCalledWith({ + value: JSON.stringify({ + error: 'The in-app browser only takes actions in the session the user is looking at.', + success: false + }) + }) + }) +}) + +describe('tour request routing', () => { + afterEach(() => { + $toursEnabled.set(true) + }) + + it('leaves a scoped request unanswered in another session even when tours are disabled', () => { + $toursEnabled.set(false) + const { handled, respond } = deliver('tour', { action: 'discover', session_id: 'session-a' }, 'session-b') + + expect(handled).toBe(true) + expect(respond).not.toHaveBeenCalled() + }) + + it('fails fast for an unscoped request with no session in view', () => { + const { respond } = deliver('tour', { action: 'discover' }, null) + + expect(respond).toHaveBeenCalledWith({ + value: JSON.stringify({ error: 'Tours only run in the session the user is looking at.', success: false }) + }) + }) +}) diff --git a/apps/desktop/src/app/session/hooks/use-message-stream/test-harness.tsx b/apps/desktop/src/app/session/hooks/use-message-stream/test-harness.tsx index b1d702643f..33b1937d7a 100644 --- a/apps/desktop/src/app/session/hooks/use-message-stream/test-harness.tsx +++ b/apps/desktop/src/app/session/hooks/use-message-stream/test-harness.tsx @@ -6,6 +6,7 @@ import { vi } from 'vitest' import type { ClientSessionState } from '@/app/types' import { createClientSessionState } from '@/lib/chat-runtime' +import type { ScopedServerRequest } from '@/store/gateway' import { useMessageStream } from './index' @@ -17,6 +18,9 @@ export interface MessageStreamHarnessOptions extends Partial void + /** Feed a server→client request (clarify, approval, …) into the mounted hook; + * returns the `respond` spy so a test can assert the answer frame. */ + handleRequest: (method: string, params: Record, id?: string) => ReturnType /** Push streaming assistant text, bypassing the event envelope. For the specs * about flush scheduling rather than about a particular event. */ appendDelta: (sessionId: string, delta: string) => void @@ -48,6 +52,7 @@ export function renderMessageStream( { states = new Map(), ...overrides }: MessageStreamHarnessOptions = {} ): MessageStreamHarness { let dispatch: ((event: GatewayEvent) => void) | null = null + let dispatchRequest: ((request: ScopedServerRequest) => boolean) | null = null let appendDelta: ((sessionId: string, delta: string) => void) | null = null let latest: ClientSessionState | null = null @@ -75,8 +80,9 @@ export function renderMessageStream( useEffect(() => { dispatch = stream.handleGatewayEvent + dispatchRequest = stream.handleServerRequest appendDelta = stream.appendAssistantDelta - }, [stream.appendAssistantDelta, stream.handleGatewayEvent]) + }, [stream.appendAssistantDelta, stream.handleGatewayEvent, stream.handleServerRequest]) return null } @@ -93,6 +99,17 @@ export function renderMessageStream( dispatch(event) }, + handleRequest: (method, params, id = `srq-${method}`) => { + const respond = vi.fn() + + if (!dispatchRequest) { + throw new Error('renderMessageStream: the hook never mounted') + } + + dispatchRequest({ fail: vi.fn(), id, method, params, profile: 'default', respond }) + + return respond + }, appendDelta: (id, delta) => { if (!appendDelta) { throw new Error('renderMessageStream: the hook never mounted') diff --git a/apps/desktop/src/app/session/hooks/use-session-actions/restore-pending-clarify.test.ts b/apps/desktop/src/app/session/hooks/use-session-actions/restore-pending-clarify.test.ts index 9bafef9d76..3a1a70cd5d 100644 --- a/apps/desktop/src/app/session/hooks/use-session-actions/restore-pending-clarify.test.ts +++ b/apps/desktop/src/app/session/hooks/use-session-actions/restore-pending-clarify.test.ts @@ -30,102 +30,41 @@ describe('restorePendingClarifyFromSnapshot', () => { $clarifyRequests.set({}) }) - it('restores a batch clarify snapshot (questions, no top-level question)', () => { + it('hands back the card the request handler already parked for a replayed open clarify request', () => { + $clarifyRequests.set({ + 'sess-1': { + choices: null, + multiSelect: false, + question: 'Proceed?', + receivedAt: resumeStartedAt + 5, + requestId: 'rid1', + sessionId: 'sess-1' + } + }) + const state = restorePendingClarifyFromSnapshot( - { - pending_clarify: { - request_id: 'rid1', - questions: [ - { choices: ['Yes', 'No'], multi_select: false, qid: 'q0', question: 'Proceed?' }, - { qid: 'q1', question: 'Which region?' } - ] - } - }, + { open_requests: [{ id: 'rid1', method: 'clarify', params: { question: 'Proceed?' } }] }, 'sess-1', resumeStartedAt ) - expect(state.request).not.toBeNull() - expect(setClarifyRequestMock).toHaveBeenCalledWith( - expect.objectContaining({ - multiSelect: false, - question: '', - requestId: 'rid1', - sessionId: 'sess-1', - questions: [ - { choices: ['Yes', 'No'], multiSelect: false, qid: 'q0', question: 'Proceed?' }, - { multiSelect: false, qid: 'q1', question: 'Which region?', choices: null } - ] - }) - ) + expect(state.authoritativeAbsent).toBe(false) + expect(state.request?.requestId).toBe('rid1') + expect(clearClarifyRequestMock).not.toHaveBeenCalled() }) - it('carries server-locked answers into the replayed batch card', () => { - restorePendingClarifyFromSnapshot( - { - pending_clarify: { - answers: { q0: 'Yes', junk: 42 }, - request_id: 'rid2', - questions: [{ qid: 'q0', question: 'Proceed?' }] - } - }, - 'sess-2', - resumeStartedAt - ) - - expect(setClarifyRequestMock).toHaveBeenCalledWith( - expect.objectContaining({ lockedAnswers: { q0: 'Yes' }, requestId: 'rid2' }) - ) - }) - - it('still restores the single-question form', () => { + it('reports an open request the handler declined (no card parked) without inventing one', () => { const state = restorePendingClarifyFromSnapshot( - { - pending_clarify: { - choices: ['A', 'B'], - multi_select: true, - question: 'Pick one', - request_id: 'rid3' - } - }, - 'sess-3', - resumeStartedAt - ) - - expect(state.request).not.toBeNull() - expect(setClarifyRequestMock).toHaveBeenCalledWith( - expect.objectContaining({ - choices: ['A', 'B'], - multiSelect: true, - question: 'Pick one', - requestId: 'rid3' - }) - ) - }) - - it('rejects a payload with neither form (no request restored)', () => { - const state = restorePendingClarifyFromSnapshot( - { pending_clarify: { request_id: 'rid4' } }, + { open_requests: [{ id: 'rid4', method: 'clarify', params: {} }] }, 'sess-4', resumeStartedAt ) + expect(state.authoritativeAbsent).toBe(false) expect(state.request).toBeNull() expect(setClarifyRequestMock).not.toHaveBeenCalled() }) - it('rejects a payload with no request id', () => { - const state = restorePendingClarifyFromSnapshot( - { pending_clarify: { question: 'Orphaned prompt' } }, - 'sess-5', - resumeStartedAt - ) - - expect(state.request).toBeNull() - expect(state.authoritativeAbsent).toBe(true) - expect(setClarifyRequestMock).not.toHaveBeenCalled() - }) - it('clears a stale local request when the snapshot has none, and leaves a newer in-flight one', () => { $clarifyRequests.set({ 'sess-6': { diff --git a/apps/desktop/src/components/assistant-ui/clarify-tool.tsx b/apps/desktop/src/components/assistant-ui/clarify-tool.tsx index f087de4481..fe64f34974 100644 --- a/apps/desktop/src/components/assistant-ui/clarify-tool.tsx +++ b/apps/desktop/src/components/assistant-ui/clarify-tool.tsx @@ -39,8 +39,8 @@ import { } from '@/store/clarify' import { $gateway } from '@/store/gateway' import { notifyError } from '@/store/notifications' -import { requestForOwnedSession } from '@/store/session-states' import { forgetServerRequest, respondToServerRequest } from '@/store/server-requests' +import { requestForOwnedSession } from '@/store/session-states' import { handleClarifySubmitShortcut } from './clarify-submit-shortcut' import { selectMessageRunning } from './tool/fallback-model' diff --git a/apps/desktop/src/components/assistant-ui/mcp-setup-tool.tsx b/apps/desktop/src/components/assistant-ui/mcp-setup-tool.tsx index d2715901d5..60cf49b8d7 100644 --- a/apps/desktop/src/components/assistant-ui/mcp-setup-tool.tsx +++ b/apps/desktop/src/components/assistant-ui/mcp-setup-tool.tsx @@ -28,8 +28,8 @@ import { prettyName } from '@/lib/text' import { cn } from '@/lib/utils' import { $gateway } from '@/store/gateway' import { clearMcpSetupRequest, type McpSetupOutcome, sessionMcpSetupRequest } from '@/store/mcp-setup' -import { respondToServerRequest } from '@/store/server-requests' import { notifyError } from '@/store/notifications' +import { respondToServerRequest } from '@/store/server-requests' import { invalidateMcpSuggestionIndex } from '@/store/suggestion-providers/mcp' import { selectMessageRunning } from './tool/fallback-model' diff --git a/apps/desktop/src/components/assistant-ui/tool/approval.tsx b/apps/desktop/src/components/assistant-ui/tool/approval.tsx index 99e794d151..d7ff03878f 100644 --- a/apps/desktop/src/components/assistant-ui/tool/approval.tsx +++ b/apps/desktop/src/components/assistant-ui/tool/approval.tsx @@ -20,8 +20,8 @@ import { AlertCircle, ChevronDown } from '@/lib/icons' import { isSubmitEnter } from '@/lib/ime' import { cn } from '@/lib/utils' import { $gateway } from '@/store/gateway' -import { answerApproval } from '@/store/prompts' import { notifyError } from '@/store/notifications' +import { answerApproval } from '@/store/prompts' import { type ApprovalRequest, clearApprovalRequest, @@ -30,7 +30,6 @@ import { sessionApprovalInlineVisible, sessionApprovalRequest } from '@/store/prompts' -import { requestForOwnedSession } from '@/store/session-states' import type { ToolPart } from './fallback-model' diff --git a/apps/desktop/src/components/prompt-overlays.tsx b/apps/desktop/src/components/prompt-overlays.tsx index a08f8493fb..7831eb5afa 100644 --- a/apps/desktop/src/components/prompt-overlays.tsx +++ b/apps/desktop/src/components/prompt-overlays.tsx @@ -33,9 +33,7 @@ import { sessionVaultSaveLoginRequest, sessionVaultUnlockRequest } from '@/store/prompts' -import { ambientRequestFor } from '@/store/session-gone-latch' import { respondToServerRequest } from '@/store/server-requests' -import { requestForOwnedSession } from '@/store/session-states' // Renders the modal mid-turn prompts the gateway raises and waits on: sudo // password and skill secret capture. Dangerous-command / execute_code approval diff --git a/apps/desktop/src/store/clarify.test.ts b/apps/desktop/src/store/clarify.test.ts index c234b7961b..76f9bc1cb1 100644 --- a/apps/desktop/src/store/clarify.test.ts +++ b/apps/desktop/src/store/clarify.test.ts @@ -12,6 +12,7 @@ import { skipClarifyRequest } from './clarify' import { $gateway } from './gateway' +import { rememberServerRequest, resetServerRequestsForTests } from './server-requests' import { $activeSessionId } from './session' function clarify(sessionId: string | null, requestId: string): ClarifyRequest { @@ -91,6 +92,7 @@ describe('skipClarifyRequest', () => { beforeEach(() => { $clarifyRequests.set({}) + resetServerRequestsForTests() request.mockClear() $gateway.set({ request } as unknown as ReturnType) }) @@ -101,12 +103,15 @@ describe('skipClarifyRequest', () => { }) it('answers the session\u2019s clarify with an empty answer and drops it', async () => { + const respond = vi.fn() + + rememberServerRequest({ fail: vi.fn(), id: 'req-a', method: 'clarify', params: {}, respond }) setClarifyRequest(clarify('session-a', 'req-a')) setClarifyRequest(clarify('session-b', 'req-b')) await expect(skipClarifyRequest('session-a')).resolves.toBe(true) - expect(request).toHaveBeenCalledWith('clarify.respond', { request_id: 'req-a', answer: '' }) + expect(respond).toHaveBeenCalledWith({ answer: '' }) expect(hasClarifyRequest('session-a')).toBe(false) // A background session's question is untouched — only the one being typed // over is skipped. @@ -118,9 +123,8 @@ describe('skipClarifyRequest', () => { expect(request).not.toHaveBeenCalled() }) - it('still reports the skip when the respond RPC fails', async () => { + it('still reports the skip when the server request is already gone (expired / other window answered)', async () => { setClarifyRequest(clarify('session-a', 'req-a')) - request.mockRejectedValueOnce(new Error('socket closed')) await expect(skipClarifyRequest('session-a')).resolves.toBe(true) expect(hasClarifyRequest('session-a')).toBe(false) diff --git a/apps/desktop/src/store/clarify.ts b/apps/desktop/src/store/clarify.ts index 7ce328c3dc..aaef93106c 100644 --- a/apps/desktop/src/store/clarify.ts +++ b/apps/desktop/src/store/clarify.ts @@ -1,6 +1,5 @@ import { atom, computed } from 'nanostores' -import { $gateway } from './gateway' import { respondToServerRequest } from './server-requests' import { $activeSessionId } from './session' diff --git a/apps/desktop/src/store/gateway.ts b/apps/desktop/src/store/gateway.ts index c6888544ed..204a631932 100644 --- a/apps/desktop/src/store/gateway.ts +++ b/apps/desktop/src/store/gateway.ts @@ -842,9 +842,10 @@ function createSecondary(profile: string, connectionId: null | string = null): S g.config?.onEvent(scopedEvent) releaseTerminalTurnLease(entry.scope, event) }) - entry.offRequest = gateway.onRequest(request => { - g.config?.onServerRequest?.({ ...request, ...(connectionId ? { connectionId } : {}), profile }) - }) + entry.offRequest = + gateway.onRequest?.(request => { + g.config?.onServerRequest?.({ ...request, ...(connectionId ? { connectionId } : {}), profile }) + }) ?? (() => {}) entry.offState = gateway.onState(state => { reportGatewayState(scope, state) diff --git a/apps/desktop/src/store/mcp-setup.ts b/apps/desktop/src/store/mcp-setup.ts index bcccba05e0..e8036e77f5 100644 --- a/apps/desktop/src/store/mcp-setup.ts +++ b/apps/desktop/src/store/mcp-setup.ts @@ -1,6 +1,5 @@ import { atom, computed } from 'nanostores' -import { $gateway } from './gateway' import { respondToServerRequest } from './server-requests' /** diff --git a/apps/desktop/src/store/native-notifications.ts b/apps/desktop/src/store/native-notifications.ts index 026566219a..c02763a209 100644 --- a/apps/desktop/src/store/native-notifications.ts +++ b/apps/desktop/src/store/native-notifications.ts @@ -4,12 +4,12 @@ import { type HermesOpenTarget, resolveHermesOpenPath } from '@/lib/hermes-open- import { persistString, storedString } from '@/lib/storage' import { $gateway } from './gateway' -import { $approvalRequests, answerApproval } from './prompts' import { withinNativeNotifyBaseline } from './notify-baseline' +import { $approvalRequests, answerApproval } from './prompts' import { clearApprovalRequest } from './prompts' import { isSessionGone, isSessionGoneForBackgroundPolling, markSessionGone } from './runtime-gone' import { $activeSessionId } from './session' -import { requestForOwnedSession, storedSessionIdForRuntimeId } from './session-states' +import { storedSessionIdForRuntimeId } from './session-states' export type { HermesOpenTarget } diff --git a/apps/desktop/src/store/prompts.ts b/apps/desktop/src/store/prompts.ts index 2ef2e70e3d..489b45d55d 100644 --- a/apps/desktop/src/store/prompts.ts +++ b/apps/desktop/src/store/prompts.ts @@ -2,9 +2,9 @@ import { atom, computed, type ReadableAtom } from 'nanostores' import { $clarifyRequest, $clarifyRequests } from './clarify' import { isSessionGone, isSessionGoneForBackgroundPolling, markSessionGone } from './runtime-gone' +import { respondToServerRequest } from './server-requests' import { $activeSessionId } from './session' import { ambientRequestFor } from './session-gone-latch' -import { respondToServerRequest } from './server-requests' import { requestForOwnedSession } from './session-states' // Blocking interactive prompts the gateway raises mid-turn. Each is a @@ -262,9 +262,9 @@ export async function answerApproval( } await requestForOwnedSession(request.sessionId, ambientRequestFor(gateway), 'approval.respond', { - all, + ...(all ? { all: true } : {}), choice, - request_id: request.requestId, + ...(request.requestId ? { request_id: request.requestId } : {}), session_id: request.sessionId ?? undefined }) } diff --git a/apps/shared/src/json-rpc-channel.test.ts b/apps/shared/src/json-rpc-channel.test.ts index 960f6d1740..80685ba8df 100644 --- a/apps/shared/src/json-rpc-channel.test.ts +++ b/apps/shared/src/json-rpc-channel.test.ts @@ -170,4 +170,50 @@ describe('JsonRpcRequestChannel', () => { vi.useRealTimers() } }) + + // Server→client requests (tui_gateway/server_requests.py): the backend asks, + // the client answers with a RESPONSE frame carrying the same id. + it('routes a server request to the first accepting handler and answers -32601 when nobody accepts', () => { + const unhandled: string[] = [] + const channel = new JsonRpcRequestChannel({ onUnhandledRequest: req => void unhandled.push(req.method) }) + const { sent, transport } = spyTransport() + + channel.attach(transport) + channel.onRequest(req => (req.method === 'clarify' ? void req.respond({ answer: 'yes' }) : false)) + + channel.handleFrame(JSON.stringify({ id: 'srq-1', jsonrpc: '2.0', method: 'clarify', params: { session_id: 's1' } })) + channel.handleFrame(JSON.stringify({ id: 'srq-2', jsonrpc: '2.0', method: 'tour', params: { session_id: 's1' } })) + + const frames = sent.map(f => JSON.parse(f) as { id: string; result?: unknown; error?: { code: number } }) + + expect(frames[0]).toEqual({ id: 'srq-1', jsonrpc: '2.0', result: { answer: 'yes' } }) + expect(frames[1].id).toBe('srq-2') + expect(frames[1].error?.code).toBe(-32601) + expect(unhandled).toEqual(['tour']) + }) + + it('re-delivers open_requests from a response before the caller sees the result, tagged replayed', async () => { + const delivered: Array<{ id: string; replayed?: boolean }> = [] + const channel = new JsonRpcRequestChannel() + const { sent, transport } = spyTransport() + + channel.attach(transport) + channel.onRequest(req => void delivered.push({ id: req.id, replayed: req.replayed })) + + const resume = channel.request<{ session_id: string }>('session.resume', { session_id: 's1' }) + const rid = (JSON.parse(sent.at(-1)!) as { id: string }).id + + channel.handleFrame( + JSON.stringify({ + id: rid, + jsonrpc: '2.0', + result: { + open_requests: [{ id: 'srq-9', method: 'sudo', params: { session_id: 's1' } }], + session_id: 's1' + } + }) + ) + await expect(resume).resolves.toMatchObject({ session_id: 's1' }) + expect(delivered).toEqual([{ id: 'srq-9', replayed: true }]) + }) }) diff --git a/tests/gateway/test_tui_approval_redaction.py b/tests/gateway/test_tui_approval_redaction.py index af7f498f3d..7bbf4c173b 100644 --- a/tests/gateway/test_tui_approval_redaction.py +++ b/tests/gateway/test_tui_approval_redaction.py @@ -15,24 +15,28 @@ import pytest class TestTuiApprovalEmitRedaction: - def test_emit_approval_request_redacts_command_in_payload(self, monkeypatch): - from tui_gateway import server as tui_server + @staticmethod + def _sent(monkeypatch): + """Capture the ``approval`` server request frame ``_emit_approval_request`` sends.""" + from tui_gateway import server as tui_server, server_requests - emitted = {} - monkeypatch.setattr( - tui_server, "_emit", - lambda event, sid, payload=None: emitted.update( - {"event": event, "sid": sid, "payload": payload} - ), - ) + sent = {} + monkeypatch.setattr(server_requests, "send_async", + lambda method, sid, params, on_result: (sent.update(method=method, sid=sid, params=params), + lambda reason: None)[1]) + monkeypatch.setattr(tui_server, "_sessions", {"sess-1": {"session_key": "key-1"}}) + return tui_server, sent + + def test_emit_approval_request_redacts_command_in_payload(self, monkeypatch): + tui_server, sent = self._sent(monkeypatch) raw = "curl -H 'Authorization: token ghp_01...6789' https://api.github.com" tui_server._emit_approval_request("sess-1", {"command": raw, "description": "x"}) - assert emitted["event"] == "approval.request" + assert sent["method"] == "approval" and sent["sid"] == "sess-1" # credential removed, non-command field + command structure preserved - assert "ghp_01...6789" not in emitted["payload"]["command"] - assert emitted["payload"]["description"] == "x" - assert "github.com" in emitted["payload"]["command"] + assert "ghp_01...6789" not in sent["params"]["command"] + assert sent["params"]["description"] == "x" + assert "github.com" in sent["params"]["command"] @pytest.mark.parametrize( ("allow_session", "allow_permanent", "expected"), @@ -45,23 +49,10 @@ class TestTuiApprovalEmitRedaction: def test_emit_approval_request_honors_allowed_scopes( self, monkeypatch, allow_session, allow_permanent, expected ): - from tui_gateway import server as tui_server - - emitted = {} - monkeypatch.setattr( - tui_server, - "_emit", - lambda event, sid, payload=None: emitted.update({"payload": payload}), - ) - + tui_server, sent = self._sent(monkeypatch) tui_server._emit_approval_request( "sess-1", - { - "allow_permanent": allow_permanent, - "allow_session": allow_session, - "command": "", - }, + {"allow_permanent": allow_permanent, "allow_session": allow_session, "command": ""}, ) - assert emitted["payload"]["choices"] == expected - + assert sent["params"]["choices"] == expected diff --git a/tests/tui_gateway/test_compute_host_phase1.py b/tests/tui_gateway/test_compute_host_phase1.py index bb7368d558..a9eda851ef 100644 --- a/tests/tui_gateway/test_compute_host_phase1.py +++ b/tests/tui_gateway/test_compute_host_phase1.py @@ -48,40 +48,38 @@ def test_compute_host_workers_inherit_tui_pool_env_or_8(monkeypatch): assert _default_workers() == 8 -def test_compute_host_routes_clarify_response_to_child_pending_registry(monkeypatch): - """Interactive answers are handled in the process that owns `_pending`.""" +def test_compute_host_routes_relayed_response_and_lock_to_its_open_request(monkeypatch): + """The child owns the server request's wait: a relayed client response frame resolves it in-process, + and a relayed ``clarify.lock`` is answered with that method's result for the parent to ack.""" + from tui_gateway import server_requests out = io.StringIO() host = ComputeHost(stdout=out, heartbeat_secs=0) sid = "host-clarify" server._sessions[sid] = {"history_lock": threading.Lock()} - calls = [] - monkeypatch.setitem( - server._methods, - "clarify.respond", - lambda rid, params: calls.append((rid, dict(params))) or {"result": {"status": "ok"}}, - ) + req = server_requests.ServerRequest(sid, "clarify", {"question": "?"}) + with server_requests._lock: + server_requests._open[req.id] = req + locks = [] + monkeypatch.setitem(server._methods, "clarify.lock", + lambda rid, params: locks.append((rid, dict(params))) or {"result": {"status": "ok", "remaining": []}}) try: - host._handle_respond( - { - "sid": sid, - "request_id": "relay-response", - "params": {"request_id": "clarify-request", "answer": "yes"}, - } - ) - assert calls == [("relay-response", {"request_id": "clarify-request", "answer": "yes"})] + host._handle_respond({"sid": sid, "request_id": "relay-lock", + "params": {"lock": {"request_id": req.id, "question_id": "q0", "answer": "a"}}}) + assert locks == [("relay-lock", {"request_id": req.id, "question_id": "q0", "answer": "a"})] + assert _json_lines(out)[-1]["response"] == {"result": {"status": "ok", "remaining": []}} + + host._handle_respond({"sid": sid, "request_id": "relay-response", + "params": {"frame": {"jsonrpc": "2.0", "id": req.id, "result": {"answer": "yes"}}}}) + assert req.answered and req.result == {"answer": "yes"} and req.event.is_set() frame = _json_lines(out)[-1] - assert frame == { - "type": "respond.ack", - "sid": sid, - "request_id": "relay-response", - "response": {"result": {"status": "ok"}}, - "host_ns": frame["host_ns"], - } + assert frame["type"] == "respond.ack" and frame["response"]["result"] == {"status": "ok"} finally: server._sessions.pop(sid, None) + server_requests.reset_for_tests() host.close() + def test_mutator_route_table_matches_prd_inventory(): assert MUTATOR_ROUTE_TABLE == { "prompt.submit": "turn-path", diff --git a/tests/tui_gateway/test_goal_command.py b/tests/tui_gateway/test_goal_command.py index 2f04458e66..674ca15f36 100644 --- a/tests/tui_gateway/test_goal_command.py +++ b/tests/tui_gateway/test_goal_command.py @@ -70,8 +70,7 @@ def server(hermes_home, monkeypatch): # _enter_buffered_busy. Clearing the per-session dicts gives the # next test a clean slate. mod._sessions.clear() - mod._pending.clear() - mod._answers.clear() + __import__("tui_gateway.server_requests", fromlist=["x"]).reset_for_tests() @pytest.fixture() diff --git a/tests/tui_gateway/test_hosted_room_member_activity_hook.py b/tests/tui_gateway/test_hosted_room_member_activity_hook.py index af73a6cbec..26521d5d88 100644 --- a/tests/tui_gateway/test_hosted_room_member_activity_hook.py +++ b/tests/tui_gateway/test_hosted_room_member_activity_hook.py @@ -57,7 +57,10 @@ def test_room_member_session_events_reach_plugins_with_room_coordinates(observer _session("room-sid", hosted=True) try: server._emit("tool.start", "room-sid", {"tool_id": "call-1", "name": "terminal", "args": {"command": "ls"}}) - server._emit("approval.request", "room-sid", {"request_id": "req-1", "command": "rm -rf build"}) + # An approval is a server→client REQUEST frame, not an event; it is member activity all the same. + from tui_gateway import server_requests + server_requests.send_async("approval", "room-sid", {"request_id": "req-1", "command": "rm -rf build"}, + lambda result: None)("test_done") server._emit("session.info", "room-sid", {"title": "chrome, not member activity"}) server._emit("tool.complete", "room-sid", {"tool_id": "call-1", "name": "terminal", "result": "ok"}) finally: @@ -71,8 +74,9 @@ def test_room_member_session_events_reach_plugins_with_room_coordinates(observer assert {k: started[k] for k in HOSTED_TASK} == HOSTED_TASK assert started["payload"]["tool_id"] == "call-1" and started["payload"]["args"] == {"command": "ls"} assert observer[1]["payload"]["request_id"] == "req-1" - # Per-session replay seq travels with the event so consumers can order/dedupe. - assert [event["seq"] for event in observer] == sorted(event["seq"] for event in observer) + # Per-session replay seq travels with the events so consumers can order/dedupe (a request frame carries none). + seqs = [event["seq"] for event in observer if event["seq"] is not None] + assert seqs == sorted(seqs) and len(seqs) == 2 def test_ordinary_session_events_never_fire_the_room_hook(observer): diff --git a/tests/tui_gateway/test_inline_rpc_gil_starvation.py b/tests/tui_gateway/test_inline_rpc_gil_starvation.py index e6cfa141df..183ab49a8c 100644 --- a/tests/tui_gateway/test_inline_rpc_gil_starvation.py +++ b/tests/tui_gateway/test_inline_rpc_gil_starvation.py @@ -56,8 +56,7 @@ def server(): mod._methods.update(methods) mod._real_stdout = real_stdout mod._sessions.clear() - mod._pending.clear() - mod._answers.clear() + __import__("tui_gateway.server_requests", fromlist=["x"]).reset_for_tests() @pytest.fixture() diff --git a/tests/tui_gateway/test_loop_command.py b/tests/tui_gateway/test_loop_command.py index 990947baa3..22af6d7f97 100644 --- a/tests/tui_gateway/test_loop_command.py +++ b/tests/tui_gateway/test_loop_command.py @@ -44,8 +44,7 @@ def server(hermes_home): mod = importlib.import_module("tui_gateway.server") yield mod mod._sessions.clear() - mod._pending.clear() - mod._answers.clear() + __import__("tui_gateway.server_requests", fromlist=["x"]).reset_for_tests() @pytest.fixture() diff --git a/tests/tui_gateway/test_protocol.py b/tests/tui_gateway/test_protocol.py index 1e32a1a4e6..b7e47cb3af 100644 --- a/tests/tui_gateway/test_protocol.py +++ b/tests/tui_gateway/test_protocol.py @@ -56,8 +56,8 @@ def server(): mod._real_stdout = real_stdout for sid in list(mod._sessions): mod._close_session_by_id(sid, end_reason="test_cleanup") - mod._pending.clear() - mod._answers.clear() + from tui_gateway import server_requests + server_requests.reset_for_tests() mod._live_transports.clear() @@ -219,42 +219,26 @@ def test_live_session_payload_replays_pending_approval(server, monkeypatch): assert replayed == first -def test_live_session_payload_replays_pending_clarify(server): - """A reattached client also receives a clarify question emitted while detached.""" +def test_live_session_payload_replays_open_requests(server): + """A reattached client receives the server→client request still blocking the session (a clarify sent + while detached), scoped to the owning runtime session only, as a snapshot not a live reference.""" + from tui_gateway import server_requests session = { - "agent": types.SimpleNamespace(), - "cols": 80, - "created_at": 1.0, - "history": [], - "history_lock": threading.Lock(), - "running": True, - "session_key": "stored-session", + "agent": types.SimpleNamespace(), "cols": 80, "created_at": 1.0, "history": [], + "history_lock": threading.Lock(), "running": True, "session_key": "stored-session", } - clarify_payload = { - "choices": ["staging", "production"], - "question": "Which deployment target?", - "request_id": "rid-clarify", - } - with server._prompt_lock: - server._pending["rid-clarify"] = ("runtime-session", threading.Event()) - server._pending_prompt_payloads["rid-clarify"] = ( - "clarify.request", - dict(clarify_payload), - ) - + req = server_requests.ServerRequest("runtime-session", "clarify", {"question": "Which?", "choices": ["a", "b"]}) + server_requests._open[req.id] = req try: payload = server._live_session_payload("runtime-session", session) other = server._live_session_payload("other-session", session) finally: - with server._prompt_lock: - server._pending.pop("rid-clarify", None) - server._pending_prompt_payloads.pop("rid-clarify", None) + server_requests._open.pop(req.id, None) - assert payload["pending_clarify"] == clarify_payload - # Snapshot, not a live reference into the registry. - assert payload["pending_clarify"] is not clarify_payload - # Scoped to the owning runtime session only. - assert "pending_clarify" not in other + assert payload["open_requests"] == [{"id": req.id, "method": "clarify", + "params": {"session_id": "runtime-session", "question": "Which?", "choices": ["a", "b"]}}] + assert payload["open_requests"][0]["params"] is not req.params + assert "open_requests" not in other def test_disable_flush_env_var_actually_wires_to_module_constant(monkeypatch): @@ -290,275 +274,145 @@ def test_emit_with_payload(capture): assert msg["params"]["payload"]["key"] == "val" -# ── Blocking prompt round-trip ─────────────────────────────────────── +# ── Server→client requests (tui_gateway/server_requests.py) ───────── -def test_block_and_respond(capture): - server, _ = capture - result = [None] - - threading.Thread( - target=lambda: result.__setitem__(0, server._block("test.prompt", "s1", {"q": "?"}, timeout=5)), - ).start() - - for _ in range(100): - if server._pending: - break - threading.Event().wait(0.01) - - rid = next(iter(server._pending)) - server._answers[rid] = "my_answer" - # _pending values are (sid, Event) tuples — unpack to set the Event - _, ev = server._pending[rid] - ev.set() - - threading.Event().wait(0.1) - assert result[0] == "my_answer" +def _frames(buf): + return [json.loads(line) for line in buf.getvalue().splitlines()] -@pytest.mark.parametrize( - "event", - ["secret.request", "sudo.request", "clarify.request", "terminal.read.request"], -) -def test_sensitive_prompt_timeout_emits_expiry(capture, event): +def _wait_open(server_requests, buf=None, timeout=2.0): + """The open request once its frame has been written (registration precedes the write).""" + deadline = time.monotonic() + timeout + while time.monotonic() < deadline: + with server_requests._lock: + req = next(iter(server_requests._open.values()), None) + if req is not None and (buf is None or req.id in buf.getvalue()): + return req + time.sleep(0.01) + raise AssertionError("server request never registered") + + +def test_server_request_round_trip_uses_response_frame(capture): + """The backend sends a real JSON-RPC request (string id, method, params.session_id) and the client's + response frame — not a method call — resolves it; the response frame itself gets no reply.""" + from tui_gateway import server_requests server, buf = capture + box = {} + thread = threading.Thread(target=lambda: box.__setitem__("r", server._ask("sudo", "s1", {}, timeout=5)), daemon=True) + thread.start() + req = _wait_open(server_requests, buf) + frame = _frames(buf)[-1] + assert frame == {"jsonrpc": "2.0", "id": req.id, "method": "sudo", "params": {"session_id": "s1"}} + assert req.id.startswith("srq-") - assert server._block(event, "s1", {}, timeout=0) == "" - - messages = [json.loads(line) for line in buf.getvalue().splitlines()] - request, expiry = [message["params"] for message in messages] - assert request["type"] == event - assert expiry["type"] == event.removesuffix(".request") + ".expire" - assert expiry["session_id"] == "s1" - assert expiry["payload"]["request_id"] == request["payload"]["request_id"] + assert server.dispatch({"jsonrpc": "2.0", "id": req.id, "result": {"value": "hunter2"}}) is None + thread.join(timeout=5) + assert box["r"] == "hunter2" + with server_requests._lock: + assert not server_requests._open -@pytest.mark.parametrize( - ("method", "value_key"), - [ - ("secret.respond", "value"), - ("sudo.respond", "password"), - ("clarify.respond", "answer"), - ("terminal.read.respond", "text"), - ], -) -def test_late_prompt_response_is_idempotent(server, method, value_key): - """All four blocking bridges tolerate a late reply after their request has - expired — the `*.respond` returns a graceful `{"status": "expired"}` instead - of the raw 4009 protocol error a client would otherwise surface verbatim.""" - response = server.handle_request( - { - "id": "late-response", - "method": method, - "params": {"request_id": "expired-request", value_key: ""}, - } - ) +@pytest.mark.parametrize("method", ["secret", "sudo", "terminal.read", "tour"]) +def test_server_request_timeout_emits_one_request_cancel(capture, method): + from tui_gateway import server_requests + server, buf = capture + assert server_requests.send(method, "s1", {}, timeout=0) is None + request, cancel = _frames(buf) + assert request["method"] == method + assert cancel["params"]["type"] == "request.cancel" + assert cancel["params"]["session_id"] == "s1" + assert cancel["params"]["payload"] == {"id": request["id"], "method": method, "reason": "timeout"} + +def test_late_response_and_lock_are_dropped_quietly(server): + """A response for an id that already ended is ignored (no reply, no error); a late clarify.lock + answers `expired` instead of a protocol error a card would surface verbatim.""" + assert server.dispatch({"jsonrpc": "2.0", "id": "srq-gone", "result": {"value": ""}}) is None + response = server.handle_request({"id": "late", "method": "clarify.lock", + "params": {"request_id": "srq-gone", "question_id": "q0", "answer": ""}}) assert response["result"] == {"status": "expired"} -# ── clarify batch (multi-question) bridge ──────────────────────────── - - -def _drain_batch_block(server, qids, timeout=5, payload=None): - """Run a batch _block on a worker thread and return (thread, result box, - emitted request payload). The caller resolves questions via - handle_request and then joins.""" +def _start_batch_clarify(server, buf, qids, timeout=None): + from tui_gateway import server_requests box = {} - - def run(): - box["answer"] = server._block( - "clarify.request", - "s1", - dict(payload or {"questions": [{"qid": q, "question": q} for q in qids]}), - timeout=timeout, - batch_qids=list(qids), - ) - - thread = threading.Thread(target=run, daemon=True) + normalized = [{"qid": q, "id": "", "question": q, "choices": None, "choices_offered": [], "multi_select": False} + for q in qids] + if timeout is not None: + server._clarify_timeout_seconds = lambda: timeout + thread = threading.Thread( + target=lambda: box.__setitem__("answer", server._clarify_block("s1", "", None, questions=normalized)), daemon=True) thread.start() - # Wait for the request to be registered so respond calls can find it. - deadline = time.monotonic() + 2 - while time.monotonic() < deadline: - with server._prompt_lock: - if server._batch_clarify: - rid = next(iter(server._batch_clarify)) - return thread, box, rid - time.sleep(0.01) - raise AssertionError("batch clarify request never registered") + return thread, box, _wait_open(server_requests, buf) -def test_clarify_batch_resolves_when_all_questions_locked(capture): +def test_clarify_batch_locks_resolve_in_order_and_keep_partial_on_timeout(capture): + """Per-question locks (clarify.lock) are editable until every qid is locked; the last lock resolves the + request with the full answer set. Only wire fields reach the renderer.""" server, buf = capture - thread, box, rid = _drain_batch_block(server, ["q0", "q1"]) - - first = server.handle_request({ - "id": "a1", "method": "clarify.respond", - "params": {"request_id": rid, "question_id": "q1", "answer": "beta"}, - }) - assert first["result"]["status"] == "ok" - assert first["result"]["remaining"] == ["q0"] - assert thread.is_alive() # one question left — still blocking - - second = server.handle_request({ - "id": "a2", "method": "clarify.respond", - "params": {"request_id": rid, "question_id": "q0", "answer": "alpha"}, - }) - assert second["result"]["status"] == "ok" - assert second["result"]["remaining"] == [] + thread, box, req = _start_batch_clarify(server, buf, ["q0", "q1"]) + sent = _frames(buf)[-1]["params"]["questions"][0] + assert set(sent) == {"qid", "question", "choices", "multi_select"} + first = server.handle_request({"id": "a1", "method": "clarify.lock", + "params": {"request_id": req.id, "question_id": "q0", "answer": "x"}}) + assert first["result"] == {"status": "ok", "remaining": ["q1"]} + redo = server.handle_request({"id": "a1b", "method": "clarify.lock", + "params": {"request_id": req.id, "question_id": "q0", "answer": "y"}}) + assert redo["result"] == {"status": "ok", "remaining": ["q1"]} + bad = server.handle_request({"id": "bad", "method": "clarify.lock", + "params": {"request_id": req.id, "question_id": "nope", "answer": ""}}) + assert bad["error"]["code"] == 4002 + assert server._open_requests("s1")[0]["params"]["answers"] == {"q0": "y"} + last = server.handle_request({"id": "a2", "method": "clarify.lock", + "params": {"request_id": req.id, "question_id": "q1", "answer": ""}}) + assert last["result"] == {"status": "ok", "remaining": []} thread.join(timeout=5) - assert not thread.is_alive() - assert json.loads(box["answer"]) == {"answers": {"q0": "alpha", "q1": "beta"}} + assert json.loads(box["answer"]) == {"answers": {"q0": "y", "q1": ""}} + + # Deadline: locked answers survive, timed_out flagged, one request.cancel. + original_timeout = server._clarify_timeout_seconds + try: + thread, box, req = _start_batch_clarify(server, buf, ["q0", "q1"], timeout=1.5) + locked = server.handle_request({"id": "b1", "method": "clarify.lock", + "params": {"request_id": req.id, "question_id": "q0", "answer": "kept"}}) + assert locked["result"]["status"] == "ok" + thread.join(timeout=5) + finally: + server._clarify_timeout_seconds = original_timeout + assert json.loads(box["answer"]) == {"answers": {"q0": "kept"}, "timed_out": True} + cancels = [f for f in _frames(buf) if f.get("method") == "event" and f["params"]["type"] == "request.cancel"] + assert [c["params"]["payload"]["id"] for c in cancels] == [req.id] -def test_clarify_batch_answer_update_overwrites_before_completion(server): - thread, box, rid = _drain_batch_block(server, ["q0", "q1"]) - - server.handle_request({ - "id": "a1", "method": "clarify.respond", - "params": {"request_id": rid, "question_id": "q0", "answer": "first"}, - }) - server.handle_request({ - "id": "a2", "method": "clarify.respond", - "params": {"request_id": rid, "question_id": "q0", "answer": "changed"}, - }) - server.handle_request({ - "id": "a3", "method": "clarify.respond", - "params": {"request_id": rid, "question_id": "q1", "answer": "done"}, - }) - - thread.join(timeout=5) - assert json.loads(box["answer"])["answers"]["q0"] == "changed" - - -def test_clarify_batch_empty_answer_is_a_locked_skip(server): - """Skipping one question locks an empty answer — it counts toward - completion instead of leaving the batch waiting.""" - thread, box, rid = _drain_batch_block(server, ["q0", "q1"]) - - server.handle_request({ - "id": "a1", "method": "clarify.respond", - "params": {"request_id": rid, "question_id": "q0", "answer": ""}, - }) - server.handle_request({ - "id": "a2", "method": "clarify.respond", - "params": {"request_id": rid, "question_id": "q1", "answer": "kept"}, - }) - - thread.join(timeout=5) - assert json.loads(box["answer"]) == {"answers": {"q0": "", "q1": "kept"}} - - -def test_clarify_batch_unknown_question_id_rejected(server): - thread, box, rid = _drain_batch_block(server, ["q0"]) - - response = server.handle_request({ - "id": "bad", "method": "clarify.respond", - "params": {"request_id": rid, "question_id": "q9", "answer": "x"}, - }) - assert response["error"]["code"] == 4002 - - server.handle_request({ - "id": "ok", "method": "clarify.respond", - "params": {"request_id": rid, "question_id": "q0", "answer": "fine"}, - }) - thread.join(timeout=5) - - -def test_clarify_batch_timeout_keeps_locked_answers(capture): - """Locked answers survive the deadline: the tool sees the partials plus - timed_out instead of an empty string.""" +def test_clarify_batch_cancel_all_is_a_response_without_answers(capture): server, buf = capture - thread, box, rid = _drain_batch_block(server, ["q0", "q1"], timeout=1) - - server.handle_request({ - "id": "a1", "method": "clarify.respond", - "params": {"request_id": rid, "question_id": "q0", "answer": "kept"}, - }) - - thread.join(timeout=10) - assert not thread.is_alive() - result = json.loads(box["answer"]) - assert result == {"answers": {"q0": "kept"}, "timed_out": True} - # The expire notification still fires for the un-finished batch. - messages = [json.loads(line) for line in buf.getvalue().splitlines()] - assert any(m["params"]["type"] == "clarify.expire" for m in messages) - - -def test_clarify_batch_cancel_all_returns_empty(server): - """A respond without question_id cancels the whole batch (Esc path).""" - thread, box, rid = _drain_batch_block(server, ["q0", "q1"]) - - server.handle_request({ - "id": "cancel", "method": "clarify.respond", - "params": {"request_id": rid, "answer": ""}, - }) - + thread, box, req = _start_batch_clarify(server, buf, ["q0", "q1"]) + server.dispatch({"jsonrpc": "2.0", "id": req.id, "result": {}}) thread.join(timeout=5) assert box["answer"] == "" -def test_clarify_batch_late_question_respond_is_idempotent(server): - response = server.handle_request({ - "id": "late", "method": "clarify.respond", - "params": {"request_id": "gone", "question_id": "q0", "answer": "x"}, - }) - assert response["result"] == {"status": "expired"} - - -def test_clarify_batch_state_cleared_after_resolution(server): - thread, box, rid = _drain_batch_block(server, ["q0"]) - server.handle_request({ - "id": "a", "method": "clarify.respond", - "params": {"request_id": rid, "question_id": "q0", "answer": "x"}, - }) - thread.join(timeout=5) - with server._prompt_lock: - assert rid not in server._batch_clarify - assert rid not in server._pending - - -def test_clarify_block_helper_builds_batch_payload(capture): - """_clarify_block forwards only wire fields (qid/question/choices/ - multi_select) — the tool-side normalized entries carry extra keys the - renderer must not see.""" +def test_clear_pending_cancels_only_that_session(capture): + from tui_gateway import server_requests server, buf = capture - normalized = [ - { - "qid": "q0", "id": "approach", "question": "Which?", - "choices": ["a (Recommended)", "b"], "choices_offered": ["a", "b"], - "multi_select": False, - }, - ] - - box = {} - - def run(): - box["answer"] = server._clarify_block("s1", "", None, questions=normalized) - - thread = threading.Thread(target=run, daemon=True) - thread.start() + a = threading.Thread(target=lambda: server_requests.send("sudo", "sid-a", {}, timeout=None), daemon=True) + b = threading.Thread(target=lambda: server_requests.send("sudo", "sid-b", {}, timeout=None), daemon=True) + a.start(); b.start() deadline = time.monotonic() + 2 - rid = None - while time.monotonic() < deadline and rid is None: - with server._prompt_lock: - rid = next(iter(server._batch_clarify), None) + while time.monotonic() < deadline and len(server_requests._open) < 2: time.sleep(0.01) - assert rid - - server.handle_request({ - "id": "a", "method": "clarify.respond", - "params": {"request_id": rid, "question_id": "q0", "answer": "a"}, - }) - thread.join(timeout=5) - - messages = [json.loads(line) for line in buf.getvalue().splitlines()] - request = messages[0]["params"] - assert request["type"] == "clarify.request" - sent = request["payload"]["questions"][0] - assert set(sent) == {"qid", "question", "choices", "multi_select"} - assert "id" not in sent and "choices_offered" not in sent + assert server._session_pending_kind("sid-a") == "sudo" + server._clear_pending("sid-a") + a.join(timeout=2) + assert not a.is_alive() and b.is_alive() + assert server._session_pending_kind("sid-a") == "" and server._session_pending_kind("sid-b") == "sudo" + server._clear_pending() + b.join(timeout=2) + assert not b.is_alive() + reasons = [f["params"]["payload"]["reason"] for f in _frames(buf) if f.get("method") == "event"] + assert reasons == ["interrupted", "shutdown"] def test_approval_pending_replays_unresolved_requests(server, monkeypatch): @@ -701,16 +555,6 @@ def test_approval_respond_4001_when_nothing_resolves(server, monkeypatch): assert response["error"]["code"] == 4001 -def test_clear_pending(server): - ev = threading.Event() - # _pending values are (sid, Event) tuples - server._pending["r1"] = ("sid-x", ev) - server._clear_pending() - - assert ev.is_set() - assert server._answers["r1"] == "" - - # ── Session lookup ─────────────────────────────────────────────────── diff --git a/tests/tui_gateway/test_review_summary_callback.py b/tests/tui_gateway/test_review_summary_callback.py index 0f4130a604..d9fe7aa68a 100644 --- a/tests/tui_gateway/test_review_summary_callback.py +++ b/tests/tui_gateway/test_review_summary_callback.py @@ -43,8 +43,7 @@ def server(): # _enter_buffered_busy. Clearing the per-session dicts gives the # next test a clean slate. mod._sessions.clear() - mod._pending.clear() - mod._answers.clear() + __import__("tui_gateway.server_requests", fromlist=["x"]).reset_for_tests() def test_init_session_attaches_background_review_callback(server, monkeypatch): diff --git a/tests/tui_gateway/test_session_control.py b/tests/tui_gateway/test_session_control.py index c9da204b4b..99541ed1fa 100644 --- a/tests/tui_gateway/test_session_control.py +++ b/tests/tui_gateway/test_session_control.py @@ -43,8 +43,7 @@ def server(hermes_home, monkeypatch): monkeypatch.setattr(mod, "_cfg_path", None) yield mod mod._sessions.clear() - mod._pending.clear() - mod._answers.clear() + __import__("tui_gateway.server_requests", fromlist=["x"]).reset_for_tests() @pytest.fixture() diff --git a/tests/tui_gateway/test_session_profile_db.py b/tests/tui_gateway/test_session_profile_db.py index f41a5873bf..1c2489bd95 100644 --- a/tests/tui_gateway/test_session_profile_db.py +++ b/tests/tui_gateway/test_session_profile_db.py @@ -64,8 +64,7 @@ def server(hermes_home): mod._methods.clear() mod._methods.update(methods) mod._sessions.clear() - mod._pending.clear() - mod._answers.clear() + __import__("tui_gateway.server_requests", fromlist=["x"]).reset_for_tests() mod._db = None diff --git a/tests/tui_gateway/test_stranded_session_adoption.py b/tests/tui_gateway/test_stranded_session_adoption.py index a1ae859e11..ed66c96c38 100644 --- a/tests/tui_gateway/test_stranded_session_adoption.py +++ b/tests/tui_gateway/test_stranded_session_adoption.py @@ -183,8 +183,7 @@ def gateway(tmp_path, monkeypatch): mod._methods.clear() mod._methods.update(methods) mod._sessions.clear() - mod._pending.clear() - mod._answers.clear() + __import__("tui_gateway.server_requests", fromlist=["x"]).reset_for_tests() mod._db = None default_db.close() diff --git a/tests/tui_gateway/test_subagent_child_mirror.py b/tests/tui_gateway/test_subagent_child_mirror.py index 5a94dc8ed0..8b1681e6ca 100644 --- a/tests/tui_gateway/test_subagent_child_mirror.py +++ b/tests/tui_gateway/test_subagent_child_mirror.py @@ -36,8 +36,7 @@ def server(): yield mod mod._sessions.clear() - mod._pending.clear() - mod._answers.clear() + __import__("tui_gateway.server_requests", fromlist=["x"]).reset_for_tests() mod._child_mirrors.clear() mod._active_child_runs.clear() diff --git a/tests/tui_gateway/test_tour_bridge_fail_fast.py b/tests/tui_gateway/test_tour_bridge_fail_fast.py index 0699895d55..290c5e0ca4 100644 --- a/tests/tui_gateway/test_tour_bridge_fail_fast.py +++ b/tests/tui_gateway/test_tour_bridge_fail_fast.py @@ -1,9 +1,9 @@ -"""A desktop client that cannot answer ``tour.request`` must not cost a full +"""A desktop client that cannot answer the ``tour`` server request must not cost a full bridge timeout per call. The renderer's handler ships in the desktop bundle; the tool is offered by the backend. An app build older than the tour tool has no branch for the event, so -nothing ever calls ``tour.respond`` and the agent blocks for the whole deadline +nothing ever answers the request and the agent blocks for the whole deadline — once per action the model tries. See tui_gateway.server._tour_request. """ @@ -23,7 +23,7 @@ def session(monkeypatch): @pytest.fixture def bridge(monkeypatch): - """Record every _block call and serve canned answers.""" + """Record every ``_ask`` (server request) call and serve canned answers.""" calls = [] def fake_block(event, sid, payload, timeout=None, **_kw): @@ -33,7 +33,7 @@ def bridge(monkeypatch): fake_block.answers = [] fake_block.calls = calls - monkeypatch.setattr(server, "_block", fake_block) + monkeypatch.setattr(server, "_ask", fake_block) return fake_block @@ -41,7 +41,7 @@ def test_first_action_is_probed_on_a_short_deadline(session, bridge): bridge.answers = [json.dumps({"success": True})] server._tour_request("s1", {"action": "targets"}) - assert bridge.calls[0]["event"] == "tour.request" + assert bridge.calls[0]["event"] == "tour" assert bridge.calls[0]["timeout"] == server._TOUR_PROBE_TIMEOUT_S assert server._TOUR_PROBE_TIMEOUT_S < server._TOUR_TIMEOUT_S diff --git a/tests/tui_gateway/test_tui_gateway_server.py b/tests/tui_gateway/test_tui_gateway_server.py index 8465eec67d..acc05d8f00 100644 --- a/tests/tui_gateway/test_tui_gateway_server.py +++ b/tests/tui_gateway/test_tui_gateway_server.py @@ -496,15 +496,17 @@ def test_compute_host_turn_end_updates_metadata_mirror(monkeypatch): server._sessions.pop("iso-sid", None) -def test_compute_host_clarify_snapshot_replays_and_proxies_batch_answers(monkeypatch): - """A host-owned clarify survives activation and receives its UI answers.""" +def test_compute_host_open_request_survives_activation_and_proxies_locks_and_responses(monkeypatch): + """A host-owned server request (batch clarify) is mirrored by the parent so `open_requests` replays it; + `clarify.lock` and the client's response frame are relayed to the child that owns the wait.""" class _Supervisor: def __init__(self): self.responses = [] def respond(self, sid, params, *, timeout=15.0): self.responses.append((sid, dict(params), timeout)) - remaining = ["q1"] if params.get("question_id") == "q0" else [] + lock = params.get("lock") or {} + remaining = ["q1"] if lock.get("question_id") == "q0" else [] return {"type": "respond.ack", "response": {"result": {"status": "ok", "remaining": remaining}}} sid = "host-clarify" @@ -515,53 +517,25 @@ def test_compute_host_clarify_snapshot_replays_and_proxies_batch_answers(monkeyp monkeypatch.setattr(server, "_get_compute_host_supervisor", lambda _cfg=None: supervisor) monkeypatch.setattr(server, "write_json", lambda _message: True) + questions = [{"qid": "q0", "question": "First?", "choices": ["a"]}, {"qid": "q1", "question": "Second?", "choices": ["b"]}] try: - server._relay_compute_host_rpc( - { - "jsonrpc": "2.0", - "method": "event", - "params": { - "type": "clarify.request", - "session_id": sid, - "payload": { - "request_id": "host-request", - "questions": [ - {"qid": "q0", "question": "First?", "choices": ["a"]}, - {"qid": "q1", "question": "Second?", "choices": ["b"]}, - ], - }, - }, - } - ) + server._relay_compute_host_rpc({"jsonrpc": "2.0", "id": "srq-host", "method": "clarify", + "params": {"session_id": sid, "questions": questions}}) activated = server._live_session_payload(sid, session) - assert activated["pending_clarify"]["request_id"] == "host-request" - - response = server.handle_request( - { - "id": "clarify-q0", - "method": "clarify.respond", - "params": {"request_id": "host-request", "question_id": "q0", "answer": "a"}, - } - ) + assert activated["open_requests"] == [{"id": "srq-host", "method": "clarify", + "params": {"session_id": sid, "questions": questions}}] + response = server.handle_request({"id": "lock-q0", "method": "clarify.lock", + "params": {"request_id": "srq-host", "question_id": "q0", "answer": "a"}}) assert response["result"] == {"status": "ok", "remaining": ["q1"]} - assert supervisor.responses == [ - (sid, {"request_id": "host-request", "question_id": "q0", "answer": "a"}, 15.0) - ] - replayed = server._live_session_payload(sid, session)["pending_clarify"] - assert replayed["answers"] == {"q0": "a"} + assert supervisor.responses == [(sid, {"lock": {"request_id": "srq-host", "question_id": "q0", "answer": "a"}}, 15.0)] + assert server._live_session_payload(sid, session)["open_requests"][0]["params"]["answers"] == {"q0": "a"} - final_response = server.handle_request( - { - "id": "clarify-q1", - "method": "clarify.respond", - "params": {"request_id": "host-request", "question_id": "q1", "answer": "b"}, - } - ) - - assert final_response["result"] == {"status": "ok", "remaining": []} - assert "pending_clarify" not in server._live_session_payload(sid, session) + # The client's response frame (cancel-all) is relayed to the child and clears the mirror. + assert server.dispatch({"jsonrpc": "2.0", "id": "srq-host", "result": {}}) is None + assert supervisor.responses[-1] == (sid, {"frame": {"jsonrpc": "2.0", "id": "srq-host", "result": {}}}, 15.0) + assert "open_requests" not in server._live_session_payload(sid, session) finally: server._sessions.pop(sid, None) @@ -13585,10 +13559,20 @@ def test_prompt_submit_row_id_accepts_full_lineage_ordinal(monkeypatch): # --------------------------------------------------------------------------- +def _open_request(sid, method="clarify"): + from tui_gateway import server_requests + req = server_requests.ServerRequest(sid, method, {}) + with server_requests._lock: + server_requests._open[req.id] = req + return req + + def test_interrupt_only_clears_own_session_pending(): - """session.interrupt on session A must NOT release pending prompts - that belong to session B.""" + """session.interrupt on session A withdraws A's open server→client requests (a request.cancel each) + and must NOT touch session B's — otherwise B's clarify/sudo/secret prompt silently resolves as if + the user cancelled it.""" import types + from tui_gateway import server_requests session_a = _session() session_a["agent"] = types.SimpleNamespace(interrupt=lambda: None) @@ -13596,71 +13580,21 @@ def test_interrupt_only_clears_own_session_pending(): session_b["agent"] = types.SimpleNamespace(interrupt=lambda: None) server._sessions["sid_a"] = session_a server._sessions["sid_b"] = session_b + req_a1, req_a2, req_b = _open_request("sid_a"), _open_request("sid_a", "sudo"), _open_request("sid_b") try: - # Simulate pending prompts on both sessions (what _block creates - # while a clarify/sudo/secret request is outstanding). - ev_a = threading.Event() - ev_b = threading.Event() - server._pending["rid-a"] = ("sid_a", ev_a) - server._pending["rid-b"] = ("sid_b", ev_b) - server._answers.clear() - - # Interrupt session A. - resp = server.handle_request( - { - "id": "1", - "method": "session.interrupt", - "params": {"session_id": "sid_a"}, - } - ) + resp = server.handle_request({"id": "1", "method": "session.interrupt", "params": {"session_id": "sid_a"}}) assert resp.get("result"), f"got error: {resp.get('error')}" - # Session A's pending must be released to empty. - assert ev_a.is_set(), "sid_a pending Event should be set after interrupt" - assert server._answers.get("rid-a") == "" - - # Session B's pending MUST remain untouched — no cross-session blast. - assert not ev_b.is_set(), ( - "CRITICAL: session.interrupt on sid_a released a pending prompt " - "belonging to sid_b — other sessions' clarify/sudo/secret " - "prompts are being silently cancelled" - ) - assert "rid-b" not in server._answers + assert req_a1.event.is_set() and req_a2.event.is_set() and not req_a1.answered + assert not req_b.event.is_set(), ( + "CRITICAL: session.interrupt on sid_a released a prompt belonging to sid_b") + assert server_requests.open_requests("sid_a") == [] + assert [r["id"] for r in server_requests.open_requests("sid_b")] == [req_b.id] finally: server._sessions.pop("sid_a", None) server._sessions.pop("sid_b", None) - server._pending.pop("rid-a", None) - server._pending.pop("rid-b", None) - server._answers.pop("rid-a", None) - server._answers.pop("rid-b", None) - - -def test_interrupt_clears_multiple_own_pending(): - """When a single session has multiple pending prompts (uncommon but - possible via nested tool calls), interrupt must release all of them.""" - import types - - sess = _session() - sess["agent"] = types.SimpleNamespace(interrupt=lambda: None) - server._sessions["sid"] = sess - - try: - ev1, ev2 = threading.Event(), threading.Event() - server._pending["r1"] = ("sid", ev1) - server._pending["r2"] = ("sid", ev2) - - resp = server.handle_request( - {"id": "1", "method": "session.interrupt", "params": {"session_id": "sid"}} - ) - assert resp.get("result") - assert ev1.is_set() and ev2.is_set() - assert server._answers.get("r1") == "" and server._answers.get("r2") == "" - finally: - server._sessions.pop("sid", None) - for key in ("r1", "r2"): - server._pending.pop(key, None) - server._answers.pop(key, None) + server_requests.reset_for_tests() def test_run_prompt_submit_registers_turn_thread_for_interrupt(monkeypatch): @@ -14295,39 +14229,16 @@ def test_wait_agent_for_prompt_expires_at_cap(monkeypatch): def test_clear_pending_without_sid_clears_all(): - """_clear_pending(None) is the shutdown path — must still release - every pending prompt regardless of owning session.""" - ev1, ev2, ev3 = threading.Event(), threading.Event(), threading.Event() - server._pending["a"] = ("sid_x", ev1) - server._pending["b"] = ("sid_y", ev2) - server._pending["c"] = ("sid_z", ev3) + """_clear_pending(None) is the process-exit path — every open request is withdrawn, and a response for + a withdrawn id is dropped quietly (no error frame back to the client).""" + from tui_gateway import server_requests + reqs = [_open_request("sid-x"), _open_request("sid-y", "sudo")] try: server._clear_pending(None) - assert ev1.is_set() and ev2.is_set() and ev3.is_set() + assert all(r.event.is_set() and not r.answered for r in reqs) + assert server.dispatch({"jsonrpc": "2.0", "id": reqs[0].id, "result": {"answer": "late"}}) is None finally: - for key in ("a", "b", "c"): - server._pending.pop(key, None) - server._answers.pop(key, None) - - -def test_respond_unpacks_sid_tuple_correctly(): - """After the (sid, Event) tuple change, _respond must still work.""" - ev = threading.Event() - server._pending["rid-x"] = ("sid_x", ev) - try: - resp = server.handle_request( - { - "id": "1", - "method": "clarify.respond", - "params": {"request_id": "rid-x", "answer": "the answer"}, - } - ) - assert resp.get("result") - assert ev.is_set() - assert server._answers.get("rid-x") == "the answer" - finally: - server._pending.pop("rid-x", None) - server._answers.pop("rid-x", None) + server_requests.reset_for_tests() # --------------------------------------------------------------------------- @@ -20703,49 +20614,44 @@ def test_speak_text_with_barge_no_monitor_when_voice_mode_off(monkeypatch): assert not listened.is_set() -def test_clarify_callback_uses_configured_timeout(monkeypatch): - """The TUI/desktop clarify bridge honors the canonical clarify timeout - (via _clarify_timeout_seconds) instead of the hardcoded _block default.""" +def _capture_server_request(monkeypatch, result): + """Stub the server-request send and capture (method, sid, params, timeout).""" + from tui_gateway import server_requests captured = {} + def fake_send(method, sid, params, *, timeout, qids=None): + captured.update(method=method, sid=sid, params=params, timeout=timeout, qids=qids) + return result + + monkeypatch.setattr(server_requests, "send", fake_send) + return captured + + +def test_clarify_callback_uses_configured_timeout(monkeypatch): + """The TUI/desktop clarify bridge sends a ``clarify`` server request with the canonical clarify timeout + (via _clarify_timeout_seconds), and returns the response's ``answer``.""" monkeypatch.setattr(server, "_clarify_timeout_seconds", lambda: 42) - - def fake_block(event, sid, payload, timeout=300): - captured.update(event=event, sid=sid, payload=payload, timeout=timeout) - return "answer" - - monkeypatch.setattr(server, "_block", fake_block) + captured = _capture_server_request(monkeypatch, {"answer": "answer"}) result = server._agent_cbs("sid-1")["clarify_callback"]("Pick one", ["a", "b"]) assert result == "answer" - assert captured["event"] == "clarify.request" + assert captured["method"] == "clarify" and captured["sid"] == "sid-1" assert captured["timeout"] == 42 - assert captured["payload"] == {"question": "Pick one", "choices": ["a", "b"]} + assert captured["params"] == {"question": "Pick one", "choices": ["a", "b"]} def test_clarify_callback_multi_select_hint(monkeypatch): - """multi_select=True adds the hint to the payload; the single-select - payload shape stays byte-identical to the pre-multi-select protocol - (older renderers must never see the extra field).""" - captured = {} - - def fake_block(event, sid, payload, timeout=300): - captured.update(payload=payload) - return "answer" - - monkeypatch.setattr(server, "_block", fake_block) + """multi_select=True adds the hint to the params; the single-select shape stays byte-identical to the + pre-multi-select protocol (older renderers must never see the extra field).""" + captured = _capture_server_request(monkeypatch, {"answer": "answer"}) cb = server._agent_cbs("sid-1")["clarify_callback"] cb("Pick many", ["a", "b"], multi_select=True) - assert captured["payload"] == { - "question": "Pick many", - "choices": ["a", "b"], - "multi_select": True, - } + assert captured["params"] == {"question": "Pick many", "choices": ["a", "b"], "multi_select": True} cb("Pick one", ["a", "b"], multi_select=False) - assert captured["payload"] == {"question": "Pick one", "choices": ["a", "b"]} + assert captured["params"] == {"question": "Pick one", "choices": ["a", "b"]} @pytest.mark.parametrize( @@ -20753,8 +20659,8 @@ def test_clarify_callback_multi_select_hint(monkeypatch): [(0, None), (-1, None), (42, 42)], ) def test_clarify_timeout_seconds_maps_non_positive_to_unlimited(monkeypatch, configured, expected): - """A ``<= 0`` clarify timeout means unlimited and reaches _block as None - (ev.wait(None) waits forever) rather than an immediate ev.wait(0) skip.""" + """A ``<= 0`` clarify timeout means unlimited and reaches the server request as None + (wait(None) waits forever) rather than an immediate wait(0) skip.""" monkeypatch.setattr("tools.clarify_gateway.get_clarify_timeout", lambda: configured) assert server._clarify_timeout_seconds() == expected diff --git a/tests/tui_gateway/test_undo_command.py b/tests/tui_gateway/test_undo_command.py index 4472cf8e2b..9e44e23e41 100644 --- a/tests/tui_gateway/test_undo_command.py +++ b/tests/tui_gateway/test_undo_command.py @@ -55,8 +55,7 @@ def server(hermes_home): mod._methods.clear() mod._methods.update(methods) mod._sessions.clear() - mod._pending.clear() - mod._answers.clear() + __import__("tui_gateway.server_requests", fromlist=["x"]).reset_for_tests() mod._db = None diff --git a/tui_gateway/AGENTS.md b/tui_gateway/AGENTS.md index c3063b42b8..6e4224b7ac 100644 --- a/tui_gateway/AGENTS.md +++ b/tui_gateway/AGENTS.md @@ -18,12 +18,20 @@ Never move agent behaviour into the renderer. ## Transport -Newline-delimited JSON-RPC over stdio: requests from Ink, events from Python. `tui_gateway/server.py` +Newline-delimited JSON-RPC over stdio, peer-to-peer: client→server method calls, server→client +**requests** (the agent asking the user something: `approval`, `clarify`, `sudo`, `secret`, `vault.*`, +`mcp.setup`, the desktop read/act bridges) and server→client `event` notifications. `tui_gateway/server.py` is the facade with the method/event catalog; methods live in `methods_*.py` siblings (`methods_config`, `methods_complete`, `methods_browser`, `methods_bot_relay`, ...), event publishing in -`event_publisher.py` / `event_replay.py`. Desktop reaches the same server over WebSocket via -`apps/shared` (`JsonRpcGatewayClient`). New RPC = a new `methods_.py` or an entry in an -existing topical sibling, registered in the table — no `if method == ...` chain (root shape rules). +`event_publisher.py` / `event_replay.py`, server→client requests in `server_requests.py` (`send()` blocks +the agent thread until the response frame with the same `srq-` id arrives; `cancel*` withdraws with a +`request.cancel` event; `open_requests(sid)` is what `session.resume` / `session.events.since` replay so a +reconnecting client re-renders the still-open questions). Desktop reaches the same server over WebSocket +via `apps/shared` (`JsonRpcGatewayClient`, `onRequest`). New RPC = a new `methods_.py` or an entry +in an existing topical sibling, registered in the table — no `if method == ...` chain (root shape rules). +New question for the user = `_ask("", sid, params, timeout)` in the emitter, a handler in +`apps/desktop/.../gateway-event/server-requests.ts` and `ui-tui/src/app/createServerRequestHandler.ts`, +and the method in `ServerRequestMap` + `apps/shared/src/gateway-events.json`. New event = a new key in `apps/shared/src/gateway-events.ts::GatewayEventMap` + `BACKEND_EVENT_NAMES` AND `apps/shared/src/gateway-events.json`; `tests/tui_gateway/test_gateway_event_contract.py` (emitter side) and `apps/shared/src/gateway-events.test.ts` (type side) both fail when either drifts. @@ -34,8 +42,8 @@ side) and `apps/shared/src/gateway-events.test.ts` (type side) both fail when ei |---|---|---| | Chat streaming | `app.tsx` + `messageLine.tsx` | `prompt.submit` → `message.delta` / `message.complete` | | Tool activity | `thinking.tsx` | `tool.start` / `tool.generating` / `tool.complete` | -| Approvals | `prompts.tsx` | `approval.request` → `approval.respond` | -| Clarify / sudo / secret | `prompts.tsx`, `maskedPrompt.tsx` | `clarify.respond`, `sudo.respond`, `secret.respond` | +| Approvals | `prompts.tsx` | server→client request `approval` → response `{choice}` | +| Clarify / sudo / secret | `prompts.tsx`, `maskedPrompt.tsx` | server→client requests `clarify` / `sudo` / `secret` (`server_requests.py`) | | Session picker | `sessionPicker.tsx` | `session.list` / `session.resume` | | Slash commands | local handler + fallthrough | `slash.exec` → `_SlashWorker`; `command.dispatch` | | Completions | `useCompletion` hook | `complete.slash`, `complete.path` | diff --git a/tui_gateway/server.py b/tui_gateway/server.py index 25e1e8e239..623203b28c 100644 --- a/tui_gateway/server.py +++ b/tui_gateway/server.py @@ -622,6 +622,11 @@ def _emit(event: str, sid: str, payload: dict | None = None) -> bool: return write_json(_event_frame(event, sid, payload)) +from tui_gateway import server_requests as _server_requests # noqa: E402 + +_server_requests.bind_sinks(lambda frame: write_json(frame), lambda event, sid, payload: _emit(event, sid, payload)) + + # Live WS peer transports (maintained by tui_gateway.ws): the only route for session-less background # events, which write_json would otherwise drop on stdio (see _broadcast_global_event). _live_transports: set[Transport] = set() diff --git a/tui_gateway/server_requests.py b/tui_gateway/server_requests.py index b98848ed99..4e47daa89f 100644 --- a/tui_gateway/server_requests.py +++ b/tui_gateway/server_requests.py @@ -68,14 +68,19 @@ class ServerRequest: _lock = threading.Lock() _open: dict[str, ServerRequest] = {} +# Frame sinks, bound by ``bind_sinks`` from server.py at import time (like the method_ctx split +# modules): importing server back from here would pick a different module object under the test +# fixtures that patch ``sys.modules`` around the server import. +_write: Callable[[dict], Any] = lambda frame: None # noqa: E731 +_emit: Callable[[str, str, dict], Any] = lambda event, sid, payload: None # noqa: E731 -def _write(frame: dict) -> None: - from tui_gateway.server import write_json - write_json(frame) + +def bind_sinks(write_json: Callable[[dict], Any], emit: Callable[[str, str, dict], Any]) -> None: + global _write, _emit + _write, _emit = write_json, emit def _emit_cancel(req: ServerRequest, reason: str) -> None: - from tui_gateway.server import _emit _emit("request.cancel", req.sid, {"id": req.id, "method": req.method, "reason": reason}) @@ -144,6 +149,14 @@ def resolve_response(frame: dict) -> bool: else: result = frame.get("result") req.result = result if isinstance(result, dict) else {} + if req.qids and "answers" in req.result: + # Batch clarify: answers locked early via clarify.lock belong to the final set even when + # the closing response only carries the tail the user answered last. + answers = req.result.get("answers") + merged = dict(req.locked) + if isinstance(answers, dict): + merged.update(answers) + req.result = {**req.result, "answers": merged} req.answered = True if req.on_result is not None: req.on_result(req.result) diff --git a/ui-tui/README.md b/ui-tui/README.md index bd1961dae8..e4b58782e3 100644 --- a/ui-tui/README.md +++ b/ui-tui/README.md @@ -229,14 +229,18 @@ Tool/status activity is shown in a live activity lane. Transcript rows stay focu ## Prompt flows -The Python gateway can pause the main loop and request structured input: +The Python gateway can pause the main loop and ask the client a question. These are JSON-RPC +**requests from the server** (string id, answered with a response frame of the same id — see +`createServerRequestHandler.ts`), not events: -- `approval.request`: allow once, allow for session, allow always, or deny -- `clarify.request`: pick from choices or type a custom answer -- `sudo.request`: masked password entry -- `secret.request`: masked value entry for a named env var +- `approval`: allow once, allow for session, allow always, or deny → `{ choice }` +- `clarify`: pick from choices or type a custom answer → `{ answer }` (batch: `{ answers }`) +- `sudo`: masked password entry → `{ value }` +- `secret`: masked value entry for a named env var → `{ value }` - `session.list`: used by `SessionPicker` for `/resume` +A withdrawn question (timeout, interrupt) arrives as a `request.cancel` event carrying its id. + These are stateful UI branches in `app.tsx`, not separate screens. ## Commands @@ -304,12 +308,7 @@ Primary event types the client handles today: | `tool.generating` | `{ name }` | | `tool.progress` | `{ name, preview }` | | `tool.complete` | `{ tool_id, name, error?, summary?, duration_s?, inline_diff?, todos? }` | -| `clarify.request` | `{ question, choices?, request_id }` | -| `approval.request` | `{ command, description, allow_permanent? }` | -| `sudo.request` | `{ request_id }` | -| `sudo.expire` | `{ request_id }` clears a timed-out sudo prompt | -| `secret.request` | `{ prompt, env_var, request_id }` | -| `secret.expire` | `{ request_id }` clears a timed-out secret prompt | +| `request.cancel` | `{ id, method, reason }` clears the withdrawn server→client request | | `background.complete` | `{ task_id, text }` | | `billing.step_up.verification` | `{ verification_url, user_code }` | | `review.summary` | `{ text }` | diff --git a/ui-tui/src/__tests__/createGatewayEventHandler.test.ts b/ui-tui/src/__tests__/createGatewayEventHandler.test.ts index 3ee8bdcaf4..0577b44848 100644 --- a/ui-tui/src/__tests__/createGatewayEventHandler.test.ts +++ b/ui-tui/src/__tests__/createGatewayEventHandler.test.ts @@ -1,7 +1,9 @@ import { beforeEach, describe, expect, it, vi } from 'vitest' import { createGatewayEventHandler } from '../app/createGatewayEventHandler.js' +import { createServerRequestHandler } from '../app/createServerRequestHandler.js' import { getOverlayState, patchOverlayState, resetOverlayState } from '../app/overlayStore.js' +import { resetServerRequestsForTests } from '../app/serverRequestStore.js' import { turnController } from '../app/turnController.js' import { getTurnState, resetTurnState } from '../app/turnStore.js' import { getUiState, patchUiState, resetUiState } from '../app/uiStore.js' @@ -58,11 +60,27 @@ const buildCtx = (appended: Msg[]) => } }) as any +/** Deliver one server→client request (`tui_gateway/server_requests.py`) to the TUI's request handler. */ +const serverRequest = (method: string, params: Record, id = `srq-${method}`) => { + const respond = vi.fn() + + const handled = createServerRequestHandler({ ringPromptBell: vi.fn(), setStatus: status => patchUiState({ status }) })({ + fail: vi.fn(), + id, + method, + params, + respond + }) + + return { handled, respond } +} + describe('createGatewayEventHandler', () => { beforeEach(() => { resetOverlayState() resetUiState() resetTurnState() + resetServerRequestsForTests() turnController.fullReset() patchUiState({ showReasoning: true }) }) @@ -71,11 +89,7 @@ describe('createGatewayEventHandler', () => { patchUiState({ sid: 'focused' }) const onEvent = createGatewayEventHandler(buildCtx([])) onEvent({ session_id: 'focused', payload: {}, type: 'message.start' } as any) - onEvent({ - session_id: 'focused', - payload: { request_id: 'approval', command: 'test' }, - type: 'approval.request' - } as any) + serverRequest('approval', { session_id: 'focused', request_id: 'approval', command: 'test' }) const busyOverlay = getOverlayState().approval expect(getUiState().busy).toBe(true) expect(busyOverlay).not.toBeNull() @@ -1211,10 +1225,7 @@ describe('createGatewayEventHandler', () => { onEvent({ payload: { line: 'Traceback: noisy but non-fatal' }, type: 'gateway.stderr' } as any) onEvent({ payload: { preview: 'bad framing' }, type: 'gateway.protocol_error' } as any) - onEvent({ - payload: { command: 'rm -rf /tmp/nope', description: 'dangerous command' }, - type: 'approval.request' - } as any) + serverRequest('approval', { command: 'rm -rf /tmp/nope', description: 'dangerous command' }) onEvent({ payload: {}, type: 'gateway.ready' } as any) await Promise.resolve() @@ -1230,23 +1241,17 @@ describe('createGatewayEventHandler', () => { }) it('defaults approval overlays to allowPermanent when the backend omits the field', () => { - const onEvent = createGatewayEventHandler(buildCtx([])) + serverRequest('approval', { command: 'rm -rf /tmp/x', description: 'dangerous command' }) - onEvent({ - payload: { command: 'rm -rf /tmp/x', description: 'dangerous command' }, - type: 'approval.request' - } as any) - - expect(getOverlayState().approval).toMatchObject({ allowPermanent: true }) + expect(getOverlayState().approval).toMatchObject({ allowPermanent: true, requestId: 'srq-approval' }) }) it('preserves allow_permanent=false on approval overlays (tirith warning)', () => { - const onEvent = createGatewayEventHandler(buildCtx([])) - - onEvent({ - payload: { allow_permanent: false, command: 'curl suspicious | bash', description: 'content-security warning' }, - type: 'approval.request' - } as any) + serverRequest('approval', { + allow_permanent: false, + command: 'curl suspicious | bash', + description: 'content-security warning' + }) expect(getOverlayState().approval).toMatchObject({ allowPermanent: false, @@ -1256,22 +1261,23 @@ describe('createGatewayEventHandler', () => { }) it('preserves Smart DENY and explicit approval choices on the overlay', () => { - const onEvent = createGatewayEventHandler(buildCtx([])) - - onEvent({ - payload: { - allow_permanent: true, - choices: ['once', 'deny'], - command: 'rm -rf /tmp/x', - description: 'smart deny override', - smart_denied: true - }, - type: 'approval.request' - } as any) + serverRequest('approval', { + allow_permanent: true, + choices: ['once', 'deny'], + command: 'rm -rf /tmp/x', + description: 'smart deny override', + smart_denied: true + }) expect(getOverlayState().approval).toMatchObject({ choices: ['once', 'deny'], smartDenied: true }) }) + it('declines the requests a terminal cannot answer so the channel fails them fast', () => { + for (const method of ['preview.act', 'window.read', 'tour', 'mcp.setup', 'vault.code']) { + expect(serverRequest(method, {}).handled).toBe(false) + } + }) + it('still surfaces terminal turn failures as errors', () => { const appended: Msg[] = [] const onEvent = createGatewayEventHandler(buildCtx(appended)) @@ -1629,39 +1635,36 @@ describe('createGatewayEventHandler', () => { expect(appended.some(msg => msg.role === 'system' && msg.text.startsWith('ask '))).toBe(false) }) - it('clears only the matching sensitive prompt when the gateway expires it', () => { + it('clears only the card whose request the gateway withdrew (request.cancel by id)', () => { const onEvent = createGatewayEventHandler(buildCtx([])) - patchOverlayState({ - secret: { envVar: 'NEW_KEY', prompt: 'Enter new key', requestId: 'secret-new' }, - sudo: { requestId: 'sudo-1' } - }) + serverRequest('secret', { env_var: 'NEW_KEY', prompt: 'Enter new key' }, 'secret-new') + serverRequest('sudo', {}, 'sudo-1') - onEvent({ payload: { request_id: 'secret-old' }, type: 'secret.expire' } as any) + onEvent({ payload: { id: 'secret-old', method: 'secret', reason: 'timeout' }, type: 'request.cancel' } as any) expect(getOverlayState().secret?.requestId).toBe('secret-new') - onEvent({ payload: { request_id: 'secret-new' }, type: 'secret.expire' } as any) + onEvent({ payload: { id: 'secret-new', method: 'secret', reason: 'timeout' }, type: 'request.cancel' } as any) expect(getOverlayState().secret).toBeNull() + expect(getOverlayState().sudo?.requestId).toBe('sudo-1') - onEvent({ payload: { request_id: 'sudo-1' }, type: 'sudo.expire' } as any) + onEvent({ payload: { id: 'sudo-1', method: 'sudo', reason: 'interrupted' }, type: 'request.cancel' } as any) expect(getOverlayState().sudo).toBeNull() }) // ── Batch (multi-question) clarify ───────────────────────────────── - it('parses a batch clarify.request into a questions overlay', () => { - const onEvent = createGatewayEventHandler(buildCtx([])) - - onEvent({ - payload: { + it('parses a batch clarify request into a questions overlay', () => { + serverRequest( + 'clarify', + { questions: [ { choices: ['a', 'b'], qid: 'q0', question: 'One?' }, { choices: null, qid: 'q1', question: 'Two?' } - ], - request_id: 'req-batch' + ] }, - type: 'clarify.request' - } as any) + 'req-batch' + ) const clarify = getOverlayState().clarify expect(clarify?.requestId).toBe('req-batch') @@ -1671,39 +1674,35 @@ describe('createGatewayEventHandler', () => { expect(clarify?.answers).toEqual({}) }) - it('seeds locked answers from a reconnect-replay batch clarify.request', () => { - const onEvent = createGatewayEventHandler(buildCtx([])) - - onEvent({ - payload: { + it('seeds locked answers from a reconnect-replayed batch clarify request', () => { + serverRequest( + 'clarify', + { answers: { q0: 'a' }, questions: [ { choices: ['a', 'b'], qid: 'q0', question: 'One?' }, { choices: null, qid: 'q1', question: 'Two?' } - ], - request_id: 'req-replay' + ] }, - type: 'clarify.request' - } as any) + 'req-replay' + ) expect(getOverlayState().clarify?.answers).toEqual({ q0: 'a' }) }) it('drops malformed batch entries and falls back to single-question shape when none survive', () => { - const onEvent = createGatewayEventHandler(buildCtx([])) - - onEvent({ - payload: { + serverRequest( + 'clarify', + { choices: ['x', 'y'], question: 'Fallback?', questions: [ { qid: '', question: 'no qid' }, { qid: 'q1', question: ' ' } - ], - request_id: 'req-bad' + ] }, - type: 'clarify.request' - } as any) + 'req-bad' + ) const clarify = getOverlayState().clarify expect(clarify?.questions).toBeUndefined() diff --git a/ui-tui/src/__tests__/useInputHandlers.test.ts b/ui-tui/src/__tests__/useInputHandlers.test.ts index 5fe403dfd8..96778fa635 100644 --- a/ui-tui/src/__tests__/useInputHandlers.test.ts +++ b/ui-tui/src/__tests__/useInputHandlers.test.ts @@ -1,6 +1,7 @@ import { describe, expect, it, vi } from 'vitest' import { getOverlayState, patchOverlayState, resetOverlayState } from '../app/overlayStore.js' +import { rememberServerRequest, resetServerRequestsForTests } from '../app/serverRequestStore.js' import { applyVoiceRecordResponse, dismissSensitivePrompt, @@ -156,31 +157,37 @@ describe('applyVoiceRecordResponse', () => { }) describe('dismissSensitivePrompt', () => { - it('clears a sudo overlay before a stale cancel RPC resolves', async () => { + const openRequest = (id: string, method: string) => { + const respond = vi.fn() + + rememberServerRequest({ fail: vi.fn(), id, method, params: {}, respond }) + + return respond + } + + it('clears a sudo overlay and answers the server request with an empty value', () => { resetOverlayState() - patchOverlayState({ sudo: { requestId: 'sudo-1' } }) - const rpc = vi.fn().mockResolvedValue(null) + resetServerRequestsForTests() + patchOverlayState({ sudo: { requestId: 'srq-sudo' } }) + const respond = openRequest('srq-sudo', 'sudo') const sys = vi.fn() - const pending = dismissSensitivePrompt(getOverlayState(), rpc, sys) + dismissSensitivePrompt(getOverlayState(), vi.fn(), sys) expect(getOverlayState().sudo).toBeNull() expect(sys).toHaveBeenCalledWith('sudo cancelled') - expect(rpc).toHaveBeenCalledWith('sudo.respond', { password: '', request_id: 'sudo-1' }) - await pending + expect(respond).toHaveBeenCalledWith({ value: '' }) }) - it('clears a secret overlay before a stale cancel RPC resolves', async () => { + it('clears a secret overlay even when its request already expired (nothing left to answer)', () => { resetOverlayState() - patchOverlayState({ secret: { envVar: 'API_KEY', prompt: 'Enter API key', requestId: 'secret-1' } }) - const rpc = vi.fn().mockResolvedValue(null) + resetServerRequestsForTests() + patchOverlayState({ secret: { envVar: 'API_KEY', prompt: 'Enter API key', requestId: 'srq-gone' } }) const sys = vi.fn() - const pending = dismissSensitivePrompt(getOverlayState(), rpc, sys) + dismissSensitivePrompt(getOverlayState(), vi.fn(), sys) expect(getOverlayState().secret).toBeNull() expect(sys).toHaveBeenCalledWith('secret entry cancelled') - expect(rpc).toHaveBeenCalledWith('secret.respond', { request_id: 'secret-1', value: '' }) - await pending }) }) diff --git a/ui-tui/src/app/createGatewayEventHandler.ts b/ui-tui/src/app/createGatewayEventHandler.ts index 5bae7abbb7..c13f1af765 100644 --- a/ui-tui/src/app/createGatewayEventHandler.ts +++ b/ui-tui/src/app/createGatewayEventHandler.ts @@ -30,8 +30,8 @@ import type { Msg, SessionInfo, SubagentProgress } from '../types.js' import { applyDelegationStatus, getDelegationState } from './delegationStore.js' import type { GatewayEventHandlerContext, NoticeLevel } from './interfaces.js' import { getOverlayState, patchOverlayState } from './overlayStore.js' -import { forgetServerRequest } from './serverRequestStore.js' import { flashGoodVibes, flashPet } from './petFlashStore.js' +import { forgetServerRequest } from './serverRequestStore.js' import { turnController } from './turnController.js' import { getTurnState } from './turnStore.js' import { getUiState, patchUiState } from './uiStore.js' diff --git a/ui-tui/src/app/useMainApp.ts b/ui-tui/src/app/useMainApp.ts index 7b8279844c..df15d49489 100644 --- a/ui-tui/src/app/useMainApp.ts +++ b/ui-tui/src/app/useMainApp.ts @@ -938,6 +938,7 @@ export function useMainApp(gw: GatewayClient) { useEffect(() => { const handler = (ev: AnyGatewayEvent) => onEventRef.current(ev) + const requestHandler = (request: ServerRequest) => { if (!onServerRequestRef.current(request)) { request.fail(JSON_RPC_METHOD_NOT_FOUND, `the terminal UI cannot answer ${request.method}`) diff --git a/website/docs/developer-guide/programmatic-integration.md b/website/docs/developer-guide/programmatic-integration.md index 449c009d79..e74662d7f3 100644 --- a/website/docs/developer-guide/programmatic-integration.md +++ b/website/docs/developer-guide/programmatic-integration.md @@ -46,8 +46,7 @@ session.create session.list session.active_list session.activate session.close session.interrupt session.history session.compress session.branch session.title session.usage session.status -clarify.respond sudo.respond secret.respond -approval.respond config.set / config.get commands.catalog +clarify.lock config.set / config.get commands.catalog command.resolve command.dispatch cli.exec reload.mcp reload.env process.stop delegation.status subagent.interrupt subagent.steer @@ -76,7 +75,20 @@ On a successful truncating submit against a durable session, the `prompt.submit` ### Events streamed back -`message.delta`, `message.complete`, `tool.start`, `tool.generating`, `tool.complete`, `approval.request`, `clarify.request`, `sudo.request`, `sudo.expire`, `secret.request`, `secret.expire`, `gateway.ready`, plus session lifecycle and error events. Expiry events carry the original `{ request_id }`; external hosts should clear only the matching pending prompt. +`message.delta`, `message.complete`, `tool.start`, `tool.generating`, `tool.complete`, `gateway.ready`, `request.cancel`, plus session lifecycle and error events. + +### Server→client requests (questions the agent asks you) + +Approvals, clarify questions, sudo/secret prompts, vault unlock, MCP setup and the desktop read/act bridges are **JSON-RPC requests from the gateway to the client**, not events. The frame carries a string id, and the client answers with a normal JSON-RPC response bearing the same id: + +``` +← {"jsonrpc":"2.0","id":"srq-7","method":"approval","params":{"session_id":"…","request_id":"…","command":"rm -rf build","description":"…"}} +→ {"jsonrpc":"2.0","id":"srq-7","result":{"choice":"once"}} +``` + +Methods: `approval` → `{choice}`; `clarify` → `{answer}` (single) or `{answers}` / `{}` cancel (batch, with `clarify.lock` to lock one answer early); `sudo`, `secret`, `vault.code`, `vault.unlock` → `{value}`; `mcp.setup` → `{result}`; `terminal.read`, `window.read`, `preview.act`, `tour` → `{value}` (JSON text). Respond with a JSON-RPC error (`-32601`) for a method your host does not implement so the agent fails fast instead of waiting out the timeout. + +When the gateway withdraws a question (timeout, interrupt, answered from another surface) it emits `request.cancel` `{ id, method, reason }`; clear only the matching prompt. `session.resume` / `session.activate` results and `session.events.since` carry `open_requests` — the still-open frames — so a reconnecting client re-renders (and can still answer) them. ### Pi-style RPC mapping @@ -94,7 +106,7 @@ Every command in the Pi-mono RPC spec ([issue #360](https://github.com/NousResea | `get_messages` | `session.history` | | `switch_session` | `session.resume` | | `fork` | `session.branch` | -| `ui_request` / `ui_response` | `clarify.respond` / `sudo.respond` / `secret.respond` / `approval.respond` | +| `ui_request` / `ui_response` | server→client requests `clarify` / `sudo` / `secret` / `approval` answered by JSON-RPC response frames | ---