Skip to content

Commit c8ff74b

Browse files
committed
Split K8sHelper, add WorkerProvisioner protocol, and add LinuxWorker implementation
1 parent af405ff commit c8ff74b

11 files changed

Lines changed: 1072 additions & 606 deletions

Sources/ContainerK8s/Commands/K8sCreate.swift

Lines changed: 75 additions & 35 deletions
Original file line numberDiff line numberDiff line change
@@ -51,6 +51,12 @@ public struct K8sCreate: AsyncParsableCommand {
5151
@Option(help: "Node image reference (default: \(K8sHelper.nodeImage))")
5252
var nodeImage: String = K8sHelper.nodeImage
5353

54+
public var worker: (any WorkerProvisioner)?
55+
56+
private enum CodingKeys: String, CodingKey {
57+
case name, remove, resourceFlags, registryFlags, imageFetchFlags, nodeImage
58+
}
59+
5460
public func run() async throws {
5561
LoggingSystem.bootstrap { _ in StderrLogHandler() }
5662
let log = Logger(label: K8sHelper.pluginName)
@@ -78,15 +84,14 @@ public struct K8sCreate: AsyncParsableCommand {
7884
try await K8sHelper.ensureImage(nodeImage: nodeImage, log: log, containerSystemConfig: containerSystemConfig)
7985

8086
let fqdn = K8sHelper.fqdn(for: name, domain: containerSystemConfig.dns.domain)
81-
let dns = Flags.DNS(domain: nil, nameservers: [], options: [], searchDomains: [])
8287

8388
let management = Flags.Management(
8489
arch: Arch.hostArchitecture().rawValue,
8590
capAdd: ["ALL"],
8691
capDrop: [],
8792
cidfile: "",
8893
detach: true,
89-
dns: dns,
94+
dns: Flags.DNS(),
9095
dnsDisabled: false,
9196
entrypoint: nil,
9297
initImage: nil,
@@ -118,7 +123,8 @@ public struct K8sCreate: AsyncParsableCommand {
118123
)
119124

120125
let updatedResource = K8sHelper.defaultedResourceFlags(resourceFlags)
121-
let processFlags = Flags.Process(cwd: nil, env: K8sHelper.nodeProxyEnv(), envFile: [], gid: nil, interactive: false, tty: false, uid: nil, ulimits: [], user: nil)
126+
var processFlags = Flags.Process()
127+
processFlags.env = K8sHelper.nodeProxyEnv()
122128

123129
var (config, kernel, initfs) = try await Utility.containerConfigFromFlags(
124130
id: name,
@@ -147,40 +153,74 @@ public struct K8sCreate: AsyncParsableCommand {
147153
initImage: initfs
148154
)
149155

150-
progress.set(description: "Starting cluster")
151-
let io = try ProcessIO.create(tty: false, interactive: false, detach: true)
152-
defer { try? io.close() }
153-
let process = try await client.bootstrap(id: name, stdio: io.stdio)
154-
try await process.start()
155-
try io.closeAfterStart()
156-
157-
progress.set(description: "Waiting for node to boot")
158-
try await K8sHelper.waitForNodeBooted(containerId: name, client: client, log: log)
159-
160-
let snapshot = try await client.get(id: name)
161-
guard let vmIP = snapshot.networks.first?.ipv4Address.address.description else {
162-
throw ContainerizationError(.internalError, message: "no VM IP for control plane \(name)")
163-
}
164-
var sans = ["127.0.0.1"]
165-
if let fqdn { sans.append(contentsOf: [vmIP, fqdn]) }
166-
167-
progress.set(description: "Running kubeadm init")
168-
try await K8sHelper.prepareNode(nodeID: name, client: client, log: log)
169-
try await K8sHelper.bootstrapControlPlane(
170-
nodeID: name, apiServerSANs: sans, advertiseAddress: vmIP,
171-
client: client, log: log)
172-
173-
progress.set(description: "Waiting for cluster to be ready")
174-
try await K8sHelper.waitForReady(containerId: name, client: client, log: log)
175-
176-
progress.set(description: "Writing kubeconfig")
156+
// From here on, clean up the cluster container if a worker step fails.
177157
do {
178-
let rawConfig = try await K8sHelper.fetchConfig(containerId: name, client: client, log: log)
179-
let kubeConfig = try await K8sHelper.transformConfig(rawConfig, containerId: name, fqdn: fqdn, client: client)
180-
try K8sHelper.mergeConfig(kubeConfig, containerId: name, setCurrentContext: true, log: log)
158+
progress.set(description: "Starting cluster")
159+
let io = try ProcessIO.create(tty: false, interactive: false, detach: true)
160+
defer { try? io.close() }
161+
let process = try await client.bootstrap(id: name, stdio: io.stdio)
162+
try await process.start()
163+
try io.closeAfterStart()
164+
165+
progress.set(description: "Waiting for node to boot")
166+
try await K8sHelper.waitForNodeBooted(containerId: name, client: client, log: log)
167+
168+
// Provision the worker before cluster init so its address can be added as a cert SAN.
169+
let workerName = "\(name)-worker-0"
170+
if let worker {
171+
progress.set(description: "Provisioning worker node")
172+
try await worker.provision(name: workerName, log: log)
173+
}
174+
175+
let snapshot = try await client.get(id: name)
176+
guard let vmIP = snapshot.networks.first?.ipv4Address.address.description else {
177+
throw ContainerizationError(.internalError, message: "no VM IP for control plane \(name)")
178+
}
179+
var sans = ["127.0.0.1"]
180+
if let fqdn { sans.append(contentsOf: [vmIP, fqdn]) }
181+
if let worker {
182+
let workerAddr = try await worker.address(name: workerName, log: log)
183+
sans.append(workerAddr)
184+
}
185+
186+
progress.set(description: "Running kubeadm init")
187+
try await K8sHelper.prepareNode(nodeID: name, client: client, log: log)
188+
try await K8sHelper.bootstrapControlPlane(
189+
nodeID: name, apiServerSANs: sans, advertiseAddress: vmIP,
190+
client: client, log: log)
191+
192+
progress.set(description: "Waiting for cluster to be ready")
193+
try await K8sHelper.waitForReady(containerId: name, client: client, log: log)
194+
195+
if let worker {
196+
let (token, caCertHash) = try await K8sHelper.createJoinToken(nodeID: name, client: client)
197+
progress.set(description: "Joining worker node")
198+
try await worker.join(
199+
name: workerName,
200+
controlPlaneEndpoint: "\(vmIP):6443",
201+
token: token,
202+
caCertHash: caCertHash,
203+
log: log)
204+
progress.set(description: "Waiting for worker node to be ready")
205+
try await worker.waitForReady(name: workerName, log: log)
206+
}
207+
208+
progress.set(description: "Writing kubeconfig")
209+
do {
210+
let rawConfig = try await K8sHelper.fetchConfig(containerId: name, client: client, log: log)
211+
let kubeConfig = try await K8sHelper.transformConfig(rawConfig, containerId: name, fqdn: fqdn, client: client)
212+
try K8sHelper.mergeConfig(kubeConfig, containerId: name, setCurrentContext: true, log: log)
213+
} catch {
214+
log.warning("failed to write kubeconfig", metadata: ["name": "\(name)", "error": "\(error)"])
215+
log.info("cluster is running; use 'container k8s write-config --name \(name)' to write the kubeconfig")
216+
}
181217
} catch {
182-
log.warning("failed to write kubeconfig", metadata: ["name": "\(name)", "error": "\(error)"])
183-
log.info("cluster is running; use 'container k8s write-config --name \(name)' to write the kubeconfig")
218+
if worker != nil {
219+
try? await client.stop(id: name)
220+
try? await client.delete(id: name)
221+
try? K8sHelper.removeConfig(containerId: name, log: log)
222+
}
223+
throw error
184224
}
185225

186226
progress.finish()

Sources/ContainerK8s/Commands/K8sDelete.swift

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,12 @@ public struct K8sDelete: AsyncParsableCommand {
3333
@Option(name: .long, help: "Cluster name (default: \(K8sHelper.defaultName))")
3434
var name: String = K8sHelper.defaultName
3535

36+
public var worker: (any WorkerProvisioner)?
37+
38+
private enum CodingKeys: String, CodingKey {
39+
case name
40+
}
41+
3642
public func run() async throws {
3743
LoggingSystem.bootstrap { _ in StderrLogHandler() }
3844
let log = Logger(label: K8sHelper.pluginName)
@@ -46,6 +52,11 @@ public struct K8sDelete: AsyncParsableCommand {
4652
}
4753
}
4854

55+
if let worker {
56+
let workerName = "\(name)-worker-0"
57+
try await worker.teardown(name: workerName, log: log)
58+
}
59+
4960
do {
5061
try? await client.stop(id: name)
5162
try await client.delete(id: name)

0 commit comments

Comments
 (0)