Files

333 lines
13 KiB
Swift
Raw Permalink Normal View History

//===----------------------------------------------------------------------===//
// 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<IO>
private let logger: Logger?
private let reason: String
private struct IO {
let from: IOCloser
let to: IOCloser
let buffer: UnsafeMutableBufferPointer<UInt8>
var closed: Bool
var registeredFd: Int32?
2026-07-17 23:48:59 -07:00
// [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
}
2026-07-17 23:48:59 -07:00
// [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()
}
2026-07-17 23:48:59 -07:00
// 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
}
2026-07-17 23:48:59 -07:00
// [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<UInt8>.allocate(capacity: Int(getpagesize()))
self.io = Mutex(
IO(
from: readFrom,
to: writeTo,
buffer: buffer,
closed: false,
2026-07-17 23:48:59 -07:00
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)
}
2026-07-17 23:48:59 -07:00
// [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")
// [Nucleic vendored patch] Never close on a bare HUP while a backpressure
// flush is in flight: the writer closing right after its final burst was
// stashed in `pending` (destination momentarily full) used to hit this
// close — whose single best-effort flush pass EAGAINed — and DROP the tail
// of the stream (the CLI's final result line). With pending outstanding,
// just return: the destination's EPOLLOUT edge flushes the backlog, the
// pump then reads the drained closed-writer pipe, observes EOF, and closes
// loss-free. If the destination dies instead, its own error/EPOLLERR wakes
// the write handler, whose failed flush closes the pair — no orphan.
if !ignoreHup && io.pending.isEmpty && !io.writeFdRegistered {
io.close(logger: self.logger)
}
return
}
2026-07-17 23:48:59 -07:00
self.pump(&io, mask: mask, ignoreHup: ignoreHup)
}
}
}
2026-07-17 23:48:59 -07:00
/// [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:
2026-07-17 23:48:59 -07:00
// 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:
2026-07-17 23:48:59 -07:00
self.logger?.error("stopping relay: short write for stdio")
io.close(logger: self.logger)
return
}
}
}
2026-07-17 23:48:59 -07:00
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