Files

519 lines
21 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 ContainerizationIO
import ContainerizationOS
import Foundation
import LCShim
import Logging
actor VsockProxy {
enum Action {
case listen
case dial
}
private enum SocketType {
case unix
case vsock
}
let id: String
private let path: URL
private let action: Action
private let port: UInt32
private let udsPerms: UInt32?
private let log: Logger?
private var listener: Socket?
private var task: Task<(), Never>?
private var connectionTasks: [UUID: Task<(), Never>] = [:]
init(
id: String,
action: Action,
port: UInt32,
path: URL,
udsPerms: UInt32?,
log: Logger? = nil
) {
self.id = id
self.action = action
self.port = port
self.path = path
self.udsPerms = udsPerms
self.log = log
}
}
extension VsockProxy {
func start() throws {
guard listener == nil else {
return
}
log?.debug(
"starting proxy",
metadata: [
"vport": "\(port)",
"uds": "\(path)",
"action": "\(action)",
])
switch action {
case .dial:
try dialHost()
case .listen:
try dialGuest()
}
}
func close() throws {
guard let listener else {
return
}
log?.debug(
"stopping proxy",
metadata: [
"vport": "\(port)",
"uds": "\(path)",
"action": "\(action)",
])
try listener.close()
for (_, t) in connectionTasks { t.cancel() }
connectionTasks.removeAll()
if action == .dial {
let fm = FileManager.default
if fm.fileExists(atPath: path.path) {
try fm.removeItem(at: path)
}
}
task?.cancel()
self.listener = nil
}
private func dialHost() throws {
let fm = FileManager.default
let parentDir = path.deletingLastPathComponent()
try fm.createDirectory(
at: parentDir,
withIntermediateDirectories: true
)
let type = try UnixType(
path: path.path,
perms: udsPerms,
unlinkExisting: true
)
let oldMask = umask(0)
defer { umask(oldMask) }
let uds = try Socket(type: type)
try uds.listen()
listener = uds
try acceptLoop(socketType: .unix)
}
private func dialGuest() throws {
let type = VsockType(
port: port,
cid: VsockType.anyCID
)
let vsock = try Socket(type: type)
try vsock.listen()
listener = vsock
try acceptLoop(socketType: .vsock)
}
private func acceptLoop(socketType: SocketType) throws {
guard let listener else {
return
}
let stream = try listener.acceptStream()
let task = Task {
do {
for try await conn in stream {
let connID = UUID()
let connTask = Task {
defer { self.connectionTasks[connID] = nil }
log?.debug(
"accepting connection",
metadata: [
"vport": "\(port)",
"uds": "\(path)",
"action": "\(action)",
"socketType": "\(socketType)",
])
do {
try await handleConn(
conn: conn,
connType: socketType
)
} catch {
self.log?.error("failed to handle connection: \(error)")
// [Nucleic vendored patch] A connection that failed before the relay
// owned it must be closed, or its fd leaks for the proxy's lifetime
// (accept vends closeOnDeinit: true, but the Socket is retained by the
// stream's yielded value until then — close deterministically).
try? conn.close()
}
}
// Safe: actor serialization ensures this runs before connTask can execute its defer.
connectionTasks[connID] = connTask
}
} catch {
self.log?.error("failed to accept connection: \(error)")
}
// [Nucleic vendored patch] If this loop ever ends while the proxy is still nominally
// running (fatal accept error), the listening socket MUST come down with it. Leaving it
// bound-but-unaccepted turned the relayed control socket into a silent black hole: every
// later client connect(2) SUCCEEDED into the kernel backlog and hung forever unanswered
// — for Nucleic, every session in the container stalling with "produced no output within
// 60s" until the VM was recreated. Closing the listener makes later connects fail fast
// (ECONNREFUSED/ENOENT), which callers surface and retry.
self.listenerLoopEnded()
}
self.task = task
}
/// [Nucleic vendored patch] The accept loop ended. If `close()` already ran (normal teardown)
/// this is a no-op; otherwise the listener died unexpectedly — tear it down so peers get
/// fail-fast refusals instead of connecting into a never-accepted backlog.
private func listenerLoopEnded() {
guard listener != nil else { return }
log?.error(
"proxy accept loop ended unexpectedly; closing listener",
metadata: [
"vport": "\(port)",
"uds": "\(path)",
"action": "\(action)",
])
try? close()
}
private func handleConn(
conn: ContainerizationOS.Socket,
connType: SocketType
) async throws {
try await withCheckedThrowingContinuation { (c: CheckedContinuation<Void, Error>) in
do {
// `relayTo` isn't used concurrently.
nonisolated(unsafe) var relayTo: ContainerizationOS.Socket
switch connType {
case .unix:
let type = VsockType(
port: port,
cid: VsockType.hostCID
)
relayTo = try Socket(
type: type,
closeOnDeinit: false
)
case .vsock:
let type = try UnixType(path: path.path)
relayTo = try Socket(
type: type,
closeOnDeinit: false
)
}
// [Nucleic vendored patch] A failed backend connect must close the socket it
// was dialing from: `closeOnDeinit` is false, so throwing out of here (the
// caller closes only `conn`) leaked one PID-1 fd per attempt while the
// backend was down — sustained control-plane churn walked toward EMFILE.
do {
try relayTo.connect()
} catch {
try? relayTo.close()
throw error
}
2026-07-17 23:48:59 -07:00
// [Nucleic vendored patch] BOTH fds must be non-blocking BEFORE either is
// registered. `Epoll.add` sets O_NONBLOCK only at registration time, and the first
// (client) registration's handler can fire — and splice toward the server fd —
// before the second (server) registration has made that fd non-blocking. A full
// destination then turned the splice into a genuinely BLOCKING call on vminitd's
// single ProcessSupervisor poller thread, freezing every exec's stdio and every
// control-plane relay in the container until the peer drained.
for fd in [conn.fileDescriptor, relayTo.fileDescriptor] {
let flags = fcntl(fd, F_GETFL)
if flags == -1 || fcntl(fd, F_SETFL, flags | O_NONBLOCK) == -1 {
self.log?.error(
"failed to set proxy fd non-blocking",
metadata: ["fd": "\(fd)", "errno": "\(errno)"])
}
}
// [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
// - the server has completely hung up or errored
// - both the client and server have half closed via:
// - read hangup on epoll
// - EOF on splice
//
// [Nucleic vendored patch] Hardened: (1) runs at most once — both fds' epoll
// handlers can reach the cleanup condition, and a second entry after a failed
// unregister would double-resume the continuation (a fatal trap in the guest's
// PID-1 agent); (2) every step is attempted independently — a thrown unregister
// used to SKIP the close(2)s, leaking both connection fds. Under control-plane
// connection churn those leaks accumulated until vminitd hit EMFILE, its control-
// socket accept loop died, and every session in the container stalled.
nonisolated(unsafe) var cleanedUp = false
let cleanup = { @Sendable [log, port, path, action] in
guard !cleanedUp else { return }
cleanedUp = true
log?.debug(
"cleaning up",
metadata: [
"vport": "\(port)",
"uds": "\(path)",
"action": "\(action)",
"toServerDone": "\(toServerDone)",
"toClientDone": "\(toClientDone)",
"clientFd": "\(conn.fileDescriptor)",
"serverFd": "\(relayTo.fileDescriptor)",
]
)
do {
try ProcessSupervisor.default.unregisterFd(conn.fileDescriptor)
} catch {
self.log?.error("Failed to unregister vsock proxy client fd: \(error)")
}
do {
try ProcessSupervisor.default.unregisterFd(relayTo.fileDescriptor)
} catch {
self.log?.error("Failed to unregister vsock proxy server fd: \(error)")
}
do {
try conn.close()
} catch {
self.log?.error("Failed to close vsock proxy client: \(error)")
}
do {
try relayTo.close()
} catch {
self.log?.error("Failed to close vsock proxy server: \(error)")
}
c.resume()
}
// [Nucleic vendored patch] These registrations were `try!` — an epoll_ctl failure
// (fd pressure, a stale registration) crashed vminitd, the VM's PID-1 agent,
// 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(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
}
}
// 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 {
// Full hangup of the client. Before tearing down, make one best-effort
// pass toward the SURVIVING peer: bytes already read off the client may
// be parked in toServer's pipe (its earlier flush EAGAINed), and the
// server is still healthy — dropping them loses a delivered-to-us
// control message. One pass only (no spin risk: a still-full server
// just returns and cleanup proceeds).
if !toServerDone {
_ = apply(Self.relayStep(&toServer, description: "client:hangup:toServer", log: self.log))
}
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()
}
}
} catch {
try? relayTo.close()
throw error
}
do {
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
}
}
// 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 {
// Mirror of the client handler: flush bytes parked toward the
// surviving client before teardown.
if !toClientDone {
_ = apply(Self.relayStep(&toClient, description: "server:hangup:toClient", log: self.log))
}
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(conn.fileDescriptor)
try? relayTo.close()
throw error
}
} catch {
c.resume(throwing: error)
}
}
}
/// [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?
) -> RelayStepOutcome {
do {
let result = try OSFile.relay(&direction)
log?.trace(
"transferred data",
metadata: [
"description": "\(description)",
"result": "\(result)",
"pendingBytes": "\(direction.pendingBytes)",
"fromFd": "\(direction.from)",
"toFd": "\(direction.to)",
]
)
switch result {
case .idle:
return .open
case .eof:
if shutdown(direction.to, Int32(SHUT_WR)) != 0 {
log?.warning(
"failed to shut down destination writes",
metadata: [
"description": "\(description)",
"errno": "\(errno)",
"fromFd": "\(direction.from)",
"toFd": "\(direction.to)",
]
)
}
return .finished
case .brokenPipe:
return .broken
}
} catch {
log?.error(
"relay failed: \(error)",
metadata: [
"description": "\(description)",
"fromFd": "\(direction.from)",
"toFd": "\(direction.to)",
]
)
return .broken
}
}
}
#endif