//===----------------------------------------------------------------------===// // 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 ContainerizationError import ContainerizationOS import Foundation import Logging import Synchronization final class IOPair: Sendable { private let io: Mutex private let logger: Logger? private let reason: String private struct IO { let from: IOCloser let to: IOCloser let buffer: UnsafeMutableBufferPointer var closed: Bool var registeredFd: Int32? // [Nucleic vendored patch] Backpressure state: bytes read from `from` that `to` couldn't // take yet (its buffer was full), plus whether `to` is currently registered for EPOLLOUT // to flush them. While `pending` is non-empty the relay reads nothing more — the source // pipe backs up and throttles the *producing process* instead of this thread. The write // fd used to be BLOCKING (only registered fds get O_NONBLOCK, and it never was), so one // slow-drained stream parked the shared ProcessSupervisor poller thread — freezing every // exec's stdio and every control-plane relay in the container at once. var pending: [UInt8] var pendingOffset: Int var writeFdRegistered: Bool func drain() { let readFrom = OSFile(fd: from.fileDescriptor) let writeTo = OSFile(fd: to.fileDescriptor) while true { let r = readFrom.read(buffer) if r.read > 0 { let view = UnsafeMutableBufferPointer( start: buffer.baseAddress, count: r.read ) let w = writeTo.write(view) if w.wrote != r.read { return } } switch r.action { case .eof, .again, .error(_): return default: break } } } mutating func close(logger: Logger?) { if self.closed { return } // [Nucleic vendored patch] Flush what we can IN ORDER: the pending backlog first, // then (only if it fully flushed) a best-effort drain of the source. Draining with // unsent pending bytes would reorder the stream. let writeTo = OSFile(fd: to.fileDescriptor) while pendingOffset < pending.count { let offset = pendingOffset let result = pending.withUnsafeMutableBufferPointer { buf in writeTo.write( UnsafeMutableBufferPointer( start: buf.baseAddress!.advanced(by: offset), count: buf.count - offset)) } if result.wrote > 0 { pendingOffset += result.wrote } if result.action != .success { break } } if pendingOffset >= pending.count { // Try and drain IO first. self.drain() } // Remove the fds from our global epoll instance first. if let fd = self.registeredFd { do { try ProcessSupervisor.default.unregisterFd(fd) } catch { logger?.error("failed to delete fd from epoll \(fd): \(error)") } self.registeredFd = nil } // [Nucleic vendored patch] The write fd may be registered for backpressure flushing. if self.writeFdRegistered { do { try ProcessSupervisor.default.unregisterFd(to.fileDescriptor) } catch { logger?.error("failed to delete write fd from epoll \(to.fileDescriptor): \(error)") } self.writeFdRegistered = false } do { try self.from.close() } catch { logger?.error("failed to close reader fd for IOPair: \(error)") } do { try self.to.close() } catch { logger?.error("failed to close writer fd for IOPair: \(error)") } self.buffer.deallocate() self.closed = true } } init( readFrom: IOCloser, writeTo: IOCloser, reason: String, logger: Logger? = nil ) { let buffer = UnsafeMutableBufferPointer.allocate(capacity: Int(getpagesize())) self.io = Mutex( IO( from: readFrom, to: writeTo, buffer: buffer, closed: false, registeredFd: nil, pending: [], pendingOffset: 0, writeFdRegistered: false )) self.reason = reason self.logger = logger } func relay(ignoreHup: Bool = false) throws { self.logger?.info("setting up relay for \(reason)") let (readFromFd, writeToFd) = self.io.withLock { io in io.registeredFd = io.from.fileDescriptor return (io.from.fileDescriptor, io.to.fileDescriptor) } // [Nucleic vendored patch] The write fd must be non-blocking BEFORE the first relay write. // `Epoll.add` only sets O_NONBLOCK on fds it registers, and the write fd is registered only // on demand (EPOLLOUT backpressure below) — so without this, the very first full-buffer // write blocked the shared poller thread. let flags = fcntl(writeToFd, F_GETFL) if flags == -1 || fcntl(writeToFd, F_SETFL, flags | O_NONBLOCK) == -1 { self.logger?.error( "failed to set relay write fd non-blocking", metadata: ["fd": "\(writeToFd)", "errno": "\(errno)"]) } try ProcessSupervisor.default.registerFd(readFromFd, mask: .input) { mask in self.io.withLock { io in if io.closed { return } if mask.isHangup && !mask.readyToRead { self.logger?.debug("received EPOLLHUP with no EPOLLIN") if !ignoreHup { io.close(logger: self.logger) } return } self.pump(&io, mask: mask, ignoreHup: ignoreHup) } } } /// [Nucleic vendored patch] One relay pass, non-blocking end to end: flush any pending /// backlog toward `to`, then (only once it's empty) drain `from`. On a full destination the /// remainder is stashed in `pending` and the write fd registered for EPOLLOUT, whose edge /// re-enters this pump — so backpressure suspends the relay instead of blocking or spinning /// the shared poller thread. Must be called with the `io` lock held. private func pump(_ io: inout IO, mask: Epoll.Mask, ignoreHup: Bool) { let readFrom = OSFile(fd: io.from.fileDescriptor) let writeTo = OSFile(fd: io.to.fileDescriptor) // Flush the pending backlog first; reads stay suspended until it clears. while io.pendingOffset < io.pending.count { let offset = io.pendingOffset let result = io.pending.withUnsafeMutableBufferPointer { buf in writeTo.write( UnsafeMutableBufferPointer( start: buf.baseAddress!.advanced(by: offset), count: buf.count - offset)) } if result.wrote > 0 { io.pendingOffset += result.wrote } switch result.action { case .success: continue case .again: self.ensureWriteRegistered(&io) return default: self.logger?.error("stopping relay: write failed during backlog flush") io.close(logger: self.logger) return } } if !io.pending.isEmpty { io.pending = [] io.pendingOffset = 0 self.unregisterWrite(&io) } // Loop so we drain fully (edge-triggered epoll requires reading until EAGAIN). while true { let r = readFrom.read(io.buffer) if r.read > 0 { let view = UnsafeMutableBufferPointer( start: io.buffer.baseAddress, count: r.read ) let w = writeTo.write(view) if w.wrote != r.read { switch w.action { case .again: // Destination full: stash the remainder and suspend reads until its // EPOLLOUT edge flushes it. (A later `read` re-reports EOF if this // chunk was the stream's last, so no EOF is lost by returning here.) io.pending = Array( UnsafeBufferPointer( start: io.buffer.baseAddress!.advanced(by: w.wrote), count: r.read - w.wrote)) io.pendingOffset = 0 self.ensureWriteRegistered(&io) return default: self.logger?.error("stopping relay: short write for stdio") io.close(logger: self.logger) return } } } switch r.action { case .error(let errno): self.logger?.error("failed with errno \(errno) while reading for fd \(io.from.fileDescriptor)") fallthrough case .eof: self.logger?.debug("closing relay for \(io.from.fileDescriptor)") io.close(logger: self.logger) return case .again: if mask.isHangup && !ignoreHup { self.logger?.error("received EPOLLHUP and EAGAIN exiting") io.close(logger: self.logger) } return default: break } } } /// [Nucleic vendored patch] Register the write fd for EPOLLOUT so the pending backlog is /// flushed when the destination drains. Registration failure closes the pair — without the /// flush wakeup the relay would hang with data stranded. Must be called with the lock held. private func ensureWriteRegistered(_ io: inout IO) { guard !io.writeFdRegistered else { return } let writeToFd = io.to.fileDescriptor do { try ProcessSupervisor.default.registerFd(writeToFd, mask: .output) { _ in self.io.withLock { io in guard !io.closed else { return } // An empty mask: HUP/EOF handling rides the read fd's own events. self.pump(&io, mask: [], ignoreHup: true) } } io.writeFdRegistered = true } catch { self.logger?.error("failed to register relay write fd for backpressure: \(error)") io.close(logger: self.logger) } } /// [Nucleic vendored patch] Drop the EPOLLOUT registration once the backlog has flushed. /// Must be called with the lock held. private func unregisterWrite(_ io: inout IO) { guard io.writeFdRegistered else { return } io.writeFdRegistered = false do { try ProcessSupervisor.default.unregisterFd(io.to.fileDescriptor) } catch { self.logger?.error("failed to unregister relay write fd: \(error)") } } func close() { self.io.withLock { io in self.logger?.info("closing relay for \(reason)") io.close(logger: self.logger) } } } #endif