Compare commits

...
4 Commits
Author SHA1 Message Date
12 05edeac562 Ignore commented-out bootptab reservations (#1358) 2026-09-30 11:32:27 +00:00
edi-oai 917c0b51ed --net-host: use VZVmnetNetworkDeviceAttachment with VMNET_HOST_MODE (#1356) 2026-09-28 22:36:09 +01:00
RKSandYibo Zhuang e3b4f71f77 fix(control-socket): recover from transient accept errors (#1351)
* fix(control-socket): recover from transient accept errors

Keep the control socket alive when Darwin reports a transient fcntl failure while accepting a client connection.

Refs #1346

* test(control-socket): use synchronous pipeline lookup

* test(control-socket): query handler on event loop

* fix(control-socket): preserve channel backpressure

* test(control-socket): exercise accept error recovery

* test(control-socket): verify accepts continue after transient errors

Check that the real server inbound stream yields a client connection after
each simulated accept error. Keep the read-routing and unrelated-error
checks, bound the wait, and close the test connections.

Adapted from the regression test suggested in:
https://github.com/openai/tart/pull/1351#issuecomment-5851285933

---------

Co-authored-by: Yibo Zhuang <yzhuang@openai.com>
2026-09-28 13:22:06 -07:00
RKS e92c3b900f Protect existing VMs during clone (#1331)
* fix(clone): protect existing VMs from overwrite

* fix(clone): preserve incomplete destination directories
2026-09-28 12:58:03 -07:00
9 changed files with 299 additions and 5 deletions
+14
View File
@@ -34,6 +34,9 @@ struct Clone: AsyncParsableCommand {
@Flag(help: "create a stacked disk that uses the source image as an immutable base")
var stacked: Bool = false
@Flag(help: "overwrite an existing local VM")
var overwrite: Bool = false
@Option(help: ArgumentHelp("limit automatic pruning to n gigabytes", valueName: "n"))
var pruneLimit: UInt = 100
@@ -52,6 +55,8 @@ struct Clone: AsyncParsableCommand {
let localStorage = try VMStorageLocal()
let remoteName = try? RemoteName(sourceName)
try rejectExistingDestination(localStorage)
if stacked {
guard remoteName != nil else {
throw ValidationError("--stacked requires a remote image")
@@ -98,6 +103,8 @@ struct Clone: AsyncParsableCommand {
let lock = try FileLock(lockURL: Config().tartHomeDir)
try lock.lock()
try rejectExistingDestination(localStorage)
let sourceState = try sourceVM.state()
let generateMAC = try localStorage.hasVMsWithMACAddress(macAddress: sourceVM.macAddress())
&& sourceState != .Suspended
@@ -151,4 +158,11 @@ struct Clone: AsyncParsableCommand {
try? tmpVMDir.removeFromDisk()
})
}
private func rejectExistingDestination(_ localStorage: VMStorageLocal) throws {
let destinationURL = localStorage.baseURL.appendingPathComponent(newName, isDirectory: true)
if !overwrite && FileManager.default.fileExists(atPath: destinationURL.path) {
throw ValidationError("VM \"\(newName)\" already exists, use --overwrite to replace it")
}
}
}
+5 -3
View File
@@ -428,7 +428,7 @@ struct Run: AsyncParsableCommand {
try vmDir.regenerateMACAddress()
}
if (netSoftnet || netHost) && isInteractiveSession() {
if netSoftnet && isInteractiveSession() {
try Softnet.configureSUIDBitIfNeeded()
}
@@ -703,9 +703,11 @@ struct Run: AsyncParsableCommand {
}
if netHost {
let config = try VMConfig.init(fromURL: vmDir.configURL)
guard #available(macOS 26, *) else {
throw ValidationError("--net-host requires macOS 26 (Tahoe) or newer")
}
return try Softnet(vmMACAddress: config.macAddress.string, extraArguments: ["--vm-net-type", "host"] + softnetExtraArguments, controlFD: netSoftnetControlFd)
return try NetworkHost()
}
if netBridged.count > 0 {
+19
View File
@@ -37,6 +37,12 @@ class ControlSocket {
do {
self.serverChannel = try await ServerBootstrap(group: eventLoopGroup)
.serverChannelInitializer { channel in
channel.pipeline.addHandler(
ControlSocketAcceptErrorHandler(),
name: "ControlSocketAcceptErrorHandler"
)
}
.bind(unixDomainSocketPath: controlSocketURL.relativePath) { childChannel in
childChannel.eventLoop.makeCompletedFuture {
return try NIOAsyncChannel<ByteBuffer, ByteBuffer>(
@@ -125,3 +131,16 @@ class ControlSocket {
return fd
}
}
private final class ControlSocketAcceptErrorHandler: ChannelInboundHandler {
typealias InboundIn = Channel
typealias InboundOut = Channel
func errorCaught(context: ChannelHandlerContext, error: Error) {
if error is NIOFcntlFailedError {
context.channel.read()
} else {
context.fireErrorCaught(error)
}
}
}
@@ -18,6 +18,11 @@ struct Bootptab {
}
for line in contents.split(whereSeparator: \.isNewline) {
// Skip comment lines
guard !line.hasPrefix("#") else {
continue
}
let fields = line.split(whereSeparator: \.isWhitespace)
// Skip lines that don't look like reservation fields
+61
View File
@@ -0,0 +1,61 @@
import Foundation
import Semaphore
import Virtualization
import vmnet
@available(macOS 26, *)
class NetworkHost: Network {
private let attachment: VZNetworkDeviceAttachment
init() throws {
var status = vmnet_return_t.VMNET_SUCCESS
guard let configuration = vmnet_network_configuration_create(.VMNET_HOST_MODE, &status) else {
throw RuntimeError.Generic("Failed to create a vmnet configuration for host-only networking: \(status)")
}
defer { Unmanaged<CFTypeRef>.fromOpaque(UnsafeRawPointer(configuration)).release() }
guard let network = vmnet_network_create(configuration, &status) else {
var message = "Failed to create a vmnet network for host-only networking: \(status)"
if status == .VMNET_NOT_AUTHORIZED {
message += ". Creating a vmnet network requires root privileges or a properly signed app with the com.apple.vm.networking entitlement."
}
throw RuntimeError.Generic(message)
}
defer { Unmanaged<CFTypeRef>.fromOpaque(UnsafeRawPointer(network)).release() }
attachment = VZVmnetNetworkDeviceAttachment(network: network)
}
func attachments() -> [VZNetworkDeviceAttachment] {
[attachment]
}
func run(_ sema: AsyncSemaphore) throws {
// no-op, only used for Softnet
}
func stop() async throws {
// no-op, only used for Softnet
}
}
extension vmnet_return_t: @retroactive CustomStringConvertible {
public var description: String {
switch self {
case .VMNET_SUCCESS: return "successfully completed"
case .VMNET_FAILURE: return "general failure"
case .VMNET_MEM_FAILURE: return "memory allocation failure"
case .VMNET_INVALID_ARGUMENT: return "invalid argument specified"
case .VMNET_SETUP_INCOMPLETE: return "interface setup is not complete"
case .VMNET_INVALID_ACCESS: return "permission denied"
case .VMNET_PACKET_TOO_BIG: return "packet size larger than MTU"
case .VMNET_BUFFER_EXHAUSTED: return "buffers exhausted in kernel"
case .VMNET_TOO_MANY_PACKETS: return "packet count exceeds limit"
case .VMNET_SHARING_SERVICE_BUSY: return "vmnet interface cannot be started as conflicting sharing service is in use"
case .VMNET_NOT_AUTHORIZED: return "the operation could not be completed due to missing authorization"
@unknown default: return "unknown vmnet status (\(rawValue))"
}
}
}
+31
View File
@@ -3,6 +3,37 @@ import Network
@testable import tart
final class BootptabTests: XCTestCase {
func testCommentedReservation() throws {
let url = FileManager.default.temporaryDirectory.appendingPathComponent(UUID().uuidString)
let contents = """
%
#client1 1 02:ab:00:01:02:03 192.168.64.2
"""
try contents.write(to: url, atomically: true, encoding: .utf8)
defer { try? FileManager.default.removeItem(at: url) }
let bootptab = try XCTUnwrap(Bootptab(url))
XCTAssertNil(try bootptab.ResolveMACAddress(
macAddress: MACAddress(fromString: "02:ab:00:01:02:03")!))
}
func testCommentedReservationWithActiveReservation() throws {
let url = FileManager.default.temporaryDirectory.appendingPathComponent(UUID().uuidString)
let contents = """
%
#client1 1 02:ab:00:01:02:03 192.168.64.2
client1 1 02:ab:00:01:02:03 192.168.64.3
"""
try contents.write(to: url, atomically: true, encoding: .utf8)
defer { try? FileManager.default.removeItem(at: url) }
let bootptab = try XCTUnwrap(Bootptab(url))
XCTAssertEqual(try bootptab.ResolveMACAddress(
macAddress: MACAddress(fromString: "02:ab:00:01:02:03")!), IPv4Address("192.168.64.3"))
}
func testResolveMACAddress() throws {
let url = FileManager.default.temporaryDirectory.appendingPathComponent(UUID().uuidString)
@@ -122,6 +122,60 @@ final class CommandBehaviorTests: XCTestCase {
XCTAssertFalse(RuntimeError.VMIsRunning("running").isFileNotFound())
}
func testCloneRejectsExistingDestinationUnlessOverwriteIsRequested() async throws {
try await withTemporaryTartHome {
let source = try makeStandaloneVM(named: "source", diskContents: "source")
let destination = try makeStandaloneVM(named: "destination", diskContents: "existing")
let command = try Clone.parseAsRoot(["source", "destination"]) as! Clone
do {
try await command.run()
XCTFail("expected cloning over an existing VM to be rejected")
} catch let error as ValidationError {
XCTAssertEqual(error.message, "VM \"destination\" already exists, use --overwrite to replace it")
}
XCTAssertEqual(try Data(contentsOf: destination.diskURL), Data("existing".utf8))
XCTAssertTrue(FileManager.default.fileExists(atPath: source.diskURL.path))
let overwriteCommand = try Clone.parseAsRoot(["--overwrite", "source", "destination"]) as! Clone
try await overwriteCommand.run()
XCTAssertEqual(try Data(contentsOf: destination.diskURL), Data("source".utf8))
}
}
func testClonePreservesIncompleteDestination() async throws {
try await withTemporaryTartHome {
_ = try makeStandaloneVM(named: "source", diskContents: "source")
let storage = try VMStorageLocal()
let destination = storage.baseURL.appendingPathComponent("destination", isDirectory: true)
try FileManager.default.createDirectory(at: destination, withIntermediateDirectories: true)
let diskURL = destination.appendingPathComponent("disk.img")
try Data("existing".utf8).write(to: diskURL)
XCTAssertFalse(storage.exists("destination"))
// Check both a local source and rejection before any remote registry access.
for source in ["source", "invalid.invalid/image:latest"] {
let command = try Clone.parseAsRoot([source, "destination"]) as! Clone
do {
try await command.run()
XCTFail("expected an incomplete destination to be preserved")
} catch let error as ValidationError {
XCTAssertEqual(error.message, "VM \"destination\" already exists, use --overwrite to replace it")
}
XCTAssertEqual(try Data(contentsOf: diskURL), Data("existing".utf8))
XCTAssertFalse(storage.exists("destination"))
}
let command = try Clone.parseAsRoot(["--overwrite", "source", "destination"]) as! Clone
try await command.run()
XCTAssertEqual(try Data(contentsOf: diskURL), Data("source".utf8))
XCTAssertTrue(storage.exists("destination"))
}
}
func testSetDiskRejectsStackedVMBeforeSavingConfig() async throws {
try await withTemporaryTartHome {
let vmDir = try VMStorageLocal().create("stacked")
@@ -208,6 +262,14 @@ final class CommandBehaviorTests: XCTestCase {
)
}
private func makeStandaloneVM(named name: String, diskContents: String) throws -> VMDirectory {
let vmDir = try VMStorageLocal().create(name)
try config().save(toURL: vmDir.configURL)
XCTAssertTrue(FileManager.default.createFile(atPath: vmDir.nvramURL.path, contents: Data()))
XCTAssertTrue(FileManager.default.createFile(atPath: vmDir.diskURL.path, contents: Data(diskContents.utf8)))
return vmDir
}
private func temporaryEntries() throws -> [URL] {
try FileManager.default.contentsOfDirectory(
at: Config().tartTmpDir,
+100
View File
@@ -1,5 +1,6 @@
import NIO
import XCTest
@testable import NIOPosix
@testable import tart
// Avoid NSObject.bind and Tart's Darwin type shadowing the system function.
@@ -30,6 +31,88 @@ final class ControlSocketTests: XCTestCase {
try await eventLoopGroup.shutdownGracefully()
}
func testAcceptErrorsRetryReadingAndForwardOtherErrors() async throws {
let temporaryDirectory = try makeTemporaryDirectory()
let originalDirectory = FileManager.default.currentDirectoryPath
defer {
FileManager.default.changeCurrentDirectoryPath(originalDirectory)
try? FileManager.default.removeItem(at: temporaryDirectory)
}
let socketURL = URL(fileURLWithPath: "control.sock", relativeTo: temporaryDirectory)
let controlSocket = try await ControlSocket(socketURL)
let channel = controlSocket.serverChannel.channel
let observations = try await channel.eventLoop.submit {
let observer = AcceptErrorObserver()
// Placing this after the recovery handler also detects context.read(), which
// bypasses downstream backpressure handlers instead of starting at the tail.
try channel.pipeline.syncOperations.addHandler(observer)
for _ in 0..<3 {
channel.pipeline.fireErrorCaught(NIOFcntlFailedError())
}
let retries = observer.readCount
let transientErrors = observer.errors.count
channel.pipeline.fireErrorCaught(ChannelError.inputClosed)
return (retries, transientErrors, observer.readCount,
observer.errors.count, observer.errors.first as? ChannelError)
}.get()
try await controlSocket.serverChannel.executeThenClose { _ in }
try await controlSocket.eventLoopGroup.shutdownGracefully()
XCTAssertEqual(observations.0, 3)
XCTAssertEqual(observations.1, 0)
XCTAssertEqual(observations.2, 3)
XCTAssertEqual(observations.3, 1)
XCTAssertEqual(observations.4, .inputClosed)
}
func testAcceptErrorsDoNotEndInboundConnections() async throws {
let temporaryDirectory = try makeTemporaryDirectory()
let originalDirectory = FileManager.default.currentDirectoryPath
defer {
FileManager.default.changeCurrentDirectoryPath(originalDirectory)
try? FileManager.default.removeItem(at: temporaryDirectory)
}
let socketURL = URL(fileURLWithPath: "control.sock", relativeTo: temporaryDirectory)
let controlSocket = try await ControlSocket(socketURL)
let serverChannel = controlSocket.serverChannel
do {
try await serverChannel.executeThenClose { inbound in
// Bound the wait if the listener stays open but stops accepting connections.
let timeout = serverChannel.channel.eventLoop.scheduleTask(in: .seconds(10)) {
serverChannel.channel.close(promise: nil)
}
defer { timeout.cancel() }
var iterator = inbound.makeAsyncIterator()
for _ in 0..<3 {
try await serverChannel.channel.eventLoop.submit {
serverChannel.channel.pipeline.fireErrorCaught(NIOFcntlFailedError())
}.get()
let clientChannel = try await ClientBootstrap(group: controlSocket.eventLoopGroup)
.connectTimeout(.seconds(5))
.connect(unixDomainSocketPath: socketURL.path)
.get()
defer { clientChannel.close(promise: nil) }
// ControlSocket.run() consumes this stream. A recoverable accept error
// must not prevent it from receiving the next connection.
let nextChannel = try await iterator.next()
let acceptedChannel = try XCTUnwrap(nextChannel, "The listener stopped delivering connections after an accept error")
try await acceptedChannel.executeThenClose { _, _ in }
}
}
} catch {
try? await controlSocket.eventLoopGroup.shutdownGracefully()
throw error
}
try await controlSocket.eventLoopGroup.shutdownGracefully()
}
func testInitializerPropagatesControlSocketCreationFailure() async throws {
let temporaryDirectory = try makeTemporaryDirectory()
let originalDirectory = FileManager.default.currentDirectoryPath
@@ -94,3 +177,20 @@ final class ControlSocketTests: XCTestCase {
return directory
}
}
private final class AcceptErrorObserver: ChannelDuplexHandler {
typealias InboundIn = Channel
typealias OutboundIn = ByteBuffer
var readCount = 0
var errors: [Error] = []
func read(context: ChannelHandlerContext) {
readCount += 1
context.read()
}
func errorCaught(context: ChannelHandlerContext, error: Error) {
errors.append(error)
}
}
+2 -2
View File
@@ -46,7 +46,7 @@ func TestClonePreservesRunningDestination(t *testing.T) {
t.Fatal(err)
}
for range 2 {
_, stderr, err := tart.Tart(t, "clone", "source", "destination")
_, stderr, err := tart.Tart(t, "clone", "--overwrite", "source", "destination")
if err == nil || !strings.Contains(strings.ToLower(stderr), "running") {
t.Fatalf("clone must reject the running destination: %v: %s", err, stderr)
}
@@ -64,7 +64,7 @@ func TestClonePreservesRunningDestination(t *testing.T) {
if err := syscall.FcntlFlock(held.Fd(), syscall.F_SETLK, &lock); err != nil {
t.Fatal(err)
}
if _, stderr, err := tart.Tart(t, "clone", "source", "destination"); err != nil {
if _, stderr, err := tart.Tart(t, "clone", "--overwrite", "source", "destination"); err != nil {
t.Fatalf("clone after shutdown: %v: %s", err, stderr)
}
currentInfo, err := os.Stat(config)