diff --git a/cmuxTests/WorkspaceForkConversationContextMenuTests.swift b/cmuxTests/WorkspaceForkConversationContextMenuTests.swift index 02361b037561..71b4a6c1f54c 100644 --- a/cmuxTests/WorkspaceForkConversationContextMenuTests.swift +++ b/cmuxTests/WorkspaceForkConversationContextMenuTests.swift @@ -1726,7 +1726,9 @@ struct WorkspaceForkConversationContextMenuTests { let loaderCallCount = OSAllocatedUnfairLock(initialState: 0) let loaderStartedCount = OSAllocatedUnfairLock(initialState: 0) let firstLoaderRelease = DispatchSemaphore(value: 0) - let secondRefreshStarted = AsyncStream.makeStream() + let firstLoaderStarted = ForkProbeTestSignal() + defer { firstLoaderRelease.signal() } + let secondRefreshStarted = ForkProbeTestSignal() let sharedIndex = SharedLiveAgentIndex( indexLoader: { let call = loaderCallCount.withLock { count in @@ -1737,6 +1739,7 @@ struct WorkspaceForkConversationContextMenuTests { count = max(count, call) } if call == 1 { + firstLoaderStarted.signal() firstLoaderRelease.wait() } let index = RestorableAgentSessionIndex.load( @@ -1779,19 +1782,17 @@ struct WorkspaceForkConversationContextMenuTests { panelId: firstPanelId ) } - for _ in 0..<1000 where loaderStartedCount.withLock({ $0 }) < 1 { - await Task.yield() - } + try #require(await firstLoaderStarted.wait(), "First index loader must start before queuing another refresh.") #expect(loaderStartedCount.withLock { $0 } >= 1) let secondRefresh = Task { - secondRefreshStarted.continuation.yield(()) + secondRefreshStarted.signal() await sharedIndex.refreshForkAvailabilityNow( workspaceId: secondWorkspaceId, panelId: secondPanelId ) } - _ = await secondRefreshStarted.stream.first(where: { _ in true }) + try #require(await secondRefreshStarted.wait(), "Expected refresh/probe must start before proceeding.") #expect(loaderStartedCount.withLock { $0 } == 1, "Concurrent refreshes share the active index loader.") secondRefresh.cancel() @@ -1931,7 +1932,7 @@ struct WorkspaceForkConversationContextMenuTests { ) } - @Test + @Test(.timeLimit(.minutes(1))) func sharedForkProbeActiveWaitPreservesLaterQueuedValidationBatches() async throws { let fm = FileManager.default let root = fm.temporaryDirectory @@ -1964,9 +1965,16 @@ struct WorkspaceForkConversationContextMenuTests { executablePath: executable.path ) let loaderCallCount = OSAllocatedUnfairLock(initialState: 0) - let blockedLoaderCount = OSAllocatedUnfairLock(initialState: 0) let releaseBlockedLoaders = DispatchSemaphore(value: 0) let activeProviderRelease = OSAllocatedUnfairLock(initialState: false) + let firstProbeStarted = ForkProbeTestSignal() + let collidingRefreshStarted = ForkProbeTestSignal() + let laterRefreshStarted = ForkProbeTestSignal() + defer { + activeProviderRelease.withLock { $0 = true } + releaseBlockedLoaders.signal() + releaseBlockedLoaders.signal() + } let firstProviderBlockCount = OSAllocatedUnfairLock(initialState: 0) let probedSessionIds = OSAllocatedUnfairLock(initialState: [String]()) let sharedIndex = SharedLiveAgentIndex( @@ -1976,7 +1984,6 @@ struct WorkspaceForkConversationContextMenuTests { return count } if call > 1 { - blockedLoaderCount.withLock { $0 += 1 } releaseBlockedLoaders.wait() } let index = RestorableAgentSessionIndex.load( @@ -2015,6 +2022,7 @@ struct WorkspaceForkConversationContextMenuTests { return count } if blockCount == 1 { + firstProbeStarted.signal() while !Task.isCancelled && !activeProviderRelease.withLock({ $0 }) { await Task.yield() } @@ -2033,35 +2041,28 @@ struct WorkspaceForkConversationContextMenuTests { panelId: firstPanelId ) } - for _ in 0..<1000 where probedSessionIds.withLock({ $0 }) != ["first-active-wait-batch"] { - await Task.yield() - } + try #require(await firstProbeStarted.wait(), "First validation must be active before queuing later batches.") #expect(probedSessionIds.withLock { $0 } == ["first-active-wait-batch"]) let collidingRefresh = Task { + collidingRefreshStarted.signal() await sharedIndex.refreshForkAvailabilityNow( workspaceId: firstWorkspaceId, panelId: firstPanelId ) } - for _ in 0..<1000 where blockedLoaderCount.withLock({ $0 }) < 1 { - await Task.yield() - } + try #require(await collidingRefreshStarted.wait(), "Queued refresh must enter before releasing loaders.") let laterRefresh = Task { + laterRefreshStarted.signal() await sharedIndex.refreshForkAvailabilityNow( workspaceId: secondWorkspaceId, panelId: secondPanelId ) } - for _ in 0..<1000 where blockedLoaderCount.withLock({ $0 }) < 2 { - await Task.yield() - } + try #require(await laterRefreshStarted.wait(), "Queued refresh must enter before releasing loaders.") releaseBlockedLoaders.signal() releaseBlockedLoaders.signal() - for _ in 0..<1000 where firstProviderBlockCount.withLock({ $0 }) < 1 { - await Task.yield() - } #expect(probedSessionIds.withLock { $0 } == ["first-active-wait-batch"]) activeProviderRelease.withLock { $0 = true } @@ -2080,7 +2081,7 @@ struct WorkspaceForkConversationContextMenuTests { #expect(forkValidationCancellationTombstoneCount(in: sharedIndex) == 0) } - @Test + @Test(.timeLimit(.minutes(1))) func cancelledSharedForkProbeRefreshPreservesSurvivingFallbackSnapshot() async throws { let fm = FileManager.default let root = fm.temporaryDirectory @@ -2105,6 +2106,9 @@ struct WorkspaceForkConversationContextMenuTests { let loaderCallCount = OSAllocatedUnfairLock(initialState: 0) let providerStartedCount = OSAllocatedUnfairLock(initialState: 0) let firstProviderRelease = OSAllocatedUnfairLock(initialState: false) + let firstProviderStarted = ForkProbeTestSignal() + let secondRefreshStarted = ForkProbeTestSignal() + defer { firstProviderRelease.withLock { $0 = true } } let secondRefreshFinished = OSAllocatedUnfairLock(initialState: false) let probedSessionIds = OSAllocatedUnfairLock(initialState: [String]()) let sharedIndex = SharedLiveAgentIndex( @@ -2125,6 +2129,7 @@ struct WorkspaceForkConversationContextMenuTests { forkSupportProvider: { snapshot, _ in providerStartedCount.withLock { $0 += 1 } probedSessionIds.withLock { $0.append(snapshot.sessionId) } + firstProviderStarted.signal() while !Task.isCancelled && !firstProviderRelease.withLock({ $0 }) { await Task.yield() } @@ -2142,12 +2147,11 @@ struct WorkspaceForkConversationContextMenuTests { fallbackSnapshot: firstFallback ) } - for _ in 0..<1000 where providerStartedCount.withLock({ $0 }) < 1 { - await Task.yield() - } + try #require(await firstProviderStarted.wait(), "First fallback probe must be active before the second refresh.") #expect(providerStartedCount.withLock { $0 } == 1) let secondRefresh = Task { + secondRefreshStarted.signal() await sharedIndex.refreshForkAvailabilityNow( workspaceId: workspaceId, panelId: panelId, @@ -2155,9 +2159,7 @@ struct WorkspaceForkConversationContextMenuTests { ) secondRefreshFinished.withLock { $0 = true } } - for _ in 0..<1000 where !secondRefreshFinished.withLock({ $0 }) { - await Task.yield() - } + try #require(await secondRefreshStarted.wait(), "Second refresh must enter before cancellation.") #expect(!secondRefreshFinished.withLock { $0 }) secondRefresh.cancel() @@ -2225,10 +2227,13 @@ struct WorkspaceForkConversationContextMenuTests { let loaderCallCount = OSAllocatedUnfairLock(initialState: 0) let loaderStartedCount = OSAllocatedUnfairLock(initialState: 0) let firstLoaderRelease = DispatchSemaphore(value: 0) + let firstLoaderStarted = ForkProbeTestSignal() + defer { firstLoaderRelease.signal() } let secondProviderRelease = OSAllocatedUnfairLock(initialState: false) + defer { secondProviderRelease.withLock { $0 = true } } let probedSessionIds = OSAllocatedUnfairLock(initialState: [String]()) - let secondRefreshStarted = AsyncStream.makeStream() - let secondProbeStarted = AsyncStream.makeStream() + let secondRefreshStarted = ForkProbeTestSignal() + let secondProbeStarted = ForkProbeTestSignal() let sharedIndex = SharedLiveAgentIndex( indexLoader: { let call = loaderCallCount.withLock { count in @@ -2239,6 +2244,7 @@ struct WorkspaceForkConversationContextMenuTests { count = max(count, call) } if call == 1 { + firstLoaderStarted.signal() firstLoaderRelease.wait() } let index = RestorableAgentSessionIndex.load( @@ -2272,7 +2278,7 @@ struct WorkspaceForkConversationContextMenuTests { forkSupportProvider: { snapshot, _ in probedSessionIds.withLock { $0.append(snapshot.sessionId) } if snapshot.sessionId == "second-replay" { - secondProbeStarted.continuation.yield(()) + secondProbeStarted.signal() while !Task.isCancelled && !secondProviderRelease.withLock({ $0 }) { await Task.yield() } @@ -2290,21 +2296,19 @@ struct WorkspaceForkConversationContextMenuTests { panelId: firstPanelId ) } - for _ in 0..<1000 where loaderStartedCount.withLock({ $0 }) < 1 { - await Task.yield() - } + try #require(await firstLoaderStarted.wait(), "First index loader must start before queuing another refresh.") #expect(loaderStartedCount.withLock { $0 } == 1) let secondRefresh = Task { - secondRefreshStarted.continuation.yield(()) + secondRefreshStarted.signal() await sharedIndex.refreshForkAvailabilityNow( workspaceId: secondWorkspaceId, panelId: secondPanelId ) } - _ = await secondRefreshStarted.stream.first(where: { _ in true }) + try #require(await secondRefreshStarted.wait(), "Expected refresh/probe must start before proceeding.") firstLoaderRelease.signal() - _ = await secondProbeStarted.stream.first(where: { _ in true }) + try #require(await secondProbeStarted.wait(), "Expected refresh/probe must start before proceeding.") #expect(probedSessionIds.withLock { $0 } == ["first-replay", "second-replay"]) secondRefresh.cancel() @@ -2399,6 +2403,12 @@ struct WorkspaceForkConversationContextMenuTests { let firstProviderRelease = OSAllocatedUnfairLock(initialState: false) let secondProviderRelease = OSAllocatedUnfairLock(initialState: false) let probedSessionIds = OSAllocatedUnfairLock(initialState: [String]()) + let firstProbeStarted = ForkProbeTestSignal() + let secondProbeStarted = ForkProbeTestSignal() + defer { + firstProviderRelease.withLock { $0 = true } + secondProviderRelease.withLock { $0 = true } + } let sharedIndex = SharedLiveAgentIndex( indexLoader: { loaderCallCount.withLock { $0 += 1 } @@ -2440,10 +2450,12 @@ struct WorkspaceForkConversationContextMenuTests { forkSupportProvider: { snapshot, _ in probedSessionIds.withLock { $0.append(snapshot.sessionId) } if snapshot.sessionId == "first-background" { + firstProbeStarted.signal() while !Task.isCancelled && !firstProviderRelease.withLock({ $0 }) { await Task.yield() } } else if snapshot.sessionId == "second-background" { + secondProbeStarted.signal() while !Task.isCancelled && !secondProviderRelease.withLock({ $0 }) { await Task.yield() } @@ -2469,27 +2481,17 @@ struct WorkspaceForkConversationContextMenuTests { } #expect(!sharedIndex.prepareForkAvailabilityProbe(workspaceId: firstWorkspaceId, panelId: firstPanelId)) - for _ in 0..<1000 where !probedSessionIds.withLock({ $0.contains("first-background") }) { - await Task.yield() - } + try #require(await firstProbeStarted.wait(), "First background probe must start before queuing the next batch.") #expect(probedSessionIds.withLock { $0 } == ["first-background"]) #expect(!sharedIndex.prepareForkAvailabilityProbe(workspaceId: secondWorkspaceId, panelId: secondPanelId)) - for _ in 0..<1000 where probedSessionIds.withLock({ $0.count }) != 1 { - await Task.yield() - } #expect(probedSessionIds.withLock { $0 } == ["first-background"]) firstProviderRelease.withLock { $0 = true } - for _ in 0..<1000 where !probedSessionIds.withLock({ $0.contains("second-background") }) { - await Task.yield() - } + try #require(await secondProbeStarted.wait(), "Second background probe must start before queuing the next batch.") #expect(probedSessionIds.withLock { $0 } == ["first-background", "second-background"]) #expect(!sharedIndex.prepareForkAvailabilityProbe(workspaceId: thirdWorkspaceId, panelId: thirdPanelId)) - for _ in 0..<1000 where probedSessionIds.withLock({ $0.count }) != 2 { - await Task.yield() - } #expect(probedSessionIds.withLock { $0 } == ["first-background", "second-background"]) secondProviderRelease.withLock { $0 = true } @@ -5147,3 +5149,22 @@ struct WorkspaceForkConversationContextMenuTests { ] } } + +// Buffered start signals establish fixture ordering without depending on executor scheduling. +private struct ForkProbeTestSignal: Sendable { + private let events = AsyncStream.makeStream() + + func signal() { + events.continuation.yield(()) + } + + func wait() async -> Bool { + let timeout = Task { + try? await Task.sleep(nanoseconds: 5_000_000_000) + guard !Task.isCancelled else { return } + events.continuation.finish() + } + defer { timeout.cancel() } + return await events.stream.first(where: { _ in true }) != nil + } +}