Merge nucleic/fuzzy-yarn-otter-2pji into dev

This commit is contained in:
2026-08-05 18:19:16 -07:00
parent 17e84c9573
commit 37c33d305b
6 changed files with 1002 additions and 198 deletions
@@ -25,6 +25,191 @@ import NIOPosix
import Synchronization
@preconcurrency import Virtualization
/// Exactly-once arbitration for Virtualization's callback-only vsock connect API. The API doesn't
/// expose a cancellation token, so a timeout settles the caller immediately and a connection which
/// arrives after that settlement is closed instead of being leaked or adopted.
private final class VZDialAttempt: @unchecked Sendable {
private enum State {
case idle
case pending(CheckedContinuation<VZVirtioSocketConnection, any Error>, Task<Void, Never>?)
case cancelled
case finished
}
private let state = Mutex<State>(.idle)
func install(_ continuation: CheckedContinuation<VZVirtioSocketConnection, any Error>) {
let wasCancelled = state.withLock { state in
switch state {
case .idle:
state = .pending(continuation, nil)
return false
case .cancelled:
state = .finished
return true
case .pending, .finished:
preconditionFailure("VZ dial continuation installed more than once")
}
}
if wasCancelled {
continuation.resume(throwing: CancellationError())
}
}
func setTimer(_ timer: Task<Void, Never>) {
let shouldCancel = state.withLock { state in
guard case .pending(let continuation, nil) = state else {
return true
}
state = .pending(continuation, timer)
return false
}
if shouldCancel {
timer.cancel()
}
}
var isPending: Bool {
state.withLock {
if case .pending = $0 { return true }
return false
}
}
func succeed(_ connection: VZVirtioSocketConnection) {
let completion = state.withLock { state -> (CheckedContinuation<VZVirtioSocketConnection, any Error>, Task<Void, Never>?)? in
guard case .pending(let continuation, let timer) = state else { return nil }
state = .finished
return (continuation, timer)
}
guard let completion else {
connection.close()
return
}
completion.1?.cancel()
// VZ's connection type isn't Sendable, but ownership transfers exactly once to the
// continuation here and this attempt never touches the successful connection again.
nonisolated(unsafe) let connection = connection
completion.0.resume(returning: connection)
}
func fail(_ error: any Error) {
let completion = state.withLock { state -> (CheckedContinuation<VZVirtioSocketConnection, any Error>, Task<Void, Never>?)? in
guard case .pending(let continuation, let timer) = state else { return nil }
state = .finished
return (continuation, timer)
}
guard let completion else { return }
completion.1?.cancel()
completion.0.resume(throwing: error)
}
func cancel() {
let completion = state.withLock { state -> (CheckedContinuation<VZVirtioSocketConnection, any Error>, Task<Void, Never>?)? in
switch state {
case .idle:
state = .cancelled
return nil
case .pending(let continuation, let timer):
state = .finished
return (continuation, timer)
case .cancelled, .finished:
return nil
}
}
guard let completion else { return }
completion.1?.cancel()
completion.0.resume(throwing: CancellationError())
}
}
/// Deadline arbiter for the whole locked dial operation. `AsyncLock` does not remove cancelled
/// waiters, so the caller is settled at the deadline while the cancelled waiter performs no connect
/// if it is eventually resumed and reaches the cancellation check inside the lock.
private final class VZLockedDialAttempt<Value: Sendable>: @unchecked Sendable {
private struct State {
var continuation: CheckedContinuation<Value, any Error>?
var operation: Task<Void, Never>?
var timer: Task<Void, Never>?
var finished = false
}
private let state = Mutex(State())
func install(_ continuation: CheckedContinuation<Value, any Error>) {
let wasCancelled = state.withLock { state in
guard !state.finished else { return true }
state.continuation = continuation
return false
}
if wasCancelled {
continuation.resume(throwing: CancellationError())
}
}
func setOperation(_ operation: Task<Void, Never>) {
let cancel = state.withLock { state in
guard !state.finished else { return true }
state.operation = operation
return false
}
if cancel { operation.cancel() }
}
func setTimer(_ timer: Task<Void, Never>) {
let cancel = state.withLock { state in
guard !state.finished else { return true }
state.timer = timer
return false
}
if cancel { timer.cancel() }
}
func succeed(_ value: Value) {
finish(.success(value), cancelOperation: false)
}
func fail(_ error: any Error, cancelOperation: Bool = false) {
finish(.failure(error), cancelOperation: cancelOperation)
}
func cancel() {
let completion = state.withLock { state -> (
CheckedContinuation<Value, any Error>?, Task<Void, Never>?, Task<Void, Never>?
)? in
guard !state.finished else { return nil }
state.finished = true
let completion = (state.continuation, state.operation, state.timer)
state.continuation = nil
state.operation = nil
state.timer = nil
return completion
}
guard let completion else { return }
completion.1?.cancel()
completion.2?.cancel()
completion.0?.resume(throwing: CancellationError())
}
private func finish(_ result: Result<Value, any Error>, cancelOperation: Bool) {
let completion = state.withLock { state -> (
CheckedContinuation<Value, any Error>, Task<Void, Never>?, Task<Void, Never>?
)? in
guard !state.finished, let continuation = state.continuation else { return nil }
state.finished = true
let completion = (continuation, state.operation, state.timer)
state.continuation = nil
state.operation = nil
state.timer = nil
return completion
}
guard let completion else { return }
completion.2?.cancel()
if cancelOperation { completion.1?.cancel() }
completion.0.resume(with: result)
}
}
public final class VZVirtualMachineInstance: Sendable {
public typealias Agent = Vminitd
@@ -40,6 +225,11 @@ public final class VZVirtualMachineInstance: Sendable {
/// The dispatch queue used for VZ operations.
public var vmQueue: DispatchQueue { queue }
/// Capture the independent, once-only stop lane for this exact VM generation.
public func emergencyStopHandle(generation: UInt64) -> EmergencyVMHandle {
EmergencyVMHandle(virtualMachine: vm, queue: queue, generation: generation)
}
/// Mutate the mount registry.
public func withMountRegistry<T: Sendable>(_ body: (inout sending [String: [AttachedFilesystem]]) throws -> sending T) rethrows -> T {
try _mounts.withLock(body)
@@ -253,14 +443,15 @@ extension VZVirtualMachineInstance: VirtualMachineInstance {
}
public func dialAgent() async throws -> Vminitd {
try await lock.withLock { _ in
let deadline = DeadlinePolicy.standard.deadline(for: .dial)
return try await withLockedDialDeadline(deadline: deadline) {
do {
let conn = try await self.vm.connect(
queue: self.queue,
port: Vminitd.port
)
let handle = try conn.dupHandle()
return try await Vminitd(connection: handle, group: self.group)
return try await self.lock.withLock { _ in
try Task.checkCancellation()
let conn = try await self.connect(port: Vminitd.port, deadline: deadline)
let handle = try conn.dupHandle()
return try await Vminitd(connection: handle, group: self.group)
}
} catch {
if let err = error as? ContainerizationError {
throw err
@@ -274,14 +465,38 @@ extension VZVirtualMachineInstance: VirtualMachineInstance {
}
}
/// Open a brand-new agent/vsock channel without entering the instance lifecycle lock.
/// Intended for independent health checks whose caller supplies a short deadline policy.
public func dialFreshAgent(deadlinePolicy: DeadlinePolicy) async throws -> Vminitd {
do {
let conn = try await self.connect(
port: Vminitd.port,
deadline: deadlinePolicy.deadline(for: .dial),
deadlinePolicy: deadlinePolicy)
let handle = try conn.dupHandle()
return try await Vminitd(
connection: handle, group: self.group, deadlinePolicy: deadlinePolicy)
} catch {
if let err = error as? ContainerizationError {
throw err
}
throw ContainerizationError(
.internalError,
message: "failed to dial fresh agent channel",
cause: error
)
}
}
public func dial(_ port: UInt32) async throws -> FileHandle {
try await lock.withLock { _ in
let deadline = DeadlinePolicy.standard.deadline(for: .dial)
return try await withLockedDialDeadline(deadline: deadline) {
do {
let conn = try await self.vm.connect(
queue: self.queue,
port: port
)
return try conn.dupHandle()
return try await self.lock.withLock { _ in
try Task.checkCancellation()
let conn = try await self.connect(port: port, deadline: deadline)
return try conn.dupHandle()
}
} catch {
if let err = error as? ContainerizationError {
throw err
@@ -295,6 +510,80 @@ extension VZVirtualMachineInstance: VirtualMachineInstance {
}
}
private func withLockedDialDeadline<Value: Sendable>(
deadline: DeadlinePolicy.Deadline,
operation: @escaping @Sendable () async throws -> Value
) async throws -> Value {
let attempt = VZLockedDialAttempt<Value>()
return try await withTaskCancellationHandler {
try await withCheckedThrowingContinuation { continuation in
attempt.install(continuation)
let operationTask = Task {
do {
attempt.succeed(try await operation())
} catch {
attempt.fail(error)
}
}
attempt.setOperation(operationTask)
let timer = Task {
do {
try await ContinuousClock().sleep(until: deadline)
} catch {
return
}
attempt.fail(
DeadlinePolicy.standard.timeoutError(for: .dial),
cancelOperation: true)
}
attempt.setTimer(timer)
}
} onCancel: {
attempt.cancel()
}
}
private func connect(
port: UInt32,
deadline: DeadlinePolicy.Deadline,
deadlinePolicy: DeadlinePolicy = .standard
) async throws -> VZVirtioSocketConnection {
let attempt = VZDialAttempt()
return try await withTaskCancellationHandler {
try await withCheckedThrowingContinuation { continuation in
attempt.install(continuation)
let timer = Task {
do {
try await ContinuousClock().sleep(until: deadline)
} catch {
return
}
attempt.fail(deadlinePolicy.timeoutError(for: .dial))
}
attempt.setTimer(timer)
self.queue.async {
guard attempt.isPending else { return }
guard let vsock = self.vm.socketDevices[0] as? VZVirtioSocketDevice else {
attempt.fail(ContainerizationError(.invalidArgument, message: "no vsock device"))
return
}
vsock.connect(toPort: port) { result in
switch result {
case .success(let connection):
attempt.succeed(connection)
case .failure(let error):
attempt.fail(error)
}
}
}
}
} onCancel: {
attempt.cancel()
}
}
public func listen(_ port: UInt32) throws -> VsockListener {
let stream = VsockListener(port: port, stopListen: self.stopListen)
let listener = VZVirtioSocketListener()