Merge nucleic/mellow-dewy-falcon-rjhr into main
This commit is contained in:
@@ -0,0 +1,821 @@
|
||||
import Foundation
|
||||
import Logging
|
||||
import RunnerCore
|
||||
import Virtualization
|
||||
|
||||
/// Everything the orchestrator tracks about one live slot.
|
||||
public struct LiveVM: Sendable {
|
||||
/// Slot index.
|
||||
public let slot: Int
|
||||
/// The ephemeral clone backing it.
|
||||
public let bundle: VMBundle
|
||||
/// The runner name registered with Gitea. Globally unique, prefixed with
|
||||
/// ``RunnerConfig/RunnerSection/namePrefix`` — this is what lets the
|
||||
/// reconcile loop tell a stale row apart from a live one.
|
||||
public let runnerName: String
|
||||
/// The guest's IP, once its DHCP lease appears.
|
||||
public var ipAddress: String?
|
||||
/// When the boot started.
|
||||
public let startedAt: Date
|
||||
|
||||
public init(
|
||||
slot: Int,
|
||||
bundle: VMBundle,
|
||||
runnerName: String,
|
||||
ipAddress: String? = nil,
|
||||
startedAt: Date
|
||||
) {
|
||||
self.slot = slot
|
||||
self.bundle = bundle
|
||||
self.runnerName = runnerName
|
||||
self.ipAddress = ipAddress
|
||||
self.startedAt = startedAt
|
||||
}
|
||||
}
|
||||
|
||||
/// The daemon: watches Gitea, boots ephemeral macOS VMs, and cleans up after
|
||||
/// them.
|
||||
///
|
||||
/// ## Loop
|
||||
///
|
||||
/// Every `pollIntervalSeconds`:
|
||||
/// 1. Fetch queued jobs (`status=queued` only — `waiting` means *blocked*).
|
||||
/// 2. Ask ``SchedulerCore/plan(state:queuedJobs:labels:maxVMs:now:jobTimeout:bootTimeout:)``
|
||||
/// what to do. The planner is pure; all I/O happens here.
|
||||
/// 3. Execute the returned actions.
|
||||
///
|
||||
/// Every `reconcileIntervalSeconds`, additionally run ``reconcileOnce()``.
|
||||
///
|
||||
/// ## Booting a slot
|
||||
///
|
||||
/// `ensureFreeSpace` → `cloneImage(named:slotMAC:)` → ``VMInstance/start(options:)``
|
||||
/// → poll `/var/db/dhcpd_leases` for the slot MAC until `bootTimeout` →
|
||||
/// ``waitForSSH(host:port:username:password:timeout:pollInterval:)`` → write the
|
||||
/// registration token into a guest file with mode `0600` → over SSH:
|
||||
///
|
||||
/// ```sh
|
||||
/// gitea-runner register --no-interactive \
|
||||
/// --instance <url> --token-file <f> \
|
||||
/// --name <prefix><uuid> --labels "macos-arm64:host" --ephemeral \
|
||||
/// && rm -f <f> \
|
||||
/// && gitea-runner daemon
|
||||
/// ```
|
||||
///
|
||||
/// The token goes through a file rather than `--token` because arguments are
|
||||
/// visible to every process on the guest, and it is deleted the instant
|
||||
/// registration returns. `--ephemeral` is **server-enforced** (Gitea 1.24+): the
|
||||
/// server hands this runner exactly one task and then deregisters it. The weaker
|
||||
/// `--once` is runner-side only and is not used.
|
||||
///
|
||||
/// When the `gitea-runner daemon` SSH command returns — which it does after the
|
||||
/// single job completes — or when `jobTimeout` elapses, the slot is torn down:
|
||||
/// force-stop the VM, delete the clone, mark the slot idle.
|
||||
///
|
||||
/// ## Reconcile
|
||||
///
|
||||
/// A VM that dies uncleanly leaves a runner row behind, and Gitea only sweeps
|
||||
/// rows at midnight — and *never* sweeps a runner that never claimed a task. So
|
||||
/// every reconcile pass lists runners and deletes any that are `ephemeral`, not
|
||||
/// `busy`, carry our name prefix, and have no live VM. The associated Running
|
||||
/// task is reaped separately by Gitea's zombie sweep (~10–15 minutes); that part
|
||||
/// is not ours to fix.
|
||||
public actor Orchestrator {
|
||||
|
||||
/// Effective configuration.
|
||||
public let config: RunnerConfig
|
||||
/// Gitea admin API client.
|
||||
public let client: GiteaClient
|
||||
/// On-disk store.
|
||||
public let store: VMStore
|
||||
/// Base image name to clone for each job.
|
||||
public let imageName: String
|
||||
|
||||
/// Structured logger.
|
||||
private let logger: Logger
|
||||
|
||||
/// The pure scheduler's state. Every mutation goes through `SchedulerCore`.
|
||||
private var state: SchedulerState
|
||||
|
||||
/// Slots that currently hold a VM, keyed by slot index.
|
||||
private var live: [Int: LiveVM] = [:]
|
||||
|
||||
/// The `VZVirtualMachine` wrapper for each live slot.
|
||||
private var instances: [Int: VMInstance] = [:]
|
||||
|
||||
/// The supervising task per slot: clone → boot → register → run → teardown.
|
||||
private var slotTasks: [Int: Task<Void, Never>] = [:]
|
||||
|
||||
/// The "the guest stopped on its own" watcher per slot.
|
||||
private var deathWatchTasks: [Int: Task<Void, Never>] = [:]
|
||||
|
||||
/// Bumped every time a slot starts a new VM, so a stale watcher from a
|
||||
/// previous occupant of the same slot cannot trigger a teardown of the
|
||||
/// current one.
|
||||
private var slotGeneration: [Int: Int] = [:]
|
||||
|
||||
/// Slots whose teardown is in flight. Guards against the SSH command
|
||||
/// returning, the death watcher firing, and the scheduler's timeout all
|
||||
/// racing to tear the same slot down.
|
||||
private var tearingDown: Set<Int> = []
|
||||
|
||||
/// Runner names minted but not yet visible in ``live`` (the window between
|
||||
/// deciding to boot and the clone finishing). The reconcile loop must not
|
||||
/// delete a row that one of these is about to create.
|
||||
private var reservedRunnerNames: Set<String> = []
|
||||
|
||||
/// Runner names whose `gitea-runner daemon` exited cleanly, meaning Gitea
|
||||
/// already deregistered them (`--ephemeral`). Teardown skips the belt-and-
|
||||
/// braces row deletion for these.
|
||||
private var completedRunnerNames: Set<String> = []
|
||||
|
||||
/// The shared registration token, resolved once and cached for the process
|
||||
/// lifetime. Never minted per VM — see ``registrationToken()``.
|
||||
private var cachedRegistrationToken: String?
|
||||
|
||||
/// The in-flight fetch of ``cachedRegistrationToken``, if any.
|
||||
///
|
||||
/// Caching the *value* alone is not enough: `Orchestrator` is an actor, so a
|
||||
/// second slot booting during the `await` on the API call would see an empty
|
||||
/// cache and mint a second token — and minting invalidates every prior token
|
||||
/// for the scope, including the one the first VM is about to use. Memoizing
|
||||
/// the task instead makes concurrent callers share one request.
|
||||
private var registrationTokenTask: Task<String, Error>?
|
||||
|
||||
/// Set once ``shutdown()`` has begun; stops new work being accepted.
|
||||
private var isShuttingDown = false
|
||||
|
||||
/// Creates an orchestrator.
|
||||
///
|
||||
/// - Parameters:
|
||||
/// - config: Validated configuration.
|
||||
/// - client: Admin-scoped Gitea client.
|
||||
/// - store: The VM store.
|
||||
/// - imageName: Base image to clone. Defaults to `default`.
|
||||
/// - logger: Structured logger.
|
||||
public init(
|
||||
config: RunnerConfig,
|
||||
client: GiteaClient,
|
||||
store: VMStore,
|
||||
imageName: String = "default",
|
||||
logger: Logger = Logger(label: "orchestrator")
|
||||
) {
|
||||
self.config = config
|
||||
self.client = client
|
||||
self.store = store
|
||||
self.imageName = imageName
|
||||
self.logger = logger
|
||||
self.state = SchedulerState(slotCount: Orchestrator.slotCount(for: config))
|
||||
}
|
||||
|
||||
/// The fixed slot count: the configured concurrency, hard-clamped to the
|
||||
/// kernel's two-guest limit.
|
||||
private static func slotCount(for config: RunnerConfig) -> Int {
|
||||
min(
|
||||
max(1, config.scheduler.maxConcurrentVMs),
|
||||
RunnerConfig.SchedulerSection.hardMaxConcurrentVMs
|
||||
)
|
||||
}
|
||||
|
||||
/// This orchestrator's slot count.
|
||||
public var slotCount: Int { state.slots.count }
|
||||
|
||||
// MARK: - Lifecycle
|
||||
|
||||
/// Runs the poll/reconcile loop until the task is cancelled.
|
||||
///
|
||||
/// On entry it purges clones orphaned by a previous crash and runs one
|
||||
/// reconcile pass, so a restart converges before it schedules anything new.
|
||||
///
|
||||
/// On cancellation it stops accepting work and calls ``shutdown()``, so
|
||||
/// `SIGTERM` from `launchd` results in guests being asked to stop rather
|
||||
/// than being killed with their filesystems dirty.
|
||||
///
|
||||
/// - Throws: Only unrecoverable errors; transient Gitea or VM failures are
|
||||
/// logged and retried on the next tick.
|
||||
public func runForever() async throws {
|
||||
try store.ensureLayout()
|
||||
|
||||
// Clones left behind by a crash are garbage: their guests are gone and
|
||||
// their runner rows, if any, are handled by the reconcile pass below.
|
||||
do {
|
||||
try store.purgeClones()
|
||||
} catch {
|
||||
logger.warning("could not purge orphaned clones", metadata: ["error": "\(error)"])
|
||||
}
|
||||
|
||||
logger.info(
|
||||
"orchestrator starting",
|
||||
metadata: [
|
||||
"image": .string(imageName),
|
||||
"slots": .stringConvertible(slotCount),
|
||||
"labels": .string(config.runner.labels.joined(separator: ",")),
|
||||
"instance": .string(config.gitea.instanceURL.absoluteString),
|
||||
]
|
||||
)
|
||||
|
||||
await reconcileOnce()
|
||||
|
||||
await withTaskGroup(of: Void.self) { group in
|
||||
group.addTask { [pollInterval = config.scheduler.pollIntervalSeconds] in
|
||||
while !Task.isCancelled {
|
||||
await self.tick()
|
||||
do {
|
||||
try await Task.sleep(for: .seconds(max(1, pollInterval)))
|
||||
} catch {
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
group.addTask { [reconcileInterval = config.scheduler.reconcileIntervalSeconds] in
|
||||
while !Task.isCancelled {
|
||||
do {
|
||||
try await Task.sleep(for: .seconds(max(1, reconcileInterval)))
|
||||
} catch {
|
||||
break
|
||||
}
|
||||
await self.reconcileOnce()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
await shutdown()
|
||||
logger.info("orchestrator stopped")
|
||||
}
|
||||
|
||||
/// Tears down every live VM and deletes their clones. Idempotent.
|
||||
public func shutdown() async {
|
||||
guard !isShuttingDown else { return }
|
||||
isShuttingDown = true
|
||||
logger.info("shutting down", metadata: ["liveVMs": .stringConvertible(live.count)])
|
||||
|
||||
// Cancel the supervising tasks first so they stop waiting on SSH, then
|
||||
// let them run their own teardown; whatever they miss we clean up below.
|
||||
let tasks = slotTasks
|
||||
slotTasks.removeAll()
|
||||
for (_, task) in tasks { task.cancel() }
|
||||
for (_, task) in tasks { await task.value }
|
||||
|
||||
for slot in live.keys.sorted() {
|
||||
await teardownSlot(slot, reason: "daemon shutdown")
|
||||
}
|
||||
}
|
||||
|
||||
/// One iteration of the poll loop: fetch, plan, execute.
|
||||
///
|
||||
/// Exposed separately so tests and `vm boot` can drive a single tick.
|
||||
///
|
||||
/// - Parameter now: Reference time, injected for testability.
|
||||
public func tick(now: Date = Date()) async {
|
||||
guard !isShuttingDown else { return }
|
||||
|
||||
let jobs: [WorkflowJob]
|
||||
do {
|
||||
jobs = try await client.listQueuedJobs()
|
||||
} catch {
|
||||
// A Gitea outage must never take the daemon down: queued jobs wait
|
||||
// up to ABANDONED_JOB_TIMEOUT (24 h), so a missed tick costs nothing.
|
||||
logger.warning("listQueuedJobs failed", metadata: ["error": .string("\(error)")])
|
||||
return
|
||||
}
|
||||
|
||||
// `waiting` means *blocked on a dependency* in Gitea's external
|
||||
// vocabulary and must never be scheduled; only `queued` is schedulable.
|
||||
let queued = jobs.filter { $0.isQueued }
|
||||
|
||||
let (newState, actions) = SchedulerCore.plan(
|
||||
state: state,
|
||||
queuedJobs: queued,
|
||||
labels: config.labelSet,
|
||||
maxVMs: slotCount,
|
||||
now: now,
|
||||
jobTimeout: TimeInterval(config.scheduler.jobTimeoutMinutes * 60),
|
||||
bootTimeout: TimeInterval(config.scheduler.bootTimeoutSeconds)
|
||||
)
|
||||
state = newState
|
||||
|
||||
for action in actions {
|
||||
switch action {
|
||||
case .none:
|
||||
continue
|
||||
|
||||
case .teardownVM(let slot, let reason):
|
||||
// Retire the slot's lifecycle task *before* tearing down, and
|
||||
// wait for it: the planner deliberately emits teardowns before
|
||||
// boots so a timed-out slot can be recycled in this same pass,
|
||||
// and `bootSlot` refuses a slot whose `slotTasks` entry is still
|
||||
// populated. The task is parked in `executor.run` on a channel we
|
||||
// are about to kill, so it would otherwise clear that entry only
|
||||
// after the boot had already been refused.
|
||||
if let task = slotTasks.removeValue(forKey: slot) {
|
||||
task.cancel()
|
||||
// Its own teardown runs to completion here, which also means
|
||||
// it cannot race a successor booted later in this pass.
|
||||
await task.value
|
||||
}
|
||||
await teardownSlot(slot, reason: reason)
|
||||
|
||||
case .bootVM(let slot, let jobHint):
|
||||
do {
|
||||
try await bootSlot(slot, jobHint: jobHint)
|
||||
} catch {
|
||||
logger.error(
|
||||
"boot refused",
|
||||
metadata: [
|
||||
"slot": .stringConvertible(slot),
|
||||
"job": .stringConvertible(jobHint),
|
||||
"error": .string("\(error)"),
|
||||
]
|
||||
)
|
||||
state = SchedulerCore.markIdle(state: state, slot: slot)
|
||||
// The job is still queued, so the ledger would never expire
|
||||
// its entry on its own and the job would never boot again.
|
||||
state = SchedulerCore.releaseJob(state: state, jobID: jobHint)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// One reconcile pass over Gitea's runner rows.
|
||||
///
|
||||
/// Deletes runners that are ephemeral, idle, ours by name prefix, and not
|
||||
/// backed by a live VM. Conservative by construction: a row we are unsure
|
||||
/// about is left alone, because deleting a live runner would fail a job.
|
||||
public func reconcileOnce() async {
|
||||
let runners: [ActionRunner]
|
||||
do {
|
||||
runners = try await client.listRunners()
|
||||
} catch {
|
||||
logger.warning("listRunners failed", metadata: ["error": .string("\(error)")])
|
||||
return
|
||||
}
|
||||
|
||||
let ours = Set(live.values.map(\.runnerName)).union(reservedRunnerNames)
|
||||
|
||||
for runner in runners {
|
||||
guard runner.isEphemeral else { continue }
|
||||
guard !runner.isBusy else { continue }
|
||||
guard RunnerNaming.hasPrefix(runner.name, prefix: config.runner.namePrefix) else { continue }
|
||||
guard !ours.contains(runner.name) else { continue }
|
||||
|
||||
do {
|
||||
try await client.deleteRunner(id: runner.id)
|
||||
logger.info(
|
||||
"reconcile: deleted orphaned runner",
|
||||
metadata: ["name": .string(runner.name), "id": .stringConvertible(runner.id)]
|
||||
)
|
||||
} catch {
|
||||
logger.warning(
|
||||
"reconcile: could not delete runner",
|
||||
metadata: ["name": .string(runner.name), "error": .string("\(error)")]
|
||||
)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// MARK: - Slot operations
|
||||
|
||||
/// Boots, provisions, registers, and then supervises one slot.
|
||||
///
|
||||
/// Returns once the slot has been handed off to its supervising task; the
|
||||
/// job itself runs asynchronously and teardown is triggered by the SSH
|
||||
/// command returning or by `jobTimeout`.
|
||||
///
|
||||
/// - Parameters:
|
||||
/// - slot: Slot index.
|
||||
/// - jobHint: The queued job that motivated this boot — a **hint** only;
|
||||
/// the server chooses which job the runner actually claims.
|
||||
public func bootSlot(_ slot: Int, jobHint: Int64) async throws {
|
||||
guard !isShuttingDown else {
|
||||
throw CoreError.provisioningFailed("shutting down; refusing to boot slot \(slot)")
|
||||
}
|
||||
guard slot >= 0, slot < slotCount else {
|
||||
throw CoreError.provisioningFailed("slot \(slot) out of range")
|
||||
}
|
||||
guard slotTasks[slot] == nil, live[slot] == nil else {
|
||||
throw CoreError.provisioningFailed("slot \(slot) is already occupied")
|
||||
}
|
||||
|
||||
// Fail fast, on the caller's turn, for the conditions that make a boot
|
||||
// pointless: no space, no image, no token source.
|
||||
try store.ensureFreeSpace(minGB: config.storage.minFreeDiskGB)
|
||||
|
||||
let runnerName = RunnerNaming.makeRunnerName(prefix: config.runner.namePrefix)
|
||||
reservedRunnerNames.insert(runnerName)
|
||||
|
||||
let generation = (slotGeneration[slot] ?? 0) + 1
|
||||
slotGeneration[slot] = generation
|
||||
state = SchedulerCore.markProvisioning(state: state, slot: slot, jobHint: jobHint, now: Date())
|
||||
|
||||
logger.info(
|
||||
"booting VM",
|
||||
metadata: [
|
||||
"slot": .stringConvertible(slot),
|
||||
"job": .stringConvertible(jobHint),
|
||||
"runner": .string(runnerName),
|
||||
]
|
||||
)
|
||||
|
||||
slotTasks[slot] = Task { [weak self] in
|
||||
guard let self else { return }
|
||||
await self.runSlotLifecycle(slot: slot, jobHint: jobHint, runnerName: runnerName, generation: generation)
|
||||
}
|
||||
}
|
||||
|
||||
/// The whole life of one slot, from clone to teardown.
|
||||
///
|
||||
/// Every failure path funnels into the same teardown, because a slot that is
|
||||
/// neither live nor idle is a slot leaked for the process's lifetime.
|
||||
private func runSlotLifecycle(slot: Int, jobHint: Int64, runnerName: String, generation: Int) async {
|
||||
let bootTimeout = Duration.seconds(max(30, config.scheduler.bootTimeoutSeconds))
|
||||
let jobTimeout = Duration.seconds(max(60, config.scheduler.jobTimeoutMinutes * 60))
|
||||
var teardownReason = "job finished"
|
||||
|
||||
do {
|
||||
let mac = try store.macAddress(forSlot: slot, slotCount: slotCount)
|
||||
// Whatever lease this MAC already holds belongs to the *previous*
|
||||
// guest on this slot — the MACs are persistent and macOS leases last
|
||||
// 24 h. `waitForLease` must not hand that address back before the new
|
||||
// guest has even brought its NIC up.
|
||||
let priorLease = DHCPLeaseParser.lease(
|
||||
forMAC: mac,
|
||||
in: DHCPLeaseParser.parseFile()
|
||||
)
|
||||
let bundle = try store.cloneImage(named: imageName, slotMAC: mac)
|
||||
live[slot] = LiveVM(slot: slot, bundle: bundle, runnerName: runnerName, startedAt: Date())
|
||||
|
||||
let instance = try VMInstance(bundle: bundle, label: "slot-\(slot)")
|
||||
instances[slot] = instance
|
||||
try await instance.start()
|
||||
|
||||
// A guest that panics, or is shut down from inside the job, must
|
||||
// land in the same teardown path as a clean finish.
|
||||
deathWatchTasks[slot] = Task { [weak self] in
|
||||
let reason = await instance.waitUntilStopped()
|
||||
await self?.vmStoppedUnexpectedly(slot: slot, generation: generation, reason: reason)
|
||||
}
|
||||
|
||||
let ip = try await waitForLease(mac: mac, timeout: bootTimeout, replacing: priorLease)
|
||||
live[slot]?.ipAddress = ip
|
||||
logger.info("guest leased address", metadata: ["slot": .stringConvertible(slot), "ip": .string(ip)])
|
||||
|
||||
try await waitForSSH(
|
||||
host: ip,
|
||||
username: config.guest.username,
|
||||
password: config.guest.password,
|
||||
timeout: bootTimeout
|
||||
)
|
||||
|
||||
let token = try await registrationToken()
|
||||
let executor = SSHExecutor(
|
||||
host: ip,
|
||||
username: config.guest.username,
|
||||
password: config.guest.password
|
||||
)
|
||||
|
||||
// With the hint: the slot is `.provisioning` right now, which carries
|
||||
// no hint to inherit, so the hintless overload would drop it — and
|
||||
// with it both the job-timeout ledger release and every log line
|
||||
// naming which job a running slot is serving.
|
||||
state = SchedulerCore.markRunning(state: state, slot: slot, jobHint: jobHint, now: Date())
|
||||
logger.info(
|
||||
"registering ephemeral runner",
|
||||
metadata: [
|
||||
"slot": .stringConvertible(slot),
|
||||
"runner": .string(runnerName),
|
||||
"job": .stringConvertible(jobHint),
|
||||
]
|
||||
)
|
||||
|
||||
let result = try await registerAndRun(
|
||||
executor: executor,
|
||||
runnerName: runnerName,
|
||||
token: token,
|
||||
timeout: jobTimeout
|
||||
)
|
||||
await executor.close()
|
||||
|
||||
if result.succeeded {
|
||||
// A clean exit means the ephemeral runner claimed its one task,
|
||||
// finished it, and was deregistered by the server.
|
||||
completedRunnerNames.insert(runnerName)
|
||||
logger.info("job finished", metadata: ["slot": .stringConvertible(slot), "runner": .string(runnerName)])
|
||||
} else {
|
||||
teardownReason = "runner exited \(result.exitCode)"
|
||||
logger.warning(
|
||||
"runner exited non-zero",
|
||||
metadata: [
|
||||
"slot": .stringConvertible(slot),
|
||||
"exit": .stringConvertible(result.exitCode),
|
||||
"stderr": .string(String(result.stderr.suffix(500))),
|
||||
]
|
||||
)
|
||||
}
|
||||
} catch is CancellationError {
|
||||
teardownReason = "cancelled"
|
||||
} catch {
|
||||
teardownReason = "\(error)"
|
||||
logger.error(
|
||||
"slot failed",
|
||||
metadata: [
|
||||
"slot": .stringConvertible(slot),
|
||||
"runner": .string(runnerName),
|
||||
"error": .string("\(error)"),
|
||||
]
|
||||
)
|
||||
}
|
||||
|
||||
await teardownSlot(slot, reason: teardownReason)
|
||||
|
||||
// The clone may never have been adopted into `live` (a MAC or clone
|
||||
// failure throws before that), in which case teardown's own removal —
|
||||
// which lives inside `if let info` — never ran. Left behind, the name
|
||||
// would sit in the reconcile loop's "ours" set for the process lifetime.
|
||||
reservedRunnerNames.remove(runnerName)
|
||||
|
||||
// A slot that did not finish a job leaves its motivating job queued, and
|
||||
// the ledger expires entries only when a job *stops* being queued — so
|
||||
// without this the job is never booted for again.
|
||||
if teardownReason != "job finished" {
|
||||
state = SchedulerCore.releaseJob(state: state, jobID: jobHint)
|
||||
}
|
||||
|
||||
slotTasks[slot] = nil
|
||||
}
|
||||
|
||||
/// Called by the death watcher when a guest stops without us asking.
|
||||
private func vmStoppedUnexpectedly(slot: Int, generation: Int, reason: VMStopReason) async {
|
||||
guard slotGeneration[slot] == generation else { return }
|
||||
guard !tearingDown.contains(slot), live[slot] != nil else { return }
|
||||
|
||||
logger.warning(
|
||||
"guest stopped unexpectedly",
|
||||
metadata: ["slot": .stringConvertible(slot), "reason": .string("\(reason)")]
|
||||
)
|
||||
slotTasks[slot]?.cancel()
|
||||
await teardownSlot(slot, reason: "guest stopped: \(reason)")
|
||||
}
|
||||
|
||||
/// Stops the VM in a slot, deletes its clone, and marks the slot idle.
|
||||
///
|
||||
/// Best-effort and never throws: teardown that could fail would leak a slot.
|
||||
///
|
||||
/// - Parameters:
|
||||
/// - slot: Slot index.
|
||||
/// - reason: Logged cause.
|
||||
public func teardownSlot(_ slot: Int, reason: String) async {
|
||||
guard !tearingDown.contains(slot) else { return }
|
||||
|
||||
let vm = instances.removeValue(forKey: slot)
|
||||
let info = live.removeValue(forKey: slot)
|
||||
guard vm != nil || info != nil else {
|
||||
state = SchedulerCore.markIdle(state: state, slot: slot)
|
||||
return
|
||||
}
|
||||
|
||||
tearingDown.insert(slot)
|
||||
slotGeneration[slot] = (slotGeneration[slot] ?? 0) + 1
|
||||
deathWatchTasks.removeValue(forKey: slot)?.cancel()
|
||||
|
||||
logger.info("tearing down slot", metadata: ["slot": .stringConvertible(slot), "reason": .string(reason)])
|
||||
|
||||
if let vm {
|
||||
// requestStop first: the guest gets a power-button press and a
|
||||
// chance to flush before we pull the plug.
|
||||
_ = await vm.requestStopThenForce(gracePeriod: .seconds(30))
|
||||
}
|
||||
|
||||
if let info {
|
||||
do {
|
||||
try store.deleteClone(info.bundle)
|
||||
} catch {
|
||||
logger.warning(
|
||||
"could not delete clone",
|
||||
metadata: ["path": .string(info.bundle.rootURL.path), "error": .string("\(error)")]
|
||||
)
|
||||
}
|
||||
|
||||
reservedRunnerNames.remove(info.runnerName)
|
||||
|
||||
// If the runner never completed a job, Gitea will keep its row
|
||||
// forever: rows are swept at midnight, and never at all for a
|
||||
// runner that claimed no task. Delete it ourselves.
|
||||
if !completedRunnerNames.contains(info.runnerName) {
|
||||
await deleteRunnerRow(named: info.runnerName)
|
||||
}
|
||||
completedRunnerNames.remove(info.runnerName)
|
||||
}
|
||||
|
||||
state = SchedulerCore.markIdle(state: state, slot: slot)
|
||||
tearingDown.remove(slot)
|
||||
}
|
||||
|
||||
/// Best-effort deletion of a runner row by name.
|
||||
private func deleteRunnerRow(named name: String) async {
|
||||
do {
|
||||
let runners = try await client.listRunners()
|
||||
guard let row = runners.first(where: { $0.name == name }) else { return }
|
||||
guard !row.isBusy else {
|
||||
// Deleting a busy runner would fail whatever job it is running;
|
||||
// leave it for the next reconcile pass.
|
||||
logger.info("runner still busy; leaving row for reconcile", metadata: ["name": .string(name)])
|
||||
return
|
||||
}
|
||||
try await client.deleteRunner(id: row.id)
|
||||
logger.info("deleted runner row", metadata: ["name": .string(name)])
|
||||
} catch {
|
||||
logger.warning(
|
||||
"could not delete runner row",
|
||||
metadata: ["name": .string(name), "error": .string("\(error)")]
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
/// Polls `/var/db/dhcpd_leases` until the slot's MAC has an address.
|
||||
///
|
||||
/// Slot MACs are persistent and macOS leases last 24 h, so the previous
|
||||
/// guest's entry for this MAC is normally still in the file when a new clone
|
||||
/// boots. `replacing` is that entry, sampled before the guest was started;
|
||||
/// the poll holds out for a lease `bootpd` wrote afterwards rather than
|
||||
/// returning an address that belongs to a VM that no longer exists.
|
||||
///
|
||||
/// The gate is deliberately soft: if no newer lease appears within half the
|
||||
/// timeout but a stale one is present, that address is used with a warning.
|
||||
/// `bootpd` overwhelmingly reissues the same address to the same MAC, and
|
||||
/// failing a boot outright over a lease record that was merely not rewritten
|
||||
/// would be worse than the stale read this guards against.
|
||||
///
|
||||
/// - Parameters:
|
||||
/// - mac: The slot's persistent MAC.
|
||||
/// - timeout: Ceiling, from ``RunnerConfig/SchedulerSection/bootTimeoutSeconds``.
|
||||
/// - replacing: The lease seen for `mac` before the guest was started.
|
||||
/// - Returns: The guest's IPv4 address.
|
||||
/// - Throws: ``CoreError/timeout(_:)``.
|
||||
public func waitForLease(
|
||||
mac: String,
|
||||
timeout: Duration,
|
||||
replacing previous: DHCPLease? = nil
|
||||
) async throws -> String {
|
||||
let start = Date()
|
||||
let deadline = start.addingTimeInterval(timeout.seconds)
|
||||
let staleFallbackAfter = start.addingTimeInterval(timeout.seconds / 2)
|
||||
|
||||
while Date() < deadline {
|
||||
try Task.checkCancellation()
|
||||
if let lease = DHCPLeaseParser.lease(forMAC: mac, in: DHCPLeaseParser.parseFile()) {
|
||||
if DHCPLeaseParser.isNewer(lease, than: previous) {
|
||||
return lease.ipAddress
|
||||
}
|
||||
if Date() >= staleFallbackAfter {
|
||||
logger.warning(
|
||||
"no fresh dhcp lease; using the previous one for this MAC",
|
||||
metadata: ["mac": .string(mac), "ip": .string(lease.ipAddress)]
|
||||
)
|
||||
return lease.ipAddress
|
||||
}
|
||||
}
|
||||
try await Task.sleep(for: .seconds(2))
|
||||
}
|
||||
throw CoreError.timeout("dhcp lease for \(mac)")
|
||||
}
|
||||
|
||||
/// Registers an ephemeral runner in the guest and starts its daemon.
|
||||
///
|
||||
/// Blocks until the daemon exits, which — because the runner is ephemeral —
|
||||
/// happens after exactly one job.
|
||||
///
|
||||
/// - Parameters:
|
||||
/// - executor: A connected guest executor.
|
||||
/// - runnerName: The unique name to register under.
|
||||
/// - token: The shared registration token.
|
||||
/// - Returns: The daemon's exit result.
|
||||
public func registerAndRun(
|
||||
executor: any GuestExecutor,
|
||||
runnerName: String,
|
||||
token: String
|
||||
) async throws -> SSHCommandResult {
|
||||
try await registerAndRun(
|
||||
executor: executor,
|
||||
runnerName: runnerName,
|
||||
token: token,
|
||||
timeout: .seconds(max(60, config.scheduler.jobTimeoutMinutes * 60))
|
||||
)
|
||||
}
|
||||
|
||||
/// ``registerAndRun(executor:runnerName:token:)`` with an explicit ceiling on
|
||||
/// how long the runner daemon may live.
|
||||
public func registerAndRun(
|
||||
executor: any GuestExecutor,
|
||||
runnerName: String,
|
||||
token: String,
|
||||
timeout: Duration
|
||||
) async throws -> SSHCommandResult {
|
||||
// Mode 0600, and removed in the same && chain below. A registration
|
||||
// token is fleet-wide; passing it as --token would publish it to every
|
||||
// process on a guest that is about to run arbitrary repository code.
|
||||
try await executor.uploadData(Data(token.utf8), remotePath: Orchestrator.tokenPath, mode: "0600")
|
||||
|
||||
// Only bare names are stored server-side; `:host` is a register-time
|
||||
// execution hint. The guest must ship no config.yaml `runner.labels`,
|
||||
// which would silently override this.
|
||||
let labels = config.labelSet.registrationArgument(schema: "host")
|
||||
let instance = config.gitea.instanceURL.absoluteString.hasSuffix("/")
|
||||
? String(config.gitea.instanceURL.absoluteString.dropLast())
|
||||
: config.gitea.instanceURL.absoluteString
|
||||
|
||||
let command = """
|
||||
export PATH=/usr/local/bin:/usr/bin:/bin:/usr/sbin:/sbin; \
|
||||
gitea-runner register --no-interactive \
|
||||
--instance \(Orchestrator.shellQuote(instance)) \
|
||||
--token-file \(Orchestrator.tokenPath) \
|
||||
--name \(Orchestrator.shellQuote(runnerName)) \
|
||||
--labels \(Orchestrator.shellQuote(labels)) \
|
||||
--ephemeral \
|
||||
&& rm -f \(Orchestrator.tokenPath) \
|
||||
&& gitea-runner daemon
|
||||
"""
|
||||
|
||||
return try await executor.run(command, timeout: timeout)
|
||||
}
|
||||
|
||||
/// Where the registration token is staged inside the guest.
|
||||
private static let tokenPath = "/tmp/.reg-token"
|
||||
|
||||
/// Single-quotes a value for `/bin/sh`.
|
||||
private static func shellQuote(_ value: String) -> String {
|
||||
"'" + value.replacingOccurrences(of: "'", with: "'\\''") + "'"
|
||||
}
|
||||
|
||||
// MARK: - Tokens
|
||||
|
||||
/// Resolves the registration token, caching it for the process lifetime.
|
||||
///
|
||||
/// Order: `registrationTokenFile`, then `registrationToken`, then — only if
|
||||
/// ``RunnerConfig/GiteaSection/fetchRegistrationTokenViaAPI`` is set —
|
||||
/// ``GiteaClient/getRegistrationToken()``.
|
||||
///
|
||||
/// - Important: Never called per VM as a way of minting a throwaway secret.
|
||||
/// Registration tokens are reusable and scope-wide, and minting a new one
|
||||
/// invalidates every prior token for that scope — including tokens held by
|
||||
/// runners registered from other hosts. One token, cached, shared.
|
||||
/// - Throws: ``CoreError/configInvalid(_:)`` when no source is available.
|
||||
public func registrationToken() async throws -> String {
|
||||
if let cachedRegistrationToken { return cachedRegistrationToken }
|
||||
|
||||
if let staticToken = try config.resolveStaticRegistrationToken(),
|
||||
!staticToken.isEmpty {
|
||||
cachedRegistrationToken = staticToken
|
||||
return staticToken
|
||||
}
|
||||
|
||||
guard config.gitea.fetchRegistrationTokenViaAPI else {
|
||||
throw CoreError.configInvalid(
|
||||
"""
|
||||
no registration token available: set gitea.registrationTokenFile or \
|
||||
gitea.registrationToken, or enable gitea.fetchRegistrationTokenViaAPI
|
||||
"""
|
||||
)
|
||||
}
|
||||
|
||||
// Share one in-flight mint between concurrent boots. Assigned before the
|
||||
// first suspension point, so the second caller cannot miss it.
|
||||
if let inFlight = registrationTokenTask {
|
||||
return try await inFlight.value
|
||||
}
|
||||
let task = Task { [client] () -> String in
|
||||
let fetched = try await client.getRegistrationToken()
|
||||
guard !fetched.isEmpty else {
|
||||
throw CoreError.configInvalid("Gitea returned an empty registration token")
|
||||
}
|
||||
return fetched
|
||||
}
|
||||
registrationTokenTask = task
|
||||
|
||||
do {
|
||||
let fetched = try await task.value
|
||||
cachedRegistrationToken = fetched
|
||||
return fetched
|
||||
} catch {
|
||||
// A failed mint must not poison every later boot.
|
||||
registrationTokenTask = nil
|
||||
throw error
|
||||
}
|
||||
}
|
||||
|
||||
/// A snapshot of live slot state, for `vm list` and diagnostics.
|
||||
public func liveVMs() -> [LiveVM] {
|
||||
live.keys.sorted().compactMap { live[$0] }
|
||||
}
|
||||
|
||||
/// A snapshot of the scheduler's slot table.
|
||||
public func slotStates() -> [VMSlot] {
|
||||
state.slots
|
||||
}
|
||||
}
|
||||
|
||||
extension Duration {
|
||||
/// This duration as a floating-point number of seconds.
|
||||
var seconds: TimeInterval {
|
||||
let components = self.components
|
||||
return TimeInterval(components.seconds) + TimeInterval(components.attoseconds) / 1e18
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user