Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
60 changes: 56 additions & 4 deletions Sources/ContainerK8s/Commands/K8sCreate.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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))")
Expand All @@ -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)
Expand All @@ -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
)
Expand All @@ -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,
Expand All @@ -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"]
Expand All @@ -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)
Expand All @@ -137,12 +172,29 @@ 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
}

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)
}
}
26 changes: 25 additions & 1 deletion Sources/ContainerK8s/Commands/K8sDelete.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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"]
)

Expand All @@ -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)
}
}
64 changes: 56 additions & 8 deletions Sources/ContainerK8s/Commands/K8sLoadImage.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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 {
Comment thread
jglogan marked this conversation as resolved.
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<String>()
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() }

Expand All @@ -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]
Expand All @@ -93,21 +141,21 @@ 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],
environment: [],
terminal: false
)
let tagProc = try await client.createProcess(
containerId: name,
containerId: target,
processId: UUID().uuidString.lowercased(),
configuration: tagConfig,
stdio: [nil, nil, nil]
Expand All @@ -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)")
}
}
}
Expand Down
22 changes: 22 additions & 0 deletions Sources/ContainerK8s/K8sHelper.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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. `<clusterName>-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] {
Comment thread
katiewasnothere marked this conversation as resolved.
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] {
Expand Down
28 changes: 28 additions & 0 deletions Sources/ContainerK8s/Provisioners/LinuxNodeProvisioner.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
7 changes: 6 additions & 1 deletion Sources/ContainerK8s/Provisioners/NodeProvisioner.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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.
///
Expand All @@ -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
}
Loading
Loading