mirror of
https://github.com/stablyai/orca.git
synced 2026-10-10 00:02:35 +00:00
fix(runtime): serve on-demand tab mounts to the phone from the PTY-grid model
A phone opening a tab not mounted since relaunch makes the desktop mount it. The freshly mounted pane is hidden, so it answers at desktop size, and both the mount-wait snapshot and the late renderer-mount-ready recovery published that raw screen and seeded the model at that size. The recovery seed now reflows the model onto the PTY grid, and both sites publish the model's snapshot through one shared step, so the phone is only served from the host model. A late recovery whose model moved past the renderer high-water while it was rebuilt is still withheld. Claude-Session: lane-host-hydrate
This commit is contained in:
@@ -0,0 +1,192 @@
|
||||
/**
|
||||
* A real OrcaRuntimeService + real legacy `terminal.subscribe` driven by a phone at 47x40, over a
|
||||
* daemon PTY reattached after a relaunch: no host model yet, the PTY at desktop size, and a desktop
|
||||
* pane that answers its serializer at desktop size because a hidden pane never refits.
|
||||
*/
|
||||
import { expect, vi } from 'vitest'
|
||||
import { OrcaRuntimeService } from './orca-runtime'
|
||||
import { RpcDispatcher } from './rpc/dispatcher'
|
||||
import type { RpcRequest } from './rpc/core'
|
||||
import { TERMINAL_METHODS } from './rpc/methods/terminal'
|
||||
import {
|
||||
TerminalStreamOpcode,
|
||||
decodeTerminalStreamFrame,
|
||||
decodeTerminalStreamJson,
|
||||
decodeTerminalStreamText
|
||||
} from '../../shared/terminal-stream-protocol'
|
||||
import { HeadlessEmulator } from '../daemon/headless-emulator'
|
||||
|
||||
const WORKTREE_ID = 'repo-1::/tmp/wt'
|
||||
export const PTY_ID = `${WORKTREE_ID}@@9c8d7e6f`
|
||||
export const DESKTOP = { cols: 200, rows: 50 }
|
||||
export const PHONE = { cols: 47, rows: 40 }
|
||||
const LONG_LINE = 'A'.repeat(120)
|
||||
export const EXPECTED_PHONE_ROWS = ['A'.repeat(47), 'A'.repeat(47), 'A'.repeat(26), '$ prompt']
|
||||
|
||||
type RuntimeInternals = {
|
||||
recordPtyWorktree: (ptyId: string, worktreeId: string, state?: { connected?: boolean }) => unknown
|
||||
issuePtyHandle: (pty: unknown) => string
|
||||
providerSnapshotPreferredPtys: Set<string>
|
||||
}
|
||||
|
||||
export function internals(runtime: OrcaRuntimeService): RuntimeInternals {
|
||||
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: test reaches protected members the runtime defines.
|
||||
return runtime as unknown as RuntimeInternals
|
||||
}
|
||||
|
||||
async function screenOn(grid: { cols: number; rows: number }): Promise<string> {
|
||||
const emulator = new HeadlessEmulator({ ...grid, scrollback: 1000 })
|
||||
try {
|
||||
await emulator.write(`${LONG_LINE}\r\n$ prompt`)
|
||||
const snapshot = emulator.getSnapshot()
|
||||
return snapshot.rehydrateSequences + snapshot.snapshotAnsi
|
||||
} finally {
|
||||
emulator.dispose()
|
||||
}
|
||||
}
|
||||
|
||||
export type PhoneSubscribeSetup = {
|
||||
paneMounted: boolean
|
||||
providerSnapshot?: boolean
|
||||
repaintOnResize?: boolean
|
||||
/** How a requested tab mount lands: at once, or when the test calls `finishMount`. */
|
||||
mount?: 'ready' | 'late'
|
||||
}
|
||||
|
||||
export function setupPhoneSubscribe(opts: PhoneSubscribeSetup) {
|
||||
const sizes = new Map([[PTY_ID, { ...DESKTOP }]])
|
||||
let paneMounted = opts.paneMounted
|
||||
let settleMount: (ready: boolean) => void = () => {}
|
||||
const mountSettled = new Promise<boolean>((resolve) => {
|
||||
settleMount = resolve
|
||||
})
|
||||
const runtime = new OrcaRuntimeService()
|
||||
// The pane orders its screen against PTY output, as a mounted desktop xterm does.
|
||||
const serializeBuffer = vi.fn(async () =>
|
||||
paneMounted
|
||||
? { data: await screenOn(DESKTOP), ...DESKTOP, seq: runtime.getPtyOutputSequence(PTY_ID) }
|
||||
: null
|
||||
)
|
||||
runtime.setPtyController({
|
||||
write: () => true,
|
||||
kill: () => true,
|
||||
getForegroundProcess: async () => null,
|
||||
getSize: (ptyId: string) => sizes.get(ptyId) ?? null,
|
||||
resize: (ptyId: string, cols: number, rows: number) => {
|
||||
sizes.set(ptyId, { cols, rows })
|
||||
if (opts.repaintOnResize) {
|
||||
// A TUI answering SIGWINCH before the subscribe reaches its own hydrate.
|
||||
runtime.onPtyData(ptyId, '\x1b[?25h', Date.now())
|
||||
}
|
||||
return true
|
||||
},
|
||||
hasRendererSerializer: () => paneMounted,
|
||||
getRendererSerializerGeneration: () => (paneMounted ? 2 : 1),
|
||||
waitForRendererSerializer: () => (opts.mount ? mountSettled : Promise.resolve(false)),
|
||||
serializeBuffer,
|
||||
// The daemon resizes its model with the PTY.
|
||||
serializeProviderBuffer: async (ptyId: string) => {
|
||||
const grid = sizes.get(ptyId) ?? DESKTOP
|
||||
return opts.providerSnapshot
|
||||
? { data: await screenOn(grid), ...grid, seq: 0, source: 'headless' as const }
|
||||
: null
|
||||
}
|
||||
})
|
||||
const finishMount = () => {
|
||||
paneMounted = true
|
||||
settleMount(true)
|
||||
}
|
||||
const requestMount = vi
|
||||
.spyOn(runtime, 'requestRendererTerminalTabMount')
|
||||
.mockImplementation(() => {
|
||||
if (opts.mount === 'ready') {
|
||||
finishMount()
|
||||
}
|
||||
return opts.mount !== undefined
|
||||
})
|
||||
const record = internals(runtime).recordPtyWorktree(PTY_ID, WORKTREE_ID, { connected: true })
|
||||
const handle = internals(runtime).issuePtyHandle(record)
|
||||
return { runtime, handle, sizes, serializeBuffer, requestMount, finishMount }
|
||||
}
|
||||
|
||||
export type PublishedSnapshot = {
|
||||
kind: string
|
||||
reason?: string
|
||||
cols: number
|
||||
rows: number
|
||||
data: string
|
||||
}
|
||||
|
||||
export function subscribePhone(runtime: OrcaRuntimeService, handle: string) {
|
||||
const frames: Uint8Array<ArrayBufferLike>[] = []
|
||||
const controller = new AbortController()
|
||||
const request: RpcRequest = {
|
||||
id: 'req-phone',
|
||||
authToken: 'tok',
|
||||
method: 'terminal.subscribe',
|
||||
params: {
|
||||
terminal: handle,
|
||||
client: { id: 'phone-1', type: 'mobile' },
|
||||
viewport: PHONE,
|
||||
capabilities: { terminalBinaryStream: 1 }
|
||||
}
|
||||
}
|
||||
const done = new RpcDispatcher({ runtime, methods: TERMINAL_METHODS }).dispatchStreaming(
|
||||
request,
|
||||
() => {},
|
||||
{
|
||||
connectionId: 'conn-phone',
|
||||
signal: controller.signal,
|
||||
sendBinary: (bytes) => {
|
||||
frames.push(bytes)
|
||||
},
|
||||
registerBinaryStreamHandler: () => () => {}
|
||||
}
|
||||
)
|
||||
/** Every complete snapshot published so far, in order. */
|
||||
const snapshots = (): PublishedSnapshot[] => {
|
||||
const published: PublishedSnapshot[] = []
|
||||
let open: PublishedSnapshot | null = null
|
||||
for (const frame of frames.flatMap((bytes) => decodeTerminalStreamFrame(bytes) ?? [])) {
|
||||
if (frame.opcode === TerminalStreamOpcode.SnapshotStart) {
|
||||
const meta = decodeTerminalStreamJson<Omit<PublishedSnapshot, 'data'>>(frame.payload)
|
||||
open = meta ? { ...meta, data: '' } : null
|
||||
} else if (frame.opcode === TerminalStreamOpcode.SnapshotChunk && open) {
|
||||
open.data += decodeTerminalStreamText(frame.payload)
|
||||
} else if (frame.opcode === TerminalStreamOpcode.SnapshotEnd && open) {
|
||||
published.push(open)
|
||||
open = null
|
||||
}
|
||||
}
|
||||
return published
|
||||
}
|
||||
const close = async () => {
|
||||
runtime.cleanupSubscription(`${handle}:phone-1`)
|
||||
controller.abort()
|
||||
await done.catch(() => {})
|
||||
}
|
||||
return { snapshots, close }
|
||||
}
|
||||
|
||||
/** What a phone xterm sized to the frame's declared grid shows after replaying it. */
|
||||
export async function paintedRows(snapshot: { cols: number; rows: number; data: string }) {
|
||||
const emulator = new HeadlessEmulator({ cols: snapshot.cols, rows: snapshot.rows, scrollback: 0 })
|
||||
try {
|
||||
await emulator.write(snapshot.data)
|
||||
return emulator.getVisibleLines().map((line) => line.trimEnd())
|
||||
} finally {
|
||||
emulator.dispose()
|
||||
}
|
||||
}
|
||||
|
||||
export async function firstSnapshot(
|
||||
runtime: OrcaRuntimeService,
|
||||
handle: string,
|
||||
timeout = 1_000
|
||||
): Promise<PublishedSnapshot> {
|
||||
const subscription = subscribePhone(runtime, handle)
|
||||
await vi.waitFor(() => expect(subscription.snapshots().length).toBeGreaterThan(0), { timeout })
|
||||
const [snapshot] = subscription.snapshots()
|
||||
await subscription.close()
|
||||
return snapshot
|
||||
}
|
||||
@@ -1,167 +1,21 @@
|
||||
/**
|
||||
* A phone subscribing to an idle reattached PTY must get its first snapshot on the phone grid.
|
||||
*
|
||||
* Harness: real OrcaRuntimeService + real legacy `terminal.subscribe`. After a relaunch the
|
||||
* reattach skips seeding the host model because a pane is mounted; the PTY then emits no byte
|
||||
* (an agent waiting for input), and the pane, hidden in another workspace, answers its
|
||||
* serializer at desktop size.
|
||||
* After a relaunch the reattach skips seeding the host model because a pane is mounted; the PTY
|
||||
* then emits no byte (an agent waiting for input), and the pane, hidden in another workspace,
|
||||
* answers its serializer at desktop size.
|
||||
*/
|
||||
import { describe, expect, it, vi } from 'vitest'
|
||||
import { OrcaRuntimeService } from './orca-runtime'
|
||||
import { RpcDispatcher } from './rpc/dispatcher'
|
||||
import type { RpcRequest } from './rpc/core'
|
||||
import { TERMINAL_METHODS } from './rpc/methods/terminal'
|
||||
import {
|
||||
TerminalStreamOpcode,
|
||||
decodeTerminalStreamFrame,
|
||||
decodeTerminalStreamJson,
|
||||
decodeTerminalStreamText
|
||||
} from '../../shared/terminal-stream-protocol'
|
||||
import { HeadlessEmulator } from '../daemon/headless-emulator'
|
||||
|
||||
const WORKTREE_ID = 'repo-1::/tmp/wt'
|
||||
const PTY_ID = `${WORKTREE_ID}@@9c8d7e6f`
|
||||
const DESKTOP = { cols: 200, rows: 50 }
|
||||
const PHONE = { cols: 47, rows: 40 }
|
||||
const LONG_LINE = 'A'.repeat(120)
|
||||
const EXPECTED_PHONE_ROWS = ['A'.repeat(47), 'A'.repeat(47), 'A'.repeat(26), '$ prompt']
|
||||
|
||||
type RuntimeInternals = {
|
||||
recordPtyWorktree: (ptyId: string, worktreeId: string, state?: { connected?: boolean }) => unknown
|
||||
issuePtyHandle: (pty: unknown) => string
|
||||
providerSnapshotPreferredPtys: Set<string>
|
||||
}
|
||||
|
||||
function internals(runtime: OrcaRuntimeService): RuntimeInternals {
|
||||
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: test reaches protected members the runtime defines.
|
||||
return runtime as unknown as RuntimeInternals
|
||||
}
|
||||
|
||||
async function screenOn(grid: { cols: number; rows: number }): Promise<string> {
|
||||
const emulator = new HeadlessEmulator({ ...grid, scrollback: 1000 })
|
||||
try {
|
||||
await emulator.write(`${LONG_LINE}\r\n$ prompt`)
|
||||
const snapshot = emulator.getSnapshot()
|
||||
return snapshot.rehydrateSequences + snapshot.snapshotAnsi
|
||||
} finally {
|
||||
emulator.dispose()
|
||||
}
|
||||
}
|
||||
|
||||
function setup(opts: {
|
||||
paneMounted: boolean
|
||||
providerSnapshot?: boolean
|
||||
repaintOnResize?: boolean
|
||||
}) {
|
||||
const sizes = new Map([[PTY_ID, { ...DESKTOP }]])
|
||||
const runtime = new OrcaRuntimeService()
|
||||
// The pane orders its screen against PTY output, as a mounted desktop xterm does.
|
||||
const serializeBuffer = vi.fn(async () =>
|
||||
opts.paneMounted
|
||||
? { data: await screenOn(DESKTOP), ...DESKTOP, seq: runtime.getPtyOutputSequence(PTY_ID) }
|
||||
: null
|
||||
)
|
||||
runtime.setPtyController({
|
||||
write: () => true,
|
||||
kill: () => true,
|
||||
getForegroundProcess: async () => null,
|
||||
getSize: (ptyId: string) => sizes.get(ptyId) ?? null,
|
||||
resize: (ptyId: string, cols: number, rows: number) => {
|
||||
sizes.set(ptyId, { cols, rows })
|
||||
if (opts.repaintOnResize) {
|
||||
// A TUI answering SIGWINCH before the subscribe reaches its own hydrate.
|
||||
runtime.onPtyData(ptyId, '\x1b[?25h', Date.now())
|
||||
}
|
||||
return true
|
||||
},
|
||||
hasRendererSerializer: () => opts.paneMounted,
|
||||
getRendererSerializerGeneration: () => 1,
|
||||
waitForRendererSerializer: async () => false,
|
||||
serializeBuffer,
|
||||
// The daemon resizes its model with the PTY.
|
||||
serializeProviderBuffer: async (ptyId: string) => {
|
||||
const grid = sizes.get(ptyId) ?? DESKTOP
|
||||
return opts.providerSnapshot
|
||||
? { data: await screenOn(grid), ...grid, seq: 0, source: 'headless' as const }
|
||||
: null
|
||||
}
|
||||
})
|
||||
const requestMount = vi.spyOn(runtime, 'requestRendererTerminalTabMount').mockReturnValue(false)
|
||||
const record = internals(runtime).recordPtyWorktree(PTY_ID, WORKTREE_ID, { connected: true })
|
||||
const handle = internals(runtime).issuePtyHandle(record)
|
||||
return { runtime, handle, sizes, serializeBuffer, requestMount }
|
||||
}
|
||||
|
||||
function subscribePhone(runtime: OrcaRuntimeService, handle: string) {
|
||||
const frames: Uint8Array<ArrayBufferLike>[] = []
|
||||
const controller = new AbortController()
|
||||
const request: RpcRequest = {
|
||||
id: 'req-phone',
|
||||
authToken: 'tok',
|
||||
method: 'terminal.subscribe',
|
||||
params: {
|
||||
terminal: handle,
|
||||
client: { id: 'phone-1', type: 'mobile' },
|
||||
viewport: PHONE,
|
||||
capabilities: { terminalBinaryStream: 1 }
|
||||
}
|
||||
}
|
||||
const done = new RpcDispatcher({ runtime, methods: TERMINAL_METHODS }).dispatchStreaming(
|
||||
request,
|
||||
() => {},
|
||||
{
|
||||
connectionId: 'conn-phone',
|
||||
signal: controller.signal,
|
||||
sendBinary: (bytes) => {
|
||||
frames.push(bytes)
|
||||
},
|
||||
registerBinaryStreamHandler: () => () => {}
|
||||
}
|
||||
)
|
||||
const decoded = () => frames.flatMap((bytes) => decodeTerminalStreamFrame(bytes) ?? [])
|
||||
const snapshot = () => {
|
||||
const start = decoded().find((frame) => frame.opcode === TerminalStreamOpcode.SnapshotStart)
|
||||
const meta = start
|
||||
? decodeTerminalStreamJson<{ cols: number; rows: number }>(start.payload)
|
||||
: null
|
||||
if (!meta) {
|
||||
return null
|
||||
}
|
||||
const data = decoded()
|
||||
.filter((frame) => frame.opcode === TerminalStreamOpcode.SnapshotChunk)
|
||||
.map((frame) => decodeTerminalStreamText(frame.payload))
|
||||
.join('')
|
||||
return { cols: meta.cols, rows: meta.rows, data }
|
||||
}
|
||||
const close = async () => {
|
||||
runtime.cleanupSubscription(`${handle}:phone-1`)
|
||||
controller.abort()
|
||||
await done.catch(() => {})
|
||||
}
|
||||
return { snapshot, close }
|
||||
}
|
||||
|
||||
/** What a phone xterm sized to the frame's declared grid shows after replaying it. */
|
||||
async function paintedRows(snapshot: { cols: number; rows: number; data: string }) {
|
||||
const emulator = new HeadlessEmulator({ cols: snapshot.cols, rows: snapshot.rows, scrollback: 0 })
|
||||
try {
|
||||
await emulator.write(snapshot.data)
|
||||
return emulator.getVisibleLines().map((line) => line.trimEnd())
|
||||
} finally {
|
||||
emulator.dispose()
|
||||
}
|
||||
}
|
||||
|
||||
async function firstSnapshot(runtime: OrcaRuntimeService, handle: string) {
|
||||
const subscription = subscribePhone(runtime, handle)
|
||||
await vi.waitFor(() => expect(subscription.snapshot()).not.toBeNull())
|
||||
const snapshot = subscription.snapshot()
|
||||
await subscription.close()
|
||||
if (!snapshot) {
|
||||
throw new Error('no snapshot frame')
|
||||
}
|
||||
return snapshot
|
||||
}
|
||||
EXPECTED_PHONE_ROWS,
|
||||
PHONE,
|
||||
PTY_ID,
|
||||
firstSnapshot,
|
||||
internals,
|
||||
paintedRows,
|
||||
setupPhoneSubscribe as setup,
|
||||
subscribePhone
|
||||
} from './mobile-phone-subscribe-test-fixture'
|
||||
|
||||
describe('phone subscribe to an idle PTY whose hidden pane sits at desktop size', () => {
|
||||
it('serves the first snapshot from a host model hydrated onto the phone grid', async () => {
|
||||
@@ -228,7 +82,9 @@ describe('phone subscribe to an idle PTY whose hidden pane sits at desktop size'
|
||||
const { runtime, handle, requestMount } = setup({ paneMounted: false })
|
||||
|
||||
const subscription = subscribePhone(runtime, handle)
|
||||
await vi.waitFor(() => expect(subscription.snapshot()).not.toBeNull(), { timeout: 5_000 })
|
||||
await vi.waitFor(() => expect(subscription.snapshots().length).toBeGreaterThan(0), {
|
||||
timeout: 5_000
|
||||
})
|
||||
await subscription.close()
|
||||
|
||||
expect(requestMount).toHaveBeenCalledTimes(1)
|
||||
|
||||
@@ -0,0 +1,53 @@
|
||||
/**
|
||||
* A phone subscribing to a tab not mounted since relaunch must get the PTY's grid, not the pane's.
|
||||
*
|
||||
* With no renderer serializer the subscribe asks the desktop to mount the tab and waits for it.
|
||||
* The freshly mounted pane is hidden, so it answers its serializer at desktop size.
|
||||
*/
|
||||
import { describe, expect, it, vi } from 'vitest'
|
||||
import {
|
||||
EXPECTED_PHONE_ROWS,
|
||||
PHONE,
|
||||
PTY_ID,
|
||||
firstSnapshot,
|
||||
paintedRows,
|
||||
setupPhoneSubscribe,
|
||||
subscribePhone
|
||||
} from './mobile-phone-subscribe-test-fixture'
|
||||
|
||||
describe('phone subscribe to a tab mounted on demand at desktop size', () => {
|
||||
it('serves the mount-wait snapshot from the model on the phone grid', async () => {
|
||||
const { runtime, handle, requestMount } = setupPhoneSubscribe({
|
||||
paneMounted: false,
|
||||
mount: 'ready'
|
||||
})
|
||||
|
||||
const snapshot = await firstSnapshot(runtime, handle)
|
||||
|
||||
expect(requestMount).toHaveBeenCalledTimes(1)
|
||||
expect({ cols: snapshot.cols, rows: snapshot.rows }).toEqual(PHONE)
|
||||
expect((await paintedRows(snapshot)).slice(0, 4)).toEqual(EXPECTED_PHONE_ROWS)
|
||||
const resubscribed = await firstSnapshot(runtime, handle)
|
||||
expect({ cols: resubscribed.cols, rows: resubscribed.rows }).toEqual(PHONE)
|
||||
expect((await paintedRows(resubscribed)).slice(0, 4)).toEqual(EXPECTED_PHONE_ROWS)
|
||||
})
|
||||
|
||||
it('publishes a late mount-ready recovery on the phone grid', async () => {
|
||||
const { runtime, handle, finishMount } = setupPhoneSubscribe({
|
||||
paneMounted: false,
|
||||
mount: 'late'
|
||||
})
|
||||
const subscription = subscribePhone(runtime, handle)
|
||||
// The bounded initial response gives up on the mount before it settles.
|
||||
await vi.waitFor(() => expect(subscription.snapshots().length).toBe(1), { timeout: 5_000 })
|
||||
|
||||
finishMount()
|
||||
await vi.waitFor(() => expect(subscription.snapshots().length).toBe(2))
|
||||
const recovery = subscription.snapshots()[1]
|
||||
await subscription.close()
|
||||
|
||||
expect(recovery).toMatchObject({ kind: 'resized', reason: 'renderer-mount-ready', ...PHONE })
|
||||
expect((await paintedRows(recovery)).slice(0, 4)).toEqual(EXPECTED_PHONE_ROWS)
|
||||
expect(runtime.hasHeadlessTerminalState(PTY_ID)).toBe(true)
|
||||
}, 10_000)
|
||||
})
|
||||
@@ -109,11 +109,11 @@ export class OrcaRuntimeWithWaitForLeafPtyId extends OrcaRuntimeWithRestoreLiveP
|
||||
oscLinks?: TerminalOscLinkRange[]
|
||||
},
|
||||
trailingOutput: { data: string; seq: number }[] = []
|
||||
): void {
|
||||
): Promise<void> {
|
||||
if (!snapshot.data) {
|
||||
return
|
||||
return Promise.resolve()
|
||||
}
|
||||
// Why: a redraw byte can create a suffix-only model before the renderer settles; replace it with the exact snapshot already sent mobile.
|
||||
// Why: a redraw byte can create a suffix-only model before the renderer settles; replace it with the renderer's full screen.
|
||||
this.providerSnapshotPreferredPtys.add(ptyId)
|
||||
this.disposeHeadlessTerminal(ptyId)
|
||||
this.seedHeadlessTerminal(
|
||||
@@ -122,11 +122,17 @@ export class OrcaRuntimeWithWaitForLeafPtyId extends OrcaRuntimeWithRestoreLiveP
|
||||
{ cols: snapshot.cols, rows: snapshot.rows },
|
||||
{ cwd: snapshot.cwd, oscLinks: snapshot.oscLinks }
|
||||
)
|
||||
// Why: a hidden pane answers at its own size; land the model on the grid later bytes paint.
|
||||
const ptyGrid = this.getTerminalSize(ptyId)
|
||||
if (ptyGrid) {
|
||||
this.resizeHeadlessTerminal(ptyId, ptyGrid.cols, ptyGrid.rows)
|
||||
}
|
||||
for (const chunk of trailingOutput) {
|
||||
this.trackHeadlessTerminalData(ptyId, chunk.data, chunk.seq)
|
||||
}
|
||||
// The seed's write chain owns subsequent live bytes; suppress on-data hydration from replacing this known-good seed.
|
||||
this.headlessHydrationState.set(ptyId, 'done')
|
||||
return this.headlessTerminals.get(ptyId)?.writeChain ?? Promise.resolve()
|
||||
}
|
||||
|
||||
waitForRendererTerminalSerializer(
|
||||
|
||||
@@ -8,6 +8,7 @@ import {
|
||||
sendSnapshotFrames,
|
||||
serializeStableMobileRendererSnapshot
|
||||
} from './terminal-snapshot-publication'
|
||||
import { seedModelFromRendererScreen } from './terminal-renderer-screen-model-seed'
|
||||
import { updateViewportForClient } from './terminal-viewport-update'
|
||||
import type {
|
||||
LegacyBinarySubscriptionState,
|
||||
@@ -47,27 +48,42 @@ export function activateLegacyBinarySubscription(
|
||||
if (recovery.seq !== runtime.getPtyOutputSequence(ptyId)) {
|
||||
return
|
||||
}
|
||||
runtime.replaceHeadlessTerminalFromRendererSnapshotForRecovery(ptyId, recovery)
|
||||
const published = await seedModelFromRendererScreen(
|
||||
runtime,
|
||||
ptyId,
|
||||
recovery,
|
||||
[],
|
||||
mobileSnapshotByteBudget(params.snapshotByteBudget, state.streamId, recoveryFrame)
|
||||
)
|
||||
// Why: bytes that landed while the model was rebuilt already went live, so a later reset would repeat them.
|
||||
if (
|
||||
state.closed ||
|
||||
!published?.data.length ||
|
||||
published.seq !== recovery.seq ||
|
||||
recovery.seq !== runtime.getPtyOutputSequence(ptyId)
|
||||
) {
|
||||
return
|
||||
}
|
||||
// Why: shipped mobile clients apply resized snapshots in place, so a blank xterm recovers without resubscribe.
|
||||
const recoveryStats = sendSnapshotFrames(state.sendFrame, {
|
||||
...recoveryFrame,
|
||||
cols: recovery.cols,
|
||||
rows: recovery.rows,
|
||||
cols: published.cols,
|
||||
rows: published.rows,
|
||||
displayMode: state.displayMode,
|
||||
source: recovery.source,
|
||||
source: published.source,
|
||||
truncated: false,
|
||||
truncatedByByteBudget: recovery.truncatedByByteBudget,
|
||||
data: recovery.data
|
||||
truncatedByByteBudget: published.truncatedByByteBudget,
|
||||
data: published.data
|
||||
})
|
||||
state.lastResizeCols = recovery.cols
|
||||
state.lastResizeCols = published.cols
|
||||
console.log('[mobile-terminal-stream] recovery snapshot', {
|
||||
terminal: params.terminal,
|
||||
streamId: state.streamId,
|
||||
reason: 'renderer-mount-ready',
|
||||
bytes: recoveryStats.bytes,
|
||||
chunks: recoveryStats.chunks,
|
||||
scrollbackRows: recovery.scrollbackRows,
|
||||
truncatedByByteBudget: recovery.truncatedByByteBudget === true
|
||||
scrollbackRows: published.scrollbackRows,
|
||||
truncatedByByteBudget: published.truncatedByByteBudget === true
|
||||
})
|
||||
})
|
||||
.catch(() => {})
|
||||
|
||||
@@ -14,6 +14,7 @@ import type {
|
||||
LegacyBinarySubscriptionState,
|
||||
TerminalSubscriptionArgs
|
||||
} from './terminal-legacy-subscription-types'
|
||||
import { seedModelFromRendererScreen } from './terminal-renderer-screen-model-seed'
|
||||
|
||||
const MOBILE_RENDERER_MOUNT_READY_TIMEOUT_MS = 3_000
|
||||
|
||||
@@ -145,17 +146,17 @@ export async function publishLegacyBinaryInitialSnapshot(
|
||||
stableRendererSnapshot?.data.length &&
|
||||
(typeof stableRendererSnapshot.seq === 'number' || state.pendingOutput.length === 0)
|
||||
) {
|
||||
serialized = stableRendererSnapshot
|
||||
const trailingOutput = state.pendingOutput.flatMap((item) => {
|
||||
const output = getOutputAfterSnapshotSeq(item, stableRendererSnapshot.seq)
|
||||
const seq = item.meta?.seq
|
||||
return output && typeof seq === 'number' ? [{ data: output.data, seq }] : []
|
||||
})
|
||||
runtime.replaceHeadlessTerminalFromRendererSnapshotForRecovery(
|
||||
const fromModel = await seedModelFromRendererScreen(
|
||||
runtime,
|
||||
ptyId,
|
||||
stableRendererSnapshot,
|
||||
trailingOutput
|
||||
state.pendingOutput,
|
||||
mobileSnapshotByteBudget(params.snapshotByteBudget, state.streamId, scrollbackFrame)
|
||||
)
|
||||
if (state.closed) {
|
||||
return
|
||||
}
|
||||
serialized = fromModel ?? stableRendererSnapshot
|
||||
}
|
||||
}
|
||||
let initialOutputOverflowed = false
|
||||
|
||||
@@ -0,0 +1,37 @@
|
||||
import type { OrcaRuntimeService } from '../../../orca-runtime'
|
||||
import {
|
||||
serializeBudgetedMobileSnapshot,
|
||||
type MobileSnapshotByteBudget
|
||||
} from './terminal-snapshot-publication'
|
||||
import { getOutputAfterSnapshotSeq } from './terminal-stream-replay'
|
||||
import type { SerializedSnapshot, TerminalOutputChunk } from './terminal-stream-types'
|
||||
|
||||
type RendererScreen = NonNullable<SerializedSnapshot>
|
||||
|
||||
/**
|
||||
* Rebuilds the host model from a renderer screen, replaying the output after its seq, and returns
|
||||
* the model's snapshot. A hidden pane answers at its own size; the model sits on the PTY grid.
|
||||
*/
|
||||
export async function seedModelFromRendererScreen(
|
||||
runtime: Pick<
|
||||
OrcaRuntimeService,
|
||||
'replaceHeadlessTerminalFromRendererSnapshotForRecovery' | 'serializeTerminalBuffer'
|
||||
>,
|
||||
ptyId: string,
|
||||
screen: RendererScreen,
|
||||
pendingOutput: readonly TerminalOutputChunk[],
|
||||
snapshotByteBudget: MobileSnapshotByteBudget | undefined
|
||||
): Promise<SerializedSnapshot> {
|
||||
const trailingOutput = pendingOutput.flatMap((item) => {
|
||||
const output = getOutputAfterSnapshotSeq(item, screen.seq)
|
||||
const seq = item.meta?.seq
|
||||
return output && typeof seq === 'number' ? [{ data: output.data, seq }] : []
|
||||
})
|
||||
await runtime.replaceHeadlessTerminalFromRendererSnapshotForRecovery(
|
||||
ptyId,
|
||||
screen,
|
||||
trailingOutput
|
||||
)
|
||||
// Why mobile: only a phone subscribe adopts a renderer screen.
|
||||
return serializeBudgetedMobileSnapshot(runtime, ptyId, true, snapshotByteBudget)
|
||||
}
|
||||
@@ -296,6 +296,7 @@ describe('terminal subscribe mount replay', () => {
|
||||
let generation = 0
|
||||
let headlessPresent = false
|
||||
const requestRendererTerminalTabMount = vi.fn(() => true)
|
||||
let model = { data: 'suffix-only redraw', cols: 80, rows: 24, seq: 1 }
|
||||
const runtime = asRuntime({
|
||||
getRuntimeId: () => 'test-runtime',
|
||||
subscribeToPtyExit: vi.fn(() => vi.fn()),
|
||||
@@ -311,7 +312,12 @@ describe('terminal subscribe mount replay', () => {
|
||||
getRendererTerminalSerializerGenerationForHandle: vi.fn(() => 0),
|
||||
getRendererTerminalSerializerGeneration: vi.fn(() => generation),
|
||||
getPtyOutputSequence: vi.fn(() => 0),
|
||||
replaceHeadlessTerminalFromRendererSnapshotForRecovery: vi.fn(),
|
||||
// The recovery seed replaces the suffix-only model with the mounted screen.
|
||||
replaceHeadlessTerminalFromRendererSnapshotForRecovery: vi.fn(
|
||||
async (_ptyId: string, snapshot: { data: string }) => {
|
||||
model = { ...model, data: snapshot.data }
|
||||
}
|
||||
),
|
||||
waitForRendererTerminalSerializer: vi.fn(async (_ptyId, afterGeneration) => {
|
||||
return generation > afterGeneration
|
||||
}),
|
||||
@@ -320,12 +326,7 @@ describe('terminal subscribe mount replay', () => {
|
||||
subscribeToTerminalData: vi.fn().mockReturnValue(vi.fn()),
|
||||
registerRemoteTerminalViewSubscriber: vi.fn(() => vi.fn()),
|
||||
readTerminal: vi.fn().mockResolvedValue({ tail: [], truncated: false }),
|
||||
serializeTerminalBuffer: vi.fn().mockResolvedValue({
|
||||
data: 'suffix-only redraw',
|
||||
cols: 80,
|
||||
rows: 24,
|
||||
seq: 1
|
||||
}),
|
||||
serializeTerminalBuffer: vi.fn(async () => model),
|
||||
// Baseline race: the pre-PTY mount settles between the attachment answer and the wait, so the
|
||||
// wait must still count that settle against the pre-mount generation.
|
||||
serializeRendererTerminalBuffer: vi
|
||||
|
||||
@@ -50,6 +50,8 @@ function subscribeMobile(pane: PaneDouble) {
|
||||
const binaryFrames: Uint8Array<ArrayBufferLike>[] = []
|
||||
const registry = createSubscriptionRegistryDouble()
|
||||
let emitData: (data: string, meta: { seq: number; rawLength: number }) => void = () => {}
|
||||
// The host model a recovery seed builds: the adopted screen plus the output it replays.
|
||||
let model: { data: string; cols: number; rows: number; seq?: number } | null = null
|
||||
const runtime = {
|
||||
getRuntimeId: () => 'test-runtime',
|
||||
subscribeToPtyExit: vi.fn(() => vi.fn()),
|
||||
@@ -60,7 +62,20 @@ function subscribeMobile(pane: PaneDouble) {
|
||||
getRendererTerminalSerializerGenerationForHandle: vi.fn(() => 1),
|
||||
getRendererTerminalSerializerGeneration: vi.fn(() => 1),
|
||||
getPtyOutputSequence: vi.fn(pane.outputSequence ?? (() => 4)),
|
||||
replaceHeadlessTerminalFromRendererSnapshotForRecovery: vi.fn(),
|
||||
replaceHeadlessTerminalFromRendererSnapshotForRecovery: vi.fn(
|
||||
async (
|
||||
_ptyId: string,
|
||||
snapshot: { data: string; seq?: number },
|
||||
trailing: { data: string; seq: number }[] = []
|
||||
) => {
|
||||
model = {
|
||||
data: snapshot.data + trailing.map((chunk) => chunk.data).join(''),
|
||||
cols: 80,
|
||||
rows: 24,
|
||||
seq: trailing.at(-1)?.seq ?? snapshot.seq
|
||||
}
|
||||
}
|
||||
),
|
||||
waitForRendererTerminalSerializer: vi.fn(pane.waitForRendererTerminalSerializer),
|
||||
handleMobileSubscribe: vi.fn().mockResolvedValue(true),
|
||||
handleMobileUnsubscribe: vi.fn(),
|
||||
@@ -72,6 +87,9 @@ function subscribeMobile(pane: PaneDouble) {
|
||||
readTerminal: vi.fn().mockResolvedValue({ tail: [], truncated: false }),
|
||||
// The restored provider snapshot wins the preference order over the live renderer.
|
||||
serializeTerminalBuffer: vi.fn(async () => {
|
||||
if (model) {
|
||||
return model
|
||||
}
|
||||
if (pane.pendingOutput) {
|
||||
emitData(pane.pendingOutput, {
|
||||
seq: pane.pendingOutputSeq ?? 3,
|
||||
@@ -171,7 +189,8 @@ describe('terminal subscribe for a pane the desktop already has mounted', () =>
|
||||
})
|
||||
|
||||
it('adopts a renderer-ordered screen even while output is pending', async () => {
|
||||
// The screen's seq (4) is an exact seam inside the pending chunk (offsets 3..5), so only `z` replays.
|
||||
// The screen's seq (4) is an exact seam inside the pending chunk (offsets 3..5), so only `z`
|
||||
// replays into the model, and the model's seq (5) keeps it from replaying to the phone again.
|
||||
const subscription = subscribeMobile({
|
||||
rendererScreen: () => 'ordered desktop prompt $ ',
|
||||
pendingOutput: 'Xz',
|
||||
@@ -185,7 +204,7 @@ describe('terminal subscribe for a pane the desktop already has mounted', () =>
|
||||
await vi.advanceTimersByTimeAsync(100)
|
||||
|
||||
expect(subscription.runtime.requestRendererTerminalTabMount).not.toHaveBeenCalled()
|
||||
expect(subscription.snapshotText()).toContain('ordered desktop prompt $ ')
|
||||
expect(subscription.snapshotText()).toContain('ordered desktop prompt $ z')
|
||||
expect(subscription.snapshotText()).not.toContain('restored provider history')
|
||||
expect(
|
||||
subscription.runtime.replaceHeadlessTerminalFromRendererSnapshotForRecovery
|
||||
@@ -194,7 +213,7 @@ describe('terminal subscribe for a pane the desktop already has mounted', () =>
|
||||
expect.objectContaining({ data: 'ordered desktop prompt $ ', seq: 4 }),
|
||||
[{ data: 'z', seq: 5 }]
|
||||
)
|
||||
expect(subscription.outputText()).toBe('z')
|
||||
expect(subscription.outputText()).toBe('')
|
||||
await subscription.close()
|
||||
})
|
||||
|
||||
|
||||
Reference in New Issue
Block a user