Skip to content
Merged
Show file tree
Hide file tree
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
Original file line number Diff line number Diff line change
Expand Up @@ -8,13 +8,16 @@ internal import CMUXDebugLog

/// Coordinates native `ghostty_surface_free` calls off the close/deinit paths.
///
/// Close/deinit frees run one at a time on a utility worker so re-entrant
/// teardown loops cannot form. Each admitted hibernation owns one independently
/// startable utility slot, so one stuck native join cannot strand another pane.
/// Deadline observers report, but never block on, stuck frees. The app constructs
/// exactly one instance and injects it through
/// Close/deinit frees run on a bounded set of utility slots so one stuck native
/// join cannot strand later closes. Each admitted hibernation owns a separate,
/// independently startable utility slot. Deadline observers report, but never
/// block on, stuck frees. The app constructs exactly one instance and injects it
/// through
/// ``TerminalSurfaceRuntimeDependencies``.
public actor TerminalSurfaceRuntimeTeardownCoordinator {
/// Maximum number of close/deinit native frees that can run concurrently.
public static let maximumConcurrentCloseTeardownCount = 2

/// Largest batch that can own independently startable native-free slots.
public static let maximumIsolatedHibernationTeardownCount = 2

Expand All @@ -27,14 +30,26 @@ public actor TerminalSurfaceRuntimeTeardownCoordinator {
#else
private var pendingReasonsById: [UUID: String] = [:]
#endif
private var queuedRequests: [TerminalSurfaceRuntimeTeardownRequest] = []
private var isWorkerRunning = false
private var queuedCloseRequests: [TerminalSurfaceRuntimeTeardownRequest] = []
private var availableCloseExecutionSlots: Set<Int>
private let closeTeardownQueues: [DispatchQueue]
private let isolatedHibernationQueues: [DispatchQueue]
private nonisolated let isolatedHibernationAdmission =
TerminalSurfaceRuntimeTeardownAdmission()

/// Creates the process's teardown coordinator.
public init() {
availableCloseExecutionSlots = Set(
0..<Self.maximumConcurrentCloseTeardownCount
)
closeTeardownQueues = (
0..<Self.maximumConcurrentCloseTeardownCount
).map { executionSlot in
DispatchQueue(
label: "com.cmux.terminal-surface-close-teardown.\(executionSlot)",
qos: .utility
)
}
isolatedHibernationQueues = (
0..<Self.maximumIsolatedHibernationTeardownCount
).map { executionSlot in
Expand Down Expand Up @@ -140,7 +155,7 @@ public actor TerminalSurfaceRuntimeTeardownCoordinator {
callbackContext: Unmanaged<GhosttySurfaceCallbackContext>?,
manualIOContext: Unmanaged<TerminalManualIOWriteBox>?,
byteTeeLease: (any TerminalByteTeeLease)?,
executionLane: TerminalSurfaceRuntimeTeardownExecutionLane = .serializedClose,
executionLane: TerminalSurfaceRuntimeTeardownExecutionLane = .boundedClose,
isolatedHibernationReservation:
TerminalSurfaceRuntimeTeardownReservation? = nil,
freeSurface: @escaping @Sendable (ghostty_surface_t) -> Void = { surface in
Expand Down Expand Up @@ -172,7 +187,7 @@ public actor TerminalSurfaceRuntimeTeardownCoordinator {

func enqueue(
_ request: TerminalSurfaceRuntimeTeardownRequest,
executionLane: TerminalSurfaceRuntimeTeardownExecutionLane = .serializedClose,
executionLane: TerminalSurfaceRuntimeTeardownExecutionLane = .boundedClose,
isolatedHibernationReservation:
TerminalSurfaceRuntimeTeardownReservation? = nil
) async {
Expand Down Expand Up @@ -209,31 +224,41 @@ public actor TerminalSurfaceRuntimeTeardownCoordinator {
isolatedHibernationReservation
)
}
case .serializedClose:
case .boundedClose:
break
}
queuedRequests.append(request)
if !isWorkerRunning {
isWorkerRunning = true
Task.detached(priority: .utility) {
while let request = await self.nextRequestForWorker() {
Task {
await self.observeTimeout(id: request.id)
}
self.freeNativeSurface(request)
await self.finishFree(request)
await self.complete(id: request.id)
queuedCloseRequests.append(request)
startAvailableCloseTeardowns()
}

private func startAvailableCloseTeardowns() {
while !queuedCloseRequests.isEmpty,
let executionSlot = availableCloseExecutionSlots.min() {
availableCloseExecutionSlots.remove(executionSlot)
let request = queuedCloseRequests.removeFirst()
Task {
await self.observeTimeout(id: request.id)
}
closeTeardownQueues[executionSlot].async {
self.freeNativeSurface(request)
Task {
await self.finishCloseTeardown(
request,
executionSlot: executionSlot
)
}
}
}
}

private func nextRequestForWorker() -> TerminalSurfaceRuntimeTeardownRequest? {
guard !queuedRequests.isEmpty else {
isWorkerRunning = false
return nil
}
return queuedRequests.removeFirst()
private func finishCloseTeardown(
_ request: TerminalSurfaceRuntimeTeardownRequest,
executionSlot: Int
) async {
await finishFree(request)
complete(id: request.id)
availableCloseExecutionSlots.insert(executionSlot)
startAvailableCloseTeardowns()
Comment thread
austinywang marked this conversation as resolved.
}

private nonisolated func freeNativeSurface(
Expand Down
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
/// Selects the ownership boundary for a native surface free.
enum TerminalSurfaceRuntimeTeardownExecutionLane: Sendable {
/// Preserves ordering for close/deinit flows that can re-enter teardown.
case serializedClose
/// Uses the bounded close/deinit pool without blocking later frees.
case boundedClose

/// Gives an explicitly owned hibernation join an independent bounded slot.
case isolatedHibernation
Expand Down
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
public import Foundation

/// An awaitable handle for one serialized native surface teardown.
/// An awaitable handle for one native surface teardown.
public struct TerminalSurfaceRuntimeTeardownTicket: Sendable {
/// Stable identity used by the surface to clear only its current teardown.
public let id: UUID
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,7 @@ public struct TerminalSurfaceRuntimeDependencies {
/// The agent-hibernation input recorder.
public let hibernationRecorder: any AgentHibernationRecording

/// The serialized native-surface free queue.
/// The bounded native-surface teardown coordinator.
public let runtimeTeardown: TerminalSurfaceRuntimeTeardownCoordinator

/// The paced native-surface creation queue for restored terminal sessions.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -326,14 +326,15 @@ extension TerminalSurface {
}
#endif

Task { @MainActor in
// Keep free behavior aligned with deinit: perform the runtime teardown on
// the next main-actor turn so SIGHUP delivery is deterministic but non-reentrant.
ghostty_surface_free(surfaceToFree)
callbackContext?.release()
manualIOContext?.release()
teeLease?.release()
}
runtimeTeardown.enqueueRuntimeTeardown(
id: id,
workspaceId: tabId,
reason: "teardown",
surface: surfaceToFree,
callbackContext: callbackContext,
manualIOContext: manualIOContext,
byteTeeLease: teeLease
)
}

/// Frees the runtime surface while keeping the model alive for an
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -120,23 +120,83 @@ private final class LifetimeRecordingByteTeeLease: TerminalByteTeeLease, @unchec
#expect(await Set(recorder.freed) == Set(surfaces.map { UInt(bitPattern: $0) }))
}

@Test func stuckCloseFreeDoesNotStrandLaterCloses() async throws {
let coordinator = TerminalSurfaceRuntimeTeardownCoordinator()
let surfaces = (0..<3).map { _ in
UnsafeMutableRawPointer.allocate(byteCount: 8, alignment: 8)
}
defer { for surface in surfaces { surface.deallocate() } }
let stuckFreeStarted = AsyncStream<Void>.makeStream()
Comment thread
austinywang marked this conversation as resolved.
let releaseStuckFree = DispatchSemaphore(value: 0)
let freedSurfaceBits = OSAllocatedUnfairLock(initialState: Set<UInt>())
defer {
releaseStuckFree.signal()
stuckFreeStarted.continuation.finish()
}

let stuckTicket = coordinator.enqueueRuntimeTeardown(
id: UUID(),
workspaceId: UUID(),
reason: "test.stuckClose",
surface: surfaces[0],
callbackContext: nil,
freeSurface: { _ in
stuckFreeStarted.continuation.yield()
_ = releaseStuckFree.wait(timeout: .distantFuture)
}
)
var stuckFreeIterator = stuckFreeStarted.stream.makeAsyncIterator()
_ = await stuckFreeIterator.next()

let laterTickets = surfaces.dropFirst().map { surface in
coordinator.enqueueRuntimeTeardown(
id: UUID(),
workspaceId: UUID(),
reason: "test.laterClose",
surface: surface,
callbackContext: nil,
freeSurface: { pointer in
let bits = UInt(bitPattern: pointer)
freedSurfaceBits.withLock {
_ = $0.insert(bits)
}
}
)
}

for ticket in laterTickets {
try #require(
await ticket.wait(timeout: .seconds(1)),
Comment thread
austinywang marked this conversation as resolved.
Comment thread
austinywang marked this conversation as resolved.
"a stuck native free stranded a later close"
)
}
#expect(await stuckTicket.wait(timeout: .zero) == false)
#expect(
freedSurfaceBits.withLock { $0 } ==
Set(surfaces.dropFirst().map { UInt(bitPattern: $0) })
)

releaseStuckFree.signal()
#expect(await stuckTicket.wait(timeout: .seconds(1)))
}

@Test func stuckHibernationFreeDoesNotStrandAnotherAdmissionOrClose() async throws {
let coordinator = TerminalSurfaceRuntimeTeardownCoordinator()
let isolatedSurface = UnsafeMutableRawPointer.allocate(byteCount: 8, alignment: 8)
let queuedIsolatedSurface = UnsafeMutableRawPointer.allocate(
byteCount: 8,
alignment: 8
)
let serializedSurface = UnsafeMutableRawPointer.allocate(byteCount: 8, alignment: 8)
let closeSurface = UnsafeMutableRawPointer.allocate(byteCount: 8, alignment: 8)
defer {
isolatedSurface.deallocate()
queuedIsolatedSurface.deallocate()
serializedSurface.deallocate()
closeSurface.deallocate()
}
let isolatedFreeStarted = AsyncStream<Void>.makeStream()
let releaseIsolatedFree = DispatchSemaphore(value: 0)
let secondIsolatedFreeCount = OSAllocatedUnfairLock(initialState: 0)
let serializedFreeCount = OSAllocatedUnfairLock(initialState: 0)
let closeFreeCount = OSAllocatedUnfairLock(initialState: 0)
defer {
releaseIsolatedFree.signal()
isolatedFreeStarted.continuation.finish()
Expand Down Expand Up @@ -181,20 +241,20 @@ private final class LifetimeRecordingByteTeeLease: TerminalByteTeeLease, @unchec
secondIsolatedFreeCount.withLock { $0 += 1 }
}
)
let serializedTicket = coordinator.enqueueRuntimeTeardown(
let closeTicket = coordinator.enqueueRuntimeTeardown(
id: UUID(),
workspaceId: UUID(),
reason: "test.serializedClose",
surface: serializedSurface,
reason: "test.close",
surface: closeSurface,
callbackContext: nil,
freeSurface: { _ in
serializedFreeCount.withLock { $0 += 1 }
closeFreeCount.withLock { $0 += 1 }
}
)

#expect(await serializedTicket.wait(timeout: .seconds(1)))
#expect(await closeTicket.wait(timeout: .seconds(1)))
#expect(await secondIsolatedTicket.wait(timeout: .seconds(1)))
#expect(serializedFreeCount.withLock { $0 } == 1)
#expect(closeFreeCount.withLock { $0 } == 1)
#expect(await isolatedTicket.wait(timeout: .zero) == false)
#expect(secondIsolatedFreeCount.withLock { $0 } == 1)

Expand All @@ -211,21 +271,49 @@ private final class LifetimeRecordingByteTeeLease: TerminalByteTeeLease, @unchec
await coordinator.cancelIsolatedHibernationTeardown(nextReservation)
}

@Test func staleIsolatedReservationFallsBackToSerializedFree() async throws {
@Test func staleIsolatedReservationFallsBackToBoundedClose() async throws {
let coordinator = TerminalSurfaceRuntimeTeardownCoordinator()
let surface = UnsafeMutableRawPointer.allocate(byteCount: 8, alignment: 8)
defer { surface.deallocate() }
let surfaces = (0..<2).map { _ in
UnsafeMutableRawPointer.allocate(byteCount: 8, alignment: 8)
}
defer { for surface in surfaces { surface.deallocate() } }
let isolatedFreeStarted = AsyncStream<Void>.makeStream()
let releaseIsolatedFree = DispatchSemaphore(value: 0)
let freeCount = OSAllocatedUnfairLock(initialState: 0)
defer {
releaseIsolatedFree.signal()
isolatedFreeStarted.continuation.finish()
}
let staleReservation = try #require(
await coordinator.reserveIsolatedHibernationTeardown()
)
await coordinator.cancelIsolatedHibernationTeardown(staleReservation)
let blockingReservation = try #require(
await coordinator.reserveIsolatedHibernationTeardown()
)
let blockingTicket = coordinator.enqueueRuntimeTeardown(
id: UUID(),
workspaceId: UUID(),
reason: "test.blockingIsolatedReservation",
surface: surfaces[0],
callbackContext: nil,
manualIOContext: nil,
byteTeeLease: nil,
executionLane: .isolatedHibernation,
isolatedHibernationReservation: blockingReservation,
freeSurface: { _ in
isolatedFreeStarted.continuation.yield()
_ = releaseIsolatedFree.wait(timeout: .distantFuture)
}
)
var isolatedFreeIterator = isolatedFreeStarted.stream.makeAsyncIterator()
_ = await isolatedFreeIterator.next()

let ticket = coordinator.enqueueRuntimeTeardown(
id: UUID(),
workspaceId: UUID(),
reason: "test.staleIsolatedReservation",
surface: surface,
surface: surfaces[1],
callbackContext: nil,
manualIOContext: nil,
byteTeeLease: nil,
Expand All @@ -238,6 +326,10 @@ private final class LifetimeRecordingByteTeeLease: TerminalByteTeeLease, @unchec

#expect(await ticket.wait(timeout: .seconds(1)))
#expect(freeCount.withLock { $0 } == 1)
#expect(await blockingTicket.wait(timeout: .zero) == false)

releaseIsolatedFree.signal()
#expect(await blockingTicket.wait(timeout: .seconds(1)))
}

@Test func byteTeeCallbackOwnerIsReleasedOnlyAfterNativeFreeReturns() async {
Expand Down
Loading
Loading