Fix the recurring session-stall pair: MainActor lock-reconcile hang + vminitd relay spin
Host (dominant): AppStore.reconcileLocks polls every 3s on the MainActor while any lock is held — effectively forever, since interrupted/errored sessions deliberately retain locks. Each pass ran heldPathDisposition's diverges() as three held×unmerged scans with two split-allocations per pathsOverlap call, pinning the main thread for tens of seconds per pass on a diverged trunk (hang-reports 2026-07-28: 100% of samples in reconcileLocks→heldPathDisposition→pathsOverlap). That froze running sessions' transcripts and starved the spawn path into the 60s "produced no output" watchdog. pathsOverlap is now allocation-free bytewise comparison with identical semantics, and divergentHeldPaths answers all three questions from one O((held+unmerged)·depth) set. Guest (persistence): VsockProxy threaded ONE offset pair through BOTH relay directions; once the EAGAIN-return backpressure patch let pending bytes persist, traffic in the other direction skewed the shared counters, made the write leg unreachable, and spun the single ProcessSupervisor poller thread forever — container-wide dead control plane until VM recreation, triggered by exactly the backpressure the host hang created. Each direction now owns its own pipe and counters (OSFile.RelayDirection), and source EOF is only surfaced after the pipe drains so SHUT_WR can't truncate a parked tail. Vendored patch docs updated (#15); inert until the initfs image is rebuilt+repointed. Co-Authored-By: Claude Fable 5 <[email protected]>
This commit is contained in:
+28
-3
@@ -233,7 +233,30 @@ rebuild whenever a guest patch changes. Built locally, not in CI: the host frame
|
|||||||
non-blocking before either is registered.
|
non-blocking before either is registered.
|
||||||
All three are guest-side and INERT until the initfs image is rebuilt (`make vminit-image`, tag
|
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.
|
`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
|
## 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 +
|
`KeychainQuery` reads: `withoutInteractiveUI` + the `errSecInteractionNotAllowed` handling +
|
||||||
the `save` duplicate retry), and patch #14 (the non-blocking/non-spinning guest I/O plane:
|
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
|
`IOPair` backpressure, the `OSFile.splice` EAGAIN return, and the `VsockProxy` pre-registration
|
||||||
non-blocking fds — all in `vminitd/`). After re-applying any
|
non-blocking fds — all in `vminitd/`), and patch #15 (the per-direction
|
||||||
`vminitd/` patch, rebuild + publish the custom init image
|
`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`.
|
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.
|
5. Update the commit hash above and in the root `Package.swift` comment.
|
||||||
6. `swift build` and run the balloon tests.
|
6. `swift build` and run the balloon tests.
|
||||||
|
|||||||
@@ -20,95 +20,103 @@ import Foundation
|
|||||||
import LCShim
|
import LCShim
|
||||||
|
|
||||||
extension OSFile {
|
extension OSFile {
|
||||||
struct SpliceFile: Sendable {
|
/// [Nucleic vendored patch] One direction of a bidirectional relay: `from` fd → transfer
|
||||||
fileprivate var file: OSFile
|
/// pipe → `to` fd. Each direction owns its OWN pipe and byte counters.
|
||||||
fileprivate var offset: Int
|
///
|
||||||
fileprivate let pipe = Pipe()
|
/// 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 {
|
/// Bytes read from the source that the destination hasn't accepted yet.
|
||||||
file.fileDescriptor
|
var pendingBytes: Int { bytesIn - bytesOut }
|
||||||
}
|
|
||||||
|
|
||||||
var reader: Int32 {
|
fileprivate var pipeReader: Int32 { pipe.fileHandleForReading.fileDescriptor }
|
||||||
pipe.fileHandleForReading.fileDescriptor
|
fileprivate var pipeWriter: Int32 { pipe.fileHandleForWriting.fileDescriptor }
|
||||||
}
|
|
||||||
|
|
||||||
var writer: Int32 {
|
init(from: Int32, to: Int32) {
|
||||||
pipe.fileHandleForWriting.fileDescriptor
|
self.from = from
|
||||||
}
|
self.to = to
|
||||||
|
|
||||||
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()
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
static func splice(from: inout SpliceFile, to: inout SpliceFile, count: Int = 1 << 16) throws -> (read: Int, wrote: Int, action: IOAction) {
|
/// The terminal state of one `relay` pass over a direction.
|
||||||
let fromOffset = from.offset
|
enum RelayResult: Sendable {
|
||||||
let toOffset = to.offset
|
/// 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 true {
|
||||||
while (from.offset - to.offset) < count {
|
// Read leg: source → pipe, until the pipe is full, the source runs dry, or EOF.
|
||||||
let toRead = count - (from.offset - to.offset)
|
var sourceDry = false
|
||||||
let bytesRead = LCShim.splice(from.fileDescriptor, nil, to.writer, nil, toRead, UInt32(bitPattern: LCShim.SPLICE_F_MOVE | LCShim.SPLICE_F_NONBLOCK))
|
if !direction.sawSourceEOF {
|
||||||
if bytesRead == -1 {
|
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 {
|
if errno != EAGAIN && errno != EIO {
|
||||||
throw POSIXError(.init(rawValue: errno)!)
|
throw POSIXError(.init(rawValue: errno)!)
|
||||||
}
|
}
|
||||||
|
sourceDry = true
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
if bytesRead == 0 {
|
if n == 0 {
|
||||||
return (0, 0, .eof)
|
direction.sawSourceEOF = true
|
||||||
}
|
|
||||||
from.offset += bytesRead
|
|
||||||
if bytesRead < toRead {
|
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
|
direction.bytesIn += n
|
||||||
|
if n < toRead { break }
|
||||||
}
|
}
|
||||||
if from.offset == to.offset {
|
|
||||||
return (from.offset - fromOffset, to.offset - toOffset, .success)
|
|
||||||
}
|
}
|
||||||
while to.offset < from.offset {
|
// Write leg: pipe → destination, until drained or the destination pushes back.
|
||||||
let toWrite = from.offset - to.offset
|
while direction.pendingBytes > 0 {
|
||||||
let bytesWrote = LCShim.splice(to.reader, nil, to.fileDescriptor, nil, toWrite, UInt32(bitPattern: LCShim.SPLICE_F_MOVE | LCShim.SPLICE_F_NONBLOCK))
|
let n = LCShim.splice(direction.pipeReader, nil, direction.to, nil, direction.pendingBytes, flags)
|
||||||
if bytesWrote == -1 {
|
if n == -1 {
|
||||||
if errno != EAGAIN && errno != EIO {
|
if errno != EAGAIN && errno != EIO {
|
||||||
throw POSIXError(.init(rawValue: errno)!)
|
throw POSIXError(.init(rawValue: errno)!)
|
||||||
}
|
}
|
||||||
// [Nucleic vendored patch] Destination full: RETURN, don't `break`. Breaking
|
// Destination full: park with the remainder in the pipe. The destination
|
||||||
// sent the outer `while true` straight back here — with the source idle and the
|
// fd's EPOLLOUT edge re-enters and resumes exactly here — never spin, and
|
||||||
// destination still full, neither leg could progress and this spun the caller's
|
// never block the shared poller thread.
|
||||||
// thread at 100% until the peer drained. The caller is an epoll handler on
|
return .idle
|
||||||
// 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 n == 0 {
|
||||||
if bytesWrote == 0 {
|
return .brokenPipe
|
||||||
return (from.offset - fromOffset, to.offset - toOffset, .brokenPipe)
|
|
||||||
}
|
|
||||||
if bytesWrote < toWrite {
|
|
||||||
break
|
|
||||||
}
|
}
|
||||||
|
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.
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -260,12 +260,19 @@ extension VsockProxy {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// `clientFile` isn't used concurrently.
|
// [Nucleic vendored patch] Each relay direction owns its own pipe and byte
|
||||||
nonisolated(unsafe) var clientFile = OSFile.SpliceFile(fd: conn.fileDescriptor)
|
// counters (see RelayDirection) — the previous shared-offset SpliceFile pair
|
||||||
nonisolated(unsafe) var eofFromClient = false
|
// let one parked direction corrupt the other's accounting and spin the poller
|
||||||
// `serverFile` isn't used concurrently.
|
// thread. Neither is used concurrently (all access is on the single
|
||||||
nonisolated(unsafe) var serverFile = OSFile.SpliceFile(fd: relayTo.fileDescriptor)
|
// ProcessSupervisor poller thread).
|
||||||
nonisolated(unsafe) var eofFromServer = false
|
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:
|
// clean up when any of these conditions apply:
|
||||||
// - the client has completely hung up or errored
|
// - the client has completely hung up or errored
|
||||||
@@ -291,20 +298,20 @@ extension VsockProxy {
|
|||||||
"vport": "\(port)",
|
"vport": "\(port)",
|
||||||
"uds": "\(path)",
|
"uds": "\(path)",
|
||||||
"action": "\(action)",
|
"action": "\(action)",
|
||||||
"eofFromClient": "\(eofFromClient)",
|
"toServerDone": "\(toServerDone)",
|
||||||
"eofFromServer": "\(eofFromServer)",
|
"toClientDone": "\(toClientDone)",
|
||||||
"clientFd": "\(clientFile.fileDescriptor)",
|
"clientFd": "\(conn.fileDescriptor)",
|
||||||
"serverFd": "\(serverFile.fileDescriptor)",
|
"serverFd": "\(relayTo.fileDescriptor)",
|
||||||
]
|
]
|
||||||
)
|
)
|
||||||
|
|
||||||
do {
|
do {
|
||||||
try ProcessSupervisor.default.unregisterFd(clientFile.fileDescriptor)
|
try ProcessSupervisor.default.unregisterFd(conn.fileDescriptor)
|
||||||
} catch {
|
} catch {
|
||||||
self.log?.error("Failed to unregister vsock proxy client fd: \(error)")
|
self.log?.error("Failed to unregister vsock proxy client fd: \(error)")
|
||||||
}
|
}
|
||||||
do {
|
do {
|
||||||
try ProcessSupervisor.default.unregisterFd(serverFile.fileDescriptor)
|
try ProcessSupervisor.default.unregisterFd(relayTo.fileDescriptor)
|
||||||
} catch {
|
} catch {
|
||||||
self.log?.error("Failed to unregister vsock proxy server fd: \(error)")
|
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,
|
// 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`;
|
// 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).
|
// `relayTo` and the first registration are released in the catch blocks below).
|
||||||
do {
|
// [Nucleic vendored patch] Interpret one relay-step outcome: returns whether
|
||||||
try ProcessSupervisor.default.registerFd(clientFile.fileDescriptor, mask: [.input, .output]) { mask in
|
// the stepped direction is now done; a broken destination ends BOTH directions
|
||||||
if mask.readyToRead && !eofFromClient {
|
// (the peer is gone). Returns rather than writing the stepped flag itself so no
|
||||||
let (fromEof, toEof) = Self.transferData(
|
// captured var is ever aliased by an inout parameter (exclusivity).
|
||||||
fromFile: &clientFile,
|
let apply = { @Sendable (outcome: RelayStepOutcome) -> Bool in
|
||||||
toFile: &serverFile,
|
switch outcome {
|
||||||
description: "readyToRead:toServer",
|
case .open: return false
|
||||||
log: self.log
|
case .finished: return true
|
||||||
)
|
case .broken:
|
||||||
eofFromClient = eofFromClient || fromEof
|
toServerDone = true
|
||||||
eofFromServer = eofFromServer || toEof
|
toClientDone = true
|
||||||
|
return true
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
if mask.readyToWrite && !eofFromServer {
|
do {
|
||||||
let (fromEof, toEof) = Self.transferData(
|
try ProcessSupervisor.default.registerFd(conn.fileDescriptor, mask: [.input, .output]) { mask in
|
||||||
fromFile: &serverFile,
|
if mask.readyToRead && !toServerDone {
|
||||||
toFile: &clientFile,
|
if apply(Self.relayStep(&toServer, description: "client:readyToRead:toServer", log: self.log)) {
|
||||||
description: "readyToWrite:toClient",
|
toServerDone = true
|
||||||
log: self.log
|
}
|
||||||
)
|
}
|
||||||
eofFromClient = eofFromClient || toEof
|
|
||||||
eofFromServer = eofFromServer || fromEof
|
// 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 {
|
if mask.isHangup {
|
||||||
eofFromClient = true
|
toServerDone = true
|
||||||
eofFromServer = true
|
toClientDone = true
|
||||||
} else if mask.isRemoteHangup && !eofFromClient {
|
} else if mask.isRemoteHangup && !toServerDone {
|
||||||
// half close, shut down client to server transfer
|
// Half close: the client sends no more. Drain the tail — relayStep
|
||||||
// we should see no more EPOLLIN events on the client fd
|
// observes the real EOF after the last buffered bytes and only then
|
||||||
// and no more EPOLLOUT events on the server fd
|
// SHUT_WRs the server, so a parked backlog is never dropped. If the
|
||||||
eofFromClient = true
|
// server is full right now the direction stays open and its EPOLLOUT
|
||||||
if shutdown(serverFile.fileDescriptor, Int32(SHUT_WR)) != 0 {
|
// edge finishes the flush.
|
||||||
self.log?.warning(
|
if apply(Self.relayStep(&toServer, description: "client:remoteHangup:toServer", log: self.log)) {
|
||||||
"failed to shut down client reads",
|
toServerDone = true
|
||||||
metadata: [
|
|
||||||
"vport": "\(self.port)",
|
|
||||||
"uds": "\(self.path)",
|
|
||||||
"errno": "\(errno)",
|
|
||||||
"eofFromClient": "\(eofFromClient)",
|
|
||||||
"eofFromServer": "\(eofFromServer)",
|
|
||||||
"clientFd": "\(clientFile.fileDescriptor)",
|
|
||||||
"serverFd": "\(serverFile.fileDescriptor)",
|
|
||||||
]
|
|
||||||
)
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
if eofFromClient && eofFromServer {
|
if toServerDone && toClientDone {
|
||||||
return cleanup()
|
return cleanup()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -384,59 +388,38 @@ extension VsockProxy {
|
|||||||
}
|
}
|
||||||
|
|
||||||
do {
|
do {
|
||||||
try ProcessSupervisor.default.registerFd(serverFile.fileDescriptor, mask: [.input, .output]) { mask in
|
try ProcessSupervisor.default.registerFd(relayTo.fileDescriptor, mask: [.input, .output]) { mask in
|
||||||
if mask.readyToRead && !eofFromServer {
|
if mask.readyToRead && !toClientDone {
|
||||||
let (fromEof, toEof) = Self.transferData(
|
if apply(Self.relayStep(&toClient, description: "server:readyToRead:toClient", log: self.log)) {
|
||||||
fromFile: &serverFile,
|
toClientDone = true
|
||||||
toFile: &clientFile,
|
}
|
||||||
description: "readyToRead:toClient",
|
|
||||||
log: self.log
|
|
||||||
)
|
|
||||||
eofFromClient = eofFromClient || toEof
|
|
||||||
eofFromServer = eofFromServer || fromEof
|
|
||||||
}
|
}
|
||||||
|
|
||||||
if mask.readyToWrite && !eofFromClient {
|
// The server drained: flush toServer's parked bytes (and whatever more the
|
||||||
let (fromEof, toEof) = Self.transferData(
|
// client has ready).
|
||||||
fromFile: &clientFile,
|
if mask.readyToWrite && !toServerDone {
|
||||||
toFile: &serverFile,
|
if apply(Self.relayStep(&toServer, description: "server:readyToWrite:toServer", log: self.log)) {
|
||||||
description: "readyToWrite:toServer",
|
toServerDone = true
|
||||||
log: self.log
|
}
|
||||||
)
|
|
||||||
eofFromClient = eofFromClient || fromEof
|
|
||||||
eofFromServer = eofFromServer || toEof
|
|
||||||
}
|
}
|
||||||
|
|
||||||
if mask.isHangup {
|
if mask.isHangup {
|
||||||
eofFromClient = true
|
toServerDone = true
|
||||||
eofFromServer = true
|
toClientDone = true
|
||||||
} else if mask.isRemoteHangup && !eofFromServer {
|
} else if mask.isRemoteHangup && !toClientDone {
|
||||||
// half close, shut down server to client transfer
|
// Half close: the server sends no more — drain the tail toward the
|
||||||
// we should see no more EPOLLIN events on the server fd
|
// client (see the client handler's mirror-image comment).
|
||||||
// and no more EPOLLOUT events on the client fd
|
if apply(Self.relayStep(&toClient, description: "server:remoteHangup:toClient", log: self.log)) {
|
||||||
eofFromServer = true
|
toClientDone = 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)",
|
|
||||||
]
|
|
||||||
)
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
if eofFromClient && eofFromServer {
|
if toServerDone && toClientDone {
|
||||||
return cleanup()
|
return cleanup()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
} catch {
|
} catch {
|
||||||
try? ProcessSupervisor.default.unregisterFd(clientFile.fileDescriptor)
|
try? ProcessSupervisor.default.unregisterFd(conn.fileDescriptor)
|
||||||
try? relayTo.close()
|
try? relayTo.close()
|
||||||
throw error
|
throw error
|
||||||
}
|
}
|
||||||
@@ -446,50 +429,65 @@ extension VsockProxy {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private static func transferData(
|
/// [Nucleic vendored patch] Outcome of one non-blocking relay pass over a direction.
|
||||||
fromFile: inout OSFile.SpliceFile,
|
enum RelayStepOutcome {
|
||||||
toFile: inout OSFile.SpliceFile,
|
/// 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,
|
description: String,
|
||||||
log: Logger?
|
log: Logger?
|
||||||
) -> (Bool, Bool) {
|
) -> RelayStepOutcome {
|
||||||
do {
|
do {
|
||||||
let (readBytes, writeBytes, action) = try OSFile.splice(from: &fromFile, to: &toFile)
|
let result = try OSFile.relay(&direction)
|
||||||
log?.trace(
|
log?.trace(
|
||||||
"transferred data",
|
"transferred data",
|
||||||
metadata: [
|
metadata: [
|
||||||
"description": "\(description)",
|
"description": "\(description)",
|
||||||
"action": "\(action)",
|
"result": "\(result)",
|
||||||
"readBytes": "\(readBytes)",
|
"pendingBytes": "\(direction.pendingBytes)",
|
||||||
"writeBytes": "\(writeBytes)",
|
"fromFd": "\(direction.from)",
|
||||||
"fromFd": "\(fromFile.fileDescriptor)",
|
"toFd": "\(direction.to)",
|
||||||
"toFd": "\(toFile.fileDescriptor)",
|
|
||||||
]
|
]
|
||||||
)
|
)
|
||||||
if action == .eof {
|
switch result {
|
||||||
// half close, shut down client to server transfer
|
case .idle:
|
||||||
// we should see no more EPOLLIN events on the client fd
|
return .open
|
||||||
// and no more EPOLLOUT events on the server fd
|
case .eof:
|
||||||
if shutdown(toFile.fileDescriptor, Int32(SHUT_WR)) != 0 {
|
if shutdown(direction.to, Int32(SHUT_WR)) != 0 {
|
||||||
log?.warning(
|
log?.warning(
|
||||||
"failed to shut down reads",
|
"failed to shut down destination writes",
|
||||||
metadata: [
|
metadata: [
|
||||||
"description": "\(description)",
|
"description": "\(description)",
|
||||||
"errno": "\(errno)",
|
"errno": "\(errno)",
|
||||||
"action": "\(action)",
|
"fromFd": "\(direction.from)",
|
||||||
"readBytes": "\(readBytes)",
|
"toFd": "\(direction.to)",
|
||||||
"writeBytes": "\(writeBytes)",
|
|
||||||
"fromFd": "\(fromFile.fileDescriptor)",
|
|
||||||
"toFd": "\(toFile.fileDescriptor)",
|
|
||||||
]
|
]
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
return (true, false)
|
return .finished
|
||||||
} else if action == .brokenPipe {
|
case .brokenPipe:
|
||||||
return (true, true)
|
return .broken
|
||||||
}
|
}
|
||||||
return (false, false)
|
|
||||||
} catch {
|
} catch {
|
||||||
return (true, true)
|
log?.error(
|
||||||
|
"relay failed: \(error)",
|
||||||
|
metadata: [
|
||||||
|
"description": "\(description)",
|
||||||
|
"fromFd": "\(direction.from)",
|
||||||
|
"toFd": "\(direction.to)",
|
||||||
|
]
|
||||||
|
)
|
||||||
|
return .broken
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user