Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
160 changes: 160 additions & 0 deletions apps/ade-cli/src/services/sync/syncHostService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11366,6 +11366,7 @@ describe("chat event replay buffer (resumable chat streams)", () => {
describe("terminal byte-offset streaming, history paging, and resize ownership", () => {
// 5000 ASCII bytes so byte offsets equal string indices in assertions.
const TRANSCRIPT_CONTENT = "0123456789".repeat(500);
const SCREEN_CSI = "\x1b[?1049h\x1b[Hhello";

beforeEach(() => {
publishMock.mockReset();
Expand Down Expand Up @@ -11411,6 +11412,11 @@ describe("terminal byte-offset streaming, history paging, and resize ownership",
const restoreDesktopSizeBySessionId = vi.fn().mockReturnValue(true);
const hasLivePty = vi.fn().mockReturnValue(true);
const writeBySessionId = vi.fn().mockReturnValue(true);
const readScreenSnapshot = vi.fn((id: string) => (
id === "session-1"
? { cols: 80, rows: 24, bufferType: "alternate" as const, serialized: SCREEN_CSI }
: null
));
let sessionAvailable = true;
const base = createHostArgs(projectRoot, []);
const host = createSyncHostService({
Expand Down Expand Up @@ -11444,6 +11450,7 @@ describe("terminal byte-offset streaming, history paging, and resize ownership",
resizeBySessionId,
restoreDesktopSizeBySessionId,
hasLivePty,
readScreenSnapshot,
enrichSessions: (rows: unknown[]) => rows,
},
} as unknown as Parameters<typeof createSyncHostService>[0]);
Expand All @@ -11457,6 +11464,7 @@ describe("terminal byte-offset streaming, history paging, and resize ownership",
resizeBySessionId,
restoreDesktopSizeBySessionId,
hasLivePty,
readScreenSnapshot,
setSessionAvailable: (available: boolean) => { sessionAvailable = available; },
};
}
Expand Down Expand Up @@ -11517,6 +11525,7 @@ describe("terminal byte-offset streaming, history paging, and resize ownership",
startOffset: 4_988,
endOffset: 5_000,
});
expect((delta.payload as { screen?: unknown }).screen).toBeUndefined();
expect(readTranscriptSnapshot).toHaveBeenCalledWith({
sessionId: "session-1",
maxBytes: 32_000,
Expand All @@ -11537,6 +11546,12 @@ describe("terminal byte-offset streaming, history paging, and resize ownership",
transcript: TRANSCRIPT_CONTENT.slice(5_000 - 1_024),
startOffset: 5_000 - 1_024,
endOffset: 5_000,
screen: {
cols: 80,
rows: 24,
bufferType: "alternate",
serialized: SCREEN_CSI,
},
});
expect((full.payload as { delta?: boolean }).delta).toBeUndefined();
expect(readTranscriptTail).not.toHaveBeenCalled();
Expand Down Expand Up @@ -11580,6 +11595,151 @@ describe("terminal byte-offset streaming, history paging, and resize ownership",
}
});

it("omits an oversized current-screen serialize instead of slicing CSI", async () => {
const { projectRoot, cleanup } = createTempProjectRoot();
const { host, readScreenSnapshot } = createTerminalHost(projectRoot);
readScreenSnapshot.mockReturnValue({
cols: 80,
rows: 24,
bufferType: "alternate",
serialized: "x".repeat(256_001),
});
let client: Awaited<ReturnType<typeof connectTerminalPeer>> | null = null;
try {
client = await connectTerminalPeer(
await host.waitUntilListening(),
host.getBootstrapToken(),
"ios-terminal-screen-cap",
);
client.ws.send(encodeSyncEnvelope({
type: "terminal_subscribe",
requestId: "sub-screen-cap",
payload: { sessionId: "session-1", maxBytes: 32_000 },
}));
const snapshot = await nextResponse(client.envelopes, "terminal_snapshot", "sub-screen-cap");
expect((snapshot.payload as { screen?: unknown }).screen).toBeUndefined();
expect((snapshot.payload as { transcript: string }).transcript).toBe(TRANSCRIPT_CONTENT);
} finally {
try { client?.ws.close(); } catch { /* ignore */ }
await host.dispose();
cleanup();
}
});

it("keeps the sync socket open when snapshot catch-up overflows unreconstructable events", async () => {
const { projectRoot, cleanup } = createTempProjectRoot();
const { host, readTranscriptSnapshot } = createTerminalHost(projectRoot);
let resolveFirstSnapshot!: (snapshot: { data: string; startOffset: number; endOffset: number }) => void;
readTranscriptSnapshot.mockImplementationOnce(() => new Promise<{
data: string;
startOffset: number;
endOffset: number;
}>((resolve) => {
resolveFirstSnapshot = resolve;
}));
let client: Awaited<ReturnType<typeof connectTerminalPeer>> | null = null;
try {
client = await connectTerminalPeer(
await host.waitUntilListening(),
host.getBootstrapToken(),
"ios-terminal-catchup-overflow",
);
client.ws.send(encodeSyncEnvelope({
type: "terminal_subscribe",
requestId: "overflow-untracked",
payload: { sessionId: "session-1", maxBytes: 32_000 },
}));
await waitForValue(
() => readTranscriptSnapshot.mock.calls.length > 0 ? true : null,
"overflow terminal snapshot capture",
);
for (let i = 0; i < 257; i += 1) {
host.handlePtyData({
sessionId: "session-1",
ptyId: "pty-1",
data: "x",
offset: null,
});
}
resolveFirstSnapshot({ data: TRANSCRIPT_CONTENT, startOffset: 0, endOffset: 5_000 });

const snapshot = await nextResponse(client.envelopes, "terminal_snapshot", "overflow-untracked");
expect(snapshot.payload).toMatchObject({
sessionId: "session-1",
transcript: TRANSCRIPT_CONTENT,
screen: { serialized: SCREEN_CSI },
});
expect(client.ws.readyState).toBe(WebSocket.OPEN);
expect(client.closeEvents).toEqual([]);
} finally {
try { client?.ws.close(); } catch { /* ignore */ }
await host.dispose();
cleanup();
}
});

it("sends the last captured transcript when capture attempts exhaust without ever failing", async () => {
const { projectRoot, cleanup } = createTempProjectRoot();
const { host, readTranscriptSnapshot } = createTerminalHost(projectRoot);
const captures: Array<(snapshot: { data: string; startOffset: number; endOffset: number }) => void> = [];
readTranscriptSnapshot.mockImplementation(() => new Promise<{
data: string;
startOffset: number;
endOffset: number;
}>((resolve) => {
captures.push(resolve);
}));
let client: Awaited<ReturnType<typeof connectTerminalPeer>> | null = null;
try {
client = await connectTerminalPeer(
await host.waitUntilListening(),
host.getBootstrapToken(),
"ios-terminal-capture-exhausted",
);
client.ws.send(encodeSyncEnvelope({
type: "terminal_subscribe",
requestId: "exhausted-recapture",
payload: { sessionId: "session-1", maxBytes: 32_000 },
}));
await waitForValue(
() => captures.length > 0 ? true : null,
"first exhausted-recapture snapshot capture",
);
host.handlePtyData({
sessionId: "session-1",
ptyId: "pty-1",
data: "x",
offset: 6_000,
});
captures[0]!({ data: "CAPTURE-1", startOffset: 0, endOffset: 5_000 });
for (let attempt = 1; attempt < 4; attempt += 1) {
await waitForValue(
() => captures.length > attempt ? true : null,
`exhausted-recapture snapshot capture ${attempt + 1}`,
);
captures[attempt]!({
data: `CAPTURE-${attempt + 1}`,
startOffset: 0,
endOffset: 5_000,
});
}

const snapshot = await nextResponse(client.envelopes, "terminal_snapshot", "exhausted-recapture");
expect(readTranscriptSnapshot).toHaveBeenCalledTimes(4);
expect(snapshot.payload).toMatchObject({
sessionId: "session-1",
transcript: "CAPTURE-4",
screen: { serialized: SCREEN_CSI },
});
expect(client.ws.readyState).toBe(WebSocket.OPEN);
expect(client.closeEvents).toEqual([]);
} finally {
try { client?.ws.close(); } catch { /* ignore */ }
await host.dispose();
cleanup();
}
});

it("replaces with a tail snapshot when sinceOffset already equals the transcript end", async () => {
const { projectRoot, cleanup } = createTempProjectRoot();
const { host, readTranscriptTail, readTranscriptSnapshot, readTranscriptRange } = createTerminalHost(projectRoot);
Expand Down
96 changes: 80 additions & 16 deletions apps/ade-cli/src/services/sync/syncHostService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -99,6 +99,7 @@ import type {
SyncTerminalInputAckPayload,
SyncTerminalInputPayload,
SyncTerminalSnapshotPayload,
SyncTerminalScreenSnapshot,
} from "../../../../desktop/src/shared/types";
import {
SYNC_APPLICATION_COMPRESSION_THRESHOLD_BYTES,
Expand Down Expand Up @@ -486,6 +487,7 @@ const MAX_TERMINAL_HISTORY_PAGE_BYTES = 524_288;
const MAX_PENDING_TERMINAL_SNAPSHOT_EVENTS = 256;
const MAX_PENDING_TERMINAL_SNAPSHOT_BYTES = 2_000_000;
const MAX_TERMINAL_SNAPSHOT_CAPTURE_ATTEMPTS = 4;
const MAX_TERMINAL_SCREEN_SERIALIZED_CHARS = 256_000;
const PEER_BACKPRESSURE_BYTES = 4 * 1024 * 1024;
const REQUIRED_SEND_MAX_BUFFERED_BYTES = 16 * 1024 * 1024;
const SEND_AND_WAIT_TIMEOUT_MS = 15_000;
Expand Down Expand Up @@ -2999,16 +3001,15 @@ export function createSyncHostService(args: SyncHostServiceArgs) {
&& peer.ws.readyState === WebSocket.OPEN;
}

function isCurrentTerminalSnapshotBarrier(
function isActiveTerminalSnapshotBarrier(
peer: PeerState,
sessionId: string,
barrier: PendingTerminalSnapshotBarrier,
lifecycleGeneration: number,
): boolean {
const currentBarrier = peer.pendingTerminalSnapshots.get(sessionId);
return isPeerLifecycleCurrent(peer, lifecycleGeneration)
&& currentBarrier === barrier
&& !barrier.failed;
&& currentBarrier === barrier;
}

function clearTerminalSnapshotBarrier(
Expand Down Expand Up @@ -3052,6 +3053,29 @@ export function createSyncHostService(args: SyncHostServiceArgs) {
barrier.queuedBytes = queuedBytes;
}

function readTerminalScreenSnapshot(sessionId: string): SyncTerminalScreenSnapshot | null {
const screen = args.ptyService.readScreenSnapshot?.(sessionId) ?? null;
if (!screen?.serialized) return null;
if (screen.serialized.length > MAX_TERMINAL_SCREEN_SERIALIZED_CHARS) return null;
if (!Number.isFinite(screen.cols) || !Number.isFinite(screen.rows)) return null;
if (screen.cols < 1 || screen.rows < 1) return null;
return {
cols: Math.floor(screen.cols),
rows: Math.floor(screen.rows),
bufferType: screen.bufferType === "alternate" ? "alternate" : "normal",
serialized: screen.serialized,
};
}

function withTerminalScreen(
sessionId: string,
snapshot: SyncTerminalSnapshotPayload,
): SyncTerminalSnapshotPayload {
if (snapshot.delta === true) return snapshot;
const screen = readTerminalScreenSnapshot(sessionId);
return screen ? { ...snapshot, screen } : snapshot;
}

function failTerminalSnapshotBarrier(
peer: PeerState,
sessionId: string,
Expand All @@ -3068,11 +3092,9 @@ export function createSyncHostService(args: SyncHostServiceArgs) {
queuedBytes: barrier.queuedBytes,
peerDeviceId: peer.metadata?.deviceId ?? peer.pairedDeviceId ?? null,
});
try {
peer.ws.close(4001, "Terminal snapshot catch-up failed");
} catch {
// The failed barrier still prevents an out-of-order or lossy flush.
}
// Do not close the controller socket. A hot Claude TUI can overflow the
// catch-up queue on LAN; tearing down sync blanks every other surface on
// that phone. The subscribe loop sends the last captured snapshot instead.
}

function enqueueTerminalSnapshotEvent(
Expand Down Expand Up @@ -7719,9 +7741,41 @@ export function createSyncHostService(args: SyncHostServiceArgs) {
: null;
let forceReplacement = false;
let barrierCompleted = false;
let lastCapturedTranscript: {
data: string;
startOffset: number | null;
endOffset: number | null;
} | null = null;
const sendReplacingSnapshot = (
session: ReturnType<typeof args.sessionService.get>,
transcriptSnapshot: {
data: string;
startOffset: number | null;
endOffset: number | null;
} | null,
): boolean => {
const snapshot = withTerminalScreen(sessionId, {
sessionId,
transcript: transcriptSnapshot?.data ?? "",
status: session?.status ?? null,
runtimeState: session?.runtimeState ?? null,
lastOutputPreview: session?.lastOutputPreview ?? null,
capturedAt: nowIso(),
startOffset: transcriptSnapshot?.startOffset ?? null,
endOffset: transcriptSnapshot?.endOffset ?? null,
live: args.ptyService.hasLivePty(sessionId),
});
return sendRequired(peer, "terminal_snapshot", snapshot, envelope.requestId);
};
try {
while (barrier.captureAttempt < MAX_TERMINAL_SNAPSHOT_CAPTURE_ATTEMPTS) {
if (!isCurrentTerminalSnapshotBarrier(peer, sessionId, barrier, lifecycleGeneration)) break;
if (!isActiveTerminalSnapshotBarrier(peer, sessionId, barrier, lifecycleGeneration)) break;
if (barrier.failed) {
if (sendReplacingSnapshot(args.sessionService.get(sessionId), lastCapturedTranscript)) {
barrierCompleted = true;
}
break;
}
barrier.captureAttempt += 1;
const session = args.sessionService.get(sessionId);
const transcriptSnapshot = session
Expand All @@ -7735,7 +7789,14 @@ export function createSyncHostService(args: SyncHostServiceArgs) {
"Sync operation aborted.",
)
: null;
if (!isCurrentTerminalSnapshotBarrier(peer, sessionId, barrier, lifecycleGeneration)) break;
if (transcriptSnapshot) lastCapturedTranscript = transcriptSnapshot;
if (!isActiveTerminalSnapshotBarrier(peer, sessionId, barrier, lifecycleGeneration)) break;
if (barrier.failed) {
if (sendReplacingSnapshot(session, lastCapturedTranscript)) {
barrierCompleted = true;
}
break;
}

const flush = planTerminalSnapshotFlush(
barrier,
Expand Down Expand Up @@ -7770,7 +7831,7 @@ export function createSyncHostService(args: SyncHostServiceArgs) {
startsAtUtf8Boundary
&& snapshotBytes.length === transcriptSnapshot.endOffset - transcriptSnapshot.startOffset
) {
snapshot = {
snapshot = withTerminalScreen(sessionId, {
sessionId,
transcript: snapshotBytes.subarray(byteStart).toString("utf8"),
status: session?.status ?? null,
Expand All @@ -7781,10 +7842,10 @@ export function createSyncHostService(args: SyncHostServiceArgs) {
endOffset: transcriptSnapshot.endOffset,
delta: true,
live: args.ptyService.hasLivePty(sessionId),
};
});
}
}
snapshot ??= {
snapshot ??= withTerminalScreen(sessionId, {
sessionId,
transcript: transcriptSnapshot?.data ?? "",
status: session?.status ?? null,
Expand All @@ -7794,7 +7855,7 @@ export function createSyncHostService(args: SyncHostServiceArgs) {
startOffset: transcriptSnapshot?.startOffset ?? null,
endOffset: transcriptSnapshot?.endOffset ?? null,
live: args.ptyService.hasLivePty(sessionId),
};
});
if (!sendRequired(peer, "terminal_snapshot", snapshot, envelope.requestId)) break;
barrierCompleted = true;
for (const event of flush.events) {
Expand All @@ -7810,10 +7871,13 @@ export function createSyncHostService(args: SyncHostServiceArgs) {
}
if (
!barrierCompleted
&& isCurrentTerminalSnapshotBarrier(peer, sessionId, barrier, lifecycleGeneration)
&& barrier.captureAttempt >= MAX_TERMINAL_SNAPSHOT_CAPTURE_ATTEMPTS
&& isActiveTerminalSnapshotBarrier(peer, sessionId, barrier, lifecycleGeneration)
&& (barrier.failed || barrier.captureAttempt >= MAX_TERMINAL_SNAPSHOT_CAPTURE_ATTEMPTS)
) {
failTerminalSnapshotBarrier(peer, sessionId, barrier, "capture_did_not_reach_stable_offset");
if (sendReplacingSnapshot(args.sessionService.get(sessionId), lastCapturedTranscript)) {
barrierCompleted = true;
}
}
} catch (error) {
if (peer.pendingTerminalSnapshots.get(sessionId) === barrier) {
Expand Down
Loading
Loading