2026-06-21 20:22:21 -07:00
//===----------------------------------------------------------------------===//
// 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 ContainerizationExtras
import ContainerizationOCI
import ContainerizationOS
import Foundation
import Logging
import Synchronization
2026-07-13 19:15:11 -07:00
// [Nucleic vendored patch] stdio-connection diagnostics. Guarded: the `os` overlay isn't importable
// under every toolchain that builds this package (e.g. the swiftly Swift used by the vminit-image CI,
// which resolves Foundation/Virtualization but not `os`), so the diagnostic degrades to a no-op there
// rather than failing the build. Local (Xcode) builds keep it.
#if canImport ( os )
import os
#endif
2026-06-21 20:22:21 -07:00
/// `LinuxProcess` represents a Linux process and is used to
/// setup and control the full lifecycle for the process.
public final class LinuxProcess : Sendable {
2026-06-25 17:37:46 -07:00
/// [Nucleic vendored patch] Diagnostic log for stdio stream-connection failures (see `setupIO`).
2026-07-13 19:15:11 -07:00
#if canImport ( os )
2026-06-25 17:37:46 -07:00
static let nucleicIOLog = os . Logger ( subsystem : "com.nucleic" , category : "container-io" )
2026-07-13 19:15:11 -07:00
#endif
2026-06-25 17:37:46 -07:00
2026-06-21 20:22:21 -07:00
/// The ID of the process. This is purely metadata for the caller.
public let id : String
/// What container owns this process (if any).
public let owningContainer : String ?
package struct StdioSetup : Sendable {
let port : UInt32
let writer : Writer
}
package struct StdioReaderSetup {
let port : UInt32
let reader : ReaderStream
}
package struct Stdio : Sendable {
let stdin : StdioReaderSetup ?
let stdout : StdioSetup ?
let stderr : StdioSetup ?
}
private struct StdioHandles : Sendable {
var stdin : FileHandle ?
var stdout : FileHandle ?
var stderr : FileHandle ?
mutating func close () throws {
if let stdin {
try stdin . close ()
stdin . readabilityHandler = nil
self . stdin = nil
}
if let stdout {
try stdout . close ()
stdout . readabilityHandler = nil
self . stdout = nil
}
if let stderr {
try stderr . close ()
stderr . readabilityHandler = nil
self . stderr = nil
}
}
}
private struct State {
var spec : ContainerizationOCI . Spec
var pid : Int32
var stdio : StdioHandles
var stdinRelay : Task < (), Never >?
var ioTracker : IoTracker ?
var deletionTask : Task < Void , Error >?
struct IoTracker {
let stream : AsyncStream < Void >
let cont : AsyncStream < Void >. Continuation
let configuredStreams : Int
}
}
/// The process ID for the container process. This will be -1
/// if the process has not been started.
public var pid : Int32 {
state . withLock { $0 . pid }
}
private let state : Mutex < State >
private let ioSetup : Stdio
private let agent : any VirtualMachineAgent
private let vm : any VirtualMachineInstance
private let ociRuntimePath : String ?
2026-06-25 17:37:46 -07:00
private let logger : Logging . Logger ? // [Nucleic vendored patch] disambiguated from os.Logger
2026-06-21 20:22:21 -07:00
private let onDelete : (@ Sendable () async -> Void )?
init (
_ id : String ,
containerID : String ? = nil ,
spec : Spec ,
io : Stdio ,
ociRuntimePath : String ?,
agent : any VirtualMachineAgent ,
vm : any VirtualMachineInstance ,
2026-06-25 17:37:46 -07:00
logger : Logging . Logger ?, // [Nucleic vendored patch] disambiguated from os.Logger
2026-06-21 20:22:21 -07:00
onDelete : (@ Sendable () async -> Void )? = nil
) {
self . id = id
self . owningContainer = containerID
self . state = Mutex < State >(. init ( spec : spec , pid : - 1 , stdio : StdioHandles ()))
self . ioSetup = io
self . agent = agent
self . ociRuntimePath = ociRuntimePath
self . vm = vm
self . logger = logger
self . onDelete = onDelete
}
}
extension LinuxProcess {
2026-07-13 18:46:48 -07:00
/// [Nucleic vendored patch] Put a connected stdio FileHandle's fd into non-blocking mode so the
/// relay's reads (``nucleicDrainNonBlocking``) can never park the shared readability queue. No-op
/// if the handle is nil.
static func nucleicSetNonBlocking ( _ handle : FileHandle ?) {
guard let fd = handle ?. fileDescriptor else { return }
let flags = fcntl ( fd , F_GETFL , 0 )
if flags >= 0 { _ = fcntl ( fd , F_SETFL , flags | O_NONBLOCK ) }
}
/// [Nucleic vendored patch] Drain `fd` (already O_NONBLOCK) without ever blocking. Returns the
/// bytes read this pass plus whether the stream hit EOF (or a hard error). On EAGAIN it returns
/// what it has with `eof == false`; the readability `DispatchSource` fires again when more data
/// arrives. Upstream read with `FileHandle.availableData`, a *blocking* read: if one exec's guest
/// stdout wedged mid-stream, that read parked Foundation's shared readability thread and
/// head-of-line-blocked EVERY other exec's stdout/stderr relay (the "one stuck session freezes the
/// others" failure). A non-blocking drain can never park that thread, so a wedged stream is
/// contained to its own exec.
static func nucleicDrainNonBlocking ( _ fd : Int32 ) -> ( data : Data , eof : Bool ) {
var out = Data ()
var buf = [ UInt8 ]( repeating : 0 , count : 64 * 1024 )
while true {
let n = buf . withUnsafeMutableBytes { read ( fd , $0 . baseAddress , $0 . count ) }
if n > 0 {
out . append ( contentsOf : buf [ 0. .< n ])
} else if n == 0 {
return ( out , true ) // EOF: guest closed the write side
} else if errno == EINTR {
continue
} else if errno == EAGAIN || errno == EWOULDBLOCK {
return ( out , false ) // drained for now; not EOF
} else {
return ( out , true ) // hard error → treat as EOF so the relay finishes
}
}
}
2026-06-21 20:22:21 -07:00
func setupIO ( listeners : [ VsockListener ?]) async throws -> [ FileHandle ?] {
let handles = try await Timeout . run ( seconds : 3 ) {
try await withThrowingTaskGroup ( of : ( Int , FileHandle ?). self ) { group in
var results = [ FileHandle ?]( repeating : nil , count : 3 )
for ( index , listener ) in listeners . enumerated () {
guard let listener else { continue }
group . addTask {
let first = await listener . first ( where : { _ in true })
try listener . finish ()
return ( index , first )
}
}
for try await ( index , fileHandle ) in group {
results [ index ] = fileHandle
}
return results
}
}
2026-06-25 17:37:46 -07:00
// [Nucleic vendored patch] Diagnostics: a configured stdio stream whose guest side never
// connected leaves its host FileHandle `nil`, so the relay / readability handler below is
// never wired — the agent's stdin is then never delivered (it hangs waiting for input) or
// its stdout is never read ("no output, just a spinner"). Log that specific failure (Console
// / `log show`, subsystem com.nucleic, category container-io) so a stall pinpoints the stream
2026-07-13 19:15:11 -07:00
// instead of proceeding silently. Log-only; behavior is unchanged. Guarded on `canImport(os)`
// (see the import) so a toolchain without the `os` overlay still builds.
#if canImport ( os )
2026-06-25 17:37:46 -07:00
let configured = [ self . ioSetup . stdin != nil , self . ioSetup . stdout != nil , self . ioSetup . stderr != nil ]
for ( index , label ) in [( 0 , "stdin" ), ( 1 , "stdout" ), ( 2 , "stderr" )] where configured [ index ] && handles [ index ] == nil {
Self . nucleicIOLog . error (
"setupIO[ \( self . id , privacy : . public ) ]: \( label , privacy : . public ) stream never connected from the guest — agent stdio will stall" )
}
2026-07-13 19:15:11 -07:00
#endif
2026-06-25 17:37:46 -07:00
2026-06-21 20:22:21 -07:00
// Note: stdin relay is started separately via startStdinRelay() after
// the process has started, to avoid a deadlock where closeStdin is
// called before the process is consuming from the pipe.
var configuredStreams = 0
let ( stream , cc ) = AsyncStream < Void >. makeStream ()
if let stdout = self . ioSetup . stdout {
configuredStreams += 1
2026-07-13 18:46:48 -07:00
// [Nucleic vendored patch] Non-blocking relay (see nucleicDrainNonBlocking): mark the
// connected fd O_NONBLOCK and drain it without a blocking read, so a wedged guest stdout
// can't head-of-line-block sibling execs' relays on Foundation's shared readability queue.
Self . nucleicSetNonBlocking ( handles [ 1 ])
2026-06-21 20:22:21 -07:00
handles [ 1 ]?. readabilityHandler = { handle in
2026-07-13 18:46:48 -07:00
let ( data , eof ) = Self . nucleicDrainNonBlocking ( handle . fileDescriptor )
if ! data . isEmpty {
do {
try stdout . writer . write ( data )
} catch {
self . logger ?. error ( "failed to write to stdout: \( error ) " )
2026-06-21 20:22:21 -07:00
}
2026-07-13 18:46:48 -07:00
}
if eof {
// The guest closed the fd it was writing into.
handles [ 1 ]?. readabilityHandler = nil
cc . yield ()
2026-06-21 20:22:21 -07:00
}
}
}
if let stderr = self . ioSetup . stderr {
configuredStreams += 1
2026-07-13 18:46:48 -07:00
// [Nucleic vendored patch] Non-blocking relay — same rationale as stdout above.
Self . nucleicSetNonBlocking ( handles [ 2 ])
2026-06-21 20:22:21 -07:00
handles [ 2 ]?. readabilityHandler = { handle in
2026-07-13 18:46:48 -07:00
let ( data , eof ) = Self . nucleicDrainNonBlocking ( handle . fileDescriptor )
if ! data . isEmpty {
do {
try stderr . writer . write ( data )
} catch {
self . logger ?. error ( "failed to write to stderr: \( error ) " )
2026-06-21 20:22:21 -07:00
}
2026-07-13 18:46:48 -07:00
}
if eof {
handles [ 2 ]?. readabilityHandler = nil
cc . yield ()
2026-06-21 20:22:21 -07:00
}
}
}
if configuredStreams > 0 {
self . state . withLock {
$0 . ioTracker = . init ( stream : stream , cont : cc , configuredStreams : configuredStreams )
}
}
return handles
}
func startStdinRelay ( handle : FileHandle ) {
guard let stdin = self . ioSetup . stdin else { return }
2026-07-28 15:35:35 -07:00
// [Nucleic vendored patch] The stdin fd is BLOCKING (only the read fds get
// O_NONBLOCK), and `FileHandle.write` parks its thread for as long as the guest
// isn't reading once the vsock buffer fills — with a large prompt line that can be
// forever, and the old direct call parked a *width-limited Swift cooperative-pool
// thread*, non-cancellably. A few wedged sessions starved the entire concurrency
// runtime: every exec's decode loop, first-output watchdogs, all of it — an
// app-wide stall. Offload each write to a per-process GCD queue so the pool thread
// suspends instead; a wedged write now costs one expendable GCD thread, and
// process deletion closing the fd still unwedges it.
let writeQueue = DispatchQueue ( label : "com.nucleic.stdin-relay" )
// The handle is used from one queue at a time (writes are serialized on writeQueue;
// the close paths run after the relay ends or via _closeStdin's own lock).
nonisolated ( unsafe ) let handle = handle
2026-06-21 20:22:21 -07:00
self . state . withLock {
$0 . stdinRelay = Task {
for await data in stdin . reader . stream () {
do {
2026-07-28 15:35:35 -07:00
try await withCheckedThrowingContinuation { ( c : CheckedContinuation < Void , Error >) in
writeQueue . async {
do {
try handle . write ( contentsOf : data )
c . resume ()
} catch {
c . resume ( throwing : error )
}
}
}
2026-06-21 20:22:21 -07:00
} catch {
self . logger ?. error ( "failed to write to stdin: \( error ) " )
break
}
}
do {
self . logger ?. debug ( "stdin relay finished, closing" )
// There's two ways we can wind up here:
//
// 1. The stream finished on its own (e.g. we wrote all the
// data) and we will close the underlying stdin in the guest below.
//
// 2. The client explicitly called closeStdin() themselves
// which will cancel this relay task AFTER actually closing
// the fds. If the client did that, then this task will be
// cancelled, and the fds are already gone so there's nothing
// for us to do.
if Task . isCancelled {
return
}
try await self . _closeStdin ()
} catch {
self . logger ?. error ( "failed to close stdin: \( error ) " )
}
}
}
}
/// Start the process.
public func start () async throws {
do {
let spec = self . state . withLock { $0 . spec }
var listeners = [ VsockListener ?]( repeating : nil , count : 3 )
if let stdin = self . ioSetup . stdin {
listeners [ 0 ] = try self . vm . listen ( stdin . port )
}
if let stdout = self . ioSetup . stdout {
listeners [ 1 ] = try self . vm . listen ( stdout . port )
}
if let stderr = self . ioSetup . stderr {
if spec . process !. terminal {
throw ContainerizationError (
. invalidArgument ,
message : "stderr should not be configured with terminal=true"
)
}
listeners [ 2 ] = try self . vm . listen ( stderr . port )
}
let t = Task {
try await self . setupIO ( listeners : listeners )
}
try await agent . createProcess (
id : self . id ,
containerID : self . owningContainer ,
stdinPort : self . ioSetup . stdin ?. port ,
stdoutPort : self . ioSetup . stdout ?. port ,
stderrPort : self . ioSetup . stderr ?. port ,
ociRuntimePath : self . ociRuntimePath ,
configuration : spec ,
options : nil
)
let result = try await t . value
2026-07-13 18:46:48 -07:00
// [Nucleic vendored patch] Atomic stdio-or-abort. If a *configured* stdio stream never
// connected from the guest (its FileHandle came back nil — the failure logged in setupIO),
// starting the process would run it with a dead stream: stdin never delivered (it hangs)
// or stdout/stderr never read ("no output, just a spinner" — the 60s stall in Nucleic
// Control). Rather than launch a black-hole process, tear the just-created exec back down
// and fail fast so the caller gets a clean, retryable start error instead of an eternal
// silent stall the watchdog has to guess at.
let configured = [
self . ioSetup . stdin != nil , self . ioSetup . stdout != nil , self . ioSetup . stderr != nil ,
]
let streamLabels = [ "stdin" , "stdout" , "stderr" ]
if let missing = ( 0. .< 3 ). first ( where : { configured [ $0 ] && result [ $0 ] == nil }) {
try ? await self . agent . deleteProcess ( id : self . id , containerID : self . owningContainer )
throw ContainerizationError (
. internalError ,
message :
"process \( self . id ) : \( streamLabels [ missing ] ) stream never connected from the guest before start; aborting so the stdio transport stall surfaces as a retryable start error"
)
}
2026-06-21 20:22:21 -07:00
let pid = try await self . agent . startProcess (
id : self . id ,
containerID : self . owningContainer
)
// Start stdin relay after process launch to avoid filling the pipe
// buffer before the process is even running.
if let stdinHandle = result [ 0 ] {
self . startStdinRelay ( handle : stdinHandle )
}
self . state . withLock {
$0 . stdio = StdioHandles (
stdin : result [ 0 ],
stdout : result [ 1 ],
stderr : result [ 2 ]
)
$0 . pid = pid
}
} catch {
if let err = error as ? ContainerizationError {
throw err
}
throw ContainerizationError (
. internalError ,
message : "failed to start process" ,
cause : error ,
)
}
}
/// Kill the process with the specified signal.
public func kill ( _ signal : Signal ) async throws {
do {
try await agent . signalProcess (
id : self . id ,
containerID : self . owningContainer ,
signal : signal . rawValue
)
} catch {
throw ContainerizationError (
. internalError ,
message : "failed to kill process" ,
cause : error
)
}
}
2026-06-25 13:55:31 -07:00
/// [Nucleic vendored patch] Deliver a signal to the whole process GROUP led by this exec'd
/// process, not just the leader. `vmexec` `setsid()`s every exec, so the process is its own
/// session/group leader and its pgid equals its pid; a negative pid makes the guest's `kill(2)`
/// target the entire group, reaching any children the agent forked (model/turn subprocesses,
/// tool shells). `kill(_:)` above signals only the leader, so a wedged child can survive a Stop
/// in a long-lived shared container — this is the group-wide counterpart. Best-effort and
/// guarded against pid ≤ 1 (a non-positive pid would target the caller's group / every process).
public func killProcessGroup ( _ signal : Signal ) async throws {
let leader = self . pid
guard leader > 1 else { return }
do {
_ = try await agent . kill ( pid : - leader , signal : signal . rawValue )
} catch {
throw ContainerizationError (
. internalError ,
message : "failed to kill process group" ,
cause : error
)
}
}
2026-06-21 20:22:21 -07:00
/// Resize the processes pty (if requested).
public func resize ( to : Terminal . Size ) async throws {
do {
try await agent . resizeProcess (
id : self . id ,
containerID : self . owningContainer ,
columns : UInt32 ( to . width ),
rows : UInt32 ( to . height )
)
} catch {
throw ContainerizationError (
. internalError ,
message : "failed to resize process" ,
cause : error
)
}
}
public func closeStdin () async throws {
do {
try await self . _closeStdin ()
self . state . withLock {
$0 . stdinRelay ?. cancel ()
}
} catch {
throw ContainerizationError (
. internalError ,
message : "failed to close stdin" ,
cause : error ,
)
}
}
func _closeStdin () async throws {
try await self . agent . closeProcessStdin (
id : self . id ,
containerID : self . owningContainer
)
}
/// Wait on the process to exit with an optional timeout. Returns the exit code of the process.
@ discardableResult
public func wait ( timeoutInSeconds : Int64 ? = nil ) async throws -> ExitStatus {
do {
let exitStatus = try await self . agent . waitProcess (
id : self . id ,
containerID : self . owningContainer ,
timeoutInSeconds : timeoutInSeconds
)
await self . waitIoComplete ()
return exitStatus
} catch {
if error is ContainerizationError {
throw error
}
throw ContainerizationError (
. internalError ,
message : "failed to wait on process" ,
cause : error
)
}
}
/// Wait until the standard output and standard error streams for the process have concluded.
private func waitIoComplete () async {
let ioTracker = self . state . withLock { $0 . ioTracker }
guard let ioTracker else {
return
}
do {
try await Timeout . run ( seconds : 3 ) {
var counter = ioTracker . configuredStreams
for await _ in ioTracker . stream {
counter -= 1
if counter == 0 {
ioTracker . cont . finish ()
break
}
}
}
} catch {
self . logger ?. error ( "timeout waiting for IO to complete for process \( id ) : \( error ) " )
}
self . state . withLock {
$0 . ioTracker = nil
}
}
/// Cleans up guest state and waits on and closes any host resources (stdio handles).
public func delete () async throws {
try await self . _delete ()
await self . onDelete ?()
}
func _delete () async throws {
let task = self . state . withLock { state in
if let existingTask = state . deletionTask {
// Deletion already in progress or finished.
return existingTask
}
let task = Task < Void , Error > {
try await self . performDeletion ()
}
state . deletionTask = task
return task
}
try await task . value
}
private func performDeletion () async throws {
do {
try await self . agent . deleteProcess (
id : self . id ,
containerID : self . owningContainer
)
} catch {
self . state . withLock {
$0 . stdinRelay ?. cancel ()
try ? $0 . stdio . close ()
}
try ? await self . agent . close ()
throw ContainerizationError (
. internalError ,
message : "failed to delete process" ,
cause : error ,
)
}
do {
try self . state . withLock {
$0 . stdinRelay ?. cancel ()
try $0 . stdio . close ()
}
} catch {
try ? await self . agent . close ()
throw ContainerizationError (
. internalError ,
message : "failed to close stdio" ,
cause : error ,
)
}
do {
try await self . agent . close ()
} catch {
throw ContainerizationError (
. internalError ,
message : "failed to close agent connection" ,
cause : error ,
)
}
}
}