feat(gateway): server→client JSON-RPC requests replace the *.request/*.respond event pairs (#110521)

The gateway asked the user questions (approval, clarify, sudo, secret,
vault, MCP setup, the desktop read/act bridges) by emitting a
`<x>.request` EVENT carrying a hand-minted request_id, blocking the
agent thread on a module dict keyed by that id, and exposing a paired
`<x>.respond` METHOD per kind — thirteen pairs, four registries
(`_pending`, `_answers`, `_batch_clarify`, `_EXPIRING_REQUESTS`) and a
per-kind reconnect snapshot (`pending_clarify` / `pending_approval`)
that only two of the thirteen kinds ever got. JSON-RPC already has the
primitive: the server sends a request frame with an id and the client
answers with a response frame bearing the same id.

`tui_gateway/server_requests.py` owns the one mechanism:

  send()          block the agent thread until the response frame
                  (`srq-<n>` ids; ints belong to the client)
  send_async()    fire-and-callback variant (bot relay)
  cancel*()       withdraw with ONE `request.cancel {id, method, reason}`
                  event (timeout / interrupt / process exit /
                  answered elsewhere) instead of per-kind *.expire
  open_requests() the still-open frames, replayed by session.resume,
                  session.activate and session.events.since so a
                  reconnecting client re-renders every kind, not two
  clarify.lock    stays a real client→server RPC (locks one batch
                  answer early); locked answers merge into the final
                  set even when the closing response carries only the
                  tail the user answered last

A client that does not implement a method answers -32601 and the agent
fails fast (the old fixed-timeout "unavailable" probes for tour/preview
still work — a wire error IS an answer). Approval: the queue entry's
settle hook withdraws the request when `/approve` from another surface,
a timeout or an interrupt resolves it first, so no window keeps a dead
card. Compute-host children own their waits; the parent mirrors their
open frames for replay and relays `clarify.lock` + response frames.

Clients: `JsonRpcRequestChannel` gains `onRequest` (unhandled → -32601,
dedup by id) and `JsonRpcGatewayClient` re-delivers `open_requests`
from the replay result. Desktop gets `gateway-event/server-requests.ts`
(one handler per method, replacing the request branches of
`input-requests.ts` / `desktop-bridge.ts`) and a `store/server-requests`
registry so every answer site calls `respondToServerRequest(id, result)`
synchronously; the TUI gets `createServerRequestHandler.ts` +
`serverRequestStore.ts`. `gateway-events.json` now pins both halves
(events + server request methods); the two contract tests check both.

Live (real stdio gateway, real `clarify_callback` on the agent thread):
before, `clarify.request` event + `clarify.respond` RPC, batch final
answers lost ('' returned); after, `{"id":"srq-…","method":"clarify"}`
frame, `session.events.since.open_requests` replays it, response frame
`{"answer":"yes"}` reaches the agent, batch lock + final response
merge to `{"q0":"1","q1":"free text"}`.
This commit is contained in:
teknium1
2026-09-14 00:24:50 -07:00
committed by Teknium
parent ebe8cda8ea
commit 9f7f2f28c0
42 changed files with 588 additions and 835 deletions
@@ -258,6 +258,7 @@ function Harness({
useGatewayBoot({
beforeConnectionSwitch,
handleGatewayEvent: () => undefined,
handleServerRequest: () => false,
onConnectionReady: () => undefined,
onGatewayReady: () => undefined,
refreshHermesConfig,
@@ -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))
@@ -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<string, unknown>) =>
act(() => stream.handleEvent({ payload, session_id: SID, type: 'clarify.request' }))
const clarifyRequest = ({ request_id, ...params }: Record<string, unknown>) =>
act(() => void stream.handleRequest('clarify', { ...params, session_id: SID }, request_id as string))
const toolStart = (payload: Record<string, unknown>) =>
act(() => stream.handleEvent({ payload, session_id: SID, type: 'tool.start' }))
@@ -34,7 +34,13 @@ const toolComplete = (payload: Record<string, unknown>) =>
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' })
@@ -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
})
})
})
})
@@ -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<string, unknown>, 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 })
})
})
})
@@ -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<Parameters<typeof u
export interface MessageStreamHarness {
/** Feed a gateway event into the mounted hook. */
handleEvent: (event: GatewayEvent) => 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<string, unknown>, id?: string) => ReturnType<typeof vi.fn>
/** 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<string, ClientSessionState>(), ...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')
@@ -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': {
@@ -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'
@@ -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'
@@ -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'
@@ -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
+7 -3
View File
@@ -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<typeof $gateway.get>)
})
@@ -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)
-1
View File
@@ -1,6 +1,5 @@
import { atom, computed } from 'nanostores'
import { $gateway } from './gateway'
import { respondToServerRequest } from './server-requests'
import { $activeSessionId } from './session'
+4 -3
View File
@@ -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)
-1
View File
@@ -1,6 +1,5 @@
import { atom, computed } from 'nanostores'
import { $gateway } from './gateway'
import { respondToServerRequest } from './server-requests'
/**
@@ -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 }
+3 -3
View File
@@ -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
})
}
+46
View File
@@ -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 }])
})
})
+20 -29
View File
@@ -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": "<write to AGENTS.md>",
},
{"allow_permanent": allow_permanent, "allow_session": allow_session, "command": "<write to AGENTS.md>"},
)
assert emitted["payload"]["choices"] == expected
assert sent["params"]["choices"] == expected
+21 -23
View File
@@ -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",
+1 -2
View File
@@ -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()
@@ -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):
@@ -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()
+1 -2
View File
@@ -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()
+125 -281
View File
@@ -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 ───────────────────────────────────────────────────
@@ -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):
+1 -2
View File
@@ -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()
+1 -2
View File
@@ -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
@@ -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()
@@ -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()
@@ -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
+69 -163
View File
@@ -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
+1 -2
View File
@@ -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
+14 -6
View File
@@ -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_<topic>.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-<n>` 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_<topic>.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("<method>", 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` |
+5
View File
@@ -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()
+17 -4
View File
@@ -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)
+10 -11
View File
@@ -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 }` |
@@ -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<string, unknown>, 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()
+19 -12
View File
@@ -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
})
})
+1 -1
View File
@@ -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'
+1
View File
@@ -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}`)
@@ -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 |
---