diff --git a/PATCHES.md b/PATCHES.md index 449b33f..43b6009 100644 --- a/PATCHES.md +++ b/PATCHES.md @@ -233,7 +233,30 @@ rebuild whenever a guest patch changes. Built locally, not in CI: the host frame non-blocking before either is registered. All three are guest-side and INERT until the initfs image is rebuilt (`make vminit-image`, tag `0.34.0-nucleic4`) and `ContainerEngine.vminitReference` is bumped after runtime validation. - Marked `[Nucleic vendored patch]`. + Marked `[Nucleic vendored patch]`. **NOTE:** the splice EAGAIN-return introduced a latent + cross-direction accounting hazard fixed by patch #15. + +15. **Per-direction relay state in `OSFile+Splice.swift` + `VsockProxy.swift` — fixes the + poller-thread spin patch #14 made reachable.** The old `SpliceFile` design threaded ONE offset + pair through BOTH directions of a proxied connection: each fd's struct served as the + read-counter for one direction and the write-counter for the other, so the loop guards compared + *differences of two directions' counters*. That was survivable only while every splice call + fully drained its transfer pipe before returning. Once patch #14's EAGAIN branch let pending + bytes persist across calls, one parked direction skewed the shared counters for the other: its + write leg's `to.offset < from.offset` guard went false with data still in the pipe, and the + outer `while true` then alternated read-EAGAIN/skip-write forever — a hard 100% spin on the + single `ProcessSupervisor` poller thread (epoll never runs again, so the event that would + un-skew the counters can never be processed). Container-wide dead control plane + frozen stdio + until the VM is recreated; triggered in practice by host-side backpressure (any slow host + reader), and bidirectional traffic on one connection (e.g. the control plane's SSE stream + + requests). Replaced `SpliceFile`/`OSFile.splice` with `OSFile.RelayDirection` (each direction + owns its OWN pipe, `bytesIn`/`bytesOut`, and `sawSourceEOF`) and `OSFile.relay`, which also + fixes a second latent defect: source EOF is now reported only after the pipe fully drains, so + the caller's SHUT_WR can never truncate a parked tail (the old read leg returned `.eof` + immediately, dropping pending bytes). `VsockProxy.handleConn` now tracks two directions with + independent done flags; a broken destination ends both. Guest-side and INERT until the initfs + image is rebuilt and repointed (mind the `vminit.ext4.reference` cache sidecar). Marked + `[Nucleic vendored patch]`. ## Re-vendoring a newer upstream commit @@ -256,8 +279,10 @@ rebuild whenever a guest patch changes. Built locally, not in CI: the host frame `KeychainQuery` reads: `withoutInteractiveUI` + the `errSecInteractionNotAllowed` handling + the `save` duplicate retry), and patch #14 (the non-blocking/non-spinning guest I/O plane: `IOPair` backpressure, the `OSFile.splice` EAGAIN return, and the `VsockProxy` pre-registration - non-blocking fds — all in `vminitd/`). After re-applying any - `vminitd/` patch, rebuild + publish the custom init image + non-blocking fds — all in `vminitd/`), and patch #15 (the per-direction + `OSFile.RelayDirection`/`OSFile.relay` rewrite + the two-direction `VsockProxy.handleConn`, + which supersede the upstream `SpliceFile`/`splice` shapes entirely — in `vminitd/`). After + re-applying any `vminitd/` patch, rebuild + publish the custom init image with `make vminit-image` + `make vminit-image-push`, and bump `ContainerEngine.vminitReference`. 5. Update the commit hash above and in the root `Package.swift` comment. 6. `swift build` and run the balloon tests. diff --git a/vminitd/Sources/VminitdCore/OSFile+Splice.swift b/vminitd/Sources/VminitdCore/OSFile+Splice.swift index 9a6f652..9e6b28a 100644 --- a/vminitd/Sources/VminitdCore/OSFile+Splice.swift +++ b/vminitd/Sources/VminitdCore/OSFile+Splice.swift @@ -20,95 +20,103 @@ import Foundation import LCShim extension OSFile { - struct SpliceFile: Sendable { - fileprivate var file: OSFile - fileprivate var offset: Int - fileprivate let pipe = Pipe() + /// [Nucleic vendored patch] One direction of a bidirectional relay: `from` fd → transfer + /// pipe → `to` fd. Each direction owns its OWN pipe and byte counters. + /// + /// The previous `SpliceFile` design threaded ONE offset pair through BOTH directions of + /// `VsockProxy`'s relay (a fd's struct served as read-counter in one direction and + /// write-counter in the other). That was survivable only while every splice call fully + /// drained its pipe before returning. Once the EAGAIN-return backpressure patch let pending + /// bytes persist across calls, one parked direction skewed the shared counters for the + /// other: its write leg's `to.offset < from.offset` guard went false with data still in the + /// pipe, and the outer loop then alternated read-EAGAIN/skip-write forever — a hard spin on + /// vminitd's single ProcessSupervisor poller thread that froze every exec's stdio and every + /// control-plane relay in the container until the VM was recreated (the persistent + /// "produced no output within 60s" / dead-control-plane state). + struct RelayDirection: Sendable { + let from: Int32 + let to: Int32 + private let pipe = Pipe() + /// Bytes spliced from `from` into the transfer pipe so far. + fileprivate var bytesIn = 0 + /// Bytes spliced from the transfer pipe into `to` so far. + fileprivate var bytesOut = 0 + /// The source reported EOF. The direction only FINISHES (`.eof`) once the pipe has + /// also drained, so the stream's tail is never dropped by an early SHUT_WR. + fileprivate var sawSourceEOF = false - var fileDescriptor: Int32 { - file.fileDescriptor - } + /// Bytes read from the source that the destination hasn't accepted yet. + var pendingBytes: Int { bytesIn - bytesOut } - var reader: Int32 { - pipe.fileHandleForReading.fileDescriptor - } + fileprivate var pipeReader: Int32 { pipe.fileHandleForReading.fileDescriptor } + fileprivate var pipeWriter: Int32 { pipe.fileHandleForWriting.fileDescriptor } - var writer: Int32 { - pipe.fileHandleForWriting.fileDescriptor - } - - init(fd: Int32) { - self.file = OSFile(fd: fd) - self.offset = 0 - } - - init(handle: FileHandle) { - self.file = OSFile(handle: handle) - self.offset = 0 - } - - init(from: OSFile, withOffset: Int = 0) { - self.file = from - self.offset = withOffset - } - - func close() throws { - try self.file.close() + init(from: Int32, to: Int32) { + self.from = from + self.to = to } } - static func splice(from: inout SpliceFile, to: inout SpliceFile, count: Int = 1 << 16) throws -> (read: Int, wrote: Int, action: IOAction) { - let fromOffset = from.offset - let toOffset = to.offset + /// The terminal state of one `relay` pass over a direction. + enum RelayResult: Sendable { + /// No more progress possible right now: the source has no data (EAGAIN) or the + /// destination is full (its EPOLLOUT edge resumes the flush of `pendingBytes`). + case idle + /// Source EOF observed AND the pipe fully drained — the direction is complete; the + /// caller should SHUT_WR the destination. + case eof + /// The destination hung up mid-write; nothing further can be delivered. + case brokenPipe + } + /// Move as much data as possible along `direction` without blocking. `count` bounds the + /// bytes buffered in the transfer pipe (it matches the default pipe capacity). + static func relay(_ direction: inout RelayDirection, count: Int = 1 << 16) throws -> RelayResult { + let flags = UInt32(bitPattern: LCShim.SPLICE_F_MOVE | LCShim.SPLICE_F_NONBLOCK) while true { - while (from.offset - to.offset) < count { - let toRead = count - (from.offset - to.offset) - let bytesRead = LCShim.splice(from.fileDescriptor, nil, to.writer, nil, toRead, UInt32(bitPattern: LCShim.SPLICE_F_MOVE | LCShim.SPLICE_F_NONBLOCK)) - if bytesRead == -1 { + // Read leg: source → pipe, until the pipe is full, the source runs dry, or EOF. + var sourceDry = false + if !direction.sawSourceEOF { + while direction.pendingBytes < count { + let toRead = count - direction.pendingBytes + let n = LCShim.splice(direction.from, nil, direction.pipeWriter, nil, toRead, flags) + if n == -1 { + if errno != EAGAIN && errno != EIO { + throw POSIXError(.init(rawValue: errno)!) + } + sourceDry = true + break + } + if n == 0 { + direction.sawSourceEOF = true + break + } + direction.bytesIn += n + if n < toRead { break } + } + } + // Write leg: pipe → destination, until drained or the destination pushes back. + while direction.pendingBytes > 0 { + let n = LCShim.splice(direction.pipeReader, nil, direction.to, nil, direction.pendingBytes, flags) + if n == -1 { if errno != EAGAIN && errno != EIO { throw POSIXError(.init(rawValue: errno)!) } - break + // Destination full: park with the remainder in the pipe. The destination + // fd's EPOLLOUT edge re-enters and resumes exactly here — never spin, and + // never block the shared poller thread. + return .idle } - if bytesRead == 0 { - return (0, 0, .eof) - } - from.offset += bytesRead - if bytesRead < toRead { - break - } - } - if from.offset == to.offset { - return (from.offset - fromOffset, to.offset - toOffset, .success) - } - while to.offset < from.offset { - let toWrite = from.offset - to.offset - let bytesWrote = LCShim.splice(to.reader, nil, to.fileDescriptor, nil, toWrite, UInt32(bitPattern: LCShim.SPLICE_F_MOVE | LCShim.SPLICE_F_NONBLOCK)) - if bytesWrote == -1 { - if errno != EAGAIN && errno != EIO { - throw POSIXError(.init(rawValue: errno)!) - } - // [Nucleic vendored patch] Destination full: RETURN, don't `break`. Breaking - // sent the outer `while true` straight back here — with the source idle and the - // destination still full, neither leg could progress and this spun the caller's - // thread at 100% until the peer drained. The caller is an epoll handler on - // vminitd's SINGLE ProcessSupervisor poller thread, so the spin froze every - // exec's stdio and every control-plane relay in the container at once (the - // all-sessions "produced no output within 60s" stall / dead control plane). - // The un-flushed bytes stay in the transfer pipe (`from.offset > to.offset` - // persists in the SpliceFiles); the destination fd is registered for EPOLLOUT, - // whose edge re-enters transferData and resumes the flush. - return (from.offset - fromOffset, to.offset - toOffset, .success) - } - to.offset += bytesWrote - if bytesWrote == 0 { - return (from.offset - fromOffset, to.offset - toOffset, .brokenPipe) - } - if bytesWrote < toWrite { - break + if n == 0 { + return .brokenPipe } + direction.bytesOut += n } + // Pipe is drained here. + if direction.sawSourceEOF { return .eof } + if sourceDry { return .idle } + // The read leg stopped only because the pipe filled (or read a full window): + // go around again — the source may still have data. } } } diff --git a/vminitd/Sources/VminitdCore/VsockProxy.swift b/vminitd/Sources/VminitdCore/VsockProxy.swift index 3e6059c..b8bbfb7 100644 --- a/vminitd/Sources/VminitdCore/VsockProxy.swift +++ b/vminitd/Sources/VminitdCore/VsockProxy.swift @@ -260,12 +260,19 @@ extension VsockProxy { } } - // `clientFile` isn't used concurrently. - nonisolated(unsafe) var clientFile = OSFile.SpliceFile(fd: conn.fileDescriptor) - nonisolated(unsafe) var eofFromClient = false - // `serverFile` isn't used concurrently. - nonisolated(unsafe) var serverFile = OSFile.SpliceFile(fd: relayTo.fileDescriptor) - nonisolated(unsafe) var eofFromServer = false + // [Nucleic vendored patch] Each relay direction owns its own pipe and byte + // counters (see RelayDirection) — the previous shared-offset SpliceFile pair + // let one parked direction corrupt the other's accounting and spin the poller + // thread. Neither is used concurrently (all access is on the single + // ProcessSupervisor poller thread). + nonisolated(unsafe) var toServer = OSFile.RelayDirection( + from: conn.fileDescriptor, to: relayTo.fileDescriptor) + nonisolated(unsafe) var toClient = OSFile.RelayDirection( + from: relayTo.fileDescriptor, to: conn.fileDescriptor) + // A direction is DONE once its source EOF fully flushed (SHUT_WR sent), its + // destination broke, or a full hangup ended the connection. + nonisolated(unsafe) var toServerDone = false + nonisolated(unsafe) var toClientDone = false // clean up when any of these conditions apply: // - the client has completely hung up or errored @@ -291,20 +298,20 @@ extension VsockProxy { "vport": "\(port)", "uds": "\(path)", "action": "\(action)", - "eofFromClient": "\(eofFromClient)", - "eofFromServer": "\(eofFromServer)", - "clientFd": "\(clientFile.fileDescriptor)", - "serverFd": "\(serverFile.fileDescriptor)", + "toServerDone": "\(toServerDone)", + "toClientDone": "\(toClientDone)", + "clientFd": "\(conn.fileDescriptor)", + "serverFd": "\(relayTo.fileDescriptor)", ] ) do { - try ProcessSupervisor.default.unregisterFd(clientFile.fileDescriptor) + try ProcessSupervisor.default.unregisterFd(conn.fileDescriptor) } catch { self.log?.error("Failed to unregister vsock proxy client fd: \(error)") } do { - try ProcessSupervisor.default.unregisterFd(serverFile.fileDescriptor) + try ProcessSupervisor.default.unregisterFd(relayTo.fileDescriptor) } catch { self.log?.error("Failed to unregister vsock proxy server fd: \(error)") } @@ -326,55 +333,52 @@ extension VsockProxy { // taking every session in the container down. Fail the one connection instead, // releasing whatever was already set up so nothing leaks (the caller closes `conn`; // `relayTo` and the first registration are released in the catch blocks below). + // [Nucleic vendored patch] Interpret one relay-step outcome: returns whether + // the stepped direction is now done; a broken destination ends BOTH directions + // (the peer is gone). Returns rather than writing the stepped flag itself so no + // captured var is ever aliased by an inout parameter (exclusivity). + let apply = { @Sendable (outcome: RelayStepOutcome) -> Bool in + switch outcome { + case .open: return false + case .finished: return true + case .broken: + toServerDone = true + toClientDone = true + return true + } + } + do { - try ProcessSupervisor.default.registerFd(clientFile.fileDescriptor, mask: [.input, .output]) { mask in - if mask.readyToRead && !eofFromClient { - let (fromEof, toEof) = Self.transferData( - fromFile: &clientFile, - toFile: &serverFile, - description: "readyToRead:toServer", - log: self.log - ) - eofFromClient = eofFromClient || fromEof - eofFromServer = eofFromServer || toEof - } - - if mask.readyToWrite && !eofFromServer { - let (fromEof, toEof) = Self.transferData( - fromFile: &serverFile, - toFile: &clientFile, - description: "readyToWrite:toClient", - log: self.log - ) - eofFromClient = eofFromClient || toEof - eofFromServer = eofFromServer || fromEof - } - - if mask.isHangup { - eofFromClient = true - eofFromServer = true - } else if mask.isRemoteHangup && !eofFromClient { - // half close, shut down client to server transfer - // we should see no more EPOLLIN events on the client fd - // and no more EPOLLOUT events on the server fd - eofFromClient = true - if shutdown(serverFile.fileDescriptor, Int32(SHUT_WR)) != 0 { - self.log?.warning( - "failed to shut down client reads", - metadata: [ - "vport": "\(self.port)", - "uds": "\(self.path)", - "errno": "\(errno)", - "eofFromClient": "\(eofFromClient)", - "eofFromServer": "\(eofFromServer)", - "clientFd": "\(clientFile.fileDescriptor)", - "serverFd": "\(serverFile.fileDescriptor)", - ] - ) + try ProcessSupervisor.default.registerFd(conn.fileDescriptor, mask: [.input, .output]) { mask in + if mask.readyToRead && !toServerDone { + if apply(Self.relayStep(&toServer, description: "client:readyToRead:toServer", log: self.log)) { + toServerDone = true } } - if eofFromClient && eofFromServer { + // The client drained: flush toClient's parked bytes (and whatever more the + // server has ready). + if mask.readyToWrite && !toClientDone { + if apply(Self.relayStep(&toClient, description: "client:readyToWrite:toClient", log: self.log)) { + toClientDone = true + } + } + + if mask.isHangup { + toServerDone = true + toClientDone = true + } else if mask.isRemoteHangup && !toServerDone { + // Half close: the client sends no more. Drain the tail — relayStep + // observes the real EOF after the last buffered bytes and only then + // SHUT_WRs the server, so a parked backlog is never dropped. If the + // server is full right now the direction stays open and its EPOLLOUT + // edge finishes the flush. + if apply(Self.relayStep(&toServer, description: "client:remoteHangup:toServer", log: self.log)) { + toServerDone = true + } + } + + if toServerDone && toClientDone { return cleanup() } } @@ -384,59 +388,38 @@ extension VsockProxy { } do { - try ProcessSupervisor.default.registerFd(serverFile.fileDescriptor, mask: [.input, .output]) { mask in - if mask.readyToRead && !eofFromServer { - let (fromEof, toEof) = Self.transferData( - fromFile: &serverFile, - toFile: &clientFile, - description: "readyToRead:toClient", - log: self.log - ) - eofFromClient = eofFromClient || toEof - eofFromServer = eofFromServer || fromEof - } - - if mask.readyToWrite && !eofFromClient { - let (fromEof, toEof) = Self.transferData( - fromFile: &clientFile, - toFile: &serverFile, - description: "readyToWrite:toServer", - log: self.log - ) - eofFromClient = eofFromClient || fromEof - eofFromServer = eofFromServer || toEof - } - - if mask.isHangup { - eofFromClient = true - eofFromServer = true - } else if mask.isRemoteHangup && !eofFromServer { - // half close, shut down server to client transfer - // we should see no more EPOLLIN events on the server fd - // and no more EPOLLOUT events on the client fd - eofFromServer = true - if shutdown(clientFile.fileDescriptor, Int32(SHUT_WR)) != 0 { - self.log?.warning( - "failed to shut down server reads", - metadata: [ - "vport": "\(self.port)", - "uds": "\(self.path)", - "errno": "\(errno)", - "eofFromClient": "\(eofFromClient)", - "eofFromServer": "\(eofFromServer)", - "clientFd": "\(clientFile.fileDescriptor)", - "serverFd": "\(serverFile.fileDescriptor)", - ] - ) + try ProcessSupervisor.default.registerFd(relayTo.fileDescriptor, mask: [.input, .output]) { mask in + if mask.readyToRead && !toClientDone { + if apply(Self.relayStep(&toClient, description: "server:readyToRead:toClient", log: self.log)) { + toClientDone = true } } - if eofFromClient && eofFromServer { + // The server drained: flush toServer's parked bytes (and whatever more the + // client has ready). + if mask.readyToWrite && !toServerDone { + if apply(Self.relayStep(&toServer, description: "server:readyToWrite:toServer", log: self.log)) { + toServerDone = true + } + } + + if mask.isHangup { + toServerDone = true + toClientDone = true + } else if mask.isRemoteHangup && !toClientDone { + // Half close: the server sends no more — drain the tail toward the + // client (see the client handler's mirror-image comment). + if apply(Self.relayStep(&toClient, description: "server:remoteHangup:toClient", log: self.log)) { + toClientDone = true + } + } + + if toServerDone && toClientDone { return cleanup() } } } catch { - try? ProcessSupervisor.default.unregisterFd(clientFile.fileDescriptor) + try? ProcessSupervisor.default.unregisterFd(conn.fileDescriptor) try? relayTo.close() throw error } @@ -446,50 +429,65 @@ extension VsockProxy { } } - private static func transferData( - fromFile: inout OSFile.SpliceFile, - toFile: inout OSFile.SpliceFile, + /// [Nucleic vendored patch] Outcome of one non-blocking relay pass over a direction. + enum RelayStepOutcome { + /// More may come (source dry, or destination full with bytes parked in the pipe). + case open + /// Source EOF fully flushed; the destination has been SHUT_WR'd. + case finished + /// The destination hung up (or the splice failed) — the connection is over. + case broken + } + + /// [Nucleic vendored patch] Run one non-blocking relay pass over `direction`. `.eof` is + /// only reported by `OSFile.relay` once the transfer pipe has fully drained, so the + /// SHUT_WR here can never truncate a parked tail. + private static func relayStep( + _ direction: inout OSFile.RelayDirection, description: String, log: Logger? - ) -> (Bool, Bool) { + ) -> RelayStepOutcome { do { - let (readBytes, writeBytes, action) = try OSFile.splice(from: &fromFile, to: &toFile) + let result = try OSFile.relay(&direction) log?.trace( "transferred data", metadata: [ "description": "\(description)", - "action": "\(action)", - "readBytes": "\(readBytes)", - "writeBytes": "\(writeBytes)", - "fromFd": "\(fromFile.fileDescriptor)", - "toFd": "\(toFile.fileDescriptor)", + "result": "\(result)", + "pendingBytes": "\(direction.pendingBytes)", + "fromFd": "\(direction.from)", + "toFd": "\(direction.to)", ] ) - if action == .eof { - // half close, shut down client to server transfer - // we should see no more EPOLLIN events on the client fd - // and no more EPOLLOUT events on the server fd - if shutdown(toFile.fileDescriptor, Int32(SHUT_WR)) != 0 { + switch result { + case .idle: + return .open + case .eof: + if shutdown(direction.to, Int32(SHUT_WR)) != 0 { log?.warning( - "failed to shut down reads", + "failed to shut down destination writes", metadata: [ "description": "\(description)", "errno": "\(errno)", - "action": "\(action)", - "readBytes": "\(readBytes)", - "writeBytes": "\(writeBytes)", - "fromFd": "\(fromFile.fileDescriptor)", - "toFd": "\(toFile.fileDescriptor)", + "fromFd": "\(direction.from)", + "toFd": "\(direction.to)", ] ) } - return (true, false) - } else if action == .brokenPipe { - return (true, true) + return .finished + case .brokenPipe: + return .broken } - return (false, false) } catch { - return (true, true) + log?.error( + "relay failed: \(error)", + metadata: [ + "description": "\(description)", + "fromFd": "\(direction.from)", + "toFd": "\(direction.to)", + ] + ) + return .broken } } }