Merge nucleic/hazy-opal-yak-cjc0 into dev

This commit is contained in:
2026-07-17 23:48:59 -07:00
parent 608e7d0450
commit 65623f1c34
5 changed files with 403 additions and 202 deletions
+168 -34
View File
@@ -33,6 +33,16 @@ final class IOPair: Sendable {
let buffer: UnsafeMutableBufferPointer<UInt8>
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)
@@ -66,10 +76,27 @@ final class IOPair: Sendable {
return
}
// Try and drain IO first.
self.drain()
// [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 fd from our global epoll instance first.
// Remove the fds from our global epoll instance first.
if let fd = self.registeredFd {
do {
try ProcessSupervisor.default.unregisterFd(fd)
@@ -78,6 +105,15 @@ final class IOPair: Sendable {
}
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()
@@ -108,7 +144,10 @@ final class IOPair: Sendable {
to: writeTo,
buffer: buffer,
closed: false,
registeredFd: nil
registeredFd: nil,
pending: [],
pendingOffset: 0,
writeFdRegistered: false
))
self.reason = reason
self.logger = logger
@@ -122,8 +161,16 @@ final class IOPair: Sendable {
return (io.from.fileDescriptor, io.to.fileDescriptor)
}
let readFrom = OSFile(fd: readFromFd)
let writeTo = OSFile(fd: writeToFd)
// [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
@@ -139,42 +186,129 @@ final class IOPair: Sendable {
return
}
// Loop so we drain fully.
while true {
let r = readFrom.read(io.buffer)
if r.read > 0 {
let view = UnsafeMutableBufferPointer(
start: io.buffer.baseAddress,
count: r.read
)
self.pump(&io, mask: mask, ignoreHup: ignoreHup)
}
}
}
let w = writeTo.write(view)
if w.wrote != r.read {
self.logger?.error("stopping relay: short write for stdio")
io.close(logger: self.logger)
return
}
}
/// [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)
switch r.action {
case .error(let errno):
self.logger?.error("failed with errno \(errno) while reading for fd \(readFromFd)")
fallthrough
case .eof:
self.logger?.debug("closing relay for \(readFromFd)")
io.close(logger: self.logger)
return
// 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:
if mask.isHangup && !ignoreHup {
self.logger?.error("received EPOLLHUP and EAGAIN exiting")
self.close()
}
// 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:
break
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)")
}
}