Bulk up docc for ContainerSandboxService. (#32)

This commit is contained in:
J Logan
2025-06-06 16:36:49 -07:00
committed by GitHub
parent 7c51835c1b
commit 6fe813a9da
3 changed files with 357 additions and 244 deletions
@@ -21,11 +21,18 @@ import ContainerizationExtras
import Foundation
import Logging
/// Track when a long running method exits, and notify the caller via a callback.
/// Track when long running work exits, and notify the caller via a callback.
public actor ExitMonitor {
/// A callback that receives the client identifier and exit code.
public typealias ExitCallback = @Sendable (String, Int32) async throws -> Void
/// A function that waits for work to complete, returning an exit code.
public typealias WaitHandler = @Sendable () async throws -> Int32
/// Create a new monitor.
///
/// - Parameters:
/// - log: The destination for log messages.
public init(log: Logger? = nil) {
self.log = log
}
@@ -34,6 +41,10 @@ public actor ExitMonitor {
private var runningTasks: [String: Task<Void, Never>] = [:]
private let log: Logger?
/// Remove tracked work from the monitor.
///
/// - Parameters:
/// - id: The client identifier for the tracked work.
public func stopTracking(id: String) async {
if let task = self.runningTasks[id] {
task.cancel()
@@ -42,6 +53,12 @@ public actor ExitMonitor {
runningTasks.removeValue(forKey: id)
}
/// Register long running work so that the monitor invokes
/// a callback when the work completes.
///
/// - Parameters:
/// - id: The client identifier for the work.
/// - onExit: The callback to invoke when the work completes.
public func registerProcess(id: String, onExit: @escaping ExitCallback) async throws {
guard self.exitCallbacks[id] == nil else {
throw ContainerizationError(.invalidState, message: "ExitMonitor already setup for process \(id)")
@@ -49,6 +66,12 @@ public actor ExitMonitor {
self.exitCallbacks[id] = onExit
}
/// Await the completion of previously registered item of work.
///
/// - Parameters:
/// - id: The client identifier for the work.
/// - waitingOn: A function that waits for the work to complete,
/// and then returns an exit code.
public func track(id: String, waitingOn: @escaping WaitHandler) async throws {
guard let onExit = self.exitCallbacks[id] else {
throw ContainerizationError(.invalidState, message: "ExitMonitor not setup for process \(id)")
@@ -18,8 +18,16 @@ import ContainerNetworkService
import ContainerXPC
import Containerization
/// Customized interface creation strategy.
/// A strategy for mapping network attachment information to a network interface.
public protocol InterfaceStrategy: Sendable {
/// Map a client network attachment request to a network interface specification.
///
/// - Parameters:
/// - attachment: General attachment information that is common
/// for all networks.
/// - additionalData: If present, attachment information that is
/// specific for the network to which the container will attach.
///
/// - Returns: An XPC message with no parameters.
func toInterface(attachment: Attachment, additionalData: XPCMessage?) throws -> Interface
}
@@ -30,6 +30,7 @@ import Logging
import struct ContainerizationOCI.Mount
import struct ContainerizationOCI.Process
/// An XPC service that manages the lifecycle of a single VM-backed container.
public actor SandboxService {
private let root: URL
private let interfaceStrategy: InterfaceStrategy
@@ -41,6 +42,13 @@ public actor SandboxService {
private var state: State = .created
private var processes: [String: ProcessInfo] = [:]
/// Create an instance with a bundle that describes the container.
///
/// - Parameters:
/// - root: The file URL for the bundle root.
/// - interfaceStrategy: The strategy for producing network interface
/// objects for each network to which the container attaches.
/// - log: The destination for log messages.
public init(root: URL, interfaceStrategy: InterfaceStrategy, log: Logger) {
self.root = root
self.interfaceStrategy = interfaceStrategy
@@ -48,6 +56,12 @@ public actor SandboxService {
self.monitor = ExitMonitor(log: log)
}
/// Start the VM and the guest agent process for a container.
///
/// - Parameters:
/// - message: An XPC message with no parameters.
///
/// - Returns: An XPC message with no parameters.
@Sendable
public func bootstrap(_ message: XPCMessage) async throws -> XPCMessage {
self.log.info("`bootstrap` xpc handler")
@@ -129,6 +143,14 @@ public actor SandboxService {
}
}
/// Start the container workload inside the virtual machine.
///
/// - Parameters:
/// - message: An XPC message with the following parameters:
/// - id: A client identifier for the process.
/// - stdio: An array of file handles for standard input, output, and error.
///
/// - Returns: An XPC message with no parameters.
@Sendable
public func startProcess(_ message: XPCMessage) async throws -> XPCMessage {
self.log.info("`start` xpc handler")
@@ -216,6 +238,305 @@ public actor SandboxService {
}
}
/// Create a process inside the virtual machine for the container.
///
/// Use this procedure to run ad hoc processes in the virtual
/// machine (`container exec`).
///
/// - Parameters:
/// - message: An XPC message with the following parameters:
/// - id: A client identifier for the process.
/// - processConfig: JSON serialization of the `ProcessConfiguration`
/// containing the process attributes.
///
/// - Returns: An XPC message with no parameters.
@Sendable
public func createProcess(_ message: XPCMessage) async throws -> XPCMessage {
log.info("`createProcess` xpc handler")
return try await self.lock.withLock { [self] _ in
switch await self.state {
case .created, .stopped(_), .starting, .stopping:
throw ContainerizationError(
.invalidState,
message: "cannot exec: container is not running"
)
case .running, .booted:
let id = try message.id()
let config = try message.processConfig()
await self.addNewProcess(id, config)
try await self.monitor.registerProcess(
id: id,
onExit: { id, code in
guard await self.processes[id] != nil else {
throw ContainerizationError(.invalidState, message: "ProcessInfo missing for process \(id)")
}
for cc in await self.waiters[id] ?? [] {
cc.resume(returning: code)
}
await self.removeWaiters(for: id)
try await self.setProcessState(id: id, state: .stopped(code))
})
return message.reply()
}
}
}
/// Return the state for the sandbox and its containers.
///
/// - Parameters:
/// - message: An XPC message with no parameters.
///
/// - Returns: An XPC message with the following parameters:
/// - snapshot: The JSON serialization of the `SandboxSnapshot`
/// that contains the state information.
@Sendable
public func state(_ message: XPCMessage) async throws -> XPCMessage {
self.log.info("`state` xpc handler")
var status: RuntimeStatus = .unknown
var networks: [Attachment] = []
var cs: ContainerSnapshot?
switch state {
case .created, .stopped(_), .starting, .booted, .stopping:
status = .stopped
case .running:
let ctr = try getContainer()
status = .running
networks = ctr.attachments
cs = ContainerSnapshot(
configuration: ctr.config,
status: RuntimeStatus.running,
networks: networks
)
}
let reply = message.reply()
try reply.setState(
.init(
status: status,
networks: networks,
containers: cs != nil ? [cs!] : []
)
)
return reply
}
/// Stop the container workload, any ad hoc processes, and the underlying
/// virtual machine.
///
/// - Parameters:
/// - message: An XPC message with the following parameters:
/// - stopOptions: JSON serialization of `ContainerStopOptions`
/// that modify stop behavior.
///
/// - Returns: An XPC message with no parameters.
@Sendable
public func stop(_ message: XPCMessage) async throws -> XPCMessage {
self.log.info("`stop` xpc handler")
let reply = try await self.lock.withLock { [self] _ in
switch await self.state {
case .stopped(_), .created, .stopping:
return message.reply()
case .starting:
throw ContainerizationError(
.invalidState,
message: "cannot stop: container is not running"
)
case .running, .booted:
let ctr = try await getContainer()
let stopOptions = try message.stopOptions()
do {
try await gracefulStopContainer(
ctr.container,
stopOpts: stopOptions
)
} catch {
log.notice("failed to stop sandbox gracefully: \(error)")
}
await setState(.stopping)
return message.reply()
}
}
do {
try await cleanupContainer()
} catch {
self.log.error("failed to cleanup container: \(error)")
}
return reply
}
/// Signal a process running in the virtual machine.
///
/// - Parameters:
/// - message: An XPC message with the following parameters:
/// - id: The process identifier.
/// - signal: The signal value.
///
/// - Returns: An XPC message with no parameters.
@Sendable
public func kill(_ message: XPCMessage) async throws -> XPCMessage {
self.log.info("`kill` xpc handler")
return try await self.lock.withLock { [self] _ in
switch await self.state {
case .created, .stopped, .starting, .booted, .stopping:
throw ContainerizationError(
.invalidState,
message: "cannot kill: container is not running"
)
case .running:
let ctr = try await getContainer()
let id = try message.id()
if id != ctr.container.id {
guard let processInfo = await self.processes[id] else {
throw ContainerizationError(.invalidState, message: "Process \(id) does not exist")
}
guard let proc = processInfo.process else {
throw ContainerizationError(.invalidState, message: "Process \(id) not started")
}
try await proc.kill(Int32(try message.signal()))
return message.reply()
}
// TODO: fix underying signal value to int64
try await ctr.container.kill(Int32(try message.signal()))
return message.reply()
}
}
}
/// Resize the terminal for a process.
///
/// - Parameters:
/// - message: An XPC message with the following parameters:
/// - id: The process identifier.
/// - width: The terminal width.
/// - height: The terminal height.
///
/// - Returns: An XPC message with no parameters.
@Sendable
public func resize(_ message: XPCMessage) async throws -> XPCMessage {
self.log.info("`resize` xpc handler")
return try await self.lock.withLock { [self] _ in
switch await self.state {
case .created, .stopped, .starting, .booted, .stopping:
throw ContainerizationError(
.invalidState,
message: "cannot resize: container is not running"
)
case .running:
let id = try message.id()
let ctr = try await getContainer()
let width = message.uint64(key: .width)
let height = message.uint64(key: .height)
if id != ctr.container.id {
guard let processInfo = await self.processes[id] else {
throw ContainerizationError(.invalidState, message: "Process \(id) does not exist")
}
guard let proc = processInfo.process else {
throw ContainerizationError(.invalidState, message: "Process \(id) not started")
}
try await proc.resize(to: .init(width: UInt16(width), height: UInt16(height)))
return message.reply()
}
try await ctr.container.resize(to: .init(width: UInt16(width), height: UInt16(height)))
return message.reply()
}
}
}
/// Wait for a process.
///
/// - Parameters:
/// - message: An XPC message with the following parameters:
/// - id: The process identifier.
///
/// - Returns: An XPC message with the following parameters:
/// - exitCode: The exit code for the process.
@Sendable
public func wait(_ message: XPCMessage) async throws -> XPCMessage {
self.log.info("`wait` xpc handler")
guard let id = message.string(key: .id) else {
throw ContainerizationError(.invalidArgument, message: "Missing id in wait xpc message")
}
let cachedCode: Int32? = try await self.lock.withLock { _ in
let ctrInfo = try await self.getContainer()
let ctr = ctrInfo.container
if id == ctr.id {
switch await self.state {
case .stopped(let code):
return code
default:
break
}
} else {
guard let processInfo = await self.processes[id] else {
throw ContainerizationError(.notFound, message: "Process with id \(id)")
}
switch processInfo.state {
case .stopped(let code):
return code
default:
break
}
}
return nil
}
if let cachedCode {
let reply = message.reply()
reply.set(key: .exitCode, value: Int64(cachedCode))
return reply
}
let exitCode = await withCheckedContinuation { cc in
// Is this safe since we are in an actor? :(
self.addWaiter(id: id, cont: cc)
}
let reply = message.reply()
reply.set(key: .exitCode, value: Int64(exitCode))
return reply
}
/// Dial a vsock port on the virtual machine.
///
/// - Parameters:
/// - message: An XPC message with the following parameters:
/// - port: The port number.
///
/// - Returns: An XPC message with the following parameters:
/// - fd: The file descriptor for the vsock.
@Sendable
public func dial(_ message: XPCMessage) async throws -> XPCMessage {
self.log.info("`dial` xpc handler")
switch self.state {
case .starting, .created, .stopped, .stopping:
throw ContainerizationError(
.invalidState,
message: "cannot dial: container is not running"
)
case .running, .booted:
let port = message.uint64(key: .port)
guard port > 0 else {
throw ContainerizationError(
.invalidArgument,
message: "no vsock port supplied for dial"
)
}
let ctr = try getContainer()
let fh = try await ctr.container.dialVsock(port: UInt32(port))
let reply = message.reply()
reply.set(key: .fd, value: fh)
return reply
}
}
private func onContainerExit(id: String, code: Int32) async throws {
self.log.info("init process exited with: \(code)")
@@ -351,246 +672,6 @@ public actor SandboxService {
return ["TERM=xterm"] + config.environment
}
@Sendable
public func createProcess(_ message: XPCMessage) async throws -> XPCMessage {
log.info("`createProcess` xpc handler")
return try await self.lock.withLock { [self] _ in
switch await self.state {
case .created, .stopped(_), .starting, .stopping:
throw ContainerizationError(
.invalidState,
message: "cannot exec: container is not running"
)
case .running, .booted:
let id = try message.id()
let config = try message.processConfig()
await self.addNewProcess(id, config)
try await self.monitor.registerProcess(
id: id,
onExit: { id, code in
guard await self.processes[id] != nil else {
throw ContainerizationError(.invalidState, message: "ProcessInfo missing for process \(id)")
}
for cc in await self.waiters[id] ?? [] {
cc.resume(returning: code)
}
await self.removeWaiters(for: id)
try await self.setProcessState(id: id, state: .stopped(code))
})
return message.reply()
}
}
}
/// Return the state for the sandbox and its containers.
@Sendable
public func state(_ message: XPCMessage) async throws -> XPCMessage {
self.log.info("`state` xpc handler")
var status: RuntimeStatus = .unknown
var networks: [Attachment] = []
var cs: ContainerSnapshot?
switch state {
case .created, .stopped(_), .starting, .booted, .stopping:
status = .stopped
case .running:
let ctr = try getContainer()
status = .running
networks = ctr.attachments
cs = ContainerSnapshot(
configuration: ctr.config,
status: RuntimeStatus.running,
networks: networks
)
}
let reply = message.reply()
try reply.setState(
.init(
status: status,
networks: networks,
containers: cs != nil ? [cs!] : []
)
)
return reply
}
/// Stop all containers inside the sandbox, aborting any processes currently
/// executing inside the container, before stopping the underlying sandbox.
@Sendable
public func stop(_ message: XPCMessage) async throws -> XPCMessage {
self.log.info("`stop` xpc handler")
let reply = try await self.lock.withLock { [self] _ in
switch await self.state {
case .stopped(_), .created, .stopping:
return message.reply()
case .starting:
throw ContainerizationError(
.invalidState,
message: "cannot stop: container is not running"
)
case .running, .booted:
let ctr = try await getContainer()
let stopOptions = try message.stopOptions()
do {
try await gracefulStopContainer(
ctr.container,
stopOpts: stopOptions
)
} catch {
log.notice("failed to stop sandbox gracefully: \(error)")
}
await setState(.stopping)
return message.reply()
}
}
do {
try await cleanupContainer()
} catch {
self.log.error("failed to cleanup container: \(error)")
}
return reply
}
@Sendable
public func kill(_ message: XPCMessage) async throws -> XPCMessage {
self.log.info("`kill` xpc handler")
return try await self.lock.withLock { [self] _ in
switch await self.state {
case .created, .stopped, .starting, .booted, .stopping:
throw ContainerizationError(
.invalidState,
message: "cannot kill: container is not running"
)
case .running:
let ctr = try await getContainer()
let id = try message.id()
if id != ctr.container.id {
guard let processInfo = await self.processes[id] else {
throw ContainerizationError(.invalidState, message: "Process \(id) does not exist")
}
guard let proc = processInfo.process else {
throw ContainerizationError(.invalidState, message: "Process \(id) not started")
}
try await proc.kill(Int32(try message.signal()))
return message.reply()
}
// TODO: fix underying signal value to int64
try await ctr.container.kill(Int32(try message.signal()))
return message.reply()
}
}
}
@Sendable
public func resize(_ message: XPCMessage) async throws -> XPCMessage {
self.log.info("`resize` xpc handler")
return try await self.lock.withLock { [self] _ in
switch await self.state {
case .created, .stopped, .starting, .booted, .stopping:
throw ContainerizationError(
.invalidState,
message: "cannot resize: container is not running"
)
case .running:
let id = try message.id()
let ctr = try await getContainer()
let width = message.uint64(key: .width)
let height = message.uint64(key: .height)
if id != ctr.container.id {
guard let processInfo = await self.processes[id] else {
throw ContainerizationError(.invalidState, message: "Process \(id) does not exist")
}
guard let proc = processInfo.process else {
throw ContainerizationError(.invalidState, message: "Process \(id) not started")
}
try await proc.resize(to: .init(width: UInt16(width), height: UInt16(height)))
return message.reply()
}
try await ctr.container.resize(to: .init(width: UInt16(width), height: UInt16(height)))
return message.reply()
}
}
}
@Sendable
public func wait(_ message: XPCMessage) async throws -> XPCMessage {
self.log.info("`wait` xpc handler")
guard let id = message.string(key: .id) else {
throw ContainerizationError(.invalidArgument, message: "Missing id in wait xpc message")
}
let cachedCode: Int32? = try await self.lock.withLock { _ in
let ctrInfo = try await self.getContainer()
let ctr = ctrInfo.container
if id == ctr.id {
switch await self.state {
case .stopped(let code):
return code
default:
break
}
} else {
guard let processInfo = await self.processes[id] else {
throw ContainerizationError(.notFound, message: "Process with id \(id)")
}
switch processInfo.state {
case .stopped(let code):
return code
default:
break
}
}
return nil
}
if let cachedCode {
let reply = message.reply()
reply.set(key: .exitCode, value: Int64(cachedCode))
return reply
}
let exitCode = await withCheckedContinuation { cc in
// Is this safe since we are in an actor? :(
self.addWaiter(id: id, cont: cc)
}
let reply = message.reply()
reply.set(key: .exitCode, value: Int64(exitCode))
return reply
}
@Sendable
public func dial(_ message: XPCMessage) async throws -> XPCMessage {
self.log.info("`dial` xpc handler")
switch self.state {
case .starting, .created, .stopped, .stopping:
throw ContainerizationError(
.invalidState,
message: "cannot dial: container is not running"
)
case .running, .booted:
let port = message.uint64(key: .port)
guard port > 0 else {
throw ContainerizationError(
.invalidArgument,
message: "no vsock port supplied for dial"
)
}
let ctr = try getContainer()
let fh = try await ctr.container.dialVsock(port: UInt32(port))
let reply = message.reply()
reply.set(key: .fd, value: fh)
return reply
}
}
private func getContainer() throws -> ContainerInfo {
guard let container else {
throw ContainerizationError(
@@ -601,7 +682,7 @@ public actor SandboxService {
return container
}
func gracefulStopContainer(_ lc: LinuxContainer, stopOpts: ContainerStopOptions) async throws {
private func gracefulStopContainer(_ lc: LinuxContainer, stopOpts: ContainerStopOptions) async throws {
// Try and gracefully shut down the process. Even if this succeeds we need to power off
// the vm, but we should try this first always.
do {
@@ -622,7 +703,7 @@ public actor SandboxService {
try await lc.stop()
}
func cleanupContainer() async throws {
private func cleanupContainer() async throws {
// Give back our lovely IP(s)
let containerInfo = try self.getContainer()
for attachment in containerInfo.attachments {
@@ -691,6 +772,7 @@ extension XPCMessage {
}
extension ContainerClient.Bundle {
/// The pathname for the workload log file.
public var containerLog: URL {
path.appendingPathComponent("stdio.log")
}