905 lines
38 KiB
Swift
905 lines
38 KiB
Swift
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:reportInterval:onAttemptFailure:)`` → 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>] = [:]
|
||
|
||
/// Why a slot's supervising task was cancelled, left for that task to find.
|
||
///
|
||
/// `Task.cancel()` carries no payload and `CancellationError` no detail, so
|
||
/// a lifecycle that catches one knows only *that* it was stopped. Everything
|
||
/// worth reading — "boot timeout: provisioning for 312s (limit 300s)" — is
|
||
/// known only to the canceller. Without this hand-off the operator sees
|
||
/// `reason=cancelled`, which names the mechanism and hides the cause.
|
||
private var slotCancelReasons: [Int: String] = [:]
|
||
|
||
/// 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 (slot, _) in tasks { slotCancelReasons[slot] = "shutting down" }
|
||
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) {
|
||
// Before the cancel, not after: the task may reach its
|
||
// `catch` the instant it is cancelled.
|
||
slotCancelReasons[slot] = reason
|
||
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
|
||
// Taken alongside the `Date()` handed to the planner, so the lifecycle
|
||
// can measure against the same deadline the planner will enforce.
|
||
let provisioningStarted = ContinuousClock.now
|
||
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,
|
||
provisioningStarted: provisioningStarted
|
||
)
|
||
}
|
||
}
|
||
|
||
/// 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.
|
||
///
|
||
/// - Parameter provisioningStarted: When the *planner's* boot deadline began
|
||
/// ticking — earlier than this function's own first instruction. Used to
|
||
/// budget the readiness waits against the deadline that will actually be
|
||
/// enforced rather than against a fresh copy of it.
|
||
private func runSlotLifecycle(
|
||
slot: Int,
|
||
jobHint: Int64,
|
||
runnerName: String,
|
||
generation: Int,
|
||
provisioningStarted: ContinuousClock.Instant
|
||
) async {
|
||
let bootTimeout = Duration.seconds(max(30, config.scheduler.bootTimeoutSeconds))
|
||
let jobTimeout = Duration.seconds(max(60, config.scheduler.jobTimeoutMinutes * 60))
|
||
var teardownReason = "job finished"
|
||
|
||
// Both readiness waits below are already supervised by the planner's
|
||
// boot deadline, which started ticking at `provisioningStarted` — before
|
||
// the clone, the boot and the DHCP lease had spent any of it. Handing
|
||
// either wait the full `bootTimeout` puts its deadline strictly *after*
|
||
// the planner's, so the planner always wins the race: this task is
|
||
// cancelled mid-wait and the specific error the wait was about to throw
|
||
// ("ssh on 192.168.65.233:22 after 41 attempts; last error: …") is
|
||
// discarded in favour of a bare cancellation. Budget from what is left
|
||
// and the wait gets to speak first.
|
||
func remainingBootBudget() -> Duration {
|
||
let spent = ContinuousClock.now - provisioningStarted
|
||
// Landing a little before the planner, so its next tick finds the
|
||
// slot already failing for a stated reason. The floor keeps an
|
||
// already-overrun budget from collapsing to zero attempts, which
|
||
// would trade one useless message for another.
|
||
return max(.seconds(15), bootTimeout - spent - .seconds(5))
|
||
}
|
||
|
||
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: remainingBootBudget(),
|
||
replacing: priorLease
|
||
)
|
||
live[slot]?.ipAddress = ip
|
||
logger.info("guest leased address", metadata: ["slot": .stringConvertible(slot), "ip": .string(ip)])
|
||
|
||
// A `Logger` is a value type, so the callback below gets its own
|
||
// copy and never touches the actor — which is what lets it be a
|
||
// plain synchronous closure called from inside the poll loop.
|
||
let log = logger
|
||
try await waitForSSH(
|
||
host: ip,
|
||
username: config.guest.username,
|
||
password: config.guest.password,
|
||
timeout: remainingBootBudget(),
|
||
onAttemptFailure: { attempt in
|
||
// A guest sharing a host with other Virtualization guests
|
||
// can take minutes to start `sshd`. Without this the wait is
|
||
// indistinguishable from a hang: the log goes quiet between
|
||
// "guest leased address" and teardown, which is exactly the
|
||
// window an operator most wants to see into.
|
||
log.info(
|
||
"waiting for guest ssh",
|
||
metadata: [
|
||
"slot": .stringConvertible(slot),
|
||
"host": .string(ip),
|
||
"attempt": .stringConvertible(attempt.attempt),
|
||
"elapsed": .string("\(attempt.elapsed.components.seconds)s"),
|
||
"error": .string(attempt.error),
|
||
]
|
||
)
|
||
}
|
||
)
|
||
|
||
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 {
|
||
// Whoever cancelled us knows why; `CancellationError` does not.
|
||
// Falling back to "cancelled" only when nobody left a note.
|
||
teardownReason = slotCancelReasons[slot] ?? "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)
|
||
}
|
||
|
||
// A note left for a cancel that arrived after the lifecycle had already
|
||
// finished on its own would otherwise be read by the *next* occupant of
|
||
// this slot, mislabelling its teardown.
|
||
slotCancelReasons[slot] = nil
|
||
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)")]
|
||
)
|
||
slotCancelReasons[slot] = "guest stopped: \(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
|
||
}
|
||
}
|