Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
119 changes: 70 additions & 49 deletions cmuxTests/WorkspaceForkConversationContextMenuTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -1726,7 +1726,9 @@ struct WorkspaceForkConversationContextMenuTests {
let loaderCallCount = OSAllocatedUnfairLock(initialState: 0)
let loaderStartedCount = OSAllocatedUnfairLock(initialState: 0)
let firstLoaderRelease = DispatchSemaphore(value: 0)
let secondRefreshStarted = AsyncStream<Void>.makeStream()
let firstLoaderStarted = ForkProbeTestSignal()
defer { firstLoaderRelease.signal() }
let secondRefreshStarted = ForkProbeTestSignal()
let sharedIndex = SharedLiveAgentIndex(
indexLoader: {
let call = loaderCallCount.withLock { count in
Expand All @@ -1737,6 +1739,7 @@ struct WorkspaceForkConversationContextMenuTests {
count = max(count, call)
}
if call == 1 {
firstLoaderStarted.signal()
firstLoaderRelease.wait()
}
let index = RestorableAgentSessionIndex.load(
Expand Down Expand Up @@ -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()
Expand Down Expand Up @@ -1931,7 +1932,7 @@ struct WorkspaceForkConversationContextMenuTests {
)
}

@Test
@Test(.timeLimit(.minutes(1)))
func sharedForkProbeActiveWaitPreservesLaterQueuedValidationBatches() async throws {
let fm = FileManager.default
let root = fm.temporaryDirectory
Expand Down Expand Up @@ -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(
Expand All @@ -1976,7 +1984,6 @@ struct WorkspaceForkConversationContextMenuTests {
return count
}
if call > 1 {
blockedLoaderCount.withLock { $0 += 1 }
releaseBlockedLoaders.wait()
}
let index = RestorableAgentSessionIndex.load(
Expand Down Expand Up @@ -2015,6 +2022,7 @@ struct WorkspaceForkConversationContextMenuTests {
return count
}
if blockCount == 1 {
firstProbeStarted.signal()
while !Task.isCancelled && !activeProviderRelease.withLock({ $0 }) {
await Task.yield()
}
Expand All @@ -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 }
Expand All @@ -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
Expand All @@ -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(
Expand All @@ -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()
}
Expand All @@ -2142,22 +2147,19 @@ 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,
fallbackSnapshot: secondFallback
)
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()
Expand Down Expand Up @@ -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<Void>.makeStream()
let secondProbeStarted = AsyncStream<Void>.makeStream()
let secondRefreshStarted = ForkProbeTestSignal()
let secondProbeStarted = ForkProbeTestSignal()
let sharedIndex = SharedLiveAgentIndex(
indexLoader: {
let call = loaderCallCount.withLock { count in
Expand All @@ -2239,6 +2244,7 @@ struct WorkspaceForkConversationContextMenuTests {
count = max(count, call)
}
if call == 1 {
firstLoaderStarted.signal()
firstLoaderRelease.wait()
}
let index = RestorableAgentSessionIndex.load(
Expand Down Expand Up @@ -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()
}
Expand All @@ -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()
Expand Down Expand Up @@ -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 }
Expand Down Expand Up @@ -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()
}
Expand All @@ -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 }
Expand Down Expand Up @@ -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<Void>.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
}
}
Loading