diff --git a/PATCHES.md b/PATCHES.md index 792babb..2d538bb 100644 --- a/PATCHES.md +++ b/PATCHES.md @@ -307,6 +307,48 @@ rebuild whenever a guest patch changes. Built locally, not in CI: the host frame checked continuation; a wedged write costs one expendable GCD thread, and process deletion closing the fd still unwedges it. Marked `[Nucleic vendored patch]`. +19. **Monotonic boundary deadlines and generation-scoped exec admission (`DeadlinePolicy.swift`, + `Vminitd.swift`, `VZVirtualMachineInstance.swift`, `LinuxContainer.swift`, `LinuxProcess.swift`).** + A central, tunable operation-class policy now supplies absolute monotonic deadlines to Phase-1 + generated vminitd gRPC calls: statistics (2s), create (10s), start (15s), process control (3s), + delete (30s), and agent graceful close (3s), with bounded filesystem/config calls and explicit + process-wait timeouts. RPC expiry cancels the call; create/start timeout races schedule bounded + deletion of any partially committed guest record. Agent close force-cancels `runConnections()` + after its grace period, closing the transport, and `LinuxProcess` teardown inherits that bound. + + A nil-timeout `waitProcess` remains the one intentionally greppable generated-call exception: + renewable leases are deferred until guest `ManagedProcess.wait()` removes cancelled waiters; + renewing today would leak one checked continuation per deadline. + + VZ vsock connects now have a three-second exactly-once callback/deadline arbiter. Timeout or task + cancellation settles the continuation once; a Virtualization connection delivered after that is + immediately closed. The same absolute deadline covers waiting for the existing VM `AsyncLock`; + a timed-out waiter is cancelled and checks cancellation before it can connect after eventually + acquiring the lock. The lock remains in place pending the lifecycle-gate migration. Container + statistics snapshots VM/ID/generation under `LinuxContainer.state` and + performs dial/RPC/close entirely unlocked. Exec now reserves `(exec ID, operation UUID, + generation)`, dials unlocked, and commits only if the reservation and public `generation` still + match; stale completions close their new agent and throw + `LinuxContainerOperationError.staleGeneration`. Generation values now come from a process-wide + atomic sequence rather than restarting at zero per object, and public `markGenerationDead()` + lets an independent retirement lane invalidate old reservations before normal state cleanup. + Stop invokes the same invalidation before boundary cleanup; reservations are removed on dial + failure, stale commit, pause, or state replacement. + +20. **Independent emergency VM stop handle (`EmergencyVMHandle.swift`, + `VZVirtualMachineInstance.swift`).** A retained, generation-labelled handle captures the raw + `VZVirtualMachine` and its required dispatch queue before higher layers expose a live container. + Its once-only `requestStop()` enqueues Virtualization's hard stop directly, without entering + `LinuxContainer.state`, the VZ instance lifecycle lock, vminitd RPCs, statistics, or the normal + engine lifecycle actor. `waitUntilStopped(deadline:)` provides a separately bounded completion + signal by observing Virtualization state on the VM queue, so callers can defer rootfs/network + reuse until the old VM is stopped and make any deadline-expiry force path explicit. + `VZVirtualMachineInstance.dialFreshAgent(deadlinePolicy:)` similarly + supplies a newly opened, caller-deadlined vsock/gRPC channel for independent health checks; it + never adopts a pooled channel or enters the instance lifecycle lock. These are generic framework + concurrency/escape primitives only; replacement quorum, admission, notification, and rate-limit + policy remain in NucleicCore. + ## Re-vendoring a newer upstream commit 1. `git clone` upstream (or copy `.build/checkouts/containerization` after bumping the URL pin diff --git a/Sources/Containerization/DeadlinePolicy.swift b/Sources/Containerization/DeadlinePolicy.swift new file mode 100644 index 0000000..5b36852 --- /dev/null +++ b/Sources/Containerization/DeadlinePolicy.swift @@ -0,0 +1,110 @@ +//===----------------------------------------------------------------------===// +// Copyright © 2025-2026 Apple Inc. and the Containerization project authors. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// https://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. +//===----------------------------------------------------------------------===// + +import ContainerizationError +import Foundation +import GRPCCore + +/// Monotonic deadlines for operations which cross a process or transport boundary. +/// +/// Keep generated gRPC calls behind ``performGRPC(operation:deadline:_:)``. Requiring both an +/// operation class and an absolute deadline makes a call which silently falls back to generated +/// default options conspicuous during review and straightforward to reject in CI. +public struct DeadlinePolicy: Sendable { + public typealias Deadline = ContinuousClock.Instant + + public enum Operation: String, Sendable { + case dial + case statistics + case createProcess + case startProcess + case processControl + case deleteProcess + case agentClose + case waitProcess + case filesystem + } + + public struct Bounds: Sendable { + public var dial: Duration = .seconds(3) + public var statistics: Duration = .seconds(2) + public var createProcess: Duration = .seconds(10) + public var startProcess: Duration = .seconds(15) + public var processControl: Duration = .seconds(3) + public var deleteProcess: Duration = .seconds(30) + public var agentClose: Duration = .seconds(3) + public var waitProcess: Duration = .seconds(30) + public var filesystem: Duration = .seconds(10) + + public init() {} + } + + public static let standard = DeadlinePolicy() + + public var bounds: Bounds + + public init(bounds: Bounds = Bounds()) { + self.bounds = bounds + } + + public func bound(for operation: Operation) -> Duration { + switch operation { + case .dial: bounds.dial + case .statistics: bounds.statistics + case .createProcess: bounds.createProcess + case .startProcess: bounds.startProcess + case .processControl: bounds.processControl + case .deleteProcess: bounds.deleteProcess + case .agentClose: bounds.agentClose + case .waitProcess: bounds.waitProcess + case .filesystem: bounds.filesystem + } + } + + public func deadline(for operation: Operation, clock: ContinuousClock = .init()) -> Deadline { + clock.now.advanced(by: bound(for: operation)) + } + + /// Invoke a generated gRPC call with the remaining part of an absolute monotonic deadline. + /// gRPC aborts the RPC when this timeout expires, so the server operation and client-side + /// continuation aren't left live after the caller receives a timeout. + public func performGRPC( + operation: Operation, + deadline: Deadline, + _ call: @Sendable (GRPCCore.CallOptions) async throws -> Result + ) async throws -> Result { + let remaining = ContinuousClock().now.duration(to: deadline) + guard remaining > .zero else { + throw timeoutError(for: operation) + } + + var options = GRPCCore.CallOptions.defaults + options.timeout = remaining + do { + return try await call(options) + } catch let error as RPCError where error.code == .deadlineExceeded { + throw timeoutError(for: operation, cause: error) + } + } + + public func timeoutError(for operation: Operation, cause: (any Error)? = nil) -> ContainerizationError { + ContainerizationError( + .timeout, + message: "\(operation.rawValue) exceeded its monotonic deadline", + cause: cause + ) + } +} diff --git a/Sources/Containerization/EmergencyVMHandle.swift b/Sources/Containerization/EmergencyVMHandle.swift new file mode 100644 index 0000000..848e869 --- /dev/null +++ b/Sources/Containerization/EmergencyVMHandle.swift @@ -0,0 +1,149 @@ +//===----------------------------------------------------------------------===// +// Copyright © 2025-2026 Apple Inc. and the Containerization project authors. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// https://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. +//===----------------------------------------------------------------------===// + +#if os(macOS) +import Foundation +import Synchronization +@preconcurrency import Virtualization + +/// A once-only escape hatch for stopping a VM when its normal lifecycle path is unavailable. +/// +/// The handle deliberately retains only Virtualization's VM object, its required dispatch queue, +/// and caller-supplied generation metadata. Requesting a stop does not enter a container state +/// gate, an instance lifecycle lock, or a guest-agent/RPC path. +public final class EmergencyVMHandle: @unchecked Sendable { + private struct State { + enum Phase { + case ready + case requested + case stopped + } + + var phase: Phase = .ready + var waiters: [UUID: CheckedContinuation] = [:] + } + + /// Stable identity for this exact retained VM, independent of caller generation numbering. + public let id = UUID() + public let generation: UInt64 + + private nonisolated(unsafe) let virtualMachine: VZVirtualMachine + private let queue: DispatchQueue + private let state = Mutex(State()) + + init(virtualMachine: VZVirtualMachine, queue: DispatchQueue, generation: UInt64) { + self.virtualMachine = virtualMachine + self.queue = queue + self.generation = generation + } + + /// Whether Virtualization has reported stop completion for this exact VM. + public var hasStopped: Bool { + state.withLock { + if case .stopped = $0.phase { return true } + return false + } + } + + /// Enqueue a hard VM stop exactly once. Returns `true` only for the caller which requested it. + /// Completion is intentionally not awaited: this API is the independent request lane used when + /// normal lifecycle tasks may themselves be unable to make progress. + @discardableResult + public func requestStop() -> Bool { + let shouldRequest = state.withLock { state in + guard case .ready = state.phase else { return false } + state.phase = .requested + return true + } + guard shouldRequest else { return false } + + queue.async { [self] in + if virtualMachine.state == .stopped { + finishStopped() + return + } + if virtualMachine.state == .stopping { + scheduleStopStateObservation() + return + } + virtualMachine.stop { [weak self] _ in + self?.scheduleStopStateObservation(immediate: true) + } + } + return true + } + + /// Wait until Virtualization reports this exact VM stopped, but never beyond `deadline`. + /// `false` is an explicit force-retirement signal for the caller; the stop request itself stays + /// independent and lock-free even if the normal instance lifecycle lane is wedged. + public func waitUntilStopped(deadline: ContinuousClock.Instant) async -> Bool { + let waiterID = UUID() + return await withCheckedContinuation { continuation in + let alreadyStopped = state.withLock { state in + guard case .stopped = state.phase else { + state.waiters[waiterID] = continuation + return false + } + return true + } + if alreadyStopped { + continuation.resume(returning: true) + return + } + + Task { [weak self] in + do { + try await ContinuousClock().sleep(until: deadline) + } catch { + // Cancellation is also a bounded failure for the waiting caller. The retained + // stop request continues independently on Virtualization's queue. + } + self?.finishWaiter(waiterID, stopped: false) + } + } + } + + private func scheduleStopStateObservation(immediate: Bool = false) { + queue.asyncAfter(deadline: .now() + (immediate ? 0 : 0.05)) { [weak self] in + guard let self else { return } + guard case .requested = state.withLock({ $0.phase }) else { return } + if virtualMachine.state == .stopped { + finishStopped() + } else { + scheduleStopStateObservation() + } + } + } + + private func finishStopped() { + let waiters = state.withLock { state -> [CheckedContinuation] in + guard case .stopped = state.phase else { + state.phase = .stopped + let waiters = Array(state.waiters.values) + state.waiters.removeAll() + return waiters + } + return [] + } + for waiter in waiters { waiter.resume(returning: true) } + } + + private func finishWaiter(_ id: UUID, stopped: Bool) { + let waiter = state.withLock { $0.waiters.removeValue(forKey: id) } + waiter?.resume(returning: stopped) + } +} +#endif diff --git a/Sources/Containerization/LinuxContainer.swift b/Sources/Containerization/LinuxContainer.swift index 5c609b9..61b69b9 100644 --- a/Sources/Containerization/LinuxContainer.swift +++ b/Sources/Containerization/LinuxContainer.swift @@ -26,11 +26,26 @@ import SystemPackage import struct ContainerizationOS.Terminal +/// Errors produced while committing generation-scoped container operations. +public enum LinuxContainerOperationError: Error, Sendable, Equatable { + /// External work completed after the container generation which reserved it was invalidated. + case staleGeneration(expected: UInt64, actual: UInt64) +} + /// `LinuxContainer` is an easy to use type for launching and managing the /// full lifecycle of a Linux container ran inside of a virtual machine. public final class LinuxContainer: Container, Sendable { public static let maxIDLength = 64 + /// Process-wide source for generation identities. Allocating from one sequence, rather than + /// restarting at zero for each `LinuxContainer` object, keeps late lifecycle evidence from an + /// old object from aliasing a replacement object's first generation. + private static let generationSequence = Atomic(0) + + private static func allocateGeneration() -> UInt64 { + generationSequence.wrappingAdd(1, ordering: .relaxed).oldValue + } + /// The identifier of the container. public let id: String @@ -135,6 +150,25 @@ public final class LinuxContainer: Container, Sendable { private let state: AsyncMutex + private let _generation: Atomic + + /// Process-wide unique identifier for the current container generation. It changes when a + /// running generation is retired, allowing callers to reject results and evidence produced by + /// an older VM instance even when its replacement uses a newly allocated container object. + public var generation: UInt64 { + _generation.load(ordering: .acquiring) + } + + /// Atomically invalidate every reservation made against the current generation. The returned + /// identity is the newly installed dead-generation sentinel; no container object can later use + /// the same value for a live generation. + @discardableResult + public func markGenerationDead() -> UInt64 { + let deadGeneration = Self.allocateGeneration() + _generation.store(deadGeneration, ordering: .releasing) + return deadGeneration + } + // Ports to be allocated from for stdio and for // unix socket relays that are sharing a guest // uds to the host. @@ -147,6 +181,11 @@ public final class LinuxContainer: Container, Sendable { private let copyQueue = DispatchQueue(label: "com.apple.containerization.copy") private enum State: Sendable { + struct ExecReservation: Sendable, Equatable { + let operationID: UUID + let generation: UInt64 + } + /// The container class has been created but no live resources are running. case initialized /// The container's virtual machine has been setup and the runtime environment has been configured. @@ -171,6 +210,7 @@ public final class LinuxContainer: Container, Sendable { let process: LinuxProcess let relayManager: UnixSocketRelayManager var vendedProcesses: [String: LinuxProcess] + var execReservations: [String: ExecReservation] let fileMountContext: FileMountContext init(_ state: CreatedState, process: LinuxProcess) { @@ -178,6 +218,7 @@ public final class LinuxContainer: Container, Sendable { self.relayManager = state.relayManager self.process = process self.vendedProcesses = [:] + self.execReservations = [:] self.fileMountContext = state.fileMountContext } @@ -186,6 +227,7 @@ public final class LinuxContainer: Container, Sendable { self.relayManager = state.relayManager self.process = state.process self.vendedProcesses = state.vendedProcesses + self.execReservations = state.execReservations self.fileMountContext = state.fileMountContext } } @@ -195,6 +237,7 @@ public final class LinuxContainer: Container, Sendable { let relayManager: UnixSocketRelayManager let process: LinuxProcess var vendedProcesses: [String: LinuxProcess] + var execReservations: [String: ExecReservation] let fileMountContext: FileMountContext init(_ state: StartedState) { @@ -202,6 +245,7 @@ public final class LinuxContainer: Container, Sendable { self.relayManager = state.relayManager self.process = state.process self.vendedProcesses = state.vendedProcesses + self.execReservations = state.execReservations self.fileMountContext = state.fileMountContext } } @@ -361,6 +405,7 @@ public final class LinuxContainer: Container, Sendable { self.logger = logger self.config = configuration self.state = AsyncMutex(.initialized) + self._generation = Atomic(Self.allocateGeneration()) self.rootfs = rootfs self.writableLayer = writableLayer } @@ -779,6 +824,10 @@ extension LinuxContainer { relayManager = createdState.relayManager } + // Invalidate unlocked operations as soon as stop has validated a live state. A failed + // stop is still the end of this usable generation: late dial completions must not commit. + self.markGenerationDead() + var firstError: Error? do { try await relayManager.stopAll() @@ -904,85 +953,101 @@ extension LinuxContainer { /// Execute a new process in the container. The process is not started after this call, and must be manually started /// via the `start` method. public func exec(_ id: String, configuration: @Sendable @escaping (inout LinuxProcessConfiguration) throws -> Void) async throws -> LinuxProcess { - try await self.state.withLock { state in - var startedState = try state.startedState("exec") - - var spec = self.generateRuntimeSpec() - var config = LinuxProcessConfiguration() - try configuration(&config) - spec.process = config.toOCI() - // [Nucleic vendored patch] Per-exec memory ceiling → the exec's OCI resources, which the - // guest applies as memory.max on this exec's own cgroup (patch #9). - if let limit = config.memoryLimitInBytes { - spec.linux?.resources?.memory?.limit = Int64(limit) - } - - let stdio = IOUtil.setup( - portAllocator: self.hostVsockPorts, - stdin: config.stdin, - stdout: config.stdout, - stderr: config.stderr - ) - let agent = try await startedState.vm.dialAgent() - let process = LinuxProcess( - id, - containerID: self.id, - spec: spec, - io: stdio, - ociRuntimePath: self.config.ociRuntimePath, - agent: agent, - vm: startedState.vm, - logger: self.logger, - onDelete: { [weak self = self] in - await self?.removeProcess(id: id) - } - ) - - startedState.vendedProcesses[id] = process - state = .started(startedState) - - return process - } + var config = LinuxProcessConfiguration() + try configuration(&config) + return try await makeExec(id, configuration: config) } /// Execute a new process in the container. The process is not started after this call, and must be manually started /// via the `start` method. public func exec(_ id: String, configuration: LinuxProcessConfiguration) async throws -> LinuxProcess { - try await self.state.withLock { - var state = try $0.startedState("exec") + try await makeExec(id, configuration: configuration) + } - var spec = self.generateRuntimeSpec() - spec.process = configuration.toOCI() - // [Nucleic vendored patch] Per-exec memory ceiling → the exec's OCI resources (see above). - if let limit = configuration.memoryLimitInBytes { - spec.linux?.resources?.memory?.limit = Int64(limit) + private struct ExecSnapshot: Sendable { + let vm: any VirtualMachineInstance + let reservation: State.ExecReservation + } + + private func makeExec(_ id: String, configuration: LinuxProcessConfiguration) async throws -> LinuxProcess { + let snapshot = try await self.state.withLock { state in + var startedState = try state.startedState("exec") + let reservation = State.ExecReservation(operationID: UUID(), generation: self.generation) + startedState.execReservations[id] = reservation + state = .started(startedState) + return ExecSnapshot(vm: startedState.vm, reservation: reservation) + } + + let agent: any VirtualMachineAgent + do { + agent = try await snapshot.vm.dialAgent() + } catch { + await releaseExecReservation(id: id, reservation: snapshot.reservation) + throw error + } + + var spec = self.generateRuntimeSpec() + spec.process = configuration.toOCI() + // [Nucleic vendored patch] Per-exec memory ceiling → the exec's OCI resources, which the + // guest applies as memory.max on this exec's own cgroup (patch #9). + if let limit = configuration.memoryLimitInBytes { + spec.linux?.resources?.memory?.limit = Int64(limit) + } + let stdio = IOUtil.setup( + portAllocator: self.hostVsockPorts, + stdin: configuration.stdin, + stdout: configuration.stdout, + stderr: configuration.stderr + ) + let process = LinuxProcess( + id, + containerID: self.id, + spec: spec, + io: stdio, + ociRuntimePath: self.config.ociRuntimePath, + agent: agent, + vm: snapshot.vm, + logger: self.logger, + onDelete: { [weak self = self] in + await self?.removeProcess(id: id) } + ) - let stdio = IOUtil.setup( - portAllocator: self.hostVsockPorts, - stdin: configuration.stdin, - stdout: configuration.stdout, - stderr: configuration.stderr - ) - let agent = try await state.vm.dialAgent() - let process = LinuxProcess( - id, - containerID: self.id, - spec: spec, - io: stdio, - ociRuntimePath: self.config.ociRuntimePath, - agent: agent, - vm: state.vm, - logger: self.logger, - onDelete: { [weak self = self] in - await self?.removeProcess(id: id) - } + let committed = await self.state.withLock { state in + guard case .started(var startedState) = state, + self.generation == snapshot.reservation.generation, + startedState.execReservations[id] == snapshot.reservation + else { + return false + } + startedState.execReservations.removeValue(forKey: id) + startedState.vendedProcesses[id] = process + state = .started(startedState) + return true + } + guard committed else { + await releaseExecReservation(id: id, reservation: snapshot.reservation) + try? await agent.close() + throw LinuxContainerOperationError.staleGeneration( + expected: snapshot.reservation.generation, + actual: self.generation ) + } + return process + } - state.vendedProcesses[id] = process - $0 = .started(state) - - return process + private func releaseExecReservation(id: String, reservation: State.ExecReservation) async { + await self.state.withLock { state in + switch state { + case .started(var startedState) where startedState.execReservations[id] == reservation: + startedState.execReservations.removeValue(forKey: id) + state = .started(startedState) + case .paused(var pausedState) where pausedState.execReservations[id] == reservation: + pausedState.execReservations.removeValue(forKey: id) + state = .paused(pausedState) + default: + break + } } } @@ -1030,21 +1095,23 @@ extension LinuxContainer { /// Get statistics for the container. public func statistics(categories: StatCategory = .all) async throws -> ContainerStatistics { - try await self.state.withLock { - let state = try $0.startedState("statistics") + let snapshot = try await self.state.withLock { state in + let startedState = try state.startedState("statistics") + return (vm: startedState.vm, containerID: self.id, generation: self.generation) + } - let stats = try await state.vm.withAgent { agent in - let allStats = try await agent.containerStatistics(containerIDs: [self.id], categories: categories) - guard let containerStats = allStats.first else { - throw ContainerizationError( - .notFound, - message: "statistics for container \(self.id) not found" - ) - } - return containerStats + return try await snapshot.vm.withAgent { agent in + let allStats = try await agent.containerStatistics( + containerIDs: [snapshot.containerID], + categories: categories + ) + guard let containerStats = allStats.first else { + throw ContainerizationError( + .notFound, + message: "statistics for container \(snapshot.containerID) not found" + ) } - - return stats + return containerStats } } diff --git a/Sources/Containerization/VZVirtualMachineInstance.swift b/Sources/Containerization/VZVirtualMachineInstance.swift index c05cfa5..327db2c 100644 --- a/Sources/Containerization/VZVirtualMachineInstance.swift +++ b/Sources/Containerization/VZVirtualMachineInstance.swift @@ -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, Task?) + case cancelled + case finished + } + + private let state = Mutex(.idle) + + func install(_ continuation: CheckedContinuation) { + 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) { + 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, Task?)? 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, Task?)? 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, Task?)? 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: @unchecked Sendable { + private struct State { + var continuation: CheckedContinuation? + var operation: Task? + var timer: Task? + var finished = false + } + + private let state = Mutex(State()) + + func install(_ continuation: CheckedContinuation) { + 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) { + 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) { + 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?, Task?, Task? + )? 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, cancelOperation: Bool) { + let completion = state.withLock { state -> ( + CheckedContinuation, Task?, Task? + )? 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(_ 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( + deadline: DeadlinePolicy.Deadline, + operation: @escaping @Sendable () async throws -> Value + ) async throws -> Value { + let attempt = VZLockedDialAttempt() + 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() diff --git a/Sources/Containerization/Vminitd.swift b/Sources/Containerization/Vminitd.swift index 57157cb..c69f266 100644 --- a/Sources/Containerization/Vminitd.swift +++ b/Sources/Containerization/Vminitd.swift @@ -33,8 +33,13 @@ public struct Vminitd: Sendable { let client: Com_Apple_Containerization_Sandbox_V3_SandboxContext.Client public let grpcClient: GRPCClient private let connectionTask: Task + private let deadlinePolicy: DeadlinePolicy - public init(connection: FileHandle, group: any EventLoopGroup) async throws { + public init( + connection: FileHandle, + group: any EventLoopGroup, + deadlinePolicy: DeadlinePolicy = .standard + ) async throws { // Configure the gRPC pipeline from inside the channel initializer — before the channel // becomes active — so no early server frames (e.g. SETTINGS) are dropped. `configure` is // supplied by `wrapping(config:serviceConfig:makeChannel:)` and must be called exactly once. @@ -54,6 +59,7 @@ public struct Vminitd: Sendable { let grpcClient = GRPCClient(transport: transport) self.grpcClient = grpcClient self.client = Com_Apple_Containerization_Sandbox_V3_SandboxContext.Client(wrapping: self.grpcClient) + self.deadlinePolicy = deadlinePolicy // Not very structured concurrency friendly, but we'd need to expose a way on the protocol to "run" the // agent otherwise, which some agents might not even need. self.connectionTask = Task { @@ -64,7 +70,41 @@ public struct Vminitd: Sendable { /// Close the connection to the guest agent. public func close() async throws { self.grpcClient.beginGracefulShutdown() - try await self.connectionTask.value + let deadline = deadlinePolicy.deadline(for: .agentClose) + do { + try await withThrowingTaskGroup(of: Void.self) { group in + group.addTask { + try await self.connectionTask.value + } + group.addTask { + try await ContinuousClock().sleep(until: deadline) + throw self.deadlinePolicy.timeoutError(for: .agentClose) + } + + do { + _ = try await group.next() + group.cancelAll() + } catch { + // Cancelling the task executing runConnections() is GRPCCore's force-shutdown + // API. It cancels in-flight RPCs and closes the underlying transport channel. + self.connectionTask.cancel() + group.cancelAll() + throw error + } + } + } catch let error as ContainerizationError where error.code == .timeout { + // The transport has already been force-shut down above. A bounded close has fulfilled + // its cleanup contract even though graceful shutdown did not finish in time. + return + } + } + + private func rpc( + operation: DeadlinePolicy.Operation, + deadline: DeadlinePolicy.Deadline, + _ call: @Sendable (GRPCCore.CallOptions) async throws -> Result + ) async throws -> Result { + try await deadlinePolicy.performGRPC(operation: operation, deadline: deadline, call) } } @@ -87,8 +127,7 @@ extension Vminitd: VirtualMachineAgent { } public func writeFile(path: String, data: Data, flags: WriteFileFlags, mode: UInt32) async throws { - _ = try await client.writeFile( - .with { + let request = Com_Apple_Containerization_Sandbox_V3_WriteFileRequest.with { $0.path = path $0.mode = mode $0.data = data @@ -97,17 +136,24 @@ extension Vminitd: VirtualMachineAgent { $0.createIfMissing = flags.create $0.createParentDirs = flags.createParentDirectories } - }) + } + let operation = DeadlinePolicy.Operation.filesystem + _ = try await rpc(operation: operation, deadline: deadlinePolicy.deadline(for: operation)) { options in + try await client.writeFile(request, options: options) + } } /// Get statistics for containers. If `containerIDs` is empty returns stats for all containers /// in the guest. If `categories` is empty, all categories are returned. public func containerStatistics(containerIDs: [String], categories: StatCategory) async throws -> [ContainerStatistics] { - let response = try await client.containerStatistics( - .with { + let request = Com_Apple_Containerization_Sandbox_V3_ContainerStatisticsRequest.with { $0.containerIds = containerIDs $0.categories = categories.toProtoCategories() - }) + } + let operation = DeadlinePolicy.Operation.statistics + let response = try await rpc(operation: operation, deadline: deadlinePolicy.deadline(for: operation)) { options in + try await client.containerStatistics(request, options: options) + } return response.containers.map { protoStats in ContainerStatistics( @@ -184,41 +230,56 @@ extension Vminitd: VirtualMachineAgent { /// Mount a filesystem in the sandbox's environment. public func mount(_ mount: ContainerizationOCI.Mount) async throws { - _ = try await client.mount( - .with { + let request = Com_Apple_Containerization_Sandbox_V3_MountRequest.with { $0.type = mount.type $0.source = mount.source $0.destination = mount.destination $0.options = mount.options - }) + } + let operation = DeadlinePolicy.Operation.filesystem + _ = try await rpc(operation: operation, deadline: deadlinePolicy.deadline(for: operation)) { options in + try await client.mount(request, options: options) + } } /// Unmount a filesystem in the sandbox's environment. public func umount(path: String, flags: Int32) async throws { - _ = try await client.umount( - .with { + let request = Com_Apple_Containerization_Sandbox_V3_UmountRequest.with { $0.path = path $0.flags = flags - }) + } + let operation = DeadlinePolicy.Operation.filesystem + _ = try await rpc(operation: operation, deadline: deadlinePolicy.deadline(for: operation)) { options in + try await client.umount(request, options: options) + } } /// Create a directory inside the sandbox's environment. public func mkdir(path: String, all: Bool, perms: UInt32) async throws { - _ = try await client.mkdir( - .with { + let request = Com_Apple_Containerization_Sandbox_V3_MkdirRequest.with { $0.path = path $0.all = all $0.perms = perms - }) + } + let operation = DeadlinePolicy.Operation.filesystem + _ = try await rpc(operation: operation, deadline: deadlinePolicy.deadline(for: operation)) { options in + try await client.mkdir(request, options: options) + } } /// Perform a filesystem operation on a path inside the sandbox's environment. public func filesystemOperation(operation: FilesystemOperation, path: String) async throws { - _ = try await client.filesystemOperation( - .with { + let request = Com_Apple_Containerization_Sandbox_V3_FilesystemOperationRequest.with { $0.operation = operation.toProtoOperation() $0.path = path - }) + } + let operationClass = DeadlinePolicy.Operation.filesystem + _ = try await rpc( + operation: operationClass, + deadline: deadlinePolicy.deadline(for: operationClass) + ) { options in + try await client.filesystemOperation(request, options: options) + } } public func createProcess( @@ -232,26 +293,34 @@ extension Vminitd: VirtualMachineAgent { options: Data? ) async throws { let enc = JSONEncoder() - _ = try await client.createProcess( - .with { - $0.id = id - if let stdinPort { - $0.stdin = stdinPort - } - if let stdoutPort { - $0.stdout = stdoutPort - } - if let stderrPort { - $0.stderr = stderrPort - } - if let containerID { - $0.containerID = containerID - } - if let ociRuntimePath { - $0.ociRuntimePath = ociRuntimePath - } - $0.configuration = try enc.encode(configuration) - }) + let request = try Com_Apple_Containerization_Sandbox_V3_CreateProcessRequest.with { + $0.id = id + if let stdinPort { + $0.stdin = stdinPort + } + if let stdoutPort { + $0.stdout = stdoutPort + } + if let stderrPort { + $0.stderr = stderrPort + } + if let containerID { + $0.containerID = containerID + } + if let ociRuntimePath { + $0.ociRuntimePath = ociRuntimePath + } + $0.configuration = try enc.encode(configuration) + } + let operation = DeadlinePolicy.Operation.createProcess + do { + _ = try await rpc(operation: operation, deadline: deadlinePolicy.deadline(for: operation)) { options in + try await client.createProcess(request, options: options) + } + } catch let error as ContainerizationError where error.code == .timeout { + cleanupTimedOutProcess(id: id, containerID: containerID, signalFirst: false) + throw error + } } @discardableResult @@ -262,7 +331,16 @@ extension Vminitd: VirtualMachineAgent { $0.containerID = containerID } } - let resp = try await client.startProcess(request) + let operation = DeadlinePolicy.Operation.startProcess + let resp: Com_Apple_Containerization_Sandbox_V3_StartProcessResponse + do { + resp = try await rpc(operation: operation, deadline: deadlinePolicy.deadline(for: operation)) { options in + try await client.startProcess(request, options: options) + } + } catch let error as ContainerizationError where error.code == .timeout { + cleanupTimedOutProcess(id: id, containerID: containerID, signalFirst: true) + throw error + } return resp.pid } @@ -274,7 +352,10 @@ extension Vminitd: VirtualMachineAgent { $0.containerID = containerID } } - _ = try await client.killProcess(request) + let operation = DeadlinePolicy.Operation.processControl + _ = try await rpc(operation: operation, deadline: deadlinePolicy.deadline(for: operation)) { options in + try await client.killProcess(request, options: options) + } } public func resizeProcess(id: String, containerID: String?, columns: UInt32, rows: UInt32) async throws { @@ -286,7 +367,10 @@ extension Vminitd: VirtualMachineAgent { $0.columns = columns $0.rows = rows } - _ = try await client.resizeProcess(request) + let operation = DeadlinePolicy.Operation.processControl + _ = try await rpc(operation: operation, deadline: deadlinePolicy.deadline(for: operation)) { options in + try await client.resizeProcess(request, options: options) + } } public func waitProcess( @@ -301,24 +385,29 @@ extension Vminitd: VirtualMachineAgent { } } - var callOpts = GRPCCore.CallOptions.defaults if let timeoutInSeconds { - callOpts.timeout = .seconds(timeoutInSeconds) - } - - do { - let resp = try await client.waitProcess(request, options: callOpts) - return ExitStatus(exitCode: resp.exitCode, exitedAt: resp.exitedAt.date) - } catch { - if let err = error as? RPCError, err.code == .deadlineExceeded { + let operation = DeadlinePolicy.Operation.waitProcess + let deadline = ContinuousClock().now.advanced(by: .seconds(timeoutInSeconds)) + do { + let resp = try await rpc(operation: operation, deadline: deadline) { options in + try await client.waitProcess(request, options: options) + } + return ExitStatus(exitCode: resp.exitCode, exitedAt: resp.exitedAt.date) + } catch let error as ContainerizationError where error.code == .timeout { throw ContainerizationError( .timeout, - message: "failed to wait for process exit within timeout of \(timeoutInSeconds!) seconds", - cause: err + message: "failed to wait for process exit within timeout of \(timeoutInSeconds) seconds", + cause: error.cause ) } - throw error } + + // Phase-2 exception, intentionally greppable: the guest's ManagedProcess.wait() does not + // remove a stored continuation when an RPC is cancelled. Renewing bounded RPCs here would + // leak one guest waiter per lease. Preserve the existing unbounded API until the guest gains + // cancellation-aware waiter removal; explicit caller timeouts above are deadline-routed. + let resp = try await client.waitProcess(request) + return ExitStatus(exitCode: resp.exitCode, exitedAt: resp.exitedAt.date) } public func deleteProcess(id: String, containerID: String?) async throws { @@ -328,16 +417,10 @@ extension Vminitd: VirtualMachineAgent { $0.containerID = containerID } } - // [Nucleic vendored patch] Bound the teardown RPC so a wedged agent channel can't hang an - // exec's cleanup forever. Nucleic fires `LinuxProcess.delete()` after every turn to reclaim - // the per-exec connection; if `deleteProcess` never returned, that reclaim task would leak - // and the connection would stay open — reintroducing the very accumulation the delete exists - // to prevent. Generous: a healthy delete returns in milliseconds, so this only trips a - // genuinely stuck channel, and `LinuxProcess.performDeletion` still closes the agent - // connection on the thrown deadline. - var callOpts = GRPCCore.CallOptions.defaults - callOpts.timeout = .seconds(30) - _ = try await client.deleteProcess(request, options: callOpts) + let operation = DeadlinePolicy.Operation.deleteProcess + _ = try await rpc(operation: operation, deadline: deadlinePolicy.deadline(for: operation)) { options in + try await client.deleteProcess(request, options: options) + } } public func closeProcessStdin(id: String, containerID: String?) async throws { @@ -347,7 +430,23 @@ extension Vminitd: VirtualMachineAgent { $0.containerID = containerID } } - _ = try await client.closeProcessStdin(request) + let operation = DeadlinePolicy.Operation.processControl + _ = try await rpc(operation: operation, deadline: deadlinePolicy.deadline(for: operation)) { options in + try await client.closeProcessStdin(request, options: options) + } + } + + /// The timed-out RPC has already been cancelled by gRPC. Reclaim any guest record which was + /// committed immediately before cancellation won the race. Cleanup has its own bounded RPCs so + /// the phase timeout can settle its caller without retaining an immortal cleanup task. + private func cleanupTimedOutProcess(id: String, containerID: String?, signalFirst: Bool) { + Task { + if signalFirst { + try? await self.signalProcess(id: id, containerID: containerID, signal: SIGKILL) + } + try? await self.deleteProcess(id: id, containerID: containerID) + try? await self.close() + } } public func up(name: String, mtu: UInt32? = nil) async throws { @@ -356,7 +455,10 @@ extension Vminitd: VirtualMachineAgent { $0.up = true if let mtu { $0.mtu = mtu } } - _ = try await client.ipLinkSet(request) + let operation = DeadlinePolicy.Operation.filesystem + _ = try await rpc(operation: operation, deadline: deadlinePolicy.deadline(for: operation)) { options in + try await client.ipLinkSet(request, options: options) + } } public func down(name: String) async throws { @@ -364,25 +466,32 @@ extension Vminitd: VirtualMachineAgent { $0.interface = name $0.up = false } - _ = try await client.ipLinkSet(request) + let operation = DeadlinePolicy.Operation.filesystem + _ = try await rpc(operation: operation, deadline: deadlinePolicy.deadline(for: operation)) { options in + try await client.ipLinkSet(request, options: options) + } } /// Get an environment variable from the sandbox's environment. public func getenv(key: String) async throws -> String { - let response = try await client.getenv( - .with { - $0.key = key - }) + let request = Com_Apple_Containerization_Sandbox_V3_GetenvRequest.with { $0.key = key } + let operation = DeadlinePolicy.Operation.filesystem + let response = try await rpc(operation: operation, deadline: deadlinePolicy.deadline(for: operation)) { options in + try await client.getenv(request, options: options) + } return response.value } /// Set an environment variable in the sandbox's environment. public func setenv(key: String, value: String) async throws { - _ = try await client.setenv( - .with { - $0.key = key - $0.value = value - }) + let request = Com_Apple_Containerization_Sandbox_V3_SetenvRequest.with { + $0.key = key + $0.value = value + } + let operation = DeadlinePolicy.Operation.filesystem + _ = try await rpc(operation: operation, deadline: deadlinePolicy.deadline(for: operation)) { options in + try await client.setenv(request, options: options) + } } } @@ -399,7 +508,10 @@ extension Vminitd { $0.mask = configuration.mask $0.flags = configuration.flags } - _ = try await client.setupEmulator(request) + let operation = DeadlinePolicy.Operation.filesystem + _ = try await rpc(operation: operation, deadline: deadlinePolicy.deadline(for: operation)) { options in + try await client.setupEmulator(request, options: options) + } } /// Sets the guest time. @@ -408,7 +520,10 @@ extension Vminitd { $0.sec = sec $0.usec = usec } - _ = try await client.setTime(request) + let operation = DeadlinePolicy.Operation.filesystem + _ = try await rpc(operation: operation, deadline: deadlinePolicy.deadline(for: operation)) { options in + try await client.setTime(request, options: options) + } } /// Set the provided sysctls inside the Sandbox's environment. @@ -416,19 +531,25 @@ extension Vminitd { let request = Com_Apple_Containerization_Sandbox_V3_SysctlRequest.with { $0.settings = settings } - _ = try await client.sysctl(request) + let operation = DeadlinePolicy.Operation.filesystem + _ = try await rpc(operation: operation, deadline: deadlinePolicy.deadline(for: operation)) { options in + try await client.sysctl(request, options: options) + } } /// Add an IP address to the sandbox's network interfaces. public func addressAdd(name: String, address: InterfaceAddress) async throws { - _ = try await client.ipAddrAdd( - .with { + let request = Com_Apple_Containerization_Sandbox_V3_IpAddrAddRequest.with { $0.interface = name $0.ipv4Address = address.ipv4Address.description if let ipv6Address = address.ipv6Address { $0.ipv6Address = ipv6Address.description } - }) + } + let operation = DeadlinePolicy.Operation.filesystem + _ = try await rpc(operation: operation, deadline: deadlinePolicy.deadline(for: operation)) { options in + try await client.ipAddrAdd(request, options: options) + } } /// Add a link-scoped route in the sandbox's environment, used to install an @@ -437,8 +558,7 @@ extension Vminitd { /// `route.ipv4Destination`/`route.ipv6Destination` carry the /// gateway address; the wire format is a CIDR string with the per-family host prefix appended. public func routeAddLink(name: String, route: LinkRoute) async throws { - _ = try await client.ipRouteAddLink( - .with { + let request = Com_Apple_Containerization_Sandbox_V3_IpRouteAddLinkRequest.with { $0.interface = name if let ipv4Destination = route.ipv4Destination { $0.dstIpv4Addr = "\(ipv4Destination.description)/32" @@ -452,26 +572,32 @@ extension Vminitd { if let ipv6Source = route.ipv6Source { $0.srcIpv6Addr = ipv6Source.description } - }) + } + let operation = DeadlinePolicy.Operation.filesystem + _ = try await rpc(operation: operation, deadline: deadlinePolicy.deadline(for: operation)) { options in + try await client.ipRouteAddLink(request, options: options) + } } /// Set the default route in the sandbox's environment. public func routeAddDefault(name: String, route: DefaultRoute) async throws { - _ = try await client.ipRouteAddDefault( - .with { + let request = Com_Apple_Containerization_Sandbox_V3_IpRouteAddDefaultRequest.with { $0.interface = name $0.ipv4Gateway = route.ipv4Gateway?.description ?? "" if let ipv6Gateway = route.ipv6Gateway { $0.ipv6Gateway = ipv6Gateway.description } - }) + } + let operation = DeadlinePolicy.Operation.filesystem + _ = try await rpc(operation: operation, deadline: deadlinePolicy.deadline(for: operation)) { options in + try await client.ipRouteAddDefault(request, options: options) + } } /// Configure DNS within the sandbox's environment. public func configureDNS(config: DNS, location: String) async throws { try config.validate() - _ = try await client.configureDns( - .with { + let request = Com_Apple_Containerization_Sandbox_V3_ConfigureDnsRequest.with { $0.location = location $0.nameservers = config.nameservers if let domain = config.domain { @@ -479,25 +605,39 @@ extension Vminitd { } $0.searchDomains = config.searchDomains $0.options = config.options - }) + } + let operation = DeadlinePolicy.Operation.filesystem + _ = try await rpc(operation: operation, deadline: deadlinePolicy.deadline(for: operation)) { options in + try await client.configureDns(request, options: options) + } } /// Configure /etc/hosts within the sandbox's environment. public func configureHosts(config: Hosts, location: String) async throws { - _ = try await client.configureHosts(config.toAgentHostsRequest(location: location)) + let operation = DeadlinePolicy.Operation.filesystem + let request = config.toAgentHostsRequest(location: location) + _ = try await rpc(operation: operation, deadline: deadlinePolicy.deadline(for: operation)) { options in + try await client.configureHosts(request, options: options) + } } /// Perform a sync call. public func sync() async throws { - _ = try await client.sync(.init()) + let operation = DeadlinePolicy.Operation.filesystem + _ = try await rpc(operation: operation, deadline: deadlinePolicy.deadline(for: operation)) { options in + try await client.sync(.init(), options: options) + } } public func kill(pid: Int32, signal: Int32) async throws -> Int32 { - let response = try await client.kill( - .with { - $0.pid = pid - $0.signal = signal - }) + let request = Com_Apple_Containerization_Sandbox_V3_KillRequest.with { + $0.pid = pid + $0.signal = signal + } + let operation = DeadlinePolicy.Operation.processControl + let response = try await rpc(operation: operation, deadline: deadlinePolicy.deadline(for: operation)) { options in + try await client.kill(request, options: options) + } return response.result } @@ -519,7 +659,10 @@ extension Vminitd { let response: Com_Apple_Containerization_Sandbox_V3_StatResponse do { - response = try await client.stat(request) + let operation = DeadlinePolicy.Operation.filesystem + response = try await rpc(operation: operation, deadline: deadlinePolicy.deadline(for: operation)) { options in + try await client.stat(request, options: options) + } } catch let error as RPCError where error.code == .notFound { throw ContainerizationError(.notFound, message: "stat: path not found '\(path.path)'", cause: error) } @@ -570,9 +713,12 @@ extension Vminitd { $0.isArchive = isArchive } - try await client.copy( - request, - onResponse: { stream in + let operation = DeadlinePolicy.Operation.filesystem + try await rpc(operation: operation, deadline: deadlinePolicy.deadline(for: operation)) { options in + try await client.copy( + request, + options: options, + onResponse: { stream in for try await response in stream.messages { if !response.error.isEmpty { throw ContainerizationError(.internalError, message: "copy: \(response.error)") @@ -587,6 +733,7 @@ extension Vminitd { } } }) + } } }