Files
abkslmandClaude Fable 5 17e84c9573 Deep-sweep fixes: cooperative-pool starvation, stdio tail loss, epoll integrity, UI pins
Host — the two remaining app-wide stall mechanisms plus main-thread pins
found by mining all nine hang reports:
- LinuxProcess.startStdinRelay wrote to a BLOCKING stdin fd on a
  width-limited cooperative-pool thread, non-cancellably; wedged guests
  starved the whole concurrency runtime (decode loops, watchdogs — an
  app-wide freeze surviving the reconcile fix). Writes now offload to a
  per-process GCD queue (vendored patch #18).
- TranscriptWriter (actor) did blocking write/fsync on the cooperative
  pool; it now runs on its own DispatchSerialQueue executor.
- UserMessageBubble's truncation probe typeset entire pasted-log-sized
  messages through CoreText per layout pass (100% main-thread pins in
  the 07-21 hang reports); certainly-long messages now skip the probe
  and render a prefix while collapsed.
- toolGroupSignature JSON-encoded every tool input in the transcript up
  to 12.5x/s on the MainActor; now a structural hash. The summary pass
  is trailing-throttled to 0.4s, and flatItems joins streaming chunks
  once instead of re-copying the prefix per delta.
- StatusFeedFetcher.parseDate allocated three formatters per call (86%
  of a pool thread in the 07-26 report); now shared statics.

Guest (vminitd) — teardown data loss and epoll registration hazards:
- IOPair no longer closes on a bare EPOLLHUP with a backpressure flush
  in flight (dropped the CLI's final output line); EPOLLOUT finishes the
  flush, then EOF closes loss-free. ManagedProcess.setExit closes only
  stdin, letting stdout/stderr self-close on EOF, with an 8s grace pass
  (patch #16).
- Epoll events carry a registration generation; the supervisor ignores
  stale events for recycled fd numbers. registerFd refuses EEXIST
  instead of clobbering the existing handler. TerminalIO's stdin relay
  writes a dup of the terminal fd so its backpressure registration
  can't collide with the stdout relay's (patch #17).
- VsockProxy flushes bytes parked toward the surviving peer on hangup,
  closes the dialing socket on a failed backend connect, and
  StandardIO/TerminalIO clean up partially-created pairs on setup
  failure (patch #16).

Full suite: 1451+292+74+20 tests, two failures — both pre-existing
environmental (MacVM base image absent on this machine; a load-flaky
liveness test that passes 3/3 in isolation).

Co-Authored-By: Claude Fable 5 <[email protected]>
2026-07-28 15:35:35 -07:00

423 lines
16 KiB
Swift

//===----------------------------------------------------------------------===//
// Copyright © 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(Linux)
import Cgroup
import Containerization
import ContainerizationError
import ContainerizationOCI
import ContainerizationOS
import Foundation
import Logging
import Synchronization
final class ManagedProcess: ContainerProcess, Sendable {
// swiftlint: disable type_name
// [Nucleic vendored patch] Sendable so the exit path's delayed grace-close can capture
// the IO existential; both conformers (StandardIO, TerminalIO) already are.
protocol IO: Sendable {
func attach(pid: Int32, fd: Int32) throws
func start(process: inout Command) throws
func resize(size: Terminal.Size) throws
func close() throws
func closeStdin() throws
func closeAfterExec() throws
}
// swiftlint: enable type_name
private struct State {
init(io: IO) {
self.io = io
}
let io: IO
var waiters: [CheckedContinuation<ContainerExitStatus, Never>] = []
var exitStatus: ContainerExitStatus? = nil
var pid: Int32?
}
private static let ackPid = "AckPid"
private static let ackConsole = "AckConsole"
let id: String
private let log: Logger
private let command: Command
private let state: Mutex<State>
private let owningPid: Int32?
// [Nucleic vendored patch] Parent cgroup for this exec's OWN per-exec child (`<parent>/<id>`);
// nil means the legacy flat layout (join the container/init cgroup via `owningPid`).
private let execCgroupParent: String?
// [Nucleic vendored patch] Hard per-exec memory.max (bytes) for this exec's cgroup; nil = none.
private let execMemoryLimitBytes: UInt64?
private let ackPipe: Pipe
private let syncPipe: Pipe
private let errorPipe: Pipe
private let terminal: Bool
private let bundle: ContainerizationOCI.Bundle
var pid: Int32? {
self.state.withLock {
$0.pid
}
}
init(
id: String,
stdio: HostStdio,
bundle: ContainerizationOCI.Bundle,
owningPid: Int32? = nil,
execCgroupParent: String? = nil, // [Nucleic vendored patch]
execMemoryLimitBytes: UInt64? = nil, // [Nucleic vendored patch]
log: Logger
) throws {
self.id = id
var log = log
log[metadataKey: "id"] = "\(id)"
self.log = log
self.owningPid = owningPid
self.execCgroupParent = execCgroupParent
self.execMemoryLimitBytes = execMemoryLimitBytes
let syncPipe = Pipe()
try syncPipe.setCloexec()
self.syncPipe = syncPipe
let ackPipe = Pipe()
try ackPipe.setCloexec()
self.ackPipe = ackPipe
let errorPipe = Pipe()
try errorPipe.setCloexec()
self.errorPipe = errorPipe
let args: [String]
if let owningPid {
args = [
"exec",
"--parent-pid",
"\(owningPid)",
"--process-path",
bundle.getExecSpecPath(id: id).path,
]
} else {
args = ["run", "--bundle-path", bundle.path.path]
}
var command = Command(
"/sbin/vmexec",
arguments: args,
extraFiles: [
syncPipe.fileHandleForWriting,
ackPipe.fileHandleForReading,
errorPipe.fileHandleForWriting,
]
)
var io: IO
if stdio.terminal {
log.info("setting up terminal I/O")
let attrs = Command.Attrs(setsid: false, setctty: false)
command.attrs = attrs
io = try TerminalIO(
stdio: stdio,
log: log
)
} else {
command.attrs = .init(setsid: false)
io = StandardIO(
stdio: stdio,
log: log
)
}
log.info("starting I/O")
// Setup IO early. We expect the host to be listening already.
try io.start(process: &command)
self.command = command
self.terminal = stdio.terminal
self.bundle = bundle
self.state = Mutex(State(io: io))
}
}
extension ManagedProcess {
/// [Nucleic vendored patch] Run the blocking start sequence OFF the cooperative executor / gRPC
/// event loop. `startBlocking()` does synchronous, potentially slow pipe reads (waiting for
/// `vmexec` to hand back the pid, then for the error pipe to close) while holding `state`'s Mutex.
/// Upstream ran that directly on the calling task, so a slow exec start parked the event loop and
/// head-of-line-blocked sibling execs' control RPCs multiplexed on the same loop (each host `exec`
/// dials its own connection, but NIO pins several connections per loop). Dispatching to a worker
/// keeps the loop responsive; the body has no `await` and `ManagedProcess` is `Sendable`, so it is
/// safe off-actor, and the per-exec Mutex still serializes only this exec's own operations.
func start() async throws -> Int32 {
try await withCheckedThrowingContinuation { (cont: CheckedContinuation<Int32, Error>) in
DispatchQueue.global(qos: .userInitiated).async {
do {
cont.resume(returning: try self.startBlocking())
} catch {
cont.resume(throwing: error)
}
}
}
}
private func startBlocking() throws -> Int32 {
do {
return try self.state.withLock {
log.info(
"starting managed process",
metadata: [
"id": "\(id)"
])
// Start the underlying process.
try command.start()
defer {
try? self.ackPipe.fileHandleForWriting.close()
try? self.syncPipe.fileHandleForReading.close()
try? self.ackPipe.fileHandleForReading.close()
try? self.syncPipe.fileHandleForWriting.close()
try? self.errorPipe.fileHandleForWriting.close()
}
// Close our side of any pipes.
try $0.io.closeAfterExec()
try self.ackPipe.fileHandleForReading.close()
try self.syncPipe.fileHandleForWriting.close()
try self.errorPipe.fileHandleForWriting.close()
let size = MemoryLayout<Int32>.size
guard let piddata = try syncPipe.fileHandleForReading.read(upToCount: size) else {
throw ContainerizationError(.internalError, message: "no PID data from sync pipe")
}
guard piddata.count == size else {
throw ContainerizationError(.internalError, message: "invalid payload")
}
let pid = piddata.withUnsafeBytes { ptr in
ptr.load(as: Int32.self)
}
log.info(
"got back pid data",
metadata: [
"pid": "\(pid)"
])
$0.pid = pid
// This should probably happen in vmexec, but we don't need to set any cgroup
// toggles so the problem is much simpler to just do it here.
if let parent = execCgroupParent {
// [Nucleic vendored patch] Per-exec cgroup: place this exec in its OWN child cgroup
// (`<parent>/<id>`) with memory.oom.group (a runaway session's OOM kills only its
// tree), a fair cpu.weight, and a pids.max fork-bomb backstop — so one session can't
// OOM-kill / starve / fork-bomb its siblings in the shared container. Best-effort:
// on ANY error fall back to the container/init cgroup so a cgroup hiccup never blocks
// the exec from starting.
do {
let execCg = Cgroup2Manager(group: URL(filePath: parent + "/" + id), logger: log)
try execCg.create()
try? execCg.setOomGroup(true)
try? execCg.setCpuWeight(100)
try? execCg.setPidsMax(4096)
if let limit = execMemoryLimitBytes, limit > 0 {
try? execCg.setMemoryMax(bytes: limit) // host-configured hard per-exec ceiling
}
try execCg.addProcess(pid: pid)
} catch {
log.error("per-exec cgroup for \(id) failed; joining the container cgroup: \(error)")
if let owningPid {
try? Cgroup2Manager.loadFromPid(pid: owningPid).addProcess(pid: pid)
}
}
} else if let owningPid {
let cgManager = try Cgroup2Manager.loadFromPid(pid: owningPid)
try cgManager.addProcess(pid: pid)
}
log.info(
"sending pid acknowledgement",
metadata: [
"pid": "\(pid)"
])
try self.ackPipe.fileHandleForWriting.write(contentsOf: Self.ackPid.data(using: .utf8)!)
if self.terminal {
log.info(
"wait for PTY FD",
metadata: [
"id": "\(id)"
])
// Wait for a new write that will contain the pty fd if we asked for one.
guard let ptyFd = try self.syncPipe.fileHandleForReading.read(upToCount: size) else {
throw ContainerizationError(
.internalError,
message: "no PTY data from sync pipe"
)
}
let fd = ptyFd.withUnsafeBytes { ptr in
ptr.load(as: Int32.self)
}
log.info(
"received PTY FD from container, attaching",
metadata: [
"id": "\(id)"
])
try $0.io.attach(pid: pid, fd: fd)
try self.ackPipe.fileHandleForWriting.write(contentsOf: Self.ackConsole.data(using: .utf8)!)
}
// Wait for the errorPipe to close (after exec).
if let errorData = try? self.errorPipe.fileHandleForReading.readToEnd(),
let errorString = String(data: errorData, encoding: .utf8),
!errorString.isEmpty
{
throw ContainerizationError(
.internalError,
message: "vmexec error: \(errorString.trimmingCharacters(in: .whitespacesAndNewlines))"
)
}
log.info(
"started managed process",
metadata: [
"pid": "\(pid)",
"id": "\(id)",
])
return pid
}
} catch {
if let errorData = try? self.errorPipe.fileHandleForReading.readToEnd(),
let errorString = String(data: errorData, encoding: .utf8),
!errorString.isEmpty
{
throw ContainerizationError(
.internalError,
message: "vmexec error: \(errorString.trimmingCharacters(in: .whitespacesAndNewlines))",
cause: error
)
}
throw error
}
}
func setExit(_ status: Int32) {
self.state.withLock { state in
self.log.info(
"managed process exit",
metadata: [
"status": "\(status)"
])
let exitStatus = ContainerExitStatus(exitCode: status, exitedAt: Date.now)
state.exitStatus = exitStatus
// [Nucleic vendored patch] Close only stdin here. Force-closing ALL stdio at
// exit raced the relay: SIGCHLD is serviced before the poller thread drains the
// pipe's final EPOLLIN edge, and the old `io.close()` drain silently dropped
// whatever the destination wouldn't take (`drain` ignores short writes) — the
// CLI's final output line, truncated at process exit. The stdout/stderr write
// ends inside vminitd were already closed by `closeAfterExec`, so the child's
// exit delivers EOF/HUP to each relay, which self-closes loss-free (the relay
// never closes with a backpressure flush in flight). The delayed full close is
// a backstop for a grandchild that inherited the pipes and never exits, or a
// host that never drains — it must never fire on a healthy teardown path.
do {
try state.io.closeStdin()
} catch {
self.log.error("failed to close stdin for exited process: \(error)")
}
let io = state.io
let log = self.log
DispatchQueue.global(qos: .utility).asyncAfter(deadline: .now() + 8) {
// Idempotent: a relay already self-closed on EOF makes this a no-op.
do {
try io.close()
} catch {
log.error("failed to close I/O for exited process (grace pass): \(error)")
}
}
for waiter in state.waiters {
waiter.resume(returning: exitStatus)
}
self.log.debug("\(state.waiters.count) managed process waiters signaled")
state.waiters.removeAll()
}
}
/// Wait on the process to exit
func wait() async -> ContainerExitStatus {
await withCheckedContinuation { cont in
self.state.withLock {
if let status = $0.exitStatus {
cont.resume(returning: status)
return
}
$0.waiters.append(cont)
}
}
}
func kill(_ signal: Int32) async throws {
try self.state.withLock {
guard let pid = $0.pid else {
throw ContainerizationError(.invalidState, message: "process PID is required")
}
guard $0.exitStatus == nil else {
return
}
self.log.info("sending signal \(signal) to process \(pid)")
guard Foundation.kill(pid, signal) == 0 else {
throw POSIXError.fromErrno()
}
}
}
func resize(size: Terminal.Size) throws {
try self.state.withLock {
guard $0.exitStatus == nil else {
return
}
try $0.io.resize(size: size)
}
}
func closeStdin() throws {
let io = self.state.withLock { $0.io }
try io.closeStdin()
}
func delete() async throws {
// vmexec doesn't require explicit cleanup - the process is cleaned up
// when it exits and IO is closed via setExit()
}
}
#endif