diff --git a/Sources/ContainerK8s/Commands/K8sCreate.swift b/Sources/ContainerK8s/Commands/K8sCreate.swift index ba6e0a2ef..9752b1b35 100644 --- a/Sources/ContainerK8s/Commands/K8sCreate.swift +++ b/Sources/ContainerK8s/Commands/K8sCreate.swift @@ -30,7 +30,7 @@ public struct K8sCreate: AsyncParsableCommand { public static let configuration = CommandConfiguration( commandName: "create", - abstract: "Create and start a local Kubernetes cluster" + abstract: "Create and start a local Kubernetes cluster and worker nodes" ) @Option(name: .long, help: "Cluster name (default: \(K8sHelper.defaultName))") @@ -54,6 +54,9 @@ public struct K8sCreate: AsyncParsableCommand { @Option(name: .long, help: "Optional path to a CNI manifest to apply, or \"NONE\" to skip installing a CNI.") var cni: String? + @Option(name: .long, help: "Number of worker nodes to create (default: 0)") + var workers: UInt = 0 + public func run() async throws { LoggingSystem.bootstrap { _ in StderrLogHandler() } let log = Logger(label: K8sHelper.pluginName) @@ -73,12 +76,14 @@ public struct K8sCreate: AsyncParsableCommand { _ = try K8sHelper.kubernetesVersion(nodeImage: nodeImage) let isTTY = isatty(FileHandle.standardError.fileDescriptor) == 1 + // TotalTasks include: + // start, kubeadm init, one per worker, wait, write kubeconfig. let progressConfig = try ProgressConfig( showSpinner: isTTY, showTasks: true, showItems: true, ignoreSmallSize: true, - totalTasks: 2, // fetch image, unpack image + totalTasks: 4 + Int(workers), clearOnFinish: isTTY, outputMode: isTTY ? .ansi : .plain ) @@ -92,7 +97,7 @@ public struct K8sCreate: AsyncParsableCommand { let provisioner = try LinuxNodeProvisioner( clusterName: name, - roles: [StandardRoles.controlPlane, StandardRoles.worker], + roles: workers == 0 ? [StandardRoles.controlPlane, StandardRoles.worker] : [StandardRoles.controlPlane], nodeImage: nodeImage, cpus: resourceFlags.cpus, memory: resourceFlags.memory, @@ -106,6 +111,7 @@ public struct K8sCreate: AsyncParsableCommand { try await provisioner.provision(name: name, log: log) let client = ContainerClient() + var workerNames: [String] = [] do { let vmIP = try await provisioner.address(name: name, log: log) var sans = ["127.0.0.1"] @@ -119,6 +125,35 @@ public struct K8sCreate: AsyncParsableCommand { cniManifestPath: cni, client: client, log: log) + if workers > 0 { + let (token, caCertHash) = try await K8sHelper.createJoinToken(nodeID: name, client: client) + let controlPlaneEndpoint = "\(vmIP):\(K8sHelper.clusterContainerPort)" + for i in 1...workers { + let workerName = "\(name)-worker-\(i)" + progress.set(description: "Joining worker \(workerName)") + let workerProvisioner = try LinuxNodeProvisioner( + clusterName: name, + roles: [StandardRoles.worker], + nodeImage: nodeImage, + cpus: resourceFlags.cpus, + memory: resourceFlags.memory, + registryScheme: registryFlags.scheme, + maxConcurrentDownloads: imageFetchFlags.maxConcurrentDownloads, + remove: remove + ) + workerNames.append(workerName) + try await workerProvisioner.provision(name: workerName, log: log) + try await workerProvisioner.join( + name: workerName, controlPlaneEndpoint: controlPlaneEndpoint, + token: token, caCertHash: caCertHash, log: log) + if skipCNI { + try await workerProvisioner.waitForRegistered(name: workerName, log: log) + } else { + try await workerProvisioner.waitForReady(name: workerName, log: log) + } + } + } + if skipCNI { progress.set(description: "Waiting for API server") try await K8sHelper.waitForAPIServer(containerId: name, client: client, log: log) @@ -137,7 +172,10 @@ public struct K8sCreate: AsyncParsableCommand { log.info("cluster is running; use 'container k8s write-config --name \(name)' to write the kubeconfig") } } catch { - try? await provisioner.teardown(name: name, log: log) + for workerName in workerNames { + await Self.teardownIfOwned(workerName, client: client, provisioner: provisioner, log: log) + } + await Self.teardownIfOwned(name, client: client, provisioner: provisioner, log: log) try? K8sHelper.removeConfig(containerId: name, log: log) throw error } @@ -145,4 +183,18 @@ public struct K8sCreate: AsyncParsableCommand { progress.finish() print(name) } + + /// Tears down `containerName` only if it's a container this plugin created + private static func teardownIfOwned( + _ containerName: String, client: ContainerClient, provisioner: LinuxNodeProvisioner, log: Logger + ) async { + guard let container = try? await client.get(id: containerName) else { return } + guard container.configuration.labels[ResourceLabelKeys.plugin] == K8sHelper.pluginName else { + log.warning( + "container exists but is not owned by this plugin, refusing teardown", + metadata: ["name": "\(containerName)"]) + return + } + try? await provisioner.teardown(name: containerName, log: log) + } } diff --git a/Sources/ContainerK8s/Commands/K8sDelete.swift b/Sources/ContainerK8s/Commands/K8sDelete.swift index bb1d28bb5..28d8c4544 100644 --- a/Sources/ContainerK8s/Commands/K8sDelete.swift +++ b/Sources/ContainerK8s/Commands/K8sDelete.swift @@ -26,7 +26,7 @@ public struct K8sDelete: AsyncParsableCommand { public static let configuration = CommandConfiguration( commandName: "delete", - abstract: "Delete a Kubernetes cluster", + abstract: "Delete a Kubernetes cluster and its nodes", aliases: ["rm"] ) @@ -46,14 +46,38 @@ public struct K8sDelete: AsyncParsableCommand { } } + let workerNames = (try? await K8sHelper.workerContainerNames(clusterName: name, client: client)) ?? [] + var failures: [(name: String, error: Error)] = [] + + for workerName in workerNames { + do { + try? await client.stop(id: workerName) + try await client.delete(id: workerName) + } catch let error as ContainerizationError where error.code == .notFound { + log.debug("worker container not found, skipping delete", metadata: ["name": "\(workerName)"]) + } catch { + log.error("failed to delete worker container", metadata: ["name": "\(workerName)", "error": "\(error)"]) + failures.append((workerName, error)) + } + } + do { try? await client.stop(id: name) try await client.delete(id: name) } catch let error as ContainerizationError where error.code == .notFound { log.debug("cluster container not found, skipping delete", metadata: ["name": "\(name)"]) + } catch { + log.error("failed to delete cluster container", metadata: ["name": "\(name)", "error": "\(error)"]) + failures.append((name, error)) } try K8sHelper.removeConfig(containerId: name, log: log) + + guard failures.isEmpty else { + let details = failures.map { "\($0.name): \($0.error)" }.joined(separator: "; ") + throw ContainerizationError(.internalError, message: "failed to delete node(s): \(details)") + } + print(name) } } diff --git a/Sources/ContainerK8s/Commands/K8sLoadImage.swift b/Sources/ContainerK8s/Commands/K8sLoadImage.swift index 0cbcb89c7..800daac69 100644 --- a/Sources/ContainerK8s/Commands/K8sLoadImage.swift +++ b/Sources/ContainerK8s/Commands/K8sLoadImage.swift @@ -46,6 +46,9 @@ public struct K8sLoadImage: AsyncParsableCommand { ) var platform: String? + @Option(name: .long, help: "Optionally load into specific nodes only") + var node: [String] = [] + public func run() async throws { LoggingSystem.bootstrap { _ in StderrLogHandler() } let log = Logger(label: K8sHelper.pluginName) @@ -65,14 +68,59 @@ public struct K8sLoadImage: AsyncParsableCommand { throw ContainerizationError(.invalidArgument, message: "\(name) is not a k8s cluster") } + let workers = try await K8sHelper.workerContainerNames(clusterName: name, client: client) + let targets = try Self.resolveTargets(clusterName: name, nodes: node, workers: workers) + log.info("Saving image", metadata: ["ref": "\(image)"]) let fq = K8sHelper.fqReference(image) let resolvedPlatform = try platform.map { try Platform(from: $0) } ?? Platform(from: "linux/\(Arch.hostArchitecture().rawValue)") try await ClientImage.save(references: [fq], out: tmpFile.string, platform: resolvedPlatform, containerSystemConfig: containerSystemConfig) - log.info("Importing image into cluster", metadata: ["target": "\(name)"]) - guard let inputHandle = FileHandle(forReadingAtPath: tmpFile.string) else { - throw ContainerizationError(.internalError, message: "failed to open image tar: \(tmpFile)") + var failures: [(node: String, error: Error)] = [] + for target in targets { + do { + try await Self.importImage( + into: target, tarPath: tmpFile.string, fq: fq, image: image, client: client, log: log) + } catch { + log.error("Failed to load image", metadata: ["target": "\(target)", "error": "\(error)"]) + failures.append((node: target, error: error)) + } + } + + guard failures.isEmpty else { + let detail = failures.map { "\($0.node): \($0.error)" }.joined(separator: "; ") + throw ContainerizationError( + .internalError, message: "failed to load image into \(failures.count) node(s): \(detail)") + } + } + + /// Resolves which node containers to load the image into: `nodes` (each validated + /// as the control plane or one of `workers`) if non-empty, otherwise the control + /// plane plus every worker. + static func resolveTargets(clusterName: String, nodes: [String], workers: [String]) throws -> [String] { + guard !nodes.isEmpty else { + return [clusterName] + workers + } + var seen = Set() + var targets: [String] = [] + for node in nodes { + guard node == clusterName || workers.contains(node) else { + throw ContainerizationError( + .invalidArgument, message: "\(node) is not a node of cluster \(clusterName)") + } + if seen.insert(node).inserted { + targets.append(node) + } + } + return targets + } + + private static func importImage( + into target: String, tarPath: String, fq: String, image: String, client: ContainerClient, log: Logger + ) async throws { + log.info("Importing image into node", metadata: ["target": "\(target)"]) + guard let inputHandle = FileHandle(forReadingAtPath: tarPath) else { + throw ContainerizationError(.internalError, message: "failed to open image tar: \(tarPath)") } defer { try? inputHandle.close() } @@ -83,7 +131,7 @@ public struct K8sLoadImage: AsyncParsableCommand { terminal: false ) let importProc = try await client.createProcess( - containerId: name, + containerId: target, processId: UUID().uuidString.lowercased(), configuration: importConfig, stdio: [inputHandle, nil, nil] @@ -93,13 +141,13 @@ public struct K8sLoadImage: AsyncParsableCommand { guard importCode == 0 else { throw ContainerizationError( .internalError, - message: "ctr import exited \(importCode) on \(name)") + message: "ctr import exited \(importCode) on \(target)") } // Tag with the fully-qualified docker.io/library/ name that kubelet expects, // but only for short (unqualified) references. if fq != image { - log.info("Tagging image for kubelet", metadata: ["short": "\(image)", "fq": "\(fq)"]) + log.info("Tagging image for kubelet", metadata: ["short": "\(image)", "fq": "\(fq)", "target": "\(target)"]) let tagConfig = ProcessConfiguration( executable: Self.ctrPath, arguments: ["--namespace", "k8s.io", "images", "tag", fq, image], @@ -107,7 +155,7 @@ public struct K8sLoadImage: AsyncParsableCommand { terminal: false ) let tagProc = try await client.createProcess( - containerId: name, + containerId: target, processId: UUID().uuidString.lowercased(), configuration: tagConfig, stdio: [nil, nil, nil] @@ -117,7 +165,7 @@ public struct K8sLoadImage: AsyncParsableCommand { guard tagCode == 0 else { throw ContainerizationError( .internalError, - message: "ctr tag exited \(tagCode) on \(name)") + message: "ctr tag exited \(tagCode) on \(target)") } } } diff --git a/Sources/ContainerK8s/K8sHelper.swift b/Sources/ContainerK8s/K8sHelper.swift index 7bfd622f8..c2cca7554 100644 --- a/Sources/ContainerK8s/K8sHelper.swift +++ b/Sources/ContainerK8s/K8sHelper.swift @@ -82,6 +82,28 @@ public struct K8sHelper { return (code, String(data: data, encoding: .utf8) ?? "") } + // MARK: - Node enumeration + + /// Names of worker containers belonging to `clusterName`, sorted (e.g. `-worker-1`). + /// Only dedicated workers are returned, never the control-plane container. With no workers + /// (`--workers 0`) the control plane doubles as the worker node and this list is empty; with + /// one or more workers the control plane is tainted and is not a combined control-plane/worker node. + static func workerContainerNames(clusterName: String, client: ContainerClient) async throws -> [String] { + let snapshots = try await client.list( + filters: ContainerListFilters(labels: [ResourceLabelKeys.plugin: pluginName]) + ) + return workerContainerNames(from: snapshots, clusterName: clusterName) + } + + /// Pure filtering logic behind `workerContainerNames(clusterName:client:)`, split out for unit testing. + static func workerContainerNames(from snapshots: [ContainerSnapshot], clusterName: String) -> [String] { + snapshots + .filter { $0.configuration.labels[ResourceLabelKeys.role] == workerRoleName } + .map { $0.configuration.id } + .filter { $0.hasPrefix("\(clusterName)-worker-") } + .sorted() + } + // MARK: - List rows static func buildK8sRows(from snapshots: [ContainerSnapshot]) -> [K8sNodeResource] { diff --git a/Sources/ContainerK8s/Provisioners/LinuxNodeProvisioner.swift b/Sources/ContainerK8s/Provisioners/LinuxNodeProvisioner.swift index c50f7d02b..77467564b 100644 --- a/Sources/ContainerK8s/Provisioners/LinuxNodeProvisioner.swift +++ b/Sources/ContainerK8s/Provisioners/LinuxNodeProvisioner.swift @@ -217,6 +217,34 @@ public struct LinuxNodeProvisioner: NodeProvisioner { } } + /// Wait until the node is registered with the API server, without requiring it to be Ready. + /// Used when no CNI is installed, since nodes stay NotReady without one. + public func waitForRegistered(name: String, log: Logger) async throws { + let timeout = 180 + let client = ContainerClient() + log.info("Waiting for node to register", metadata: ["node": "\(name)"]) + for attempt in 1...timeout { + let code: Int32 + do { + code = try await K8sHelper.runProbe( + client: client, + containerId: clusterName, + arguments: ["get", "node/\(name)", "--request-timeout=2s"]) + } catch { + throw ContainerizationError( + .internalError, + message: "node \(name) stopped unexpectedly while waiting for registration: \(error)") + } + if code == 0 { return } + if attempt == timeout { + throw ContainerizationError( + .timeout, + message: "node \(name) did not register within \(timeout * 2)s") + } + try await Task.sleep(for: .seconds(2)) + } + } + public func teardown(name: String, log: Logger) async throws { let client = ContainerClient() do { diff --git a/Sources/ContainerK8s/Provisioners/NodeProvisioner.swift b/Sources/ContainerK8s/Provisioners/NodeProvisioner.swift index 706910824..3b2ecc2d2 100644 --- a/Sources/ContainerK8s/Provisioners/NodeProvisioner.swift +++ b/Sources/ContainerK8s/Provisioners/NodeProvisioner.swift @@ -33,7 +33,8 @@ public struct StandardRoles { /// /// For a **worker** node, the caller additionally calls: /// 4. `join` — run kubeadm join with the bootstrap token and CA cert hash -/// 5. `waitForReady` — poll until the node is registered and Ready in the cluster +/// 5. `waitForReady` — poll until the node is registered and Ready in the cluster, or +/// `waitForRegistered` when no CNI is installed (nodes stay NotReady without one) /// /// `K8sDelete` calls `teardown` before removing cluster containers. /// @@ -60,6 +61,10 @@ public protocol NodeProvisioner: Sendable { /// Poll until the node with `name` is registered and Ready in the cluster. func waitForReady(name: String, log: Logger) async throws + /// Poll until the node with `name` is registered in the cluster, without requiring it to be Ready. + /// Used instead of `waitForReady` when no CNI is installed. + func waitForRegistered(name: String, log: Logger) async throws + /// Remove the machine identified by `name`. func teardown(name: String, log: Logger) async throws } diff --git a/Tests/IntegrationTests/K8s/TestK8sCreateCleanupSerial.swift b/Tests/IntegrationTests/K8s/TestK8sCreateCleanupSerial.swift new file mode 100644 index 000000000..ec385fdfa --- /dev/null +++ b/Tests/IntegrationTests/K8s/TestK8sCreateCleanupSerial.swift @@ -0,0 +1,75 @@ +//===----------------------------------------------------------------------===// +// Copyright © 2026 Apple Inc. and the container 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. +//===----------------------------------------------------------------------===// + +import ContainerTestSupport +import Foundation +import Testing + +@Suite(.serialized) +struct TestK8sCreateCleanupSerial { + + /// If a container occupying the control-plane name already exists and isn't owned by + /// the k8s plugin, `k8s create` must fail (name collision) without deleting it. + @Test func testFailedCreateDoesNotTeardownUnrelatedControlPlaneNameContainer() async throws { + try await ContainerFixture.with { f in + let name = "k8s-\(f.testID)" + + try await f.doLongRun(name: name, autoRemove: false, waitUntilRunning: true) + f.addCleanup { + try? f.doStop(name) + try? f.doRemove(name, force: true) + } + let preExistingId = try f.getContainerId(name) + + try f.restoreWarmupImage(.kindestNodeV1_35_5) + let result = try f.run(["k8s", "create", "--name", name]) + print("[k8s-create-cleanup] k8s create stderr: \(result.error)") + #expect(result.status != 0) + + // The pre-existing, unrelated container must survive untouched. + #expect(try f.getContainerId(name) == preExistingId) + } + } + + /// If a container occupying a would-be worker name already exists and isn't owned by + /// the k8s plugin, the control-plane node that *was* legitimately created should be + /// torn down on failure, but the unrelated worker-named container must be left alone. + @Test func testFailedCreateDoesNotTeardownUnrelatedWorkerNameContainer() async throws { + try await ContainerFixture.with { f in + let name = "k8s-\(f.testID)" + let workerName = "\(name)-worker-1" + f.addCleanup { _ = try? f.run(["k8s", "delete", "--name", name]) } + + try await f.doLongRun(name: workerName, autoRemove: false, waitUntilRunning: true) + f.addCleanup { + try? f.doStop(workerName) + try? f.doRemove(workerName, force: true) + } + let preExistingId = try f.getContainerId(workerName) + + try f.restoreWarmupImage(.kindestNodeV1_35_5) + let result = try f.run(["k8s", "create", "--name", name, "--workers", "1"]) + print("[k8s-create-cleanup] k8s create stderr: \(result.error)") + #expect(result.status != 0) + + // The unrelated worker-named container must survive untouched. + #expect(try f.getContainerId(workerName) == preExistingId) + + // The control-plane node this run legitimately created should have been torn down. + #expect(throws: (any Error).self) { try f.getContainerStatus(name) } + } + } +} diff --git a/Tests/IntegrationTests/K8s/TestK8sMultiNodeSerial.swift b/Tests/IntegrationTests/K8s/TestK8sMultiNodeSerial.swift new file mode 100644 index 000000000..6f135e2ba --- /dev/null +++ b/Tests/IntegrationTests/K8s/TestK8sMultiNodeSerial.swift @@ -0,0 +1,122 @@ +//===----------------------------------------------------------------------===// +// Copyright © 2026 Apple Inc. and the container 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. +//===----------------------------------------------------------------------===// + +import ContainerTestSupport +import Foundation +import Testing + +@Suite(.serialized) +struct TestK8sMultiNodeSerial { + + @discardableResult + private func kubectl(_ f: ContainerFixture, node: String, args: [String]) throws -> (output: String, status: Int32) { + print("[k8s-multi] kubectl \(args.joined(separator: " ")) (node: \(node))") + let result = try f.run(["exec", node, "kubectl"] + args) + return (result.output, result.status) + } + + // "kubectl get nodes --no-headers" lines look like: + // + private func nodeRows(_ output: String) -> [(name: String, status: String, roles: String)] { + output.split(separator: "\n").compactMap { line in + let fields = line.split(separator: " ", omittingEmptySubsequences: true).map(String.init) + guard fields.count >= 3 else { return nil } + return (name: fields[0], status: fields[1], roles: fields[2]) + } + } + + @Test func testCreateWithWorkersRegistersAllNodes() async throws { + try await ContainerFixture.with { f in + let name = "k8s-\(f.testID)" + let workerNames = [1, 2].map { "\(name)-worker-\($0)" } + f.addCleanup { _ = try? f.run(["k8s", "delete", "--name", name]) } + + try f.restoreWarmupImage(.kindestNodeV1_35_5) + let result = try f.run(["k8s", "create", "--name", name, "--workers", "2"]) + if result.status != 0 { + print("[k8s-multi] k8s create stderr: \(result.error)") + f.dumpNodeDiagnostics(node: name) + } + try result.check() + #expect(result.output.contains(name)) + + // Control plane and both worker containers should all be running. + #expect(try f.getContainerStatus(name) == "running") + for workerName in workerNames { + #expect(try f.getContainerStatus(workerName) == "running") + } + + // The control plane and both workers must be registered and Ready in-cluster. + let (nodesOutput, nodesStatus) = try kubectl(f, node: name, args: ["get", "nodes", "--no-headers"]) + #expect(nodesStatus == 0) + let rows = nodeRows(nodesOutput) + #expect(rows.count == 3) + + let controlPlaneRow = rows.first { $0.name == name } + #expect(controlPlaneRow != nil) + #expect(controlPlaneRow?.roles == "control-plane") + #expect(controlPlaneRow?.status == "Ready") + + for workerName in workerNames { + let workerRow = rows.first { $0.name == workerName } + #expect(workerRow != nil) + #expect(workerRow?.roles == "") + #expect(workerRow?.status == "Ready") + } + + // container k8s list should surface every node under the same cluster. + let listResult = try f.run(["k8s", "list"]) + #expect(listResult.status == 0) + #expect(listResult.output.contains(name)) + for workerName in workerNames { + #expect(listResult.output.contains(workerName)) + } + } + } + + @Test func testDeleteRemovesAllWorkerContainers() async throws { + try await ContainerFixture.with { f in + let name = "k8s-\(f.testID)" + let workerNames = [1, 2].map { "\(name)-worker-\($0)" } + f.addCleanup { _ = try? f.run(["k8s", "delete", "--name", name]) } + + try f.restoreWarmupImage(.kindestNodeV1_35_5) + let createResult = try f.run(["k8s", "create", "--name", name, "--workers", "2"]) + if createResult.status != 0 { + print("[k8s-multi] k8s create stderr: \(createResult.error)") + f.dumpNodeDiagnostics(node: name) + } + try createResult.check() + + for workerName in workerNames { + #expect(try f.getContainerStatus(workerName) == "running") + } + + let deleteResult = try f.run(["k8s", "delete", "--name", name]) + try deleteResult.check() + + // The control plane and every worker container should be gone, not just the + // control plane — this is exactly what `K8sDelete`'s worker enumeration covers. + #expect(throws: (any Error).self) { try f.getContainerStatus(name) } + for workerName in workerNames { + #expect(throws: (any Error).self) { try f.getContainerStatus(workerName) } + } + + let listResult = try f.run(["k8s", "list"]) + #expect(!listResult.output.contains(name)) + } + } +} diff --git a/Tests/K8sPluginTests/K8sListTests.swift b/Tests/K8sPluginTests/K8sListTests.swift index 4400a9866..0f3776cae 100644 --- a/Tests/K8sPluginTests/K8sListTests.swift +++ b/Tests/K8sPluginTests/K8sListTests.swift @@ -127,6 +127,36 @@ struct K8sNodeRowTests { } } +// MARK: - K8sHelper.workerContainerNames + +@Suite("K8sHelper.workerContainerNames") +struct WorkerContainerNamesTests { + @Test func emptyInputProducesNoNames() { + #expect(K8sHelper.workerContainerNames(from: [], clusterName: "dev").isEmpty) + } + + @Test func excludesControlPlane() throws { + let cp = try makeControlPlane("dev") + let names = K8sHelper.workerContainerNames(from: [cp], clusterName: "dev") + #expect(names.isEmpty) + } + + @Test func includesMatchingWorkers() throws { + let cp = try makeControlPlane("dev") + let w1 = try makeWorker("dev-worker-1") + let w2 = try makeWorker("dev-worker-2") + let names = K8sHelper.workerContainerNames(from: [cp, w2, w1], clusterName: "dev") + #expect(names == ["dev-worker-1", "dev-worker-2"]) + } + + @Test func excludesWorkersFromOtherClusters() throws { + let w1 = try makeWorker("dev-worker-1") + let other = try makeWorker("dev2-worker-1") + let names = K8sHelper.workerContainerNames(from: [w1, other], clusterName: "dev") + #expect(names == ["dev-worker-1"]) + } +} + // MARK: - buildK8sRows ordering and grouping @Suite("K8sHelper.buildK8sRows") diff --git a/Tests/K8sPluginTests/K8sLoadImageTests.swift b/Tests/K8sPluginTests/K8sLoadImageTests.swift new file mode 100644 index 000000000..b214e486f --- /dev/null +++ b/Tests/K8sPluginTests/K8sLoadImageTests.swift @@ -0,0 +1,92 @@ +//===----------------------------------------------------------------------===// +// Copyright © 2026 Apple Inc. and the container 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. +//===----------------------------------------------------------------------===// + +import ContainerizationError +import Foundation +import Testing + +@testable import ContainerK8s + +// MARK: - K8sLoadImage flag parsing + +@Suite("K8sLoadImage --node flag") +struct K8sLoadImageNodeFlagTests { + @Test func nodeDefaultsToEmptyWhenNotProvided() throws { + let command = try K8sLoadImage.parse(["my-app:latest"]) + #expect(command.node.isEmpty) + } + + @Test func nodeCapturesProvidedValue() throws { + let command = try K8sLoadImage.parse(["--node", "k8s-dev-worker-1", "my-app:latest"]) + #expect(command.node == ["k8s-dev-worker-1"]) + } + + @Test func nodeCapturesMultipleProvidedValues() throws { + let command = try K8sLoadImage.parse([ + "--node", "k8s-dev-worker-1", "--node", "k8s-dev-worker-2", "my-app:latest", + ]) + #expect(command.node == ["k8s-dev-worker-1", "k8s-dev-worker-2"]) + } +} + +// MARK: - K8sLoadImage.resolveTargets + +@Suite("K8sLoadImage.resolveTargets") +struct ResolveTargetsTests { + @Test func noNodesReturnsControlPlaneAndAllWorkers() throws { + let targets = try K8sLoadImage.resolveTargets( + clusterName: "k8s-dev", nodes: [], workers: ["k8s-dev-worker-1", "k8s-dev-worker-2"]) + #expect(targets == ["k8s-dev", "k8s-dev-worker-1", "k8s-dev-worker-2"]) + } + + @Test func noNodesWithNoWorkersReturnsOnlyControlPlane() throws { + let targets = try K8sLoadImage.resolveTargets(clusterName: "k8s-dev", nodes: [], workers: []) + #expect(targets == ["k8s-dev"]) + } + + @Test func nodeMatchingControlPlaneReturnsJustControlPlane() throws { + let targets = try K8sLoadImage.resolveTargets( + clusterName: "k8s-dev", nodes: ["k8s-dev"], workers: ["k8s-dev-worker-1"]) + #expect(targets == ["k8s-dev"]) + } + + @Test func nodeMatchingWorkerReturnsJustThatWorker() throws { + let targets = try K8sLoadImage.resolveTargets( + clusterName: "k8s-dev", nodes: ["k8s-dev-worker-2"], workers: ["k8s-dev-worker-1", "k8s-dev-worker-2"]) + #expect(targets == ["k8s-dev-worker-2"]) + } + + @Test func multipleNodesReturnsEachInOrder() throws { + let targets = try K8sLoadImage.resolveTargets( + clusterName: "k8s-dev", nodes: ["k8s-dev-worker-2", "k8s-dev"], + workers: ["k8s-dev-worker-1", "k8s-dev-worker-2"]) + #expect(targets == ["k8s-dev-worker-2", "k8s-dev"]) + } + + @Test func duplicateNodesAreDeduplicated() throws { + let targets = try K8sLoadImage.resolveTargets( + clusterName: "k8s-dev", nodes: ["k8s-dev-worker-1", "k8s-dev-worker-1"], + workers: ["k8s-dev-worker-1"]) + #expect(targets == ["k8s-dev-worker-1"]) + } + + @Test func unknownNodeThrowsInvalidArgument() throws { + #expect(throws: ContainerizationError.self) { + try K8sLoadImage.resolveTargets( + clusterName: "k8s-dev", nodes: ["not-a-node"], workers: ["k8s-dev-worker-1"]) + } + } +} diff --git a/docs/command-reference.md b/docs/command-reference.md index 3d7cce714..08fe8f9e0 100644 --- a/docs/command-reference.md +++ b/docs/command-reference.md @@ -1621,7 +1621,7 @@ container system property list --format json ## Kubernetes Cluster Management -`container k8s` manages local single-node Kubernetes clusters backed by container VMs. Each cluster runs a Kubernetes control-plane node inside a container using `kindest/node` and `kubeadm`. +`container k8s` manages local Kubernetes clusters backed by container VMs. Each cluster runs a Kubernetes control-plane node inside a container using `kindest/node` and `kubeadm`, with optional additional worker nodes. > [!IMPORTANT] > The `k8s` command is an experimental feature and its subcommands and options are subject to change. @@ -1633,7 +1633,7 @@ Creates and starts a local Kubernetes cluster. Pulls the node image if needed, r **Usage** ```bash -container k8s create [--name ] [--node-image ] [--cni ] [--rm] [] [--debug] +container k8s create [--name ] [--node-image ] [--cni ] [--workers ] [--rm] [] [--debug] ``` **Options** @@ -1641,6 +1641,7 @@ container k8s create [--name ] [--node-image ] [--cni ] [--rm * `--name `: Cluster name (default: `k8s-dev`) * `--node-image `: Node image reference (default: `docker.io/kindest/node:v1.35.5`) * `--cni `: Optional path to a CNI manifest to apply, or `NONE` (case-insensitive) to skip installing a CNI. If not provided, the bundled kindnet CNI is used. With `NONE`, the command returns without waiting for nodes to become `Ready`, since that requires a CNI; apply your own afterward with `kubectl apply`. +* `--workers `: Number of worker nodes to create (default: `0`, meaning the control-plane node also acts as a worker) * `--rm`: Remove the cluster container after it stops **Resource Options** @@ -1671,13 +1672,18 @@ container k8s create --name temp-cluster --rm # create a cluster using a custom CNI manifest instead of the bundled kindnet container k8s create --cni ./my-cni.yaml +<<<<<<< HEAD # create a cluster with no CNI installed container k8s create --cni NONE +======= +# create a cluster with a control plane and 3 worker nodes +container k8s create --name my-cluster --workers 3 +>>>>>>> 5a62027a (Add support for multi node k8s clusters) ``` ### `container k8s delete (rm)` -Stops and deletes a Kubernetes cluster container and removes its entry from `~/.kube/config`. +Stops and deletes a Kubernetes cluster container, including any worker nodes, and removes its entry from `~/.kube/config`. **Usage** @@ -1719,12 +1725,12 @@ container k8s ls ### `container k8s load-image` -Exports an image from the local `container` image store and imports it into the cluster's containerd (in the `k8s.io` namespace) so that Kubernetes can schedule pods that reference it. +Exports an image from the local `container` image store and imports it into every node's containerd (in the `k8s.io` namespace) so that Kubernetes can schedule pods that reference it on any node. **Usage** ```bash -container k8s load-image [--name ] [--platform ] [--debug] +container k8s load-image [--name ] [--platform ] [--node ...] [--debug] ``` **Arguments** @@ -1735,18 +1741,25 @@ container k8s load-image [--name ] [--platform ] [--debu * `--name `: Cluster name (default: `k8s-dev`) * `--platform `: Platform of the image variant to load from a multi-arch image (format: os/arch[/variant], default: `linux/`). Use this when the local store contains a multi-arch manifest list and you want to select a specific variant. +* `--node `: Load into specific nodes instead of every node in the cluster. Repeat to target multiple nodes. **Examples** ```bash -# load an image into the default cluster +# load an image into every node of the default cluster container k8s load-image my-app:latest -# load an image into a named cluster +# load an image into every node of a named cluster container k8s load-image --name my-cluster my-app:latest # load the amd64 variant of a multi-arch image container k8s load-image --platform linux/amd64 my-app:latest + +# load an image into a single worker node only +container k8s load-image --node k8s-dev-worker-1 my-app:latest + +# load an image into specific worker nodes only +container k8s load-image --node k8s-dev-worker-1 --node k8s-dev-worker-2 my-app:latest ``` ### `container k8s write-config` diff --git a/docs/how-to.md b/docs/how-to.md index 13eff14ed..f2eb2632f 100644 --- a/docs/how-to.md +++ b/docs/how-to.md @@ -20,5 +20,5 @@ - [Logs](./logs.md) — container output, VM boot logs, and the `container` system's own logs. - [`config.toml` reference](./container-system-config.md) — every configuration key, its default, and how to view your merged configuration. - [Container machines](./container-machine.md) — persistent Linux environments built from OCI images, with your home directory mounted in and the filesystem surviving stop/start. -- [Kubernetes clusters](./kubernetes.md) — run local single-node Kubernetes clusters for development and testing, load your own images, and test deployments before production. +- [Kubernetes clusters](./kubernetes.md) — run local Kubernetes clusters, single-node or multi-node, for development and testing, load your own images, and test deployments before production. - [Shell completions](./shell-completions.md) — generate and install completion scripts for `zsh`, `bash`, and `fish`. diff --git a/docs/kubernetes.md b/docs/kubernetes.md index 36fb2b841..2e7882fa8 100644 --- a/docs/kubernetes.md +++ b/docs/kubernetes.md @@ -80,6 +80,17 @@ container k8s create --name high-resource --cpus 4 --memory 8g By default, clusters use 1/4 of your host's CPUs (minimum 2) and 1/4 of your host's memory (minimum 2GB). +### Multi-node clusters + +By default, `container k8s create` creates a single node that acts as both control plane and worker. Use `--workers` to add dedicated worker nodes instead: + +```bash +# Create a cluster with a control plane and 3 worker nodes +container k8s create --name my-cluster --workers 3 +``` + +When `--workers` is greater than `0`, the control-plane node is tainted so pods only schedule onto worker nodes, matching typical multi-node cluster behavior. Worker nodes share the same `--node-image`, `--cpus`, and `--memory` settings as the control plane, and are named `-worker-`. + ### Access clusters with kubectl Once a cluster is created, `kubectl` works normally: @@ -105,8 +116,14 @@ Build an image and load it into your cluster: # Build a local image container build -t my-app:latest . -# Load the image into the cluster +# Load the image into every node of the cluster container k8s load-image my-app:latest + +# Load the image into a single node instead +container k8s load-image --node k8s-dev-worker-1 my-app:latest + +# Load the image into specific nodes only +container k8s load-image --node k8s-dev-worker-1 --node k8s-dev-worker-2 my-app:latest ``` The image is placed in the `k8s.io` namespace, making it available for pod scheduling: