mirror of
https://github.com/stablyai/orca.git
synced 2026-09-29 16:02:50 +00:00
fix(wsl): coalesce relay identity reads across polling
This commit is contained in:
@@ -22,6 +22,10 @@ export {
|
||||
RETIRED_PANE_FENCES_MAX
|
||||
} from './server/server-constants'
|
||||
export { isValidPaneKey }
|
||||
export {
|
||||
getAgentHookRemotePostCount,
|
||||
resetAgentHookRemotePostCount
|
||||
} from './server/server-ingest-remote'
|
||||
|
||||
/** Public composition seam for the loopback hook listener and relay status adapter. */
|
||||
export class AgentHookServer extends AgentHookServerLifecycle {}
|
||||
|
||||
@@ -19,6 +19,19 @@ import type { AgentHookEventPayload } from '../../../shared/agent-hook-listener/
|
||||
import { isValidPiProviderSessionOnly } from './server-status-identity'
|
||||
import { AgentHookServerIngestTerminal } from './server-ingest-terminal'
|
||||
|
||||
// Diagnostics seam for the WSL hooks-off contract. This counts relay posts
|
||||
// reaching the host boundary; local loopback hook traffic is intentionally not
|
||||
// included.
|
||||
let agentHookRemotePostCount = 0
|
||||
|
||||
export function getAgentHookRemotePostCount(): number {
|
||||
return agentHookRemotePostCount
|
||||
}
|
||||
|
||||
export function resetAgentHookRemotePostCount(): void {
|
||||
agentHookRemotePostCount = 0
|
||||
}
|
||||
|
||||
export abstract class AgentHookServerIngestRemote extends AgentHookServerIngestTerminal {
|
||||
/** Ingest a payload from the relay JSON-RPC channel (not the local HTTP server); connectionId is stamped here. Main is still the SSH trust boundary, so re-run the canonical normalizer before caching. */
|
||||
ingestRemote(
|
||||
@@ -49,6 +62,7 @@ export abstract class AgentHookServerIngestRemote extends AgentHookServerIngestT
|
||||
},
|
||||
connectionId: string | null
|
||||
): void {
|
||||
agentHookRemotePostCount++
|
||||
// Why: wire crosses a trust boundary — re-check/trim so an empty connectionId can't poison caches.
|
||||
if (connectionId !== null && typeof connectionId !== 'string') {
|
||||
return
|
||||
|
||||
@@ -31,6 +31,18 @@ type GuestInstallState = {
|
||||
lastInstallAt?: number
|
||||
}
|
||||
|
||||
// Diagnostics seam for the hooks-off contract: relay residency must remain
|
||||
// live while no guest hook installation work is attempted.
|
||||
let wslRelayGuestHookInstallCount = 0
|
||||
|
||||
export function getWslRelayGuestHookInstallCount(): number {
|
||||
return wslRelayGuestHookInstallCount
|
||||
}
|
||||
|
||||
export function resetWslRelayGuestHookInstallCount(): void {
|
||||
wslRelayGuestHookInstallCount = 0
|
||||
}
|
||||
|
||||
export async function runWslRelayGuestInstall(
|
||||
deps: GuestInstallDeps,
|
||||
state: GuestInstallState,
|
||||
@@ -43,6 +55,7 @@ export async function runWslRelayGuestInstall(
|
||||
) {
|
||||
return
|
||||
}
|
||||
wslRelayGuestHookInstallCount++
|
||||
state.lastInstallAt = Date.now()
|
||||
await installWslGuestHooks({
|
||||
mux,
|
||||
|
||||
@@ -19,7 +19,7 @@ export async function readWslRelayProcessIdentity(options: {
|
||||
ensure: (distro: string) => Promise<void> | void
|
||||
getState: (distro: string) => { phase: string; mux?: SshChannelMultiplexer } | undefined
|
||||
disposed: boolean
|
||||
requestOptions?: { signal?: AbortSignal; timeoutMs?: number }
|
||||
requestOptions?: { signal?: AbortSignal; timeoutMs?: number; stableUntilReset?: boolean }
|
||||
}): Promise<WslRelayIdentityResult[]> {
|
||||
const unavailable = (reason: string): WslRelayIdentityResult[] =>
|
||||
options.anchors.map(() => ({ status: 'unverifiable' as const, reason, capturedAgeMs: 0 }))
|
||||
|
||||
@@ -18,8 +18,11 @@ import {
|
||||
import {
|
||||
getWslRelayIdentityRpcCount,
|
||||
resetWslRelayIdentityRpcCount,
|
||||
getWslRelayGuestHookInstallCount,
|
||||
resetWslRelayGuestHookInstallCount,
|
||||
WslHookRelayManager
|
||||
} from './wsl-hook-relay-manager'
|
||||
import { getAgentHookRemotePostCount, resetAgentHookRemotePostCount } from './server'
|
||||
import { FAILURE_COOLDOWN_BASE_MS, type WslHookRelayManagerDeps } from './wsl-hook-relay-deps'
|
||||
import {
|
||||
AGENT_HOOK_INSTALL_PLUGINS_METHOD,
|
||||
@@ -450,6 +453,8 @@ describe('WslHookRelayManager', () => {
|
||||
|
||||
it('serves identity reads through connectWslRelayState with hooks disabled', async () => {
|
||||
resetWslRelayIdentityRpcCount()
|
||||
resetWslRelayGuestHookInstallCount()
|
||||
resetAgentHookRemotePostCount()
|
||||
const waitForSentinel = vi.fn(async () => {
|
||||
const transport = guestTransport()
|
||||
const guest = harnesses.at(-1)!.guestDispatcher
|
||||
@@ -482,6 +487,8 @@ describe('WslHookRelayManager', () => {
|
||||
{ status: 'live', processName: 'claude' }
|
||||
])
|
||||
expect(getWslRelayIdentityRpcCount()).toBe(1)
|
||||
expect(getWslRelayGuestHookInstallCount()).toBe(0)
|
||||
expect(getAgentHookRemotePostCount()).toBe(0)
|
||||
expect(deps.installHooks).not.toHaveBeenCalled()
|
||||
manager.disposeAll()
|
||||
})
|
||||
|
||||
@@ -24,6 +24,10 @@ import {
|
||||
wslRelayIdentityRpcCount
|
||||
} from './wsl-hook-relay-identity'
|
||||
import { wslHookRelayStateKey } from './wsl-hook-relay-state-key'
|
||||
import {
|
||||
getWslRelayGuestHookInstallCount,
|
||||
resetWslRelayGuestHookInstallCount
|
||||
} from './wsl-hook-relay-guest-install'
|
||||
import {
|
||||
sanitizeWslHookInstanceKey,
|
||||
WSL_RELAY_HOOKS_SET_ENABLED_METHOD
|
||||
@@ -44,6 +48,8 @@ export function resetWslRelayIdentityRpcCount(): void {
|
||||
resetIdentityCounter()
|
||||
}
|
||||
|
||||
export { getWslRelayGuestHookInstallCount, resetWslRelayGuestHookInstallCount }
|
||||
|
||||
export class WslHookRelayManager {
|
||||
private deps: WslHookRelayManagerDeps
|
||||
private recovery: WslRelayRecovery
|
||||
@@ -155,7 +161,7 @@ export class WslHookRelayManager {
|
||||
readProcessIdentity = (
|
||||
distro: string,
|
||||
anchors: Parameters<typeof readWslRelayProcessIdentity>[0]['anchors'],
|
||||
options?: { signal?: AbortSignal; timeoutMs?: number }
|
||||
options?: { signal?: AbortSignal; timeoutMs?: number; stableUntilReset?: boolean }
|
||||
) =>
|
||||
readWslRelayProcessIdentity({
|
||||
distro,
|
||||
|
||||
@@ -71,6 +71,7 @@ export abstract class DaemonPtySessionInventory extends DaemonPtyProcessInspecti
|
||||
sessions.map((session) => session.wslShellAnchor!),
|
||||
{
|
||||
signal: opts?.signal,
|
||||
stableUntilReset: true,
|
||||
timeoutMs:
|
||||
opts?.deadlineMs === undefined
|
||||
? undefined
|
||||
|
||||
@@ -20,12 +20,14 @@ import { resolveWslSessionContext } from './wsl-session-context'
|
||||
import { resolveSafePtyDefaultCwd } from '../providers/pty-default-cwd'
|
||||
import { resolveUnixShellPath } from '../providers/local-pty-utils'
|
||||
import type { PtySpawnOptions, PtySpawnResult } from '../providers/types'
|
||||
import { wslRelayIdentityReader } from '../providers/wsl-relay-identity-reader'
|
||||
import { injectHistoryEnv, injectWslFishHistoryEnv, logHistoryInjection } from '../terminal-history'
|
||||
import { addWslEnvKeys } from '../wsl-env'
|
||||
|
||||
export abstract class DaemonPtySessionSpawn extends DaemonPtySpawnResult {
|
||||
async spawn(opts: PtySpawnOptions): Promise<PtySpawnResult> {
|
||||
const spawnOpts = this.withHistoryIsolation(opts)
|
||||
wslRelayIdentityReader.reset()
|
||||
const sessionId = spawnOpts.sessionId ?? mintPtySessionId(spawnOpts.worktreeId)
|
||||
const operation: PendingDaemonSpawnOperation = {
|
||||
exitsBySessionId: new Map(),
|
||||
|
||||
@@ -122,6 +122,7 @@ import {
|
||||
ptyWslShellAnchors
|
||||
} from './local-pty-provider-state'
|
||||
import { wslRelayIdentityReader } from './wsl-relay-identity-reader'
|
||||
import { wslHookRelayManager } from '../agent-hooks/wsl-hook-relay-manager'
|
||||
import {
|
||||
applyLocalPtyProviderMockDefaults,
|
||||
createLocalPtyMockProcess,
|
||||
@@ -250,6 +251,44 @@ describe('LocalPtyProvider', () => {
|
||||
readBatch.mockRestore()
|
||||
}
|
||||
})
|
||||
|
||||
it('coalesces repeated WSL list-process bursts until a PTY event resets identity', async () => {
|
||||
Object.defineProperty(process, 'platform', { configurable: true, value: 'win32' })
|
||||
const spawned = await provider.spawn({ cols: 80, rows: 24, cwd: '/tmp/wsl-owned-cwd' })
|
||||
const anchor = {
|
||||
distro: 'Ubuntu',
|
||||
bootId: '11111111-1111-1111-1111-111111111111',
|
||||
shellPid: 100,
|
||||
shellStartTime: 1,
|
||||
tty: '/dev/pts/1'
|
||||
}
|
||||
ptyWslDistroById.set(spawned.id, 'Ubuntu')
|
||||
ptyWslShellAnchors.set(spawned.id, anchor)
|
||||
wslRelayIdentityReader.reset()
|
||||
const identityRequest = vi
|
||||
.spyOn(wslHookRelayManager, 'readProcessIdentity')
|
||||
.mockResolvedValue([
|
||||
{
|
||||
status: 'unverifiable' as const,
|
||||
reason: 'relay_unavailable',
|
||||
capturedAgeMs: 0
|
||||
}
|
||||
])
|
||||
try {
|
||||
await Promise.all(Array.from({ length: 5 }, () => provider.listProcesses()))
|
||||
expect(identityRequest).toHaveBeenCalledOnce()
|
||||
|
||||
await Promise.all(Array.from({ length: 5 }, () => provider.listProcesses()))
|
||||
expect(identityRequest).toHaveBeenCalledOnce()
|
||||
|
||||
wslRelayIdentityReader.reset()
|
||||
await provider.listProcesses()
|
||||
expect(identityRequest).toHaveBeenCalledTimes(2)
|
||||
} finally {
|
||||
identityRequest.mockRestore()
|
||||
clearPtyState(spawned.id)
|
||||
}
|
||||
})
|
||||
})
|
||||
|
||||
describe('getDefaultShell', () => {
|
||||
|
||||
@@ -41,6 +41,7 @@ import {
|
||||
type ShellStartupIdentityScanState
|
||||
} from '../shell-startup-identity-scanner'
|
||||
import type { PtySpawnOptions, PtySpawnResult } from './types'
|
||||
import { wslRelayIdentityReader } from './wsl-relay-identity-reader'
|
||||
|
||||
export function activateLocalPtySession(args: {
|
||||
id: string
|
||||
@@ -55,6 +56,8 @@ export function activateLocalPtySession(args: {
|
||||
}): PtySpawnResult {
|
||||
const { id, incarnationId, spawn, getOptions, plan, env, proc, spawnedWslDistro } = args
|
||||
createPtyPhysicalExit(id)
|
||||
// A newly created PTY can introduce a fresh WSL shell identity.
|
||||
wslRelayIdentityReader.reset()
|
||||
ptyReportsChildExitStatus.set(id, args.reportsChildExitStatus)
|
||||
ptyProcesses.set(id, proc)
|
||||
ptyInitialCwd.set(id, plan.cwd)
|
||||
@@ -78,13 +81,20 @@ export function activateLocalPtySession(args: {
|
||||
)
|
||||
ptyLoadGeneration.set(id, getLoadGeneration())
|
||||
ptyIncarnations.set(id, incarnationId)
|
||||
// Foreground identity is event-invalidated: once output arrives, the next
|
||||
// inventory may observe an agent transition instead of reusing the old read.
|
||||
const resetWslIdentityCache = (): void => {
|
||||
wslRelayIdentityReader.reset()
|
||||
}
|
||||
getOptions().onSpawned?.(id, incarnationId)
|
||||
|
||||
const emitIngressData = (emission: PtyIngressEmission): void => {
|
||||
const sequenceChars = emission.rawEndSeq - emission.rawStartSeq
|
||||
if (emission.transformed || sequenceChars !== emission.data.length) {
|
||||
resetWslIdentityCache()
|
||||
getOptions().onData?.(id, emission.data, Date.now(), sequenceChars, true)
|
||||
} else {
|
||||
resetWslIdentityCache()
|
||||
getOptions().onData?.(id, emission.data, Date.now())
|
||||
}
|
||||
for (const cb of dataListeners) {
|
||||
|
||||
@@ -148,6 +148,9 @@ export async function listLocalPtyProcesses(opts?: {
|
||||
.filter((anchor): anchor is NonNullable<typeof anchor> => anchor !== undefined)
|
||||
const results = await wslRelayIdentityReader.readBatch(distro, anchors, {
|
||||
signal: opts?.signal,
|
||||
// A process inventory is reused until a local PTY event resets the
|
||||
// reader; polling itself is not evidence that guest state changed.
|
||||
stableUntilReset: true,
|
||||
timeoutMs:
|
||||
opts?.deadlineMs === undefined ? undefined : Math.max(1, opts.deadlineMs - Date.now())
|
||||
})
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
import { afterEach, describe, expect, it, vi } from 'vitest'
|
||||
|
||||
import { createWslRelayIdentityReader } from './wsl-relay-identity-reader'
|
||||
import type { WslRelayIdentityResult } from '../../shared/wsl-hook-relay-contract'
|
||||
|
||||
describe('wsl relay identity reader cost invariant', () => {
|
||||
afterEach(() => {
|
||||
@@ -31,4 +32,166 @@ describe('wsl relay identity reader cost invariant', () => {
|
||||
await expect(reader.read('Ubuntu', anchor)).resolves.toMatchObject({ capturedAgeMs: 260 })
|
||||
expect(readProcessIdentity).toHaveBeenCalledOnce()
|
||||
})
|
||||
|
||||
it('keeps inventory reads stable until a PTY event resets the cache', async () => {
|
||||
vi.useFakeTimers()
|
||||
vi.setSystemTime(0)
|
||||
const readProcessIdentity = vi.fn().mockResolvedValue([
|
||||
{
|
||||
status: 'live' as const,
|
||||
processName: 'claude',
|
||||
anchor: {
|
||||
distro: 'Ubuntu',
|
||||
bootId: '11111111-1111-1111-1111-111111111111',
|
||||
shellPid: 1,
|
||||
shellStartTime: 1,
|
||||
tty: '/dev/pts/1'
|
||||
},
|
||||
capturedAgeMs: 0
|
||||
}
|
||||
])
|
||||
const reader = createWslRelayIdentityReader({ readProcessIdentity })
|
||||
const anchor = {
|
||||
distro: 'Ubuntu',
|
||||
bootId: '11111111-1111-1111-1111-111111111111',
|
||||
shellPid: 1,
|
||||
shellStartTime: 1,
|
||||
tty: '/dev/pts/1'
|
||||
}
|
||||
|
||||
await reader.readBatch('Ubuntu', [anchor], { stableUntilReset: true })
|
||||
vi.setSystemTime(60_000)
|
||||
await reader.readBatch('Ubuntu', [anchor], { stableUntilReset: true })
|
||||
expect(readProcessIdentity).toHaveBeenCalledOnce()
|
||||
|
||||
reader.reset()
|
||||
await reader.readBatch('Ubuntu', [anchor], { stableUntilReset: true })
|
||||
expect(readProcessIdentity).toHaveBeenCalledTimes(2)
|
||||
})
|
||||
|
||||
it('coalesces list-session bursts and reuses the capture until reset', async () => {
|
||||
vi.useFakeTimers()
|
||||
vi.setSystemTime(0)
|
||||
const readProcessIdentity = vi.fn().mockResolvedValue([
|
||||
{
|
||||
status: 'live' as const,
|
||||
processName: 'claude',
|
||||
anchor: {
|
||||
distro: 'Ubuntu',
|
||||
bootId: '11111111-1111-1111-1111-111111111111',
|
||||
shellPid: 1,
|
||||
shellStartTime: 1,
|
||||
tty: '/dev/pts/1'
|
||||
},
|
||||
capturedAgeMs: 0
|
||||
}
|
||||
])
|
||||
const reader = createWslRelayIdentityReader({ readProcessIdentity })
|
||||
const anchor = {
|
||||
distro: 'Ubuntu',
|
||||
bootId: '11111111-1111-1111-1111-111111111111',
|
||||
shellPid: 1,
|
||||
shellStartTime: 1,
|
||||
tty: '/dev/pts/1'
|
||||
}
|
||||
|
||||
await Promise.all(
|
||||
Array.from({ length: 5 }, () =>
|
||||
reader.readBatch('Ubuntu', [anchor], { stableUntilReset: true })
|
||||
)
|
||||
)
|
||||
expect(readProcessIdentity).toHaveBeenCalledOnce()
|
||||
|
||||
vi.setSystemTime(60_000)
|
||||
await Promise.all(
|
||||
Array.from({ length: 5 }, () =>
|
||||
reader.readBatch('Ubuntu', [anchor], { stableUntilReset: true })
|
||||
)
|
||||
)
|
||||
expect(readProcessIdentity).toHaveBeenCalledOnce()
|
||||
|
||||
reader.reset()
|
||||
await reader.readBatch('Ubuntu', [anchor], { stableUntilReset: true })
|
||||
expect(readProcessIdentity).toHaveBeenCalledTimes(2)
|
||||
})
|
||||
|
||||
it('does not republish an in-flight capture after an event reset', async () => {
|
||||
let resolveCapture!: (result: WslRelayIdentityResult[]) => void
|
||||
const replacementResult: WslRelayIdentityResult[] = [
|
||||
{
|
||||
status: 'live',
|
||||
processName: 'claude',
|
||||
anchor: {
|
||||
distro: 'Ubuntu',
|
||||
bootId: '11111111-1111-1111-1111-111111111111',
|
||||
shellPid: 1,
|
||||
shellStartTime: 1,
|
||||
tty: '/dev/pts/1'
|
||||
},
|
||||
capturedAgeMs: 0
|
||||
}
|
||||
]
|
||||
const readProcessIdentity = vi.fn(
|
||||
() =>
|
||||
new Promise<WslRelayIdentityResult[]>((resolve) => {
|
||||
resolveCapture = resolve
|
||||
})
|
||||
)
|
||||
const reader = createWslRelayIdentityReader({ readProcessIdentity })
|
||||
const anchor = {
|
||||
distro: 'Ubuntu',
|
||||
bootId: '11111111-1111-1111-1111-111111111111',
|
||||
shellPid: 1,
|
||||
shellStartTime: 1,
|
||||
tty: '/dev/pts/1'
|
||||
}
|
||||
const first = reader.readBatch('Ubuntu', [anchor], { stableUntilReset: true })
|
||||
expect(readProcessIdentity).toHaveBeenCalledOnce()
|
||||
reader.reset()
|
||||
resolveCapture(replacementResult)
|
||||
await first
|
||||
|
||||
const second = reader.readBatch('Ubuntu', [anchor], { stableUntilReset: true })
|
||||
// The second request is intentionally independent of the stale first one.
|
||||
resolveCapture(replacementResult)
|
||||
await second
|
||||
expect(readProcessIdentity).toHaveBeenCalledTimes(2)
|
||||
})
|
||||
|
||||
it('keeps a stable snapshot marked stable after a foreground read', async () => {
|
||||
vi.useFakeTimers()
|
||||
vi.setSystemTime(0)
|
||||
const readProcessIdentity = vi.fn().mockResolvedValue([
|
||||
{
|
||||
status: 'live' as const,
|
||||
processName: 'claude',
|
||||
anchor: {
|
||||
distro: 'Ubuntu',
|
||||
bootId: '11111111-1111-1111-1111-111111111111',
|
||||
shellPid: 1,
|
||||
shellStartTime: 1,
|
||||
tty: '/dev/pts/1'
|
||||
},
|
||||
capturedAgeMs: 0
|
||||
}
|
||||
])
|
||||
const reader = createWslRelayIdentityReader({ readProcessIdentity })
|
||||
const anchor = {
|
||||
distro: 'Ubuntu',
|
||||
bootId: '11111111-1111-1111-1111-111111111111',
|
||||
shellPid: 1,
|
||||
shellStartTime: 1,
|
||||
tty: '/dev/pts/1'
|
||||
}
|
||||
|
||||
await reader.readBatch('Ubuntu', [anchor], { stableUntilReset: true })
|
||||
vi.setSystemTime(60_000)
|
||||
await reader.readBatch('Ubuntu', [anchor])
|
||||
vi.setSystemTime(120_000)
|
||||
await reader.readBatch('Ubuntu', [anchor], { stableUntilReset: true })
|
||||
expect(readProcessIdentity).toHaveBeenCalledTimes(2)
|
||||
vi.setSystemTime(180_000)
|
||||
await reader.readBatch('Ubuntu', [anchor], { stableUntilReset: true })
|
||||
expect(readProcessIdentity).toHaveBeenCalledTimes(2)
|
||||
})
|
||||
})
|
||||
|
||||
@@ -16,7 +16,7 @@ export type WslRelayIdentityReader = {
|
||||
readBatch: (
|
||||
distro: string,
|
||||
anchors: readonly WslShellProcessAnchor[],
|
||||
options?: { signal?: AbortSignal; timeoutMs?: number }
|
||||
options?: { signal?: AbortSignal; timeoutMs?: number; stableUntilReset?: boolean }
|
||||
) => Promise<WslRelayIdentityResult[]>
|
||||
reset: () => void
|
||||
}
|
||||
@@ -26,21 +26,40 @@ export function createWslRelayIdentityReader(
|
||||
): WslRelayIdentityReader {
|
||||
const cache = new Map<
|
||||
string,
|
||||
{ at: number; anchors: readonly WslShellProcessAnchor[]; results: WslRelayIdentityResult[] }
|
||||
{
|
||||
at: number
|
||||
anchors: readonly WslShellProcessAnchor[]
|
||||
results: WslRelayIdentityResult[]
|
||||
stableUntilReset: boolean
|
||||
}
|
||||
>()
|
||||
const pending = new Map<
|
||||
string,
|
||||
Promise<{ anchors: readonly WslShellProcessAnchor[]; results: WslRelayIdentityResult[] }>
|
||||
{
|
||||
promise: Promise<{
|
||||
anchors: readonly WslShellProcessAnchor[]
|
||||
results: WslRelayIdentityResult[]
|
||||
}>
|
||||
stableUntilReset: boolean
|
||||
resetGeneration: number
|
||||
}
|
||||
>()
|
||||
let resetGeneration = 0
|
||||
const readBatch = async (
|
||||
distro: string,
|
||||
anchors: readonly WslShellProcessAnchor[],
|
||||
options?: { signal?: AbortSignal; timeoutMs?: number }
|
||||
options?: { signal?: AbortSignal; timeoutMs?: number; stableUntilReset?: boolean }
|
||||
): Promise<WslRelayIdentityResult[]> => {
|
||||
const key = distro.trim().toLowerCase()
|
||||
const now = Date.now()
|
||||
const prior = cache.get(key)
|
||||
if (prior && now - prior.at < 500) {
|
||||
if (
|
||||
prior &&
|
||||
((options?.stableUntilReset === true && prior.stableUntilReset) || now - prior.at < 500)
|
||||
) {
|
||||
if (options?.stableUntilReset === true) {
|
||||
prior.stableUntilReset = true
|
||||
}
|
||||
const byAnchor = new Map(
|
||||
prior.anchors.map((anchor, index) => [JSON.stringify(anchor), prior.results[index]!])
|
||||
)
|
||||
@@ -56,7 +75,10 @@ export function createWslRelayIdentityReader(
|
||||
}
|
||||
const active = pending.get(key)
|
||||
if (active) {
|
||||
const result = await active
|
||||
if (options?.stableUntilReset === true) {
|
||||
active.stableUntilReset = true
|
||||
}
|
||||
const result = await active.promise
|
||||
const byAnchor = new Map(
|
||||
result.anchors.map((anchor, index) => [JSON.stringify(anchor), result.results[index]!])
|
||||
)
|
||||
@@ -64,16 +86,27 @@ export function createWslRelayIdentityReader(
|
||||
return anchors.map((anchor) => byAnchor.get(JSON.stringify(anchor))!)
|
||||
}
|
||||
}
|
||||
const request = manager
|
||||
.readProcessIdentity(distro, anchors, options)
|
||||
.then((results) => ({ anchors, results }))
|
||||
pending.set(key, request)
|
||||
const requestGeneration = resetGeneration
|
||||
const activeRequest = {
|
||||
stableUntilReset: options?.stableUntilReset === true,
|
||||
resetGeneration: requestGeneration,
|
||||
promise: manager
|
||||
.readProcessIdentity(distro, anchors, options)
|
||||
.then((results) => ({ anchors, results }))
|
||||
}
|
||||
pending.set(key, activeRequest)
|
||||
try {
|
||||
const result = await request
|
||||
cache.set(key, { at: Date.now(), ...result })
|
||||
const result = await activeRequest.promise
|
||||
if (activeRequest.resetGeneration === resetGeneration) {
|
||||
cache.set(key, {
|
||||
at: Date.now(),
|
||||
stableUntilReset: activeRequest.stableUntilReset || prior?.stableUntilReset === true,
|
||||
...result
|
||||
})
|
||||
}
|
||||
return result.results
|
||||
} finally {
|
||||
if (pending.get(key) === request) {
|
||||
if (pending.get(key) === activeRequest) {
|
||||
pending.delete(key)
|
||||
}
|
||||
}
|
||||
@@ -82,6 +115,7 @@ export function createWslRelayIdentityReader(
|
||||
read: async (distro, anchor, options) => (await readBatch(distro, [anchor], options))[0]!,
|
||||
readBatch,
|
||||
reset: () => {
|
||||
resetGeneration += 1
|
||||
cache.clear()
|
||||
pending.clear()
|
||||
}
|
||||
|
||||
@@ -67,11 +67,13 @@ export class OrcaRuntimeWithRefreshFloatingWorkspacePtyLiveness extends OrcaRunt
|
||||
const binding = persistedBindingByPtyId.get(ptyId)
|
||||
if (!pty && binding) {
|
||||
// Why: a live daemon PTY restored from disk needs its pane identity before mobile can issue a safe handle.
|
||||
pty = this.recordPtyWorktree(ptyId, FLOATING_TERMINAL_WORKTREE_ID, {
|
||||
connected: true,
|
||||
tabId: binding.tabId,
|
||||
paneKey: binding.paneKey
|
||||
})
|
||||
pty = this.withPtyLivenessRefreshBookkeeping(() =>
|
||||
this.recordPtyWorktree(ptyId, FLOATING_TERMINAL_WORKTREE_ID, {
|
||||
connected: true,
|
||||
tabId: binding.tabId,
|
||||
paneKey: binding.paneKey
|
||||
})
|
||||
)
|
||||
}
|
||||
if (pty) {
|
||||
pty.connected = true
|
||||
|
||||
@@ -9,6 +9,36 @@ export class OrcaRuntimeWithRefreshPtyWorktreeRecordsFromController extends Orca
|
||||
targetWorktreeId: string | null = null,
|
||||
deadline?: number,
|
||||
signal?: AbortSignal
|
||||
): Promise<Set<string> | null> {
|
||||
const refreshKey = targetWorktreeId === null ? 'aggregate' : `target:${targetWorktreeId}`
|
||||
const pending = this.ptyLivenessRefreshPromises.get(refreshKey)
|
||||
if (pending) {
|
||||
return pending
|
||||
}
|
||||
const invalidationSequence = this.ptyLivenessInvalidationSequence
|
||||
const refresh = this.refreshPtyWorktreeRecordsFromControllerUncoalesced(
|
||||
resolvedWorktrees,
|
||||
targetWorktreeId,
|
||||
deadline,
|
||||
signal,
|
||||
invalidationSequence
|
||||
)
|
||||
this.ptyLivenessRefreshPromises.set(refreshKey, refresh)
|
||||
try {
|
||||
return await refresh
|
||||
} finally {
|
||||
if (this.ptyLivenessRefreshPromises.get(refreshKey) === refresh) {
|
||||
this.ptyLivenessRefreshPromises.delete(refreshKey)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private async refreshPtyWorktreeRecordsFromControllerUncoalesced(
|
||||
resolvedWorktrees: ResolvedWorktree[],
|
||||
targetWorktreeId: string | null,
|
||||
deadline?: number,
|
||||
signal?: AbortSignal,
|
||||
invalidationSequence = this.ptyLivenessInvalidationSequence
|
||||
): Promise<Set<string> | null> {
|
||||
this.ptyLivenessRefreshInProgress += 1
|
||||
try {
|
||||
@@ -20,7 +50,11 @@ export class OrcaRuntimeWithRefreshPtyWorktreeRecordsFromController extends Orca
|
||||
false,
|
||||
signal ? { signal } : undefined
|
||||
)
|
||||
if (inventory) {
|
||||
if (
|
||||
inventory &&
|
||||
targetWorktreeId === null &&
|
||||
invalidationSequence === this.ptyLivenessInvalidationSequence
|
||||
) {
|
||||
this.ptyLivenessRefreshRequired = false
|
||||
}
|
||||
return inventory ? new Set(inventory.livePtyIds) : null
|
||||
|
||||
+13
-11
@@ -217,17 +217,19 @@ export class OrcaRuntimeWithRefreshPtyWorktreeRecordsWithControllerInventory ext
|
||||
}
|
||||
this.restoredOrchestrationAuthorityByPtyId.delete(session.id)
|
||||
if (worktreeId) {
|
||||
const pty = this.recordPtyWorktree(session.id, worktreeId, {
|
||||
connected: true,
|
||||
...(session.incarnationId ? { incarnationId: session.incarnationId } : {}),
|
||||
agentSessionOwners: session.incarnationId ? (session.agentSessionOwners ?? []) : [],
|
||||
...(session.wslDistro !== undefined
|
||||
? { isWsl: Boolean(session.wslDistro), wslDistro: session.wslDistro }
|
||||
: {}),
|
||||
...(restoresExactSurface
|
||||
? { tabId: persistedSurface.tabId, paneKey: persistedSurface.paneKey }
|
||||
: {})
|
||||
})
|
||||
const pty = this.withPtyLivenessRefreshBookkeeping(() =>
|
||||
this.recordPtyWorktree(session.id, worktreeId, {
|
||||
connected: true,
|
||||
...(session.incarnationId ? { incarnationId: session.incarnationId } : {}),
|
||||
agentSessionOwners: session.incarnationId ? (session.agentSessionOwners ?? []) : [],
|
||||
...(session.wslDistro !== undefined
|
||||
? { isWsl: Boolean(session.wslDistro), wslDistro: session.wslDistro }
|
||||
: {}),
|
||||
...(restoresExactSurface
|
||||
? { tabId: persistedSurface.tabId, paneKey: persistedSurface.paneKey }
|
||||
: {})
|
||||
})
|
||||
)
|
||||
if (restoresExactSurface && controllerIdentity) {
|
||||
this.rememberRestoredOrchestrationAuthority(
|
||||
pty,
|
||||
|
||||
@@ -342,9 +342,31 @@ export class OrcaRuntimeWithRuntimeId {
|
||||
|
||||
protected ptyLivenessRefreshInProgress = 0
|
||||
|
||||
// Concurrent catalog requests share one event-invalidated controller census;
|
||||
// without this, each request observes the stale flag before the first one clears it.
|
||||
protected ptyLivenessRefreshPromises = new Map<string, Promise<Set<string> | null>>()
|
||||
|
||||
// Monotonic event fence prevents an in-flight census from clearing a newer invalidation.
|
||||
protected ptyLivenessInvalidationSequence = 0
|
||||
|
||||
// Reconciliation bookkeeping must not look like a new PTY event.
|
||||
protected ptyLivenessRefreshBookkeeping = false
|
||||
|
||||
protected invalidatePtyLivenessSnapshot(): void {
|
||||
if (this.ptyLivenessRefreshInProgress === 0) {
|
||||
this.ptyLivenessRefreshRequired = true
|
||||
if (this.ptyLivenessRefreshBookkeeping) {
|
||||
return
|
||||
}
|
||||
this.ptyLivenessRefreshRequired = true
|
||||
this.ptyLivenessInvalidationSequence += 1
|
||||
}
|
||||
|
||||
protected withPtyLivenessRefreshBookkeeping<T>(callback: () => T): T {
|
||||
const prior = this.ptyLivenessRefreshBookkeeping
|
||||
this.ptyLivenessRefreshBookkeeping = true
|
||||
try {
|
||||
return callback()
|
||||
} finally {
|
||||
this.ptyLivenessRefreshBookkeeping = prior
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2,6 +2,7 @@ import { describe, expect, it, vi } from 'vitest'
|
||||
import {
|
||||
OrcaRuntimeService,
|
||||
registerSshGitProvider,
|
||||
resetPlatform,
|
||||
setPlatform,
|
||||
worktreePathComparison
|
||||
} from '../orca-runtime-test-mocks.spec'
|
||||
@@ -11,6 +12,24 @@ import {
|
||||
getWslRelayIdentityRpcCount,
|
||||
resetWslRelayIdentityRpcCount
|
||||
} from '../../agent-hooks/wsl-hook-relay-manager'
|
||||
import { listLocalPtyProcesses } from '../../providers/local-pty-session-operations'
|
||||
import {
|
||||
clearPtyState,
|
||||
ptyInitialCwd,
|
||||
ptyProcesses,
|
||||
ptyWslDistroById,
|
||||
ptyWslShellAnchors,
|
||||
ptyWorktreeId
|
||||
} from '../../providers/local-pty-provider-state'
|
||||
import { wslRelayIdentityReader } from '../../providers/wsl-relay-identity-reader'
|
||||
|
||||
function deferred<T>(): { promise: Promise<T>; resolve: (value: T) => void } {
|
||||
let resolve!: (value: T) => void
|
||||
const promise = new Promise<T>((next) => {
|
||||
resolve = next
|
||||
})
|
||||
return { promise, resolve }
|
||||
}
|
||||
|
||||
describe('OrcaRuntimeService', () => {
|
||||
it('does not read guest inventories on idle worktree poll ticks', async () => {
|
||||
@@ -42,6 +61,138 @@ describe('OrcaRuntimeService', () => {
|
||||
await runtime.getWorktreePs()
|
||||
expect(relayProcessReads).toHaveBeenCalledOnce()
|
||||
})
|
||||
|
||||
it('coalesces concurrent and repeated worktree polls after one reconciliation', async () => {
|
||||
resetWslRelayIdentityRpcCount()
|
||||
const relayInventory =
|
||||
deferred<{ id: string; cwd: string; title: string; worktreeId: string }[]>()
|
||||
const relayProcessReads = vi.fn(() => relayInventory.promise)
|
||||
const runtime = new OrcaRuntimeService(store)
|
||||
runtime.setPtyController({
|
||||
write: () => true,
|
||||
kill: () => true,
|
||||
getForegroundProcess: async () => null,
|
||||
listProcesses: relayProcessReads
|
||||
})
|
||||
syncSinglePty(runtime)
|
||||
|
||||
const concurrentPolls = [
|
||||
runtime.getWorktreePs(),
|
||||
runtime.getWorktreePs(),
|
||||
runtime.getWorktreePs()
|
||||
]
|
||||
await vi.waitFor(() => expect(relayProcessReads).toHaveBeenCalledOnce())
|
||||
relayInventory.resolve([
|
||||
{
|
||||
id: 'pty-1',
|
||||
cwd: '/tmp/worktree-a',
|
||||
title: 'shell',
|
||||
worktreeId: 'repo-1::/tmp/worktree-a'
|
||||
}
|
||||
])
|
||||
await Promise.all(concurrentPolls)
|
||||
|
||||
const identityRpcCountAfterFirstReconciliation = getWslRelayIdentityRpcCount()
|
||||
expect(identityRpcCountAfterFirstReconciliation).toBe(0)
|
||||
for (let poll = 0; poll < 5; poll += 1) {
|
||||
await runtime.getWorktreePs()
|
||||
}
|
||||
expect(relayProcessReads).toHaveBeenCalledOnce()
|
||||
expect(getWslRelayIdentityRpcCount()).toBe(identityRpcCountAfterFirstReconciliation)
|
||||
})
|
||||
|
||||
it('keeps WSL inventory reads stable across repeated worktree polls', async () => {
|
||||
setPlatform('win32')
|
||||
resetWslRelayIdentityRpcCount()
|
||||
const ptyId = 'wsl-pty-1'
|
||||
const anchor = {
|
||||
distro: 'Ubuntu',
|
||||
bootId: '11111111-1111-1111-1111-111111111111',
|
||||
shellPid: 1,
|
||||
shellStartTime: 1,
|
||||
tty: '/dev/pts/1'
|
||||
}
|
||||
const guestIdentityReads = vi.spyOn(wslRelayIdentityReader, 'readBatch').mockResolvedValue([
|
||||
{
|
||||
status: 'unverifiable' as const,
|
||||
reason: 'relay_unavailable',
|
||||
capturedAgeMs: 0
|
||||
}
|
||||
])
|
||||
const fakePty = { process: 'bash' }
|
||||
ptyProcesses.set(ptyId, fakePty as never)
|
||||
ptyInitialCwd.set(ptyId, '/tmp/worktree-a')
|
||||
ptyWorktreeId.set(ptyId, 'repo-1::/tmp/worktree-a')
|
||||
ptyWslDistroById.set(ptyId, 'Ubuntu')
|
||||
ptyWslShellAnchors.set(ptyId, anchor)
|
||||
const guestReads = vi.fn()
|
||||
const runtime = new OrcaRuntimeService(store)
|
||||
runtime.setPtyController({
|
||||
write: () => true,
|
||||
kill: () => true,
|
||||
getForegroundProcess: async () => null,
|
||||
listProcesses: async () => {
|
||||
guestReads()
|
||||
return listLocalPtyProcesses()
|
||||
}
|
||||
})
|
||||
syncSinglePty(runtime)
|
||||
|
||||
try {
|
||||
await runtime.getWorktreePs()
|
||||
const identityRpcCountAfterFirstReconciliation = getWslRelayIdentityRpcCount()
|
||||
const guestReadsAfterFirstReconciliation = guestReads.mock.calls.length
|
||||
expect(guestIdentityReads).toHaveBeenCalledOnce()
|
||||
expect(guestReadsAfterFirstReconciliation).toBe(1)
|
||||
for (let poll = 0; poll < 5; poll += 1) {
|
||||
await runtime.getWorktreePs()
|
||||
}
|
||||
expect(guestReads).toHaveBeenCalledTimes(guestReadsAfterFirstReconciliation)
|
||||
expect(guestIdentityReads).toHaveBeenCalledOnce()
|
||||
expect(getWslRelayIdentityRpcCount()).toBe(identityRpcCountAfterFirstReconciliation)
|
||||
} finally {
|
||||
guestIdentityReads.mockRestore()
|
||||
clearPtyState(ptyId)
|
||||
resetPlatform()
|
||||
}
|
||||
})
|
||||
|
||||
it('does not let targeted liveness satisfy a pending aggregate poll', async () => {
|
||||
const relayInventory =
|
||||
deferred<{ id: string; cwd: string; title: string; worktreeId: string }[]>()
|
||||
const relayProcessReads = vi.fn(() => relayInventory.promise)
|
||||
const runtime = new OrcaRuntimeService(store)
|
||||
runtime.setPtyController({
|
||||
write: () => true,
|
||||
kill: () => true,
|
||||
getForegroundProcess: async () => null,
|
||||
listProcesses: relayProcessReads
|
||||
})
|
||||
syncSinglePty(runtime)
|
||||
const refresh = (
|
||||
runtime as unknown as {
|
||||
refreshPtyWorktreeRecordsFromController: (
|
||||
worktrees: readonly unknown[],
|
||||
targetWorktreeId: string
|
||||
) => Promise<Set<string> | null>
|
||||
}
|
||||
).refreshPtyWorktreeRecordsFromController
|
||||
|
||||
const targeted = refresh.call(runtime, [], 'repo-1::/tmp/worktree-a')
|
||||
await vi.waitFor(() => expect(relayProcessReads).toHaveBeenCalledOnce())
|
||||
relayInventory.resolve([
|
||||
{
|
||||
id: 'pty-1',
|
||||
cwd: '/tmp/worktree-a',
|
||||
title: 'shell',
|
||||
worktreeId: 'repo-1::/tmp/worktree-a'
|
||||
}
|
||||
])
|
||||
await targeted
|
||||
|
||||
await runtime.getWorktreePs()
|
||||
expect(relayProcessReads).toHaveBeenCalledTimes(2)
|
||||
})
|
||||
it('keeps pinned and unread worktrees when active rows fill the mobile summary limit', async () => {
|
||||
setPlatform('win32')
|
||||
const remoteRepo = {
|
||||
|
||||
Reference in New Issue
Block a user