Files
orca/src/relay/workspace-session-handler.ts
T
Neil 510305e574 fix(relay): signal capacity loss instead of dropping, hanging, or truncating (#17870)
Three failures with one shape: a payload past a fixed capacity was met with
silence, with a wait that never ends, or with a prefix presented as a whole.

**The workspace snapshot was silently dropped.** `workspace.changed` carries the
tab/session list, and a snapshot past the producer frame capacity (12288 B on a
Node <=21 remote) was dropped with only a relay stderr line, so the client kept a
stale list forever. The relay now publishes per client and, for a client whose
sink refused the frame, sends a compact `workspace.stale` marker on the control
lane; the client re-reads through `workspace.get`, whose lane is budgeted in
megabytes rather than in one producer frame. A new JSON-RPC notification rather
than a new field on `workspace.changed`: `normalizeSnapshot(undefined, ns)` yields
revision 0 and an empty session, so a Rule-1 field would make an old client
replace its tab list with nothing — worse than the drop. An old client ignores the
unknown method and is exactly where it is today. The marker retention/retry
machinery is extracted from the `fs.changed` overflow path and shared by both.

**The Windows upload hung, and the fix for it could truncate.** `#16432` was
attributed to `[Console]::In.ReadToEnd()` materializing the base64 bundle. That is
not what the reporter measured: he also measured
`new IO.StreamReader([Console]::OpenStandardInput())` — an incremental reader —
hanging at 1 MB. The limit is in the stdin the host hands PowerShell over a
non-pty ssh exec, not in the string the script builds.

- `uploadFileViaSystemSsh` — the user file-import path — was piping a whole file
  into one Windows stdin, unchunked and untimed. That is the path large files
  take; it now chunks into 32 KB writes and bounds each wait.
- The Windows directory upload reuses that single-file path rather than repeating
  a weaker copy of chunk-read + write-buffer; the `ino`/`dev` TOCTOU verification
  comes with it.
- A Windows write needing more than one exec lands on a `.orca-partial` staging
  path and is published by rename, so a failed chunk cannot leave a truncated
  artifact under the real name. `exclusive` is enforced once at the rename, not on
  the first chunk, where a retry met its own leftovers.
- The mkdir batch reads stdin through the stream reader the reporter measured
  surviving 50 KB, not `[Console]::In`, which he measured wedging at that size.
- `waitForChannelClose` takes an optional bound. A wedged PowerShell stays alive
  at idle CPU and never closes, so without one the promise is simply never
  settled and the caller waits forever with no error to show.

**Quick Open showed a prefix as the whole workspace.** The mechanism "a full page
means there is more" only works if the caller named the cap, and the failing UI
named none — it hardcoded `truncated: false`. Quick Open now names
`QUICK_OPEN_LISTING_MAX_RESULTS` on both the Electron IPC hop and the runtime-RPC
hop (the field #17954 added to `files.listAll`), and reads a full page as
truncation. The local hop honours the cap too, which it previously ignored.

Rebase note on `fs.listFiles`: an earlier revision of this work also clamped the
host unconditionally, and #17934 escalated an uncapped request to an explicit
error. #17954 has since landed and made an oversized reply streamable, which
removes the premise — the host no longer has to choose between a prefix and a
refusal, so it returns the whole listing when no limit is named and only clamps a
limit it was given. Keeping either would have regressed #17954 and hard-failed
three in-tree callers that deliberately pass no options
(`runtime-file-commands-search-runtime-files.ts:81`,
`filesystem-read-handlers.ts:125`, `runtime-file-commands-constructor.ts:41`).
2026-09-02 15:42:08 -07:00

193 lines
5.9 KiB
TypeScript

import { existsSync, mkdirSync, readFileSync, renameSync, writeFileSync } from 'node:fs'
import { homedir } from 'node:os'
import { dirname, join } from 'node:path'
import type { RelayDispatcher } from './dispatcher'
import { publishWorkspaceSnapshotChange } from './workspace-snapshot-publication'
type RemoteWorkspaceSnapshot = {
namespace: string
revision: number
updatedAt: number
schemaVersion: number
session: Record<string, unknown>
}
type ConnectedClient = {
clientId: string
name: string
lastSeenAt: number
}
type PatchResult =
| { ok: true; snapshot: RemoteWorkspaceSnapshot }
| {
ok: false
reason: 'stale-revision' | 'unavailable'
snapshot?: RemoteWorkspaceSnapshot
message?: string
}
const SNAPSHOT_SCHEMA_VERSION = 1
const PRESENCE_TTL_MS = 45_000
function emptySession(): Record<string, unknown> {
return {
activeRepoId: null,
activeWorktreeId: null,
activeTabId: null,
tabsByWorktree: {},
terminalLayoutsByTabId: {}
}
}
function sanitizeNamespace(namespace: unknown): string {
const raw = typeof namespace === 'string' && namespace.trim() ? namespace.trim() : 'default'
return raw.replace(/[^a-zA-Z0-9._-]/g, '_').slice(0, 160) || 'default'
}
function sanitizeClientName(value: string): string {
return value.replace(/\s+/g, ' ').trim().slice(0, 80)
}
function sanitizeClientId(value: string): string {
return value.trim().slice(0, 200)
}
export class WorkspaceSessionHandler {
private readonly clientsByNamespace = new Map<string, Map<string, ConnectedClient>>()
constructor(
private dispatcher: RelayDispatcher,
private baseDir = join(homedir(), '.orca', 'sessions')
) {
this.dispatcher.onRequest('workspace.get', (params) => this.get(params))
this.dispatcher.onRequest('workspace.patch', (params) => this.patch(params))
this.dispatcher.onRequest('workspace.presence', (params) => this.presence(params))
}
private snapshotPath(namespace: string): string {
return join(this.baseDir, `${namespace}.json`)
}
private read(namespace: string): RemoteWorkspaceSnapshot {
const path = this.snapshotPath(namespace)
if (!existsSync(path)) {
return {
namespace,
revision: 0,
updatedAt: 0,
schemaVersion: SNAPSHOT_SCHEMA_VERSION,
session: emptySession()
}
}
try {
const parsed = JSON.parse(readFileSync(path, 'utf-8')) as Partial<RemoteWorkspaceSnapshot>
return {
namespace,
revision:
typeof parsed.revision === 'number' && Number.isFinite(parsed.revision)
? parsed.revision
: 0,
updatedAt:
typeof parsed.updatedAt === 'number' && Number.isFinite(parsed.updatedAt)
? parsed.updatedAt
: 0,
schemaVersion:
typeof parsed.schemaVersion === 'number' && Number.isFinite(parsed.schemaVersion)
? parsed.schemaVersion
: SNAPSHOT_SCHEMA_VERSION,
session:
parsed.session && typeof parsed.session === 'object' && !Array.isArray(parsed.session)
? (parsed.session as Record<string, unknown>)
: emptySession()
}
} catch {
return {
namespace,
revision: 0,
updatedAt: 0,
schemaVersion: SNAPSHOT_SCHEMA_VERSION,
session: emptySession()
}
}
}
private write(snapshot: RemoteWorkspaceSnapshot): void {
const path = this.snapshotPath(snapshot.namespace)
mkdirSync(dirname(path), { recursive: true, mode: 0o700 })
const tmpPath = `${path}.tmp`
writeFileSync(tmpPath, JSON.stringify(snapshot, null, 2), { mode: 0o600 })
renameSync(tmpPath, path)
}
private async get(params: Record<string, unknown>): Promise<RemoteWorkspaceSnapshot> {
return this.read(sanitizeNamespace(params.namespace))
}
private async patch(params: Record<string, unknown>): Promise<PatchResult> {
const namespace = sanitizeNamespace(params.namespace)
const current = this.read(namespace)
const baseRevision = Number(params.baseRevision)
if (Number.isFinite(baseRevision) && baseRevision !== current.revision) {
return { ok: false, reason: 'stale-revision', snapshot: current }
}
const patch = params.patch as { kind?: unknown; session?: unknown } | undefined
if (
!patch ||
patch.kind !== 'replace-session' ||
!patch.session ||
typeof patch.session !== 'object' ||
Array.isArray(patch.session)
) {
return { ok: false, reason: 'unavailable', message: 'Invalid workspace patch' }
}
const snapshot: RemoteWorkspaceSnapshot = {
namespace,
revision: current.revision + 1,
updatedAt: Date.now(),
schemaVersion: SNAPSHOT_SCHEMA_VERSION,
session: patch.session as Record<string, unknown>
}
this.write(snapshot)
publishWorkspaceSnapshotChange(
this.dispatcher,
{
namespace,
snapshot,
sourceClientId: typeof params.clientId === 'string' ? params.clientId : undefined
},
namespace
)
return { ok: true, snapshot }
}
private async presence(params: Record<string, unknown>): Promise<{ clients: ConnectedClient[] }> {
const namespace = sanitizeNamespace(params.namespace)
const clientId = typeof params.clientId === 'string' ? sanitizeClientId(params.clientId) : ''
const name = typeof params.clientName === 'string' ? sanitizeClientName(params.clientName) : ''
const clients = this.clientsByNamespace.get(namespace) ?? new Map<string, ConnectedClient>()
this.clientsByNamespace.set(namespace, clients)
const now = Date.now()
for (const [id, client] of clients) {
if (now - client.lastSeenAt > PRESENCE_TTL_MS) {
clients.delete(id)
}
}
if (clientId) {
clients.set(clientId, {
clientId,
name: name || 'Unknown device',
lastSeenAt: now
})
}
return {
clients: Array.from(clients.values()).sort((a, b) => b.lastSeenAt - a.lastSeenAt)
}
}
}