Skip to content
Merged
2 changes: 1 addition & 1 deletion src/main/daemon/daemon-pty-adapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1938,7 +1938,7 @@ export class DaemonPtyAdapter implements IPtyProvider {
}
}

// Why final=true not teardown: clean disconnect needs the full-depth snapshot as the restore source, but the
// Why final=true not teardown: clean disconnect needs the full daemon-window snapshot as the restore source, but the
// detached daemon's PTYs keep running for warm reattach, so shell-ready scanner state must stay intact.
private async checkpointAllSessions(): Promise<void> {
const completed = await this.checkpointSessions(this.activeSessionIds, { final: true })
Expand Down
8 changes: 8 additions & 0 deletions src/main/daemon/daemon-pty-router-history-handoff.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import { DaemonPtyRouter } from './daemon-pty-router'
import { DaemonServer } from './daemon-server'
import { getDaemonSocketPath } from './daemon-spawner'
import { TERMINAL_HISTORY_INLINE_SEED_CODE_UNITS } from './terminal-history-seed-chunks'
import { DAEMON_SESSION_SCROLLBACK_ROWS } from './daemon-session-scrollback-window'
import type { DaemonFileLog } from './daemon-file-log'
import type { SubprocessHandle } from './session'

Expand Down Expand Up @@ -93,14 +94,21 @@ describe('DaemonPtyRouter history handoff', () => {
}
await Promise.all([legacyServer?.shutdown(), currentServer?.shutdown()])
rmSync(testDir, { recursive: true, force: true })
vi.unstubAllEnvs()
vi.restoreAllMocks()
})

it('moves a slept v29 session to v30 with its full large history', async () => {
const sessionId = 'large-legacy-session'
const marker = 'V29-HISTORY-HANDOFF-MARKER'
// Why: this fixture plays an OLD daemon binary, whose sessions retained ~5000 rows. New code
// applies the flat window at session create, which would shrink the history below the chunked
// seed threshold and silently skip the transfer path this test exists to cover.
vi.stubEnv('ORCA_DAEMON_SESSION_SCROLLBACK_ROWS', '5000')
await legacyAdapter.spawn({ sessionId, cols: 400, rows: 24 })
legacySubprocess.emitData(`${'x'.repeat(TERMINAL_HISTORY_INLINE_SEED_CODE_UNITS + 1)}${marker}`)
// Keep only the played legacy session deep; the receiving daemon must use today's flat window.
vi.stubEnv('ORCA_DAEMON_SESSION_SCROLLBACK_ROWS', String(DAEMON_SESSION_SCROLLBACK_ROWS))
router = new DaemonPtyRouter({ current: currentAdapter, legacy: [legacyAdapter] })
await router.discoverLegacySessions()
const client = (
Expand Down
135 changes: 135 additions & 0 deletions src/main/daemon/daemon-server-attachment-lifecycle.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,135 @@
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
import type { Socket } from 'node:net'
import { mkdtempSync, rmSync } from 'node:fs'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { DaemonClient } from './client'
import { DaemonServer } from './daemon-server'
import { getDaemonSocketPath } from './daemon-spawner'
import type { SubprocessHandle } from './session'
import type { CreateOrAttachOptions, CreateOrAttachResult } from './terminal-host'

function createMockSubprocess(): SubprocessHandle {
let onExit: ((code: number) => void) | undefined
return {
pid: 55_555,
getForegroundProcess: () => null,
write: vi.fn(),
resize: vi.fn(),
kill: vi.fn(() => onExit?.(0)),
forceKill: vi.fn(() => onExit?.(137)),
signal: vi.fn(),
onData: vi.fn(),
onExit(callback) {
onExit = callback
},
dispose: vi.fn()
}
}

type DaemonAttachmentPrivate = {
host: {
detach: (sessionId: string, token: symbol) => void
detachClients: (attachments: readonly { sessionId: string; token: symbol }[]) => void
createOrAttach: (options: CreateOrAttachOptions) => Promise<CreateOrAttachResult>
sessions: Map<string, { hasAttachedClients: boolean }>
}
clients: Map<string, { streamSocket: Socket | null }>
streamClientIdBySessionId: Map<string, string>
attachTokenBySessionId: Map<string, symbol>
}

describe('DaemonServer attachment lifecycle', () => {
let directory: string
let server: DaemonServer
let client: DaemonClient

beforeEach(async () => {
directory = mkdtempSync(join(tmpdir(), 'daemon-attachment-lifecycle-'))
const socketPath = getDaemonSocketPath(directory)
const tokenPath = join(directory, 'daemon.token')
server = new DaemonServer({
socketPath,
tokenPath,
spawnSubprocess: createMockSubprocess
})
await server.start()
client = new DaemonClient({ socketPath, tokenPath })
await client.ensureConnected()
})

afterEach(async () => {
client.disconnect()
await server.shutdown()
rmSync(directory, { recursive: true, force: true })
})

async function createSession(sessionId: string): Promise<void> {
await client.request('createOrAttach', { sessionId, cols: 80, rows: 24 })
}

it('parks attached sessions when their stream transport closes', async () => {
await createSession('transport-owned-session')
const daemon = server as unknown as DaemonAttachmentPrivate
const detachClients = vi.spyOn(daemon.host, 'detachClients')
const streamSocket = [...daemon.clients.values()][0]?.streamSocket

streamSocket?.destroy()

await vi.waitFor(() => expect(detachClients).toHaveBeenCalledOnce())
expect(detachClients.mock.calls[0]?.[0]).toEqual([
{ sessionId: 'transport-owned-session', token: expect.any(Symbol) }
])
expect(daemon.host.sessions.get('transport-owned-session')?.hasAttachedClients).toBe(false)
expect(daemon.streamClientIdBySessionId.has('transport-owned-session')).toBe(false)
expect(daemon.attachTokenBySessionId.has('transport-owned-session')).toBe(false)
})

it('parks only the requesting client attachment on explicit detach', async () => {
await createSession('explicitly-detached-session')
const daemon = server as unknown as DaemonAttachmentPrivate
const detach = vi.spyOn(daemon.host, 'detach')

await client.request('detach', { sessionId: 'explicitly-detached-session' })

expect(detach).toHaveBeenCalledWith('explicitly-detached-session', expect.any(Symbol))
expect(daemon.host.sessions.get('explicitly-detached-session')?.hasAttachedClients).toBe(false)
expect(daemon.streamClientIdBySessionId.has('explicitly-detached-session')).toBe(false)
expect(daemon.attachTokenBySessionId.has('explicitly-detached-session')).toBe(false)
})

it('parks an attachment completed after its transport already closed', async () => {
const daemon = server as unknown as DaemonAttachmentPrivate
const originalCreateOrAttach = daemon.host.createOrAttach.bind(daemon.host)
let attachmentCreated!: () => void
let finishRequest!: () => void
const created = new Promise<void>((resolve) => {
attachmentCreated = resolve
})
const gate = new Promise<void>((resolve) => {
finishRequest = resolve
})
vi.spyOn(daemon.host, 'createOrAttach').mockImplementation(async (options) => {
const result = await originalCreateOrAttach(options)
attachmentCreated()
await gate
return result
})
const request = client
.request('createOrAttach', { sessionId: 'close-race-session', cols: 80, rows: 24 })
.catch(() => undefined)
await created

client.disconnect()
await vi.waitFor(() => expect(daemon.clients.size).toBe(0))
expect(daemon.host.sessions.get('close-race-session')?.hasAttachedClients).toBe(true)
finishRequest()

await request
await vi.waitFor(() =>
expect(daemon.host.sessions.get('close-race-session')?.hasAttachedClients).toBe(false)
)
expect(daemon.streamClientIdBySessionId.has('close-race-session')).toBe(false)
expect(daemon.attachTokenBySessionId.has('close-race-session')).toBe(false)
})
})
57 changes: 55 additions & 2 deletions src/main/daemon/daemon-server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -174,6 +174,7 @@ export class DaemonServer {
}
})
private streamClientIdBySessionId = new Map<string, string>()
private attachTokenBySessionId = new Map<string, symbol>()
private lastInputAtBySessionId = new Map<string, number>()
private pendingPtySpawnPreparations = new Map<string, Set<PendingPtySpawnPreparation>>()
private historySeedTransfers = new TerminalHistorySeedTransferRegistry()
Expand Down Expand Up @@ -217,7 +218,17 @@ export class DaemonServer {
this.onAuthenticatedClientPair = opts.onAuthenticatedClientPair ?? (() => {})
this.host = new TerminalHost({
spawnSubprocess: opts.spawnSubprocess,
...(opts.onPtySessionExit ? { onSessionReaped: opts.onPtySessionExit } : {})
// Why host-level and not the attach callback: a session whose client transport already dropped
// has no attachment left to fire exit bookkeeping, and the daemon must still notice it can idle.
onSessionReaped: (sessionId) => {
this.streamClientIdBySessionId.delete(sessionId)
this.attachTokenBySessionId.delete(sessionId)
this.lastInputAtBySessionId.delete(sessionId)
this.transientFactRelay.onSessionExit(sessionId)
this.streamDataBatcher.refreshSessionDroppability(sessionId)
opts.onPtySessionExit?.(sessionId)
this.reevaluateIdleShutdown()
}
})
this.ptySpawnHealthCheck = opts.ptySpawnHealthCheck ?? checkPtySpawnHealth
this.preparePtySpawn = opts.preparePtySpawn ?? (() => Promise.resolve())
Expand Down Expand Up @@ -681,6 +692,7 @@ export class DaemonServer {
this.historySeedTransfers.clearOwner(clientId)
const wasFullyAuthenticated = client.authenticatedPairEstablished
this.streamDataBatcher.clear(clientId)
this.detachClientSessions(clientId)
client.streamSocket?.destroy()
this.clients.delete(clientId)
this.recordFullyAuthenticatedDisconnect(wasFullyAuthenticated)
Expand Down Expand Up @@ -718,6 +730,7 @@ export class DaemonServer {
// Why: a preflight that outlives its output channel would create an unattached daemon PTY.
this.cancelPendingPtySpawnPreparationsForClient(client.clientId)
this.streamDataBatcher.clear(client.clientId)
this.detachClientSessions(client.clientId)
client.streamSocket = null
}

Expand Down Expand Up @@ -843,6 +856,36 @@ export class DaemonServer {
}
}

private detachClientSessions(clientId: string): void {
const attachments: { sessionId: string; token: symbol }[] = []
for (const [sessionId, attachedClientId] of this.streamClientIdBySessionId) {
if (attachedClientId !== clientId) {
continue
}
const token = this.attachTokenBySessionId.get(sessionId)
if (token) {
attachments.push({ sessionId, token })
}
this.streamClientIdBySessionId.delete(sessionId)
this.attachTokenBySessionId.delete(sessionId)
}
if (attachments.length > 0) {
this.host.detachClients(attachments)
}
}

private detachSessionForClient(sessionId: string, clientId: string): void {
if (this.streamClientIdBySessionId.get(sessionId) !== clientId) {
return
}
const token = this.attachTokenBySessionId.get(sessionId)
if (token) {
this.host.detach(sessionId, token)
}
this.streamClientIdBySessionId.delete(sessionId)
this.attachTokenBySessionId.delete(sessionId)
}

private async routeRequest(clientId: string, request: DaemonRequest): Promise<unknown> {
const client = this.clients.get(clientId)

Expand Down Expand Up @@ -966,6 +1009,7 @@ export class DaemonServer {
this.transientFactRelay.onSessionExit(routedSessionId)
this.streamDataBatcher.refreshSessionDroppability(routedSessionId)
this.streamClientIdBySessionId.delete(routedSessionId)
this.attachTokenBySessionId.delete(routedSessionId)
this.lastInputAtBySessionId.delete(routedSessionId)
this.reevaluateIdleShutdown()
}
Expand All @@ -976,7 +1020,16 @@ export class DaemonServer {
this.reevaluateIdleShutdown()
}
routedSessionId = result.agentSessionEnsure?.owner.ptyId ?? p.sessionId
if (
this.clients.get(clientId) !== client ||
!client.authenticatedPairEstablished ||
client.streamSocket === null
) {
this.host.detach(routedSessionId, result.attachToken)
throw new TerminalAttachCanceledError(routedSessionId)
}
this.streamClientIdBySessionId.set(routedSessionId, clientId)
this.attachTokenBySessionId.set(routedSessionId, result.attachToken)
this.streamDataBatcher.refreshSessionDroppability(routedSessionId)
// Why an attach-time marker: background resync can precede this attach, so scan suppression must start at the new stream's head.
if (this.transientFactRelay.isBackgrounded(routedSessionId)) {
Expand Down Expand Up @@ -1123,7 +1176,7 @@ export class DaemonServer {
return {}

case 'detach':
// Note: detach token handling simplified — full impl would track tokens per client
this.detachSessionForClient(request.payload.sessionId, clientId)
this.log.log('session-detached', {
sessionId: request.payload.sessionId
})
Expand Down
94 changes: 94 additions & 0 deletions src/main/daemon/daemon-session-scrollback-window.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,94 @@
/**
* OOM regression: a daemon owning 100+ terminals retained ~5000 rows of grid per session with no
* bound, grew to ~1.9 GB, and was killed under system memory pressure — losing every session it
* owned. Sessions now retain a flat window; deep scrolling on an open terminal is the renderer's
* live buffer, and the window is what a rebuild (reload/remount/restart/remote attach) restores.
*/
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
import { TerminalHost } from './terminal-host'
import type { SubprocessHandle } from './session'
import {
DAEMON_SESSION_SCROLLBACK_ROWS,
resolveDaemonSessionScrollbackRows
} from './daemon-session-scrollback-window'

describe('resolveDaemonSessionScrollbackRows', () => {
it('defaults to the flat window', () => {
expect(resolveDaemonSessionScrollbackRows({} as NodeJS.ProcessEnv)).toBe(
DAEMON_SESSION_SCROLLBACK_ROWS
)
})

it('accepts the inclusive override bounds and rejects everything outside them', () => {
for (const raw of ['100', '2500', '5000']) {
const env = { ORCA_DAEMON_SESSION_SCROLLBACK_ROWS: raw } as NodeJS.ProcessEnv
expect(resolveDaemonSessionScrollbackRows(env)).toBe(Number(raw))
}
// Why bounded: 0 loses the visible screen's context; huge values silently reintroduce the
// unbounded retention this window exists to prevent.
for (const raw of ['0', '50', '99', '5001', '50000', '-1', '3.5', 'nonsense', '']) {
const env = { ORCA_DAEMON_SESSION_SCROLLBACK_ROWS: raw } as NodeJS.ProcessEnv
expect(resolveDaemonSessionScrollbackRows(env)).toBe(DAEMON_SESSION_SCROLLBACK_ROWS)
}
})
})

describe('daemon session scrollback window', () => {
let host: TerminalHost
let dataCb: ((data: string) => void) | null

function createMockSubprocess(): SubprocessHandle {
let onExitCb: ((code: number) => void) | null = null
return {
pid: 4242,
getForegroundProcess: vi.fn(() => null),
write: vi.fn(),
resize: vi.fn(),
kill: vi.fn(() => {
setTimeout(() => onExitCb?.(0), 1)
}),
forceKill: vi.fn(() => onExitCb?.(137)),
signal: vi.fn(),
onData(cb) {
dataCb = cb
},
onExit(cb) {
onExitCb = cb
},
dispose: vi.fn()
} as SubprocessHandle
}

beforeEach(() => {
dataCb = null
host = new TerminalHost({ spawnSubprocess: () => createMockSubprocess() })
})

afterEach(async () => {
await host.dispose()
})

it('caps retained rows at the window while keeping the newest content', async () => {
await host.createOrAttach({
sessionId: 'windowed',
cols: 80,
rows: 24,
streamClient: { onData: vi.fn(), onExit: vi.fn() }
})
const total = DAEMON_SESSION_SCROLLBACK_ROWS + 500
for (let i = 1; i <= total; i += 1) {
dataCb?.(`LINE_${String(i).padStart(5, '0')}\r\n`)
}
await vi.waitFor(() => {
const snapshot = host.getSnapshot('windowed')
expect(snapshot?.snapshotAnsi ?? snapshot?.scrollbackAnsi).toBeTruthy()
const text = `${snapshot?.scrollbackAnsi ?? ''}${snapshot?.snapshotAnsi ?? ''}`
// Newest line always present; the oldest has scrolled past the window.
expect(text).toContain(`LINE_${String(total).padStart(5, '0')}`)
expect(text).not.toContain('LINE_00001')
const retainedMatches = text.match(/LINE_\d{5}/g) ?? []
expect(retainedMatches.length).toBeLessThanOrEqual(DAEMON_SESSION_SCROLLBACK_ROWS + 24)
expect(retainedMatches.length).toBeGreaterThan(DAEMON_SESSION_SCROLLBACK_ROWS - 50)
})
})
})
Loading
Loading