fix: hydrate worker records and resume fences atomically

This commit is contained in:
Jinwoo-H
2026-09-07 19:51:36 -04:00
parent 67209085ad
commit 7f36880759
10 changed files with 198 additions and 52 deletions
+5 -3
View File
@@ -25,12 +25,13 @@
"coverageNotes": "Deterministic local/SSH partition and mixed-source tests; hidden desktop restart/reload and retirement E2Es. No live independently versioned paired-server or physical Windows/Linux claim.",
"motivatingLinks": ["https://github.com/stablyai/orca/pull/19344"],
"invariant": "A runtime-authored settled-worker fence survives old writers, absent old-source authority, overlapping hydration, and sibling host write failure; authoritative scoped retirement preserves sibling fences.",
"oracle": "Nine round-2 reviewer probes plus publication/read-boundary tests; each of the six round-3 fix mutations must fail. Restart/reload does not resume a fenced worker, while explicit retirement permits resume.",
"oracle": "Nine round-2 reviewer probes plus atomic pending-hydration and publication/read-boundary tests; queue-only-fence and stale-read mutations must fail alongside the six round-3 fix mutations. Restart/reload does not resume a fenced worker, while explicit retirement permits resume.",
"commands": [
"ORCA_BACKGROUND_LAUNCH=1 npx vitest run --config config/vitest.config.ts src/renderer/src/lib/runtime-session-application.test.ts src/main/persistence/loading-store/store-runtime-authored-session-fences.test.ts src/main/runtime/runtime-legacy-worker-session-publication.test.ts src/main/runtime/runtime-legacy-worker-host-write-publication.test.ts src/main/runtime/runtime-legacy-worker-recovery-publication-order.test.ts src/renderer/src/lib/legacy-worker-old-source-resume.test.ts src/renderer/src/lib/runtime-session-refresh-hydration-order.test.ts src/renderer/src/store/slices/agent-sleep-resume-fence-races.test.ts src/renderer/src/store/slices/workspace-scoped-resume-fence-hydration.test.ts src/renderer/src/components/terminal-pane/pty-connection/sleeping-record-access-recordless-fence.test.ts",
"ORCA_BACKGROUND_LAUNCH=1 npx vitest run --config config/vitest.config.ts src/renderer/src/lib/runtime-session-pending-hydration.test.ts src/renderer/src/lib/runtime-session-application.test.ts src/main/persistence/loading-store/store-runtime-authored-session-fences.test.ts src/main/runtime/runtime-legacy-worker-session-publication.test.ts src/main/runtime/runtime-legacy-worker-host-write-publication.test.ts src/main/runtime/runtime-legacy-worker-recovery-publication-order.test.ts src/renderer/src/lib/legacy-worker-old-source-resume.test.ts src/renderer/src/lib/runtime-session-refresh-hydration-order.test.ts src/renderer/src/store/slices/agent-sleep-resume-fence-races.test.ts src/renderer/src/store/slices/workspace-scoped-resume-fence-hydration.test.ts src/renderer/src/components/terminal-pane/pty-connection/sleeping-record-access-recordless-fence.test.ts",
"ORCA_BACKGROUND_LAUNCH=1 SKIP_BUILD=1 npx playwright test tests/e2e/settled-worker-resume-fence-restart.spec.ts tests/e2e/completed-worker-retirement-resume.spec.ts --config tests/playwright.config.ts --project=electron-headless --workers=1"
],
"testFiles": [
"src/renderer/src/lib/runtime-session-pending-hydration.test.ts",
"src/renderer/src/lib/runtime-session-application.test.ts",
"src/main/persistence/loading-store/store-runtime-authored-session-fences.test.ts",
"src/main/runtime/runtime-legacy-worker-session-publication.test.ts",
@@ -49,7 +50,8 @@
"file": "src/renderer/src/lib/runtime-session-application.test.ts",
"assertions": [
"applies startup runtime fields at read time, before delayed catalog hydration",
"orders a scoped retirement after an in-flight invalidation and retains its sibling"
"orders a scoped retirement after an in-flight invalidation and retains its sibling",
"publishes sleeping records and their authority in a single subscriber notification"
]
}
],
@@ -136,3 +136,66 @@ it('an authoritative empty set retires a legacy flag while an absent set preserv
store.getState().hydrateWorkspaceSession({ ...oldSession, legacyWorkerResumeFencesByPaneKey: {} })
expect(store.getState().legacyWorkerResumeFencesByPaneKey).toEqual({})
})
it('publishes sleeping records and their authority in a single subscriber notification', () => {
const session = {
...getDefaultWorkspaceSession(),
legacyWorkerResumeFencesByPaneKey: { [pane]: true as const },
sleepingAgentSessionsByPaneKey: {
[pane]: {
paneKey: pane,
tabId: 'target',
worktreeId: 'folder:legacy',
agent: 'codex' as const,
providerSession: { key: 'session_id' as const, id: 'provider' },
prompt: '',
state: 'working' as const,
capturedAt: 1,
updatedAt: 1,
origin: 'live' as const
}
}
}
const observed: boolean[][] = []
const unsubscribe = useAppStore.subscribe((state) => {
observed.push([
!!state.sleepingAgentSessionsByPaneKey[pane],
!!state.legacyWorkerResumeFencesByPaneKey[pane]
])
})
try {
useAppStore.getState().hydrateWorkspaceSession(session, {
additionalValidWorkspaceKeys: ['folder:legacy']
})
expect(observed).toEqual([[true, true]])
} finally {
unsubscribe()
}
})
it.each(['resolve', 'reject'] as const)(
'a later refresh can retire hydration after the older read %ss',
async (outcome) => {
let settle!: () => void
const read = vi
.fn()
.mockImplementationOnce(
() =>
new Promise<Record<string, true>>((resolve, reject) => {
settle = () => (outcome === 'resolve' ? resolve({}) : reject(new Error('offline')))
})
)
.mockResolvedValue({})
vi.stubGlobal('window', { api: { app: { getLegacyWorkerResumeFences: read } } })
const pending = refreshLegacyWorkerResumeFences()
useAppStore.getState().hydrateWorkspaceSession({
...getDefaultWorkspaceSession(),
legacyWorkerResumeFencesByPaneKey: { [pane]: true }
})
settle()
await pending
expect(useAppStore.getState().legacyWorkerResumeFencesByPaneKey).toEqual({ [pane]: true })
await refreshLegacyWorkerResumeFences()
expect(useAppStore.getState().legacyWorkerResumeFencesByPaneKey).toEqual({})
}
)
@@ -3,18 +3,20 @@ import type { WorkspaceSessionState } from '../../../shared/workspace-session-st
import { parsePaneKey, parseLegacyNumericPaneKey } from '../../../shared/stable-pane-id'
import type { TerminalStoreGet, TerminalStoreSet } from '../store/terminals/terminal-state'
const pending: (() => void)[] = []
const pendingReads: (() => void)[] = []
let hydratedDuringRead = new Map<TerminalStoreGet, ReadonlySet<string> | null>()
let reading = false
const appliedSessions = new WeakMap<TerminalStoreGet, WeakSet<WorkspaceSessionState>>()
function drainApplications(): void {
reading = false
while (!reading && pending.length > 0) {
pending.shift()!()
hydratedDuringRead.clear()
while (!reading && pendingReads.length > 0) {
pendingReads.shift()!()
}
}
// Reads and synchronous hydration share this lane, including the read before startup's catalog wait.
// Serialize runtime reads; synchronous hydration supersedes their overlapping snapshot scopes.
export function readAndApplyRuntimeSession<T>(
read: () => Promise<T>,
apply: (value: T) => void
@@ -22,6 +24,7 @@ export function readAndApplyRuntimeSession<T>(
return new Promise((resolve, reject) => {
const start = (): void => {
reading = true
hydratedDuringRead = new Map()
void (async () => read())().then(
(value) => {
try {
@@ -40,7 +43,7 @@ export function readAndApplyRuntimeSession<T>(
)
}
if (reading) {
pending.push(start)
pendingReads.push(start)
} else {
start()
}
@@ -53,6 +56,35 @@ export function applyReadRuntimeSession(
get: TerminalStoreGet,
targetTabIds?: ReadonlySet<string>
): void {
const fields = runtimeSessionFields(session, get, targetTabIds)
const superseded = hydratedDuringRead.get(get)
if (superseded === null) {
return
}
if (superseded) {
const current = get().legacyWorkerResumeFencesByPaneKey
fields.legacyWorkerResumeFencesByPaneKey = {
...Object.fromEntries(
Object.entries(fields.legacyWorkerResumeFencesByPaneKey).filter(
([key]) => !paneInScope(key, superseded)
)
),
...Object.fromEntries(Object.entries(current).filter(([key]) => paneInScope(key, superseded)))
}
}
set(fields)
}
function paneInScope(key: string, tabIds: ReadonlySet<string>): boolean {
const tabId = parsePaneKey(key)?.tabId ?? parseLegacyNumericPaneKey(key)?.tabId
return tabId !== undefined && tabIds.has(tabId)
}
function runtimeSessionFields(
session: WorkspaceSessionState,
get: TerminalStoreGet,
targetTabIds?: ReadonlySet<string>
): { legacyWorkerResumeFencesByPaneKey: Record<string, true> } {
let applied = appliedSessions.get(get)
if (!applied) {
applied = new WeakSet()
@@ -66,33 +98,35 @@ export function applyReadRuntimeSession(
...current,
...readWorkspaceSessionResumeFences(session)
}
const inScope = (key: string): boolean => {
const tabId = parsePaneKey(key)?.tabId ?? parseLegacyNumericPaneKey(key)?.tabId
return tabId !== undefined && targetTabIds!.has(tabId)
}
set({
return {
legacyWorkerResumeFencesByPaneKey: targetTabIds
? {
...Object.fromEntries(Object.entries(current).filter(([key]) => !inScope(key))),
...Object.fromEntries(Object.entries(fences).filter(([key]) => inScope(key)))
...Object.fromEntries(
Object.entries(current).filter(([key]) => !paneInScope(key, targetTabIds))
),
...Object.fromEntries(
Object.entries(fences).filter(([key]) => paneInScope(key, targetTabIds))
)
}
: fences
})
}
}
export function hydrateRuntimeSessionFields(
session: WorkspaceSessionState,
set: TerminalStoreSet,
get: TerminalStoreGet,
targetTabIds?: ReadonlySet<string>
): void {
): { legacyWorkerResumeFencesByPaneKey: Record<string, true> } {
if (appliedSessions.get(get)?.has(session)) {
return
return { legacyWorkerResumeFencesByPaneKey: get().legacyWorkerResumeFencesByPaneKey }
}
const apply = (): void => applyReadRuntimeSession(session, set, get, targetTabIds)
if (reading) {
pending.push(apply)
} else {
apply()
const previous = hydratedDuringRead.get(get)
// Track snapshot scopes, never policy values, until the older read settles.
hydratedDuringRead.set(
get,
!targetTabIds || previous === null ? null : new Set([...(previous ?? []), ...targetTabIds])
)
}
return runtimeSessionFields(session, get, targetTabIds)
}
@@ -0,0 +1,66 @@
import { afterEach, expect, it, vi } from 'vitest'
import { useAppStore } from '@/store'
import { getDefaultWorkspaceSession } from '../../../shared/constants'
import { folderWorkspaceKey } from '../../../shared/workspace-scope'
import { refreshLegacyWorkerResumeFences } from './legacy-worker-resume-fence-refresh'
import { resumeSleepingAgentSessionsForWorktree } from './resume-sleeping-agent-session'
const initial = useAppStore.getState()
afterEach(() => {
useAppStore.setState(initial, true)
vi.unstubAllGlobals()
})
it.each(['canonical', 'legacy'] as const)(
'cannot auto-resume %s protected records while hydration waits behind a read',
async (kind) => {
let resolve!: (value: Record<string, true>) => void
vi.stubGlobal('window', {
api: {
app: {
getLegacyWorkerResumeFences: () =>
new Promise<Record<string, true>>((done) => {
resolve = done
})
}
}
})
const pending = refreshLegacyWorkerResumeFences()
const paneKey = 'tab-legacy:11111111-2222-4333-8444-555555555555'
const record = {
paneKey,
tabId: 'tab-legacy',
worktreeId: 'folder:legacy',
agent: 'claude' as const,
providerSession: { key: 'session_id' as const, id: 'session-legacy' },
prompt: 'continue legacy work',
state: 'working' as const,
capturedAt: 1,
updatedAt: 1,
origin: 'live' as const,
...(kind === 'legacy'
? { automaticResumeBlockedBy: 'legacy-orchestration-worker' as const }
: {})
}
useAppStore.getState().hydrateWorkspaceSession(
{
...getDefaultWorkspaceSession(),
sleepingAgentSessionsByPaneKey: { [paneKey]: record },
...(kind === 'canonical'
? { legacyWorkerResumeFencesByPaneKey: { [paneKey]: true as const } }
: {})
},
{ additionalValidWorkspaceKeys: [folderWorkspaceKey('legacy')] }
)
const count = resumeSleepingAgentSessionsForWorktree('folder:legacy')
const observed = {
count,
tabs: useAppStore.getState().tabsByWorktree['folder:legacy']?.map((tab) => tab.id),
claims: useAppStore.getState().automaticAgentResumeClaimsByTabId,
recordStillPresent: !!useAppStore.getState().sleepingAgentSessionsByPaneKey[paneKey]
}
resolve({})
await pending
console.log(kind, observed)
expect(count).toBe(0)
expect(useAppStore.getState().sleepingAgentSessionsByPaneKey[paneKey]).toBe(record)
}
)
@@ -61,7 +61,7 @@ describe('resume fences during asynchronous sleep writes', () => {
{ providerSession: { key: 'session_id', id: 'session-1' } }
)
mockApi.pty.kill.mockImplementationOnce(async () => {
store.getState().setLegacyWorkerResumeFences({ [pane]: true })
store.setState({ legacyWorkerResumeFencesByPaneKey: { [pane]: true } })
expect(store.getState().legacyWorkerResumeFencesByPaneKey[pane]).toBe(true)
})
await store.getState().shutdownWorktreeTerminals(wt, { keepIdentifiers: true })
@@ -98,7 +98,7 @@ describe('resume fences during asynchronous sleep writes', () => {
{ providerSession: { key: 'session_id', id: 'session-1' } }
)
mockApi.pty.kill.mockImplementationOnce(async () => {
store.getState().setLegacyWorkerResumeFences({ [pane]: true })
store.setState({ legacyWorkerResumeFencesByPaneKey: { [pane]: true } })
expect(store.getState().legacyWorkerResumeFencesByPaneKey[pane]).toBe(true)
throw new Error('kill_failed')
})
@@ -58,7 +58,7 @@ describe('a resume fence for a pane with no sleeping record', () => {
it('is retired only by the runtime replacing the set', () => {
const store = fencedStore()
store.getState().setLegacyWorkerResumeFences({})
store.setState({ legacyWorkerResumeFencesByPaneKey: {} })
expect(store.getState().legacyWorkerResumeFencesByPaneKey).toEqual({})
})
@@ -102,7 +102,7 @@ describe('a resume fence for a pane with no sleeping record', () => {
it('strips a stale projection once the runtime retires the fence', () => {
const store = fencedStore()
store.getState().captureAllSleepingAgentSessions('quit')
store.getState().setLegacyWorkerResumeFences({})
store.setState({ legacyWorkerResumeFencesByPaneKey: {} })
const projected = buildSleepingAgentSessionData(store.getState())
@@ -1,5 +1,4 @@
import type { SleepingAgentSessionRecord } from '../../../../shared/agent-session-resume'
import { sameFenceSet } from './legacy-worker-resume-fences'
import type { AgentStatusSlice } from './agent-status-slice-contract'
import type { AgentStatusRuntime } from './agent-status-runtime'
import { collectSleepingAgentSessionRecordsForWorktree } from './agent-status-recovery-collection'
@@ -23,7 +22,6 @@ export function createAgentStatusRecoveryActions(
| 'captureAllSleepingAgentSessions'
| 'clearSleepingAgentSession'
| 'clearSleepingAgentSessionsByPaneKey'
| 'setLegacyWorkerResumeFences'
| 'clearSleepingAgentSessionsByWorktree'
| 'pruneSleepingAgentSessions'
> {
@@ -110,14 +108,6 @@ export function createAgentStatusRecoveryActions(
clearSleepingAgentSession: (paneKey) => clearSleepingAgentSessionsByPaneKey([paneKey]),
clearSleepingAgentSessionsByPaneKey,
setLegacyWorkerResumeFences: (fences) => {
set((s) =>
sameFenceSet(s.legacyWorkerResumeFencesByPaneKey, fences)
? s
: { legacyWorkerResumeFencesByPaneKey: fences }
)
},
clearSleepingAgentSessionsByWorktree: (worktreeId) => {
set((s) => {
let changed = false
@@ -162,8 +162,6 @@ export type AgentStatusSlice = {
captureAllSleepingAgentSessions: (mode: AllAgentSessionCaptureMode) => void
clearSleepingAgentSession: (paneKey: string) => void
clearSleepingAgentSessionsByPaneKey: (paneKeys: readonly string[]) => void
/** Install main's fenced-pane set wholesale. The only writer, and level-triggered. */
setLegacyWorkerResumeFences: (fences: Record<string, true>) => void
clearSleepingAgentSessionsByWorktree: (worktreeId: string) => void
pruneSleepingAgentSessions: (validWorktreeIds: Set<string>) => void
@@ -1,9 +0,0 @@
/** Main rewrites the whole set on every recovery pass, so compare by content: an identity check
* would republish an unchanged set and wake every session subscriber. */
export function sameFenceSet(current: Record<string, true>, next: Record<string, true>): boolean {
const nextKeys = Object.keys(next)
return (
Object.keys(current).length === nextKeys.length &&
nextKeys.every((paneKey) => current[paneKey] === true)
)
}
@@ -40,7 +40,6 @@ export function createWorkspaceTerminalHydrationActions(
])
)
: undefined
hydrateRuntimeSessionFields(session, set, get, targetTabIds)
const ownershipTransferTabIds = options?.replaceWorkspaceKeys
? new Set(
options.replaceWorkspaceKeys.flatMap((workspaceKey) =>
@@ -225,9 +224,12 @@ export function createWorkspaceTerminalHydrationActions(
validTabIds
})
}
return options?.replaceWorkspaceKeys
? targetScopedWorkspaceHydrationPatch(s, hydrated, session, options)
: hydrated
return {
...(options?.replaceWorkspaceKeys
? targetScopedWorkspaceHydrationPatch(s, hydrated, session, options)
: hydrated),
...hydrateRuntimeSessionFields(session, get, targetTabIds)
}
})
for (const [tabId, transfers] of ownershipTransfersByTabId) {
transferNormalizedTerminalLayoutPtyOwnership(get(), tabId, transfers)