mirror of
https://github.com/stablyai/orca.git
synced 2026-09-30 00:03:15 +00:00
fix(runtime): stop a first status publication retiring in-flight worktree scans
A paired runtime host's first status publication counted as a connection change, advancing the connection generation. Worktree scans already in flight against that same connection were then discarded, so the sidebar showed a strict subset of the host's worktrees until an unrelated refresh. Two independent defects, both fixed: - `connectionChanged` conflated "no entry yet" with "recorded unreachable". Only the latter is a reconnect. The provider-session bump keeps the broader predicate, since a first publication is a real session start for integration-readiness caches. - A stale-generation result was thrown away with no retry, so even a genuine mid-flight reconnect silently dropped completed work. The scan is now re-read once against the new generation.
This commit is contained in:
@@ -0,0 +1,75 @@
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import { create } from 'zustand'
|
||||
import type { RuntimeStatus } from '../../../../shared/runtime-types'
|
||||
import {
|
||||
clearRuntimeEnvironmentConnectionGenerationsForTests,
|
||||
createRuntimeStatusSlice,
|
||||
getRuntimeEnvironmentConnectionGeneration,
|
||||
type RuntimeStatusSlice
|
||||
} from './runtime-status'
|
||||
|
||||
vi.mock('sonner', () => ({
|
||||
toast: { warning: vi.fn(), dismiss: vi.fn() }
|
||||
}))
|
||||
|
||||
function createSliceStore() {
|
||||
return create<RuntimeStatusSlice>()((...a) => ({
|
||||
...createRuntimeStatusSlice(...(a as unknown as Parameters<typeof createRuntimeStatusSlice>))
|
||||
}))
|
||||
}
|
||||
|
||||
function makeStatus(runtimeId: string): RuntimeStatus {
|
||||
return {
|
||||
runtimeId,
|
||||
rendererGraphEpoch: 0,
|
||||
graphStatus: 'ready',
|
||||
authoritativeWindowId: null,
|
||||
liveTabCount: 0,
|
||||
liveLeafCount: 0,
|
||||
runtimeProtocolVersion: 3,
|
||||
minCompatibleRuntimeClientVersion: 3
|
||||
} as RuntimeStatus
|
||||
}
|
||||
|
||||
beforeEach(() => {
|
||||
clearRuntimeEnvironmentConnectionGenerationsForTests()
|
||||
vi.stubGlobal('window', { api: {}, dispatchEvent: vi.fn() })
|
||||
})
|
||||
|
||||
afterEach(() => {
|
||||
vi.unstubAllGlobals()
|
||||
})
|
||||
|
||||
describe('first runtime status publication', () => {
|
||||
it('does not advance the connection generation when no entry existed', () => {
|
||||
// Regression (#19241): the first status for a paired environment counted as a
|
||||
// connection change, so the generation fence retired worktree scans already in
|
||||
// flight against that same connection. Those repos stayed absent from the sidebar
|
||||
// until an unrelated refresh happened to run.
|
||||
const store = createSliceStore()
|
||||
const before = getRuntimeEnvironmentConnectionGeneration('env-a')
|
||||
expect(store.getState().runtimeStatusByEnvironmentId.get('env-a')).toBeUndefined()
|
||||
|
||||
store
|
||||
.getState()
|
||||
.setRuntimeEnvironmentStatus('env-a', { status: makeStatus('runtime-a'), checkedAt: 1 })
|
||||
|
||||
expect(getRuntimeEnvironmentConnectionGeneration('env-a')).toBe(before)
|
||||
expect(store.getState().runtimeStatusByEnvironmentId.get('env-a')?.connectionGeneration).toBe(
|
||||
before
|
||||
)
|
||||
})
|
||||
|
||||
it('still advances when a recorded-unreachable host comes back', () => {
|
||||
// The other reading of the old `previous?.status == null`: this one IS a reconnect.
|
||||
const store = createSliceStore()
|
||||
store.getState().setRuntimeEnvironmentStatus('env-a', { status: null, checkedAt: 1 })
|
||||
const before = getRuntimeEnvironmentConnectionGeneration('env-a')
|
||||
|
||||
store
|
||||
.getState()
|
||||
.setRuntimeEnvironmentStatus('env-a', { status: makeStatus('runtime-a'), checkedAt: 2 })
|
||||
|
||||
expect(getRuntimeEnvironmentConnectionGeneration('env-a')).toBe(before + 1)
|
||||
})
|
||||
})
|
||||
@@ -181,7 +181,8 @@ describe('runtime-status slice', () => {
|
||||
|
||||
const map = store.getState().runtimeStatusByEnvironmentId
|
||||
expect(map.size).toBe(1)
|
||||
expect(map.get('env-a')).toEqual({ status: null, checkedAt: 5, connectionGeneration: 1 })
|
||||
// Generation 0: a first publication is not a reconnect, and going offline never bumps.
|
||||
expect(map.get('env-a')).toEqual({ status: null, checkedAt: 5, connectionGeneration: 0 })
|
||||
})
|
||||
|
||||
it('retains a learned paired device id after disconnecting a legacy environment', () => {
|
||||
@@ -302,7 +303,8 @@ describe('runtime-status slice', () => {
|
||||
})
|
||||
|
||||
expect(store.getState().runtimeStatusByEnvironmentId).not.toBe(before)
|
||||
expect(store.getState().runtimeStatusByEnvironmentId.get('env-a')?.connectionGeneration).toBe(2)
|
||||
// The first publication holds 0; only the runtime-id change advances it.
|
||||
expect(store.getState().runtimeStatusByEnvironmentId.get('env-a')?.connectionGeneration).toBe(1)
|
||||
})
|
||||
|
||||
it('does not toast when the first probe finds a saved server offline', () => {
|
||||
@@ -525,14 +527,16 @@ describe('runtime-status slice', () => {
|
||||
status: makeStatus({ runtimeId: 'runtime-a' }),
|
||||
checkedAt: 2
|
||||
})
|
||||
expect(store.getState().runtimeStatusByEnvironmentId.get('env-a')?.connectionGeneration).toBe(1)
|
||||
// Neither the first publication nor a stable re-poll is a connection change.
|
||||
expect(store.getState().runtimeStatusByEnvironmentId.get('env-a')?.connectionGeneration).toBe(0)
|
||||
|
||||
store.getState().setRuntimeEnvironmentStatus('env-a', { status: null, checkedAt: 3 })
|
||||
store.getState().setRuntimeEnvironmentStatus('env-a', {
|
||||
status: makeStatus({ runtimeId: 'runtime-a' }),
|
||||
checkedAt: 4
|
||||
})
|
||||
expect(store.getState().runtimeStatusByEnvironmentId.get('env-a')?.connectionGeneration).toBe(2)
|
||||
// Offline -> online is a real reconnect, so recovery still advances.
|
||||
expect(store.getState().runtimeStatusByEnvironmentId.get('env-a')?.connectionGeneration).toBe(1)
|
||||
})
|
||||
|
||||
it('keeps stored and canonical generations aligned after same-id re-pairing', () => {
|
||||
@@ -552,7 +556,9 @@ describe('runtime-status slice', () => {
|
||||
expect(store.getState().runtimeStatusByEnvironmentId.get('env-a')?.connectionGeneration).toBe(
|
||||
getRuntimeEnvironmentConnectionGeneration('env-a')
|
||||
)
|
||||
expect(getRuntimeEnvironmentConnectionGeneration('env-a')).toBe(3)
|
||||
// The re-pair itself advanced the generation and dropped the entry; the first
|
||||
// publication under the new pairing must not advance it a second time.
|
||||
expect(getRuntimeEnvironmentConnectionGeneration('env-a')).toBe(1)
|
||||
})
|
||||
|
||||
it('invalidates provider state only when the active runtime session changes', () => {
|
||||
|
||||
@@ -173,9 +173,17 @@ export const createRuntimeStatusSlice: StateCreator<AppState, [], [], RuntimeSta
|
||||
}
|
||||
set((s) => {
|
||||
const sessionEnded = status.status === null && previous?.status != null
|
||||
const connectionChanged =
|
||||
// A reachable answer where we held none (never asked, or recorded unreachable) or
|
||||
// where the runtime id moved starts a runtime session.
|
||||
const runtimeSessionStarted =
|
||||
status.status !== null &&
|
||||
(previous?.status == null || previous.status.runtimeId !== status.status.runtimeId)
|
||||
// Why narrower than the session start: a first publication has no prior connection to
|
||||
// differ from, so it is not a reconnect. Advancing the generation there retires reads
|
||||
// already issued against this very connection — a startup worktree scan that had
|
||||
// already answered was discarded, leaving those repos absent until an unrelated
|
||||
// refresh (#19241).
|
||||
const connectionChanged = runtimeSessionStarted && previous !== undefined
|
||||
const activeEnvironmentId = s.settings?.activeRuntimeEnvironmentId?.trim()
|
||||
const connectionGeneration = connectionChanged
|
||||
? runtimeStatusConnectionGeneration.advanceRuntimeEnvironmentConnectionGeneration(
|
||||
@@ -186,7 +194,9 @@ export const createRuntimeStatusSlice: StateCreator<AppState, [], [], RuntimeSta
|
||||
runtimeStatusConnectionGeneration.getRuntimeEnvironmentConnectionGeneration(
|
||||
environmentId
|
||||
))
|
||||
if (activeEnvironmentId === environmentId && (sessionEnded || connectionChanged)) {
|
||||
// Why the session flag and not `connectionChanged`: integration-readiness caches key
|
||||
// off the runtime session, for which a first publication is a real transition.
|
||||
if (activeEnvironmentId === environmentId && (sessionEnded || runtimeSessionStarted)) {
|
||||
bumpProviderRuntimeSessionGeneration()
|
||||
}
|
||||
const nextEntry = { ...status, connectionGeneration }
|
||||
|
||||
@@ -0,0 +1,133 @@
|
||||
import { beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import type { AppState } from '../types'
|
||||
import type { RuntimeEnvironmentCallRequest } from '../../runtime/runtime-compatibility-test-fixture'
|
||||
import { advanceRuntimeEnvironmentConnectionGeneration } from './runtime-status-connection-generation'
|
||||
import { makeDetectedResult } from './worktrees-detected-listing-fixtures'
|
||||
import { makeWorktree } from './worktrees-slice-test-fixtures'
|
||||
import {
|
||||
createTestStore,
|
||||
resetRemoteRuntimeMocks,
|
||||
resetWorktreeSliceModuleMemory,
|
||||
runtimeEnvironmentCall
|
||||
} from './worktrees-slice-test-harness'
|
||||
|
||||
vi.mock('sonner', () => ({
|
||||
toast: {
|
||||
warning: vi.fn(),
|
||||
info: vi.fn(),
|
||||
success: vi.fn(),
|
||||
error: vi.fn(),
|
||||
dismiss: vi.fn()
|
||||
}
|
||||
}))
|
||||
|
||||
const ENV = 'env-remote'
|
||||
const REPO_IDS = ['repo1', 'repo2', 'repo3', 'repo4', 'repo5'] as const
|
||||
|
||||
function seedRuntimeRepos(store: ReturnType<typeof createTestStore>) {
|
||||
store.setState({
|
||||
repos: REPO_IDS.map((id) => ({
|
||||
id,
|
||||
path: `C:/remote/${id}`,
|
||||
executionHostId: `runtime:${ENV}`
|
||||
})),
|
||||
settings: { activeRuntimeEnvironmentId: ENV },
|
||||
hasHydratedWorktreePurge: true
|
||||
} as unknown as Partial<AppState>)
|
||||
}
|
||||
|
||||
function detectedFor(repoId: string) {
|
||||
return makeDetectedResult(repoId, [
|
||||
makeWorktree({
|
||||
id: `${repoId}::/remote/${repoId}/wt`,
|
||||
repoId,
|
||||
path: `/remote/${repoId}/wt`,
|
||||
branch: 'refs/heads/main'
|
||||
})
|
||||
])
|
||||
}
|
||||
|
||||
function repoOf(args: RuntimeEnvironmentCallRequest): string {
|
||||
return (args.params as { repo: string }).repo
|
||||
}
|
||||
|
||||
function detectedListReply(repoId: string) {
|
||||
return {
|
||||
id: 'rpc',
|
||||
ok: true,
|
||||
result: detectedFor(repoId),
|
||||
_meta: { runtimeId: 'runtime-remote' }
|
||||
}
|
||||
}
|
||||
|
||||
beforeEach(() => {
|
||||
resetWorktreeSliceModuleMemory()
|
||||
vi.clearAllMocks()
|
||||
resetRemoteRuntimeMocks()
|
||||
})
|
||||
|
||||
describe('fetchAllWorktrees across a runtime connection-generation change', () => {
|
||||
it('publishes every repo when nothing perturbs the connection', async () => {
|
||||
const store = createTestStore()
|
||||
seedRuntimeRepos(store)
|
||||
runtimeEnvironmentCall.mockImplementation((args: RuntimeEnvironmentCallRequest) =>
|
||||
detectedListReply(repoOf(args))
|
||||
)
|
||||
|
||||
await store.getState().fetchAllWorktrees()
|
||||
|
||||
expect(Object.keys(store.getState().worktreesByRepo).sort()).toEqual([...REPO_IDS])
|
||||
})
|
||||
|
||||
it('re-reads instead of dropping the repos still in flight when the generation moves', async () => {
|
||||
// Regression (#19241): scans run concurrently across a host's repos, so one generation
|
||||
// bump discarded every repo still outstanding. Their rows then never reached the
|
||||
// sidebar — a strict subset of the host's worktrees, stable until an unrelated refresh.
|
||||
const store = createTestStore()
|
||||
seedRuntimeRepos(store)
|
||||
const parked = new Map<string, () => void>()
|
||||
let bumped = false
|
||||
runtimeEnvironmentCall.mockImplementation(async (args: RuntimeEnvironmentCallRequest) => {
|
||||
const repo = repoOf(args)
|
||||
// Park every repo but the first, so the bump lands with four scans outstanding.
|
||||
if (repo !== 'repo1' && !bumped) {
|
||||
await new Promise<void>((resolve) => parked.set(repo, resolve))
|
||||
}
|
||||
return detectedListReply(repo)
|
||||
})
|
||||
|
||||
const fetching = store.getState().fetchAllWorktrees()
|
||||
await vi.waitFor(() => expect(parked.size).toBe(REPO_IDS.length - 1))
|
||||
bumped = true
|
||||
advanceRuntimeEnvironmentConnectionGeneration(ENV)
|
||||
for (const release of parked.values()) {
|
||||
release()
|
||||
}
|
||||
await fetching
|
||||
|
||||
expect(Object.keys(store.getState().worktreesByRepo).sort()).toEqual([...REPO_IDS])
|
||||
for (const repoId of REPO_IDS) {
|
||||
expect(store.getState().worktreesByRepo[repoId]).toHaveLength(1)
|
||||
}
|
||||
})
|
||||
|
||||
it('gives up after one re-read so a churning connection cannot stall the caller', async () => {
|
||||
const store = createTestStore()
|
||||
seedRuntimeRepos(store)
|
||||
const attemptsByRepo = new Map<string, number>()
|
||||
runtimeEnvironmentCall.mockImplementation((args: RuntimeEnvironmentCallRequest) => {
|
||||
const repo = repoOf(args)
|
||||
attemptsByRepo.set(repo, (attemptsByRepo.get(repo) ?? 0) + 1)
|
||||
// Bump on every answer: the retried read is stale again the moment it lands.
|
||||
advanceRuntimeEnvironmentConnectionGeneration(ENV)
|
||||
return detectedListReply(repo)
|
||||
})
|
||||
|
||||
await store.getState().fetchAllWorktrees()
|
||||
|
||||
expect(Object.keys(store.getState().worktreesByRepo)).toEqual([])
|
||||
for (const repoId of REPO_IDS) {
|
||||
expect(attemptsByRepo.get(repoId)).toBe(2)
|
||||
}
|
||||
})
|
||||
})
|
||||
@@ -35,6 +35,15 @@ const runtimeDetectedWorktreeRefreshesInFlight = new Map<
|
||||
Promise<DetectedWorktreeListResult>
|
||||
>()
|
||||
|
||||
const STALE_RUNTIME_GENERATION_ERROR = 'runtime_environment_generation_changed'
|
||||
// Why exactly one: a second stale answer means the connection is still churning, and
|
||||
// retrying into that would stall the caller instead of letting it fail visibly.
|
||||
const STALE_RUNTIME_GENERATION_RETRIES = 1
|
||||
|
||||
export function isStaleRuntimeGenerationError(error: unknown): boolean {
|
||||
return error instanceof Error && error.message === STALE_RUNTIME_GENERATION_ERROR
|
||||
}
|
||||
|
||||
export const detectedWorktreeRefreshLeaseRegistry = createDetectedWorktreeRefreshLeaseRegistry({
|
||||
startProviderRequest: startDetectedWorktreeProviderRequest,
|
||||
cancelProviderRequest: async (request) => {
|
||||
@@ -133,58 +142,83 @@ export function normalizeNotAdmittedProviderResult(
|
||||
}
|
||||
}
|
||||
|
||||
async function listDetectedWorktreesForRuntimeRepoOnce(
|
||||
settings: AppState['settings'],
|
||||
repoId: string,
|
||||
options: DetectedWorktreeRefreshOptions,
|
||||
environmentId: string
|
||||
): Promise<DetectedWorktreeRefreshOutcome> {
|
||||
// Why recomputed per attempt: the key embeds both generations, so a retry after a
|
||||
// reconnect must not join the superseded connection's in-flight scan.
|
||||
const key = detectedWorktreeRefreshKey(settings, repoId, options)
|
||||
const connectionGeneration = getEnvironmentSshStateGeneration(environmentId)
|
||||
const runtimeConnectionGeneration = getRuntimeEnvironmentConnectionGeneration(environmentId)
|
||||
let refresh = runtimeDetectedWorktreeRefreshesInFlight.get(key)
|
||||
if (!refresh) {
|
||||
refresh = listDetectedWorktreesForRepo(settings, repoId, {
|
||||
reuseRecentCompatibilityFailure: options.reuseRecentCompatibilityFailure
|
||||
})
|
||||
runtimeDetectedWorktreeRefreshesInFlight.set(key, refresh)
|
||||
}
|
||||
try {
|
||||
const result = await refresh
|
||||
if (
|
||||
getEnvironmentSshStateGeneration(environmentId) !== connectionGeneration ||
|
||||
getRuntimeEnvironmentConnectionGeneration(environmentId) !== runtimeConnectionGeneration
|
||||
) {
|
||||
throw new Error(STALE_RUNTIME_GENERATION_ERROR)
|
||||
}
|
||||
// Why (#10562): the scan coalesces, but teardown must not — each caller carries
|
||||
// its own known-id snapshot and purges its own state, so a caller that joined
|
||||
// an in-flight scan would otherwise purge without ever stopping those terminals.
|
||||
await teardownMissingWorktreeTerminalsBestEffort(
|
||||
settings,
|
||||
repoId,
|
||||
options.connectionId,
|
||||
options.knownWorktreeIds,
|
||||
result
|
||||
)
|
||||
return {
|
||||
status: 'admitted',
|
||||
result,
|
||||
executionHostId: options.executionHostId,
|
||||
runtimeAuthority: {
|
||||
environmentId,
|
||||
connectionGeneration,
|
||||
runtimeConnectionGeneration
|
||||
}
|
||||
}
|
||||
} finally {
|
||||
if (runtimeDetectedWorktreeRefreshesInFlight.get(key) === refresh) {
|
||||
runtimeDetectedWorktreeRefreshesInFlight.delete(key)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
export async function listDetectedWorktreesForRepoCoalesced(
|
||||
settings: AppState['settings'],
|
||||
repoId: string,
|
||||
options: DetectedWorktreeRefreshOptions
|
||||
): Promise<DetectedWorktreeRefreshOutcome> {
|
||||
const key = detectedWorktreeRefreshKey(settings, repoId, options)
|
||||
const target = getActiveRuntimeTarget(settings)
|
||||
if (target.kind === 'environment') {
|
||||
const connectionGeneration = getEnvironmentSshStateGeneration(target.environmentId)
|
||||
const runtimeConnectionGeneration = getRuntimeEnvironmentConnectionGeneration(
|
||||
target.environmentId
|
||||
)
|
||||
let refresh = runtimeDetectedWorktreeRefreshesInFlight.get(key)
|
||||
if (!refresh) {
|
||||
refresh = listDetectedWorktreesForRepo(settings, repoId, {
|
||||
reuseRecentCompatibilityFailure: options.reuseRecentCompatibilityFailure
|
||||
})
|
||||
runtimeDetectedWorktreeRefreshesInFlight.set(key, refresh)
|
||||
}
|
||||
try {
|
||||
const result = await refresh
|
||||
if (
|
||||
getEnvironmentSshStateGeneration(target.environmentId) !== connectionGeneration ||
|
||||
getRuntimeEnvironmentConnectionGeneration(target.environmentId) !==
|
||||
runtimeConnectionGeneration
|
||||
) {
|
||||
throw new Error('runtime_environment_generation_changed')
|
||||
}
|
||||
// Why (#10562): the scan coalesces, but teardown must not — each caller carries
|
||||
// its own known-id snapshot and purges its own state, so a caller that joined
|
||||
// an in-flight scan would otherwise purge without ever stopping those terminals.
|
||||
await teardownMissingWorktreeTerminalsBestEffort(
|
||||
settings,
|
||||
repoId,
|
||||
options.connectionId,
|
||||
options.knownWorktreeIds,
|
||||
result
|
||||
)
|
||||
return {
|
||||
status: 'admitted',
|
||||
result,
|
||||
executionHostId: options.executionHostId,
|
||||
runtimeAuthority: {
|
||||
environmentId: target.environmentId,
|
||||
connectionGeneration,
|
||||
runtimeConnectionGeneration
|
||||
for (let attempt = 0; ; attempt += 1) {
|
||||
try {
|
||||
return await listDetectedWorktreesForRuntimeRepoOnce(
|
||||
settings,
|
||||
repoId,
|
||||
options,
|
||||
target.environmentId
|
||||
)
|
||||
} catch (err) {
|
||||
// Why re-read instead of surfacing: the fence proves this answer predates the
|
||||
// current connection, not that the repo has no worktrees. Callers drop the repo
|
||||
// on any throw, so a discarded scan left those rows absent until an unrelated
|
||||
// refresh happened to run (#19241).
|
||||
if (attempt >= STALE_RUNTIME_GENERATION_RETRIES || !isStaleRuntimeGenerationError(err)) {
|
||||
throw err
|
||||
}
|
||||
}
|
||||
} finally {
|
||||
if (runtimeDetectedWorktreeRefreshesInFlight.get(key) === refresh) {
|
||||
runtimeDetectedWorktreeRefreshesInFlight.delete(key)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user