mirror of
https://github.com/cirruslabs/tart.git
synced 2026-10-01 19:51:10 +02:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
05edeac562 | ||
|
|
917c0b51ed | ||
|
|
e3b4f71f77 | ||
|
|
e92c3b900f |
@@ -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")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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))"
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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,
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user