Nucleic: Gitea Runner macOS VM Support
This commit is contained in:
@@ -0,0 +1,733 @@
|
||||
import Foundation
|
||||
|
||||
/// The runner's on-disk configuration, loaded from
|
||||
/// `~/.config/gitea-macos-runner/config.json`.
|
||||
///
|
||||
/// Every section has defaults, and decoding tolerates missing keys, so a minimal
|
||||
/// config only needs `gitea.instanceURL` plus a way to obtain tokens. See
|
||||
/// `Resources/config.example.json` for an annotated full example.
|
||||
public struct RunnerConfig: Codable, Sendable, Equatable {
|
||||
|
||||
// MARK: - Sections
|
||||
|
||||
/// How to reach the Gitea instance and how to authenticate to it.
|
||||
public struct GiteaSection: Codable, Sendable, Equatable {
|
||||
/// Base URL of the Gitea instance, e.g. `https://gitea.example.com`.
|
||||
/// Paths are appended to this, so a trailing slash is harmless.
|
||||
public var instanceURL: URL
|
||||
|
||||
/// A Gitea admin API token, inline. Used for the admin Actions endpoints
|
||||
/// (job listing, runner listing/deletion, registration-token minting).
|
||||
/// Prefer ``adminTokenFile`` so the secret is not world-readable in JSON.
|
||||
///
|
||||
/// - Important: Exactly one of this and ``adminTokenFile`` must be set.
|
||||
/// ``RunnerConfig/validated()`` rejects both-set and neither-set alike;
|
||||
/// a stale inline token sitting beside a live token file is exactly the
|
||||
/// ambiguity that produces a baffling 401 at 3am.
|
||||
public var adminToken: String?
|
||||
|
||||
/// Path to a file whose (trimmed) contents are the admin API token.
|
||||
/// Tilde-expanded.
|
||||
///
|
||||
/// - Important: Exactly one of this and ``adminToken`` must be set — see
|
||||
/// that property. This one does *not* silently win over an inline
|
||||
/// value; setting both is a validation error.
|
||||
public var adminTokenFile: String?
|
||||
|
||||
/// The shared runner registration token, inline.
|
||||
///
|
||||
/// - Important: Registration tokens are **reusable** and **scoped**.
|
||||
/// Minting a new token for a scope invalidates all prior tokens of that
|
||||
/// scope, so per-VM tokens must never be pre-generated. One shared
|
||||
/// token serves the whole fleet. See docs/DESIGN.md, Verified Fact 4.
|
||||
public var registrationToken: String?
|
||||
|
||||
/// Path to a file whose (trimmed) contents are the registration token.
|
||||
/// Tilde-expanded. Takes precedence over ``registrationToken``.
|
||||
public var registrationTokenFile: String?
|
||||
|
||||
/// When no static registration token is configured, fetch one from
|
||||
/// `POST /api/v1/admin/actions/runners/registration-token`.
|
||||
///
|
||||
/// Defaults to `false` because that endpoint effectively returns the
|
||||
/// *existing* active token for the scope, and any implementation change
|
||||
/// that made it mint a fresh one would invalidate tokens held by runners
|
||||
/// registered elsewhere.
|
||||
public var fetchRegistrationTokenViaAPI: Bool
|
||||
|
||||
public init(
|
||||
instanceURL: URL,
|
||||
adminToken: String? = nil,
|
||||
adminTokenFile: String? = nil,
|
||||
registrationToken: String? = nil,
|
||||
registrationTokenFile: String? = nil,
|
||||
fetchRegistrationTokenViaAPI: Bool = false
|
||||
) {
|
||||
self.instanceURL = instanceURL
|
||||
self.adminToken = adminToken
|
||||
self.adminTokenFile = adminTokenFile
|
||||
self.registrationToken = registrationToken
|
||||
self.registrationTokenFile = registrationTokenFile
|
||||
self.fetchRegistrationTokenViaAPI = fetchRegistrationTokenViaAPI
|
||||
}
|
||||
}
|
||||
|
||||
/// Identity and provenance of the runners registered inside each guest.
|
||||
public struct RunnerSection: Codable, Sendable, Equatable {
|
||||
/// Bare label names this host serves. Matched case-sensitively against a
|
||||
/// job's `labels` (i.e. its `runs-on:`). The `:host` schema suffix is
|
||||
/// added only when calling `gitea-runner register`.
|
||||
public var labels: [String]
|
||||
|
||||
/// Prefix for generated runner names. Must be distinctive enough that the
|
||||
/// reconcile loop can tell our stale rows from other runners'.
|
||||
public var namePrefix: String
|
||||
|
||||
/// Template for the `gitea-runner` release asset to install in the guest.
|
||||
/// `{version}` is substituted with ``version``.
|
||||
public var runnerDownloadURL: String
|
||||
|
||||
/// The `gitea-runner` version to install (v3.x; the binary was renamed
|
||||
/// from `act_runner`, and now lives at `gitea.com/gitea/runner`).
|
||||
public var version: String
|
||||
|
||||
public init(
|
||||
labels: [String] = ["macos-arm64"],
|
||||
namePrefix: String = "macos-vm-",
|
||||
runnerDownloadURL: String = RunnerSection.defaultDownloadURLTemplate,
|
||||
version: String = "3.0.2"
|
||||
) {
|
||||
self.labels = labels
|
||||
self.namePrefix = namePrefix
|
||||
self.runnerDownloadURL = runnerDownloadURL
|
||||
self.version = version
|
||||
}
|
||||
|
||||
/// Default release-asset URL template for the darwin/arm64 build.
|
||||
public static let defaultDownloadURLTemplate =
|
||||
"https://gitea.com/gitea/runner/releases/download/v{version}/gitea-runner-{version}-darwin-arm64"
|
||||
|
||||
/// ``runnerDownloadURL`` with `{version}` substituted.
|
||||
public var resolvedDownloadURL: URL {
|
||||
get throws {
|
||||
let substituted = runnerDownloadURL.replacingOccurrences(of: "{version}", with: version)
|
||||
guard let url = URL(string: substituted), url.scheme != nil else {
|
||||
throw CoreError.configInvalid(
|
||||
"runner.runnerDownloadURL does not form a valid URL: \(substituted)")
|
||||
}
|
||||
return url
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Polling cadence, concurrency, and the timeouts that bound a stuck VM.
|
||||
public struct SchedulerSection: Codable, Sendable, Equatable {
|
||||
/// How many macOS guests may run at once.
|
||||
///
|
||||
/// - Important: Hard-clamped to 2 by ``RunnerConfig/validated()``. Apple's
|
||||
/// kernel enforces a limit of two concurrent macOS VMs per host; a third
|
||||
/// `start()` fails with `VZError.virtualMachineLimitExceeded`.
|
||||
public var maxConcurrentVMs: Int
|
||||
|
||||
/// Seconds between queued-job polls.
|
||||
public var pollIntervalSeconds: Int
|
||||
|
||||
/// Seconds between reconcile passes that sweep orphaned runner rows.
|
||||
public var reconcileIntervalSeconds: Int
|
||||
|
||||
/// Wall-clock ceiling on a single job before its VM is torn down.
|
||||
public var jobTimeoutMinutes: Int
|
||||
|
||||
/// Ceiling on boot + DHCP lease + SSH readiness before a slot is
|
||||
/// declared dead and recycled.
|
||||
public var bootTimeoutSeconds: Int
|
||||
|
||||
public init(
|
||||
maxConcurrentVMs: Int = 2,
|
||||
pollIntervalSeconds: Int = 5,
|
||||
reconcileIntervalSeconds: Int = 300,
|
||||
jobTimeoutMinutes: Int = 120,
|
||||
bootTimeoutSeconds: Int = 300
|
||||
) {
|
||||
self.maxConcurrentVMs = maxConcurrentVMs
|
||||
self.pollIntervalSeconds = pollIntervalSeconds
|
||||
self.reconcileIntervalSeconds = reconcileIntervalSeconds
|
||||
self.jobTimeoutMinutes = jobTimeoutMinutes
|
||||
self.bootTimeoutSeconds = bootTimeoutSeconds
|
||||
}
|
||||
|
||||
/// The absolute cap on concurrent macOS guests, enforced by the kernel.
|
||||
public static let hardMaxConcurrentVMs = 2
|
||||
}
|
||||
|
||||
/// Shape of each guest VM and the credentials used to reach it over SSH.
|
||||
///
|
||||
/// - Note: These credentials only ever exist on the NAT network between the
|
||||
/// host and its own ephemeral guests. They are not secrets in any
|
||||
/// meaningful sense, but they are also why the NAT attachment (rather than
|
||||
/// bridged networking) is not optional.
|
||||
public struct GuestSection: Codable, Sendable, Equatable {
|
||||
/// The admin account created by Setup Assistant automation.
|
||||
public var username: String
|
||||
/// That account's password, also used for SSH password auth.
|
||||
public var password: String
|
||||
/// Virtual CPUs per guest.
|
||||
public var cpuCount: Int
|
||||
/// RAM per guest, in gibibytes.
|
||||
public var memoryGB: Int
|
||||
/// Backing disk size per guest, in gibibytes. Sparse (ASIF) where
|
||||
/// available, so this is a ceiling rather than an allocation.
|
||||
public var diskGB: Int
|
||||
|
||||
public init(
|
||||
username: String = "admin",
|
||||
password: String = "admin",
|
||||
cpuCount: Int = 4,
|
||||
memoryGB: Int = 8,
|
||||
diskGB: Int = 64
|
||||
) {
|
||||
self.username = username
|
||||
self.password = password
|
||||
self.cpuCount = cpuCount
|
||||
self.memoryGB = memoryGB
|
||||
self.diskGB = diskGB
|
||||
}
|
||||
}
|
||||
|
||||
/// Where images, clones, IPSWs, and host state live on disk.
|
||||
public struct StorageSection: Codable, Sendable, Equatable {
|
||||
/// Root of the store. Tilde-expanded.
|
||||
///
|
||||
/// - Important: Clones are made with APFS copy-on-write, which requires
|
||||
/// source and destination on the *same volume*. Keep base images and
|
||||
/// ephemeral clones under one root.
|
||||
public var storeDir: String
|
||||
|
||||
/// Refuse to clone a new VM when the store volume has less than this
|
||||
/// much free space. CoW clones start near-free but grow as the guest
|
||||
/// writes, so a floor well above one clone's nominal size is prudent.
|
||||
public var minFreeDiskGB: Int
|
||||
|
||||
public init(
|
||||
storeDir: String = "~/Library/Application Support/gitea-macos-runner",
|
||||
minFreeDiskGB: Int = 20
|
||||
) {
|
||||
self.storeDir = storeDir
|
||||
self.minFreeDiskGB = minFreeDiskGB
|
||||
}
|
||||
}
|
||||
|
||||
// MARK: - Stored properties
|
||||
|
||||
public var gitea: GiteaSection
|
||||
public var runner: RunnerSection
|
||||
public var scheduler: SchedulerSection
|
||||
public var guest: GuestSection
|
||||
public var storage: StorageSection
|
||||
|
||||
public init(
|
||||
gitea: GiteaSection,
|
||||
runner: RunnerSection = .init(),
|
||||
scheduler: SchedulerSection = .init(),
|
||||
guest: GuestSection = .init(),
|
||||
storage: StorageSection = .init()
|
||||
) {
|
||||
self.gitea = gitea
|
||||
self.runner = runner
|
||||
self.scheduler = scheduler
|
||||
self.guest = guest
|
||||
self.storage = storage
|
||||
}
|
||||
|
||||
// MARK: - Defaults
|
||||
|
||||
/// A configuration with every default applied and a placeholder instance URL.
|
||||
/// Used by `config init` to seed a new file, and by tests.
|
||||
public static var `default`: RunnerConfig {
|
||||
RunnerConfig(gitea: GiteaSection(instanceURL: URL(string: "https://gitea.example.com")!))
|
||||
}
|
||||
|
||||
/// The conventional config path, `~/.config/gitea-macos-runner/config.json`,
|
||||
/// tilde-expanded.
|
||||
public static var defaultPath: String {
|
||||
expandTilde("~/.config/gitea-macos-runner/config.json")
|
||||
}
|
||||
|
||||
// MARK: - Loading & validation
|
||||
|
||||
/// Loads and validates a configuration from a JSON file.
|
||||
///
|
||||
/// - Parameter path: Filesystem path; tilde-expanded. Defaults to
|
||||
/// ``defaultPath``.
|
||||
/// - Returns: A validated configuration.
|
||||
/// - Throws: ``CoreError/configInvalid(_:)`` if the file is missing,
|
||||
/// unparseable, or fails ``validated()``.
|
||||
public static func load(from path: String = RunnerConfig.defaultPath) throws -> RunnerConfig {
|
||||
let expanded = expandTilde(path)
|
||||
|
||||
guard FileManager.default.fileExists(atPath: expanded) else {
|
||||
throw CoreError.configInvalid("no configuration file at \(expanded)")
|
||||
}
|
||||
|
||||
let data: Data
|
||||
do {
|
||||
data = try Data(contentsOf: URL(fileURLWithPath: expanded))
|
||||
} catch {
|
||||
throw CoreError.configInvalid("cannot read \(expanded): \(error.localizedDescription)")
|
||||
}
|
||||
|
||||
let decoded: RunnerConfig
|
||||
do {
|
||||
decoded = try JSONDecoder().decode(RunnerConfig.self, from: data)
|
||||
} catch let error as DecodingError {
|
||||
throw CoreError.configInvalid("\(expanded): \(RunnerConfig.describe(error))")
|
||||
} catch {
|
||||
throw CoreError.configInvalid("\(expanded): \(error.localizedDescription)")
|
||||
}
|
||||
|
||||
return try decoded.validated()
|
||||
}
|
||||
|
||||
/// Renders a `DecodingError` as something an operator can act on, since the
|
||||
/// default description is a multi-line dump of the underlying context.
|
||||
private static func describe(_ error: DecodingError) -> String {
|
||||
func keyPath(_ context: DecodingError.Context) -> String {
|
||||
let path = context.codingPath.map(\.stringValue).joined(separator: ".")
|
||||
return path.isEmpty ? "<root>" : path
|
||||
}
|
||||
switch error {
|
||||
case .keyNotFound(let key, let context):
|
||||
let parent = keyPath(context)
|
||||
return "missing required key `\(key.stringValue)`"
|
||||
+ (parent == "<root>" ? "" : " under `\(parent)`")
|
||||
case .typeMismatch(let type, let context):
|
||||
return "key `\(keyPath(context))` has the wrong type (expected \(type))"
|
||||
case .valueNotFound(let type, let context):
|
||||
return "key `\(keyPath(context))` is null (expected \(type))"
|
||||
case .dataCorrupted(let context):
|
||||
let path = keyPath(context)
|
||||
return path == "<root>"
|
||||
? "not valid JSON (\(context.debugDescription))"
|
||||
: "key `\(path)` is malformed (\(context.debugDescription))"
|
||||
@unknown default:
|
||||
return "\(error)"
|
||||
}
|
||||
}
|
||||
|
||||
/// Writes this configuration as pretty-printed JSON, creating parent
|
||||
/// directories as needed.
|
||||
///
|
||||
/// - Parameter path: Destination; tilde-expanded.
|
||||
public func save(to path: String) throws {
|
||||
let expanded = RunnerConfig.expandTilde(path)
|
||||
let url = URL(fileURLWithPath: expanded)
|
||||
|
||||
let encoder = JSONEncoder()
|
||||
encoder.outputFormatting = [.prettyPrinted, .sortedKeys, .withoutEscapingSlashes]
|
||||
|
||||
do {
|
||||
try FileManager.default.createDirectory(
|
||||
at: url.deletingLastPathComponent(),
|
||||
withIntermediateDirectories: true)
|
||||
var data = try encoder.encode(self)
|
||||
data.append(0x0A) // trailing newline, so the file is diff-friendly
|
||||
try data.write(to: url, options: .atomic)
|
||||
} catch {
|
||||
throw CoreError.configInvalid("cannot write \(expanded): \(error.localizedDescription)")
|
||||
}
|
||||
}
|
||||
|
||||
/// Writes the annotated example configuration shipped in `Resources/`, or —
|
||||
/// when that resource is not reachable — this configuration serialized by
|
||||
/// ``save(to:)``.
|
||||
///
|
||||
/// `config init` uses this so a fresh install lands an operator on the
|
||||
/// commented example rather than a bare JSON dump.
|
||||
///
|
||||
/// - Parameters:
|
||||
/// - path: Destination; tilde-expanded.
|
||||
/// - exampleContents: The example document, if the caller could load it.
|
||||
/// - overwrite: When `false` (the default) an existing file is left alone.
|
||||
/// - Returns: `true` if a file was written, `false` if one already existed.
|
||||
@discardableResult
|
||||
public func writeExample(
|
||||
to path: String,
|
||||
exampleContents: String? = nil,
|
||||
overwrite: Bool = false
|
||||
) throws -> Bool {
|
||||
let expanded = RunnerConfig.expandTilde(path)
|
||||
if !overwrite, FileManager.default.fileExists(atPath: expanded) {
|
||||
return false
|
||||
}
|
||||
guard let example = exampleContents else {
|
||||
try save(to: expanded)
|
||||
return true
|
||||
}
|
||||
let url = URL(fileURLWithPath: expanded)
|
||||
do {
|
||||
try FileManager.default.createDirectory(
|
||||
at: url.deletingLastPathComponent(),
|
||||
withIntermediateDirectories: true)
|
||||
try Data(example.utf8).write(to: url, options: .atomic)
|
||||
} catch {
|
||||
throw CoreError.configInvalid("cannot write \(expanded): \(error.localizedDescription)")
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
/// Returns a normalized copy, or throws describing what is wrong.
|
||||
///
|
||||
/// Normalization clamps ``SchedulerSection/maxConcurrentVMs`` into
|
||||
/// `1...2` and expands tildes in path-bearing fields. Validation rejects a
|
||||
/// non-http(s) instance URL, an empty label list, an empty name prefix,
|
||||
/// non-positive intervals or timeouts, a guest with fewer than 1 CPU or less
|
||||
/// than 1 GB of RAM, and a configuration with no way to obtain either token.
|
||||
///
|
||||
/// - Returns: The normalized configuration.
|
||||
/// - Throws: ``CoreError/configInvalid(_:)``.
|
||||
public func validated() throws -> RunnerConfig {
|
||||
var c = self
|
||||
|
||||
// --- gitea.instanceURL ------------------------------------------------
|
||||
let scheme = c.gitea.instanceURL.scheme?.lowercased()
|
||||
guard scheme == "http" || scheme == "https" else {
|
||||
throw CoreError.configInvalid(
|
||||
"gitea.instanceURL must be an http:// or https:// URL, got \"\(c.gitea.instanceURL.absoluteString)\"")
|
||||
}
|
||||
guard let host = c.gitea.instanceURL.host, !host.isEmpty else {
|
||||
throw CoreError.configInvalid(
|
||||
"gitea.instanceURL has no host: \"\(c.gitea.instanceURL.absoluteString)\"")
|
||||
}
|
||||
|
||||
// --- admin token: exactly one source ----------------------------------
|
||||
//
|
||||
// Both-set is rejected rather than silently preferring one, because a
|
||||
// stale inline token sitting next to a live token file is precisely the
|
||||
// kind of ambiguity that produces a baffling 401 at 3am.
|
||||
let inlineAdmin = RunnerConfig.nonEmpty(c.gitea.adminToken)
|
||||
let fileAdmin = RunnerConfig.nonEmpty(c.gitea.adminTokenFile)
|
||||
switch (inlineAdmin, fileAdmin) {
|
||||
case (nil, nil):
|
||||
throw CoreError.configInvalid(
|
||||
"no admin API token configured: set exactly one of gitea.adminToken or gitea.adminTokenFile")
|
||||
case (.some, .some):
|
||||
throw CoreError.configInvalid(
|
||||
"gitea.adminToken and gitea.adminTokenFile are both set: use exactly one")
|
||||
default:
|
||||
break
|
||||
}
|
||||
c.gitea.adminToken = inlineAdmin
|
||||
c.gitea.adminTokenFile = fileAdmin.map(RunnerConfig.expandTilde)
|
||||
|
||||
// --- registration token: at least one source --------------------------
|
||||
//
|
||||
// Unlike the admin token, a file and an inline value are not mutually
|
||||
// exclusive here (the file wins); what is rejected is having no source
|
||||
// at all with the API fallback switched off.
|
||||
let inlineReg = RunnerConfig.nonEmpty(c.gitea.registrationToken)
|
||||
let fileReg = RunnerConfig.nonEmpty(c.gitea.registrationTokenFile)
|
||||
if inlineReg == nil, fileReg == nil, !c.gitea.fetchRegistrationTokenViaAPI {
|
||||
throw CoreError.configInvalid(
|
||||
"no runner registration token configured: set gitea.registrationTokenFile "
|
||||
+ "(or gitea.registrationToken), or set gitea.fetchRegistrationTokenViaAPI to true")
|
||||
}
|
||||
c.gitea.registrationToken = inlineReg
|
||||
c.gitea.registrationTokenFile = fileReg.map(RunnerConfig.expandTilde)
|
||||
|
||||
// --- runner -----------------------------------------------------------
|
||||
let labels = c.runner.labels.map { $0.trimmingCharacters(in: .whitespaces) }
|
||||
guard !labels.isEmpty else {
|
||||
throw CoreError.configInvalid("runner.labels must not be empty")
|
||||
}
|
||||
if labels.contains(where: \.isEmpty) {
|
||||
throw CoreError.configInvalid("runner.labels contains an empty label name")
|
||||
}
|
||||
// Bare names only: the `:schema` suffix belongs on the `register
|
||||
// --labels` argument, never in stored config, and Gitea reports bare
|
||||
// names on jobs — so a configured "macos-arm64:host" would never match.
|
||||
if let schemed = labels.first(where: { $0.contains(":") }) {
|
||||
throw CoreError.configInvalid(
|
||||
"runner.labels must contain bare names only, but \"\(schemed)\" carries a ':schema' suffix; "
|
||||
+ "the schema is appended automatically at registration time")
|
||||
}
|
||||
c.runner.labels = labels
|
||||
|
||||
let prefix = c.runner.namePrefix.trimmingCharacters(in: .whitespaces)
|
||||
guard !prefix.isEmpty else {
|
||||
throw CoreError.configInvalid("runner.namePrefix must not be empty")
|
||||
}
|
||||
c.runner.namePrefix = prefix
|
||||
|
||||
guard !c.runner.version.trimmingCharacters(in: .whitespaces).isEmpty else {
|
||||
throw CoreError.configInvalid("runner.version must not be empty")
|
||||
}
|
||||
c.runner.version = c.runner.version.trimmingCharacters(in: .whitespaces)
|
||||
_ = try c.runner.resolvedDownloadURL
|
||||
|
||||
// --- scheduler --------------------------------------------------------
|
||||
//
|
||||
// Clamped rather than rejected: Apple's kernel caps concurrent macOS
|
||||
// guests at two, and that is not a limit a config file gets to negotiate.
|
||||
c.scheduler.maxConcurrentVMs = min(
|
||||
max(c.scheduler.maxConcurrentVMs, 1),
|
||||
SchedulerSection.hardMaxConcurrentVMs)
|
||||
|
||||
guard c.scheduler.pollIntervalSeconds > 0 else {
|
||||
throw CoreError.configInvalid("scheduler.pollIntervalSeconds must be greater than 0")
|
||||
}
|
||||
guard c.scheduler.reconcileIntervalSeconds > 0 else {
|
||||
throw CoreError.configInvalid("scheduler.reconcileIntervalSeconds must be greater than 0")
|
||||
}
|
||||
guard c.scheduler.jobTimeoutMinutes > 0 else {
|
||||
throw CoreError.configInvalid("scheduler.jobTimeoutMinutes must be greater than 0")
|
||||
}
|
||||
guard c.scheduler.bootTimeoutSeconds > 0 else {
|
||||
throw CoreError.configInvalid("scheduler.bootTimeoutSeconds must be greater than 0")
|
||||
}
|
||||
|
||||
// --- guest ------------------------------------------------------------
|
||||
guard !c.guest.username.trimmingCharacters(in: .whitespaces).isEmpty else {
|
||||
throw CoreError.configInvalid("guest.username must not be empty")
|
||||
}
|
||||
// SSH password auth is the only channel into the guest, and an empty
|
||||
// password would leave the boot hanging at authentication with no
|
||||
// diagnostic worth reading.
|
||||
guard !c.guest.password.isEmpty else {
|
||||
throw CoreError.configInvalid("guest.password must not be empty")
|
||||
}
|
||||
guard c.guest.cpuCount >= 1 else {
|
||||
throw CoreError.configInvalid("guest.cpuCount must be at least 1")
|
||||
}
|
||||
guard c.guest.memoryGB >= 1 else {
|
||||
throw CoreError.configInvalid("guest.memoryGB must be at least 1")
|
||||
}
|
||||
guard c.guest.diskGB >= 1 else {
|
||||
throw CoreError.configInvalid("guest.diskGB must be at least 1")
|
||||
}
|
||||
|
||||
// --- storage ----------------------------------------------------------
|
||||
let storeDir = c.storage.storeDir.trimmingCharacters(in: .whitespaces)
|
||||
guard !storeDir.isEmpty else {
|
||||
throw CoreError.configInvalid("storage.storeDir must not be empty")
|
||||
}
|
||||
c.storage.storeDir = RunnerConfig.expandTilde(storeDir)
|
||||
guard c.storage.minFreeDiskGB >= 0 else {
|
||||
throw CoreError.configInvalid("storage.minFreeDiskGB must not be negative")
|
||||
}
|
||||
|
||||
return c
|
||||
}
|
||||
|
||||
/// Trims a string and maps `""` to `nil`, so an empty JSON value reads as
|
||||
/// "not configured" rather than as a zero-length token.
|
||||
private static func nonEmpty(_ value: String?) -> String? {
|
||||
guard let trimmed = value?.trimmingCharacters(in: .whitespacesAndNewlines),
|
||||
!trimmed.isEmpty
|
||||
else { return nil }
|
||||
return trimmed
|
||||
}
|
||||
|
||||
/// The admin API token, resolved from ``GiteaSection/adminTokenFile`` (read
|
||||
/// and trimmed) or ``GiteaSection/adminToken``.
|
||||
///
|
||||
/// On a configuration that has been through ``validated()`` exactly one of
|
||||
/// those is set, so the file-first order here never actually chooses between
|
||||
/// two live values.
|
||||
///
|
||||
/// - Returns: The token, or `nil` when neither source is configured.
|
||||
public func resolveAdminToken() throws -> String? {
|
||||
if let path = RunnerConfig.nonEmpty(gitea.adminTokenFile) {
|
||||
return try RunnerConfig.readTokenFile(path, describedAs: "gitea.adminTokenFile")
|
||||
}
|
||||
return RunnerConfig.nonEmpty(gitea.adminToken)
|
||||
}
|
||||
|
||||
/// The registration token from static configuration only — file first, then
|
||||
/// inline value. Returns `nil` when the caller must fall back to the API
|
||||
/// (see ``GiteaSection/fetchRegistrationTokenViaAPI``).
|
||||
public func resolveStaticRegistrationToken() throws -> String? {
|
||||
if let path = RunnerConfig.nonEmpty(gitea.registrationTokenFile) {
|
||||
return try RunnerConfig.readTokenFile(path, describedAs: "gitea.registrationTokenFile")
|
||||
}
|
||||
return RunnerConfig.nonEmpty(gitea.registrationToken)
|
||||
}
|
||||
|
||||
/// Reads a secret from a file: tilde-expanded, trimmed of surrounding
|
||||
/// whitespace and newlines (an `echo`-written token file always has one).
|
||||
///
|
||||
/// - Throws: ``CoreError/configInvalid(_:)`` when the file is missing,
|
||||
/// unreadable, not UTF-8, or empty once trimmed.
|
||||
private static func readTokenFile(_ path: String, describedAs key: String) throws -> String {
|
||||
let expanded = expandTilde(path)
|
||||
guard FileManager.default.fileExists(atPath: expanded) else {
|
||||
throw CoreError.configInvalid("\(key): no such file: \(expanded)")
|
||||
}
|
||||
let data: Data
|
||||
do {
|
||||
data = try Data(contentsOf: URL(fileURLWithPath: expanded))
|
||||
} catch {
|
||||
throw CoreError.configInvalid("\(key): cannot read \(expanded): \(error.localizedDescription)")
|
||||
}
|
||||
guard let text = String(data: data, encoding: .utf8) else {
|
||||
throw CoreError.configInvalid("\(key): \(expanded) is not valid UTF-8")
|
||||
}
|
||||
let token = text.trimmingCharacters(in: .whitespacesAndNewlines)
|
||||
guard !token.isEmpty else {
|
||||
throw CoreError.configInvalid("\(key): \(expanded) is empty")
|
||||
}
|
||||
return token
|
||||
}
|
||||
|
||||
/// Whether a token file is readable by users other than its owner.
|
||||
///
|
||||
/// Permissions are deliberately **not** enforced — refusing to start because
|
||||
/// a file is `0644` would be a poor trade on a single-user CI Mac — but
|
||||
/// `doctor` surfaces this as a warning.
|
||||
///
|
||||
/// - Parameter path: Path to check; tilde-expanded.
|
||||
/// - Returns: `true` when group or other bits are set, `false` when the file
|
||||
/// is owner-only, and `nil` when the mode cannot be read.
|
||||
public static func tokenFileIsGroupOrWorldReadable(_ path: String) -> Bool? {
|
||||
let expanded = expandTilde(path)
|
||||
guard
|
||||
let attrs = try? FileManager.default.attributesOfItem(atPath: expanded),
|
||||
let mode = attrs[.posixPermissions] as? NSNumber
|
||||
else { return nil }
|
||||
return (mode.int16Value & 0o077) != 0
|
||||
}
|
||||
|
||||
/// Paths of configured token files whose permissions are looser than `0600`.
|
||||
/// Empty when everything is owner-only or nothing is file-backed.
|
||||
public var insecureTokenFilePaths: [String] {
|
||||
[gitea.adminTokenFile, gitea.registrationTokenFile]
|
||||
.compactMap { RunnerConfig.nonEmpty($0) }
|
||||
.filter { RunnerConfig.tokenFileIsGroupOrWorldReadable($0) == true }
|
||||
}
|
||||
|
||||
/// ``StorageSection/storeDir`` with `~` expanded, as a `URL`.
|
||||
public var storeDirectoryURL: URL {
|
||||
URL(fileURLWithPath: RunnerConfig.expandTilde(storage.storeDir), isDirectory: true)
|
||||
}
|
||||
|
||||
/// The label set used for job matching.
|
||||
public var labelSet: LabelSet {
|
||||
LabelSet(runner.labels)
|
||||
}
|
||||
|
||||
// MARK: - Helpers
|
||||
|
||||
/// Expands a leading `~` or `~/` to the current user's home directory.
|
||||
///
|
||||
/// `NSString.expandingTildeInPath` is used rather than `FileManager`'s
|
||||
/// deprecated home lookup so the behaviour matches the shell.
|
||||
///
|
||||
/// - Parameter path: A possibly tilde-prefixed path.
|
||||
/// - Returns: An absolute path.
|
||||
public static func expandTilde(_ path: String) -> String {
|
||||
(path as NSString).expandingTildeInPath
|
||||
}
|
||||
|
||||
// MARK: - Codable
|
||||
|
||||
private enum CodingKeys: String, CodingKey {
|
||||
case gitea, runner, scheduler, guest, storage
|
||||
}
|
||||
|
||||
/// Decodes a configuration, substituting section defaults for absent keys.
|
||||
public init(from decoder: Decoder) throws {
|
||||
let c = try decoder.container(keyedBy: CodingKeys.self)
|
||||
self.gitea = try c.decode(GiteaSection.self, forKey: .gitea)
|
||||
self.runner = try c.decodeIfPresent(RunnerSection.self, forKey: .runner) ?? .init()
|
||||
self.scheduler = try c.decodeIfPresent(SchedulerSection.self, forKey: .scheduler) ?? .init()
|
||||
self.guest = try c.decodeIfPresent(GuestSection.self, forKey: .guest) ?? .init()
|
||||
self.storage = try c.decodeIfPresent(StorageSection.self, forKey: .storage) ?? .init()
|
||||
}
|
||||
}
|
||||
|
||||
// MARK: - Tolerant section decoding
|
||||
|
||||
extension RunnerConfig.GiteaSection {
|
||||
private enum CodingKeys: String, CodingKey {
|
||||
case instanceURL, adminToken, adminTokenFile
|
||||
case registrationToken, registrationTokenFile, fetchRegistrationTokenViaAPI
|
||||
}
|
||||
|
||||
public init(from decoder: Decoder) throws {
|
||||
let c = try decoder.container(keyedBy: CodingKeys.self)
|
||||
self.instanceURL = try c.decode(URL.self, forKey: .instanceURL)
|
||||
self.adminToken = try c.decodeIfPresent(String.self, forKey: .adminToken)
|
||||
self.adminTokenFile = try c.decodeIfPresent(String.self, forKey: .adminTokenFile)
|
||||
self.registrationToken = try c.decodeIfPresent(String.self, forKey: .registrationToken)
|
||||
self.registrationTokenFile = try c.decodeIfPresent(String.self, forKey: .registrationTokenFile)
|
||||
self.fetchRegistrationTokenViaAPI =
|
||||
try c.decodeIfPresent(Bool.self, forKey: .fetchRegistrationTokenViaAPI) ?? false
|
||||
}
|
||||
}
|
||||
|
||||
extension RunnerConfig.RunnerSection {
|
||||
private enum CodingKeys: String, CodingKey {
|
||||
case labels, namePrefix, runnerDownloadURL, version
|
||||
}
|
||||
|
||||
public init(from decoder: Decoder) throws {
|
||||
let d = RunnerConfig.RunnerSection()
|
||||
let c = try decoder.container(keyedBy: CodingKeys.self)
|
||||
self.labels = try c.decodeIfPresent([String].self, forKey: .labels) ?? d.labels
|
||||
self.namePrefix = try c.decodeIfPresent(String.self, forKey: .namePrefix) ?? d.namePrefix
|
||||
self.runnerDownloadURL =
|
||||
try c.decodeIfPresent(String.self, forKey: .runnerDownloadURL) ?? d.runnerDownloadURL
|
||||
self.version = try c.decodeIfPresent(String.self, forKey: .version) ?? d.version
|
||||
}
|
||||
}
|
||||
|
||||
extension RunnerConfig.SchedulerSection {
|
||||
private enum CodingKeys: String, CodingKey {
|
||||
case maxConcurrentVMs, pollIntervalSeconds, reconcileIntervalSeconds
|
||||
case jobTimeoutMinutes, bootTimeoutSeconds
|
||||
}
|
||||
|
||||
public init(from decoder: Decoder) throws {
|
||||
let d = RunnerConfig.SchedulerSection()
|
||||
let c = try decoder.container(keyedBy: CodingKeys.self)
|
||||
self.maxConcurrentVMs =
|
||||
try c.decodeIfPresent(Int.self, forKey: .maxConcurrentVMs) ?? d.maxConcurrentVMs
|
||||
self.pollIntervalSeconds =
|
||||
try c.decodeIfPresent(Int.self, forKey: .pollIntervalSeconds) ?? d.pollIntervalSeconds
|
||||
self.reconcileIntervalSeconds =
|
||||
try c.decodeIfPresent(Int.self, forKey: .reconcileIntervalSeconds) ?? d.reconcileIntervalSeconds
|
||||
self.jobTimeoutMinutes =
|
||||
try c.decodeIfPresent(Int.self, forKey: .jobTimeoutMinutes) ?? d.jobTimeoutMinutes
|
||||
self.bootTimeoutSeconds =
|
||||
try c.decodeIfPresent(Int.self, forKey: .bootTimeoutSeconds) ?? d.bootTimeoutSeconds
|
||||
}
|
||||
}
|
||||
|
||||
extension RunnerConfig.GuestSection {
|
||||
private enum CodingKeys: String, CodingKey {
|
||||
case username, password, cpuCount, memoryGB, diskGB
|
||||
}
|
||||
|
||||
public init(from decoder: Decoder) throws {
|
||||
let d = RunnerConfig.GuestSection()
|
||||
let c = try decoder.container(keyedBy: CodingKeys.self)
|
||||
self.username = try c.decodeIfPresent(String.self, forKey: .username) ?? d.username
|
||||
self.password = try c.decodeIfPresent(String.self, forKey: .password) ?? d.password
|
||||
self.cpuCount = try c.decodeIfPresent(Int.self, forKey: .cpuCount) ?? d.cpuCount
|
||||
self.memoryGB = try c.decodeIfPresent(Int.self, forKey: .memoryGB) ?? d.memoryGB
|
||||
self.diskGB = try c.decodeIfPresent(Int.self, forKey: .diskGB) ?? d.diskGB
|
||||
}
|
||||
}
|
||||
|
||||
extension RunnerConfig.StorageSection {
|
||||
private enum CodingKeys: String, CodingKey {
|
||||
case storeDir, minFreeDiskGB
|
||||
}
|
||||
|
||||
public init(from decoder: Decoder) throws {
|
||||
let d = RunnerConfig.StorageSection()
|
||||
let c = try decoder.container(keyedBy: CodingKeys.self)
|
||||
self.storeDir = try c.decodeIfPresent(String.self, forKey: .storeDir) ?? d.storeDir
|
||||
self.minFreeDiskGB =
|
||||
try c.decodeIfPresent(Int.self, forKey: .minFreeDiskGB) ?? d.minFreeDiskGB
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,110 @@
|
||||
import Foundation
|
||||
|
||||
/// The single error domain shared by every layer of the runner.
|
||||
///
|
||||
/// Host-side (`RunnerHost`) code wraps Virtualization.framework's `VZError` into
|
||||
/// these cases rather than propagating it, so the CLI only ever has to render one
|
||||
/// error type. `unimplemented` exists so that skeleton bodies can `throw` instead
|
||||
/// of trapping in code paths where a trap would take down the daemon.
|
||||
public enum CoreError: Error, Sendable {
|
||||
/// A code path that has not been written yet.
|
||||
case unimplemented
|
||||
|
||||
/// The on-disk configuration is missing, malformed, or internally inconsistent.
|
||||
/// The payload is a human-readable explanation suitable for printing to stderr.
|
||||
case configInvalid(String)
|
||||
|
||||
/// The Gitea API returned a non-2xx status.
|
||||
/// - Parameters:
|
||||
/// - status: The HTTP status code.
|
||||
/// - message: The response body (truncated) or a decoded API error message.
|
||||
case gitea(status: Int, message: String)
|
||||
|
||||
/// An SSH session could not be established, authenticated, or the remote
|
||||
/// command exited non-zero when a zero exit was required.
|
||||
case sshFailed(String)
|
||||
|
||||
/// A bounded wait elapsed. The payload names what was being waited on
|
||||
/// (for example `"dhcp lease for aa:bb:cc:dd:ee:ff"` or `"ssh on 192.168.64.7"`).
|
||||
case timeout(String)
|
||||
|
||||
/// A required external tool or file was absent (`diskutil`, an IPSW, the
|
||||
/// `gitea-runner` release asset, …).
|
||||
case notFound(String)
|
||||
|
||||
/// The host cannot run VMs: wrong architecture, unsupported macOS, missing
|
||||
/// `com.apple.security.virtualization` entitlement, or a locked login keychain.
|
||||
case hostUnsupported(String)
|
||||
|
||||
/// Apple's kernel-enforced limit of two concurrent macOS guests was hit.
|
||||
/// Surfaced distinctly because it is transient and the scheduler retries.
|
||||
case vmLimitExceeded
|
||||
|
||||
/// A VM bundle on disk is missing files or has an unreadable `config.json`.
|
||||
case bundleCorrupt(String)
|
||||
|
||||
/// Not enough free space on the store volume to safely clone or grow a VM.
|
||||
/// - Parameters:
|
||||
/// - requiredGB: The configured floor.
|
||||
/// - availableGB: What the volume actually has.
|
||||
case insufficientDiskSpace(requiredGB: Int, availableGB: Int)
|
||||
|
||||
/// A subprocess (`diskutil`, `codesign`, `security`, …) exited non-zero.
|
||||
case processFailed(command: String, exitCode: Int32, output: String)
|
||||
|
||||
/// The image build or provisioning pipeline failed at a named stage.
|
||||
case provisioningFailed(String)
|
||||
}
|
||||
|
||||
extension CoreError: CustomStringConvertible {
|
||||
/// A one-line, user-facing rendering of the error.
|
||||
public var description: String {
|
||||
switch self {
|
||||
case .unimplemented:
|
||||
return "not implemented"
|
||||
|
||||
case .configInvalid(let detail):
|
||||
return "invalid configuration: \(detail)"
|
||||
|
||||
case .gitea(let status, let message):
|
||||
let trimmed = message.trimmingCharacters(in: .whitespacesAndNewlines)
|
||||
return trimmed.isEmpty
|
||||
? "gitea API error (HTTP \(status))"
|
||||
: "gitea API error (HTTP \(status)): \(trimmed)"
|
||||
|
||||
case .sshFailed(let detail):
|
||||
return "ssh failed: \(detail)"
|
||||
|
||||
case .timeout(let what):
|
||||
return "timed out waiting for \(what)"
|
||||
|
||||
case .notFound(let what):
|
||||
return "not found: \(what)"
|
||||
|
||||
case .hostUnsupported(let detail):
|
||||
return "host cannot run VMs: \(detail)"
|
||||
|
||||
case .vmLimitExceeded:
|
||||
return "macOS guest limit reached (Apple allows at most 2 concurrent VMs per host)"
|
||||
|
||||
case .bundleCorrupt(let detail):
|
||||
return "VM bundle is corrupt: \(detail)"
|
||||
|
||||
case .insufficientDiskSpace(let requiredGB, let availableGB):
|
||||
return "insufficient disk space: need \(requiredGB) GB free, have \(availableGB) GB"
|
||||
|
||||
case .processFailed(let command, let exitCode, let output):
|
||||
let trimmed = output.trimmingCharacters(in: .whitespacesAndNewlines)
|
||||
return trimmed.isEmpty
|
||||
? "`\(command)` exited \(exitCode)"
|
||||
: "`\(command)` exited \(exitCode): \(trimmed)"
|
||||
|
||||
case .provisioningFailed(let stage):
|
||||
return "provisioning failed: \(stage)"
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
extension CoreError: LocalizedError {
|
||||
public var errorDescription: String? { description }
|
||||
}
|
||||
@@ -0,0 +1,246 @@
|
||||
import Foundation
|
||||
|
||||
/// One entry from macOS's `/var/db/dhcpd_leases`.
|
||||
///
|
||||
/// The Virtualization NAT attachment hands guests addresses from the host's
|
||||
/// built-in `bootpd`, which records each lease in that file. There is no API for
|
||||
/// this, so parsing the file keyed by the guest's MAC is how we learn a VM's IP.
|
||||
public struct DHCPLease: Sendable, Equatable {
|
||||
/// The guest's advertised hostname (`name=` in the lease block). Often the
|
||||
/// guest's local hostname, sometimes absent.
|
||||
public let name: String?
|
||||
/// The leased IPv4 address, e.g. `192.168.64.7`.
|
||||
public let ipAddress: String
|
||||
/// The hardware address, **normalized**: lowercase, colon-separated, each
|
||||
/// octet zero-padded to two hex digits, with the `1,` type prefix stripped.
|
||||
public let hwAddress: String
|
||||
/// Lease expiry, parsed from the `lease=` hex epoch, when present.
|
||||
public let leaseExpiry: Date?
|
||||
|
||||
public init(name: String?, ipAddress: String, hwAddress: String, leaseExpiry: Date?) {
|
||||
self.name = name
|
||||
self.ipAddress = ipAddress
|
||||
self.hwAddress = hwAddress
|
||||
self.leaseExpiry = leaseExpiry
|
||||
}
|
||||
}
|
||||
|
||||
/// Parser for `/var/db/dhcpd_leases`.
|
||||
///
|
||||
/// ## File format
|
||||
///
|
||||
/// A sequence of brace-delimited blocks of `key=value` lines:
|
||||
///
|
||||
/// ```
|
||||
/// {
|
||||
/// name=macos-guest
|
||||
/// ip_address=192.168.64.7
|
||||
/// hw_address=1,aa:bb:c:dd:ee:ff
|
||||
/// identifier=1,aa:bb:c:dd:ee:ff
|
||||
/// lease=0x67a1b2c3
|
||||
/// }
|
||||
/// ```
|
||||
///
|
||||
/// Two details bite:
|
||||
///
|
||||
/// 1. `hw_address` carries a leading hardware-type prefix (`1,` for Ethernet)
|
||||
/// that is not part of the MAC.
|
||||
/// 2. Octets are **not zero-padded** — `aa:bb:c:dd:ee:ff` is the same address
|
||||
/// that `VZMACAddress.string` renders as `aa:bb:0c:dd:ee:ff`. Comparing raw
|
||||
/// strings silently fails to match; both sides must be normalized.
|
||||
///
|
||||
/// Blocks accumulate: a MAC can appear more than once as leases are renewed or
|
||||
/// reissued, so lookups take the **newest** lease (latest `leaseExpiry`, falling
|
||||
/// back to last-in-file when expiry is missing).
|
||||
///
|
||||
/// - Note: macOS's DHCP lease time is 24 hours. That is exactly why clones must
|
||||
/// reuse a small set of **persistent per-slot MACs** rather than randomizing a
|
||||
/// MAC per VM: a randomized fleet would fill this file with day-long stale
|
||||
/// leases and exhaust the NAT subnet.
|
||||
public enum DHCPLeaseParser {
|
||||
/// The canonical path of the lease database.
|
||||
public static let defaultPath = "/var/db/dhcpd_leases"
|
||||
|
||||
/// Parses the whole file.
|
||||
///
|
||||
/// Malformed blocks are skipped rather than throwing — the file is written
|
||||
/// by another process and may be observed mid-write.
|
||||
///
|
||||
/// - Parameter text: The file's contents.
|
||||
/// - Returns: Leases in file order.
|
||||
public static func parse(_ text: String) -> [DHCPLease] {
|
||||
var leases: [DHCPLease] = []
|
||||
var fields: [String: String] = [:]
|
||||
var inBlock = false
|
||||
|
||||
for rawLine in text.split(separator: "\n", omittingEmptySubsequences: false) {
|
||||
let line = rawLine.trimmingCharacters(in: .whitespaces)
|
||||
if line.isEmpty { continue }
|
||||
|
||||
if line.hasPrefix("{") {
|
||||
// A `{` while already inside a block means the previous one was
|
||||
// truncated (the file is written by bootpd and can be observed
|
||||
// mid-write). Drop it and start over rather than merging.
|
||||
inBlock = true
|
||||
fields = [:]
|
||||
continue
|
||||
}
|
||||
|
||||
if line.hasPrefix("}") {
|
||||
if inBlock, let lease = makeLease(from: fields) { leases.append(lease) }
|
||||
inBlock = false
|
||||
fields = [:]
|
||||
continue
|
||||
}
|
||||
|
||||
guard inBlock, let separator = line.firstIndex(of: "=") else { continue }
|
||||
let key = line[line.startIndex..<separator].trimmingCharacters(in: .whitespaces).lowercased()
|
||||
let value = line[line.index(after: separator)...].trimmingCharacters(in: .whitespaces)
|
||||
if key.isEmpty { continue }
|
||||
fields[key] = value
|
||||
}
|
||||
|
||||
return leases
|
||||
}
|
||||
|
||||
/// Builds a lease from one block's `key=value` pairs, or `nil` when the block
|
||||
/// lacks the two fields that make it useful (an address and a MAC we can
|
||||
/// normalize). Never throws: a half-written block is simply not a lease.
|
||||
private static func makeLease(from fields: [String: String]) -> DHCPLease? {
|
||||
guard
|
||||
let ip = fields["ip_address"], !ip.isEmpty,
|
||||
let rawMAC = fields["hw_address"] ?? fields["identifier"],
|
||||
let mac = normalizeMAC(rawMAC)
|
||||
else { return nil }
|
||||
|
||||
let name = fields["name"].flatMap { $0.isEmpty ? nil : $0 }
|
||||
return DHCPLease(
|
||||
name: name,
|
||||
ipAddress: ip,
|
||||
hwAddress: mac,
|
||||
leaseExpiry: fields["lease"].flatMap(parseLeaseTime)
|
||||
)
|
||||
}
|
||||
|
||||
/// Parses a `lease=` value. `bootpd` writes a hex epoch (`0x66b2c0de`), but
|
||||
/// a plain decimal epoch has been observed too, so both are accepted.
|
||||
private static func parseLeaseTime(_ raw: String) -> Date? {
|
||||
let text = raw.trimmingCharacters(in: .whitespaces).lowercased()
|
||||
guard !text.isEmpty else { return nil }
|
||||
|
||||
let seconds: UInt64?
|
||||
if text.hasPrefix("0x") {
|
||||
seconds = UInt64(text.dropFirst(2), radix: 16)
|
||||
} else {
|
||||
seconds = UInt64(text, radix: 10)
|
||||
}
|
||||
|
||||
guard let seconds else { return nil }
|
||||
return Date(timeIntervalSince1970: TimeInterval(seconds))
|
||||
}
|
||||
|
||||
/// Reads and parses the lease database from disk.
|
||||
///
|
||||
/// - Parameter path: Defaults to ``defaultPath``.
|
||||
/// - Returns: Leases, or `[]` when the file does not exist yet (no guest has
|
||||
/// ever leased an address).
|
||||
public static func parseFile(at path: String = DHCPLeaseParser.defaultPath) -> [DHCPLease] {
|
||||
guard let text = try? String(contentsOfFile: path, encoding: .utf8) else { return [] }
|
||||
return parse(text)
|
||||
}
|
||||
|
||||
/// Finds the current IP for a MAC.
|
||||
///
|
||||
/// Both `mac` and each lease's `hwAddress` are normalized before comparison.
|
||||
///
|
||||
/// - Parameters:
|
||||
/// - mac: The guest's MAC, in any common rendering.
|
||||
/// - leases: Leases from ``parse(_:)``.
|
||||
/// - Returns: The newest matching lease's IP, or `nil`.
|
||||
public static func ipAddress(forMAC mac: String, in leases: [DHCPLease]) -> String? {
|
||||
lease(forMAC: mac, in: leases)?.ipAddress
|
||||
}
|
||||
|
||||
/// Finds the newest lease for a MAC.
|
||||
///
|
||||
/// Callers that must distinguish a *fresh* lease from the 24 h-old one the
|
||||
/// slot's previous guest left behind need the whole record, not just its
|
||||
/// address — see ``isNewer(_:than:)``.
|
||||
///
|
||||
/// - Parameters:
|
||||
/// - mac: The guest's MAC, in any common rendering.
|
||||
/// - leases: Leases from ``parse(_:)``.
|
||||
/// - Returns: The newest matching lease, or `nil`.
|
||||
public static func lease(forMAC mac: String, in leases: [DHCPLease]) -> DHCPLease? {
|
||||
guard let wanted = normalizeMAC(mac) else { return nil }
|
||||
|
||||
var best: DHCPLease?
|
||||
for lease in leases where lease.hwAddress == wanted {
|
||||
guard let current = best else {
|
||||
best = lease
|
||||
continue
|
||||
}
|
||||
// Newest expiry wins; a missing expiry sorts oldest. `>=` means that
|
||||
// among equally-dated (or equally-undated) duplicates the last block
|
||||
// in the file wins, which is the one bootpd wrote most recently.
|
||||
let candidate = lease.leaseExpiry ?? .distantPast
|
||||
let incumbent = current.leaseExpiry ?? .distantPast
|
||||
if candidate >= incumbent { best = lease }
|
||||
}
|
||||
|
||||
return best
|
||||
}
|
||||
|
||||
/// Whether `candidate` is a lease `bootpd` wrote *after* `previous`.
|
||||
///
|
||||
/// Slot MACs are persistent and macOS leases live 24 h, so a MAC almost
|
||||
/// always still has its previous guest's entry when the next clone boots.
|
||||
/// A caller that accepted the first entry it saw would hand out a stale
|
||||
/// address and then spend the whole boot timeout SSHing at nothing.
|
||||
///
|
||||
/// `bootpd` rewrites the block — bumping `lease=` — whenever it hands the
|
||||
/// address out again, so a strictly later expiry means a new lease. A
|
||||
/// changed address means the same thing. With no `previous` (first boot on
|
||||
/// this MAC) anything counts as new.
|
||||
///
|
||||
/// - Parameters:
|
||||
/// - candidate: The lease just read from the file.
|
||||
/// - previous: The lease observed before the guest was started.
|
||||
/// - Returns: `true` when `candidate` may be used.
|
||||
public static func isNewer(_ candidate: DHCPLease, than previous: DHCPLease?) -> Bool {
|
||||
guard let previous else { return true }
|
||||
if candidate.ipAddress != previous.ipAddress { return true }
|
||||
guard let previousExpiry = previous.leaseExpiry else { return true }
|
||||
guard let candidateExpiry = candidate.leaseExpiry else { return false }
|
||||
return candidateExpiry > previousExpiry
|
||||
}
|
||||
|
||||
/// Normalizes a MAC to lowercase, colon-separated, zero-padded octets.
|
||||
///
|
||||
/// Accepts an optional `<type>,` prefix (as written by `bootpd`), and
|
||||
/// tolerates `-` separators.
|
||||
///
|
||||
/// - Parameter raw: For example `1,aa:bb:c:dd:ee:ff` or `AA-BB-0C-DD-EE-FF`.
|
||||
/// - Returns: For example `aa:bb:0c:dd:ee:ff`, or `nil` if unparseable.
|
||||
public static func normalizeMAC(_ raw: String) -> String? {
|
||||
var text = raw.trimmingCharacters(in: .whitespaces)
|
||||
|
||||
// `bootpd` prefixes the hardware type: `1,` for Ethernet.
|
||||
if let comma = text.lastIndex(of: ",") {
|
||||
text = String(text[text.index(after: comma)...])
|
||||
}
|
||||
text = text.replacingOccurrences(of: "-", with: ":")
|
||||
|
||||
let octets = text.split(separator: ":", omittingEmptySubsequences: false)
|
||||
guard octets.count == 6 else { return nil }
|
||||
|
||||
var normalized: [String] = []
|
||||
normalized.reserveCapacity(6)
|
||||
for octet in octets {
|
||||
guard (1...2).contains(octet.count), octet.allSatisfy(\.isHexDigit) else { return nil }
|
||||
normalized.append(String(repeating: "0", count: 2 - octet.count) + octet.lowercased())
|
||||
}
|
||||
|
||||
return normalized.joined(separator: ":")
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,357 @@
|
||||
import Foundation
|
||||
|
||||
#if canImport(FoundationNetworking)
|
||||
import FoundationNetworking
|
||||
#endif
|
||||
|
||||
/// The HTTP seam under ``GiteaClient``.
|
||||
///
|
||||
/// Everything network-facing goes through this protocol so tests can supply a
|
||||
/// canned transport without a live Gitea instance.
|
||||
public protocol HTTPTransport: Sendable {
|
||||
/// Performs a request.
|
||||
///
|
||||
/// - Parameter request: A fully-formed request, including auth headers.
|
||||
/// - Returns: The response body and its HTTP status code.
|
||||
/// - Throws: Transport-level errors only; a non-2xx status is *not* an error
|
||||
/// here — ``GiteaClient`` maps that to ``CoreError/gitea(status:message:)``.
|
||||
func send(_ request: URLRequest) async throws -> (Data, Int)
|
||||
}
|
||||
|
||||
/// The production transport, backed by `URLSession`.
|
||||
public struct URLSessionTransport: HTTPTransport {
|
||||
/// The underlying session.
|
||||
public let session: URLSession
|
||||
|
||||
/// Creates a transport.
|
||||
///
|
||||
/// - Parameter session: Defaults to an ephemeral session with a 30 s request
|
||||
/// timeout, so a hung Gitea cannot stall the poll loop.
|
||||
public init(session: URLSession = URLSessionTransport.makeDefaultSession()) {
|
||||
self.session = session
|
||||
}
|
||||
|
||||
/// Builds the default ephemeral session.
|
||||
public static func makeDefaultSession() -> URLSession {
|
||||
let cfg = URLSessionConfiguration.ephemeral
|
||||
cfg.timeoutIntervalForRequest = 30
|
||||
cfg.timeoutIntervalForResource = 60
|
||||
return URLSession(configuration: cfg)
|
||||
}
|
||||
|
||||
public func send(_ request: URLRequest) async throws -> (Data, Int) {
|
||||
// `dataTask` + a continuation rather than `session.data(for:)`, because
|
||||
// the async URLSession API is not uniformly available in
|
||||
// swift-corelibs-foundation, and RunnerCore must build on Linux.
|
||||
try await withCheckedThrowingContinuation { (continuation: CheckedContinuation<(Data, Int), Error>) in
|
||||
let task = session.dataTask(with: request) { data, response, error in
|
||||
if let error {
|
||||
continuation.resume(throwing: error)
|
||||
return
|
||||
}
|
||||
guard let http = response as? HTTPURLResponse else {
|
||||
continuation.resume(
|
||||
throwing: CoreError.gitea(status: 0, message: "no HTTP response"))
|
||||
return
|
||||
}
|
||||
continuation.resume(returning: (data ?? Data(), http.statusCode))
|
||||
}
|
||||
task.resume()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// A thin, typed client for the subset of Gitea's admin Actions API this daemon
|
||||
/// needs.
|
||||
///
|
||||
/// All endpoints used here are **admin**-scoped, so the token must belong to a
|
||||
/// Gitea administrator. `doctor` verifies that by calling ``listRunners()``.
|
||||
public struct GiteaClient: Sendable {
|
||||
/// Instance base URL, e.g. `https://gitea.example.com`.
|
||||
public let baseURL: URL
|
||||
/// Admin API token, sent as `Authorization: token <value>`.
|
||||
public let token: String
|
||||
/// The HTTP seam.
|
||||
public let transport: any HTTPTransport
|
||||
|
||||
/// Creates a client.
|
||||
///
|
||||
/// - Parameters:
|
||||
/// - baseURL: Instance base URL; a trailing slash is tolerated.
|
||||
/// - token: Admin API token.
|
||||
/// - transport: Defaults to ``URLSessionTransport``.
|
||||
public init(baseURL: URL, token: String, transport: any HTTPTransport = URLSessionTransport()) {
|
||||
self.baseURL = baseURL
|
||||
self.token = token
|
||||
self.transport = transport
|
||||
}
|
||||
|
||||
// MARK: - Endpoints
|
||||
|
||||
/// Lists jobs currently waiting for a runner.
|
||||
///
|
||||
/// `GET /api/v1/admin/actions/jobs?status=queued&limit=<limit>`
|
||||
///
|
||||
/// A queued job stays queued until a matching runner claims it, or until
|
||||
/// Gitea's `ABANDONED_JOB_TIMEOUT` (default 24 h, swept every 6 h) expires
|
||||
/// it. There is therefore no urgency risk in a 5-second poll.
|
||||
///
|
||||
/// - Parameter limit: Page size. The scheduler only ever needs a handful.
|
||||
/// - Returns: The queued jobs, oldest-first as Gitea returns them.
|
||||
/// - Throws: ``CoreError/gitea(status:message:)`` on a non-2xx response.
|
||||
public func listQueuedJobs(limit: Int = 50) async throws -> [WorkflowJob] {
|
||||
// `status=queued` and nothing else. Gitea's `convertToInternal` maps
|
||||
// "queued" onto StatusWaiting ("ready, waiting for a runner") and maps
|
||||
// "waiting" onto StatusBlocked ("blocked on a dependency") — so asking
|
||||
// for "waiting" would return exactly the jobs that must not be booted.
|
||||
let request = try makeRequest(
|
||||
method: "GET",
|
||||
path: "/api/v1/admin/actions/jobs",
|
||||
query: [
|
||||
URLQueryItem(name: "status", value: "queued"),
|
||||
URLQueryItem(name: "limit", value: String(max(limit, 1))),
|
||||
])
|
||||
let response = try await send(request, as: WorkflowJobsResponse.self)
|
||||
return response.items
|
||||
}
|
||||
|
||||
/// Lists every registered runner on the instance.
|
||||
///
|
||||
/// `GET /api/v1/admin/actions/runners?page=<n>&limit=<limit>`
|
||||
///
|
||||
/// Used by the reconcile loop and by `doctor` (as an admin-scope probe).
|
||||
///
|
||||
/// Paginated deliberately: an unpaginated request returns only Gitea's
|
||||
/// default first page, and reconcile is precisely the thing that stops an
|
||||
/// instance from accumulating orphan rows. Missing rows past the page
|
||||
/// boundary would let the leak accelerate — each undeleted row pushes more
|
||||
/// rows out of view — and it would silently no-op the targeted cleanup that
|
||||
/// looks a single runner up by name.
|
||||
///
|
||||
/// - Parameters:
|
||||
/// - limit: Page size.
|
||||
/// - maxPages: Defensive ceiling, so a server that ignores `page` cannot
|
||||
/// spin this forever.
|
||||
/// - Returns: Every runner across all fetched pages, in server order.
|
||||
public func listRunners(limit: Int = 50, maxPages: Int = 50) async throws -> [ActionRunner] {
|
||||
let pageSize = max(limit, 1)
|
||||
var all: [ActionRunner] = []
|
||||
|
||||
for page in 1...max(maxPages, 1) {
|
||||
let request = try makeRequest(
|
||||
method: "GET",
|
||||
path: "/api/v1/admin/actions/runners",
|
||||
query: [
|
||||
URLQueryItem(name: "page", value: String(page)),
|
||||
URLQueryItem(name: "limit", value: String(pageSize)),
|
||||
])
|
||||
let response = try await send(request, as: RunnersResponse.self)
|
||||
all.append(contentsOf: response.items)
|
||||
// Terminate on the server's own total rather than on a short page:
|
||||
// Gitea clamps `limit` to its configured maximum, so a page shorter
|
||||
// than the one we asked for is not evidence that it is the last.
|
||||
if response.items.isEmpty { break }
|
||||
if let total = response.totalCount, all.count >= total { break }
|
||||
}
|
||||
|
||||
return all
|
||||
}
|
||||
|
||||
/// Deletes a runner row.
|
||||
///
|
||||
/// `DELETE /api/v1/admin/actions/runners/{id}`
|
||||
///
|
||||
/// Needed because a VM that dies uncleanly leaves its row behind: Gitea only
|
||||
/// sweeps runner rows at midnight, and never sweeps a runner that claimed no
|
||||
/// task. A 404 is treated as success (someone else already removed it).
|
||||
///
|
||||
/// - Parameter id: The runner id.
|
||||
public func deleteRunner(id: Int64) async throws {
|
||||
let request = try makeRequest(
|
||||
method: "DELETE",
|
||||
path: "/api/v1/admin/actions/runners/\(id)")
|
||||
// Gitea answers 204. 200 is accepted for tolerance, and 404 counts as
|
||||
// success: the reconcile loop's only goal is that the row be gone, and
|
||||
// it races with Gitea's own midnight sweep and with `--ephemeral`
|
||||
// auto-deregistration.
|
||||
try await sendIgnoringBody(request, acceptingStatuses: [200, 202, 204, 404])
|
||||
}
|
||||
|
||||
/// Returns the instance-scoped runner registration token.
|
||||
///
|
||||
/// `POST /api/v1/admin/actions/runners/registration-token`
|
||||
///
|
||||
/// - Warning: In current Gitea this returns the *existing* active token for
|
||||
/// the scope rather than minting a new one — but the semantics of "mint"
|
||||
/// are that a new token **invalidates all prior tokens of that scope**.
|
||||
/// Never call this per VM as a way of getting a throwaway secret; call it
|
||||
/// once and cache. Prefer seeding the token server-side via
|
||||
/// `GITEA_RUNNER_REGISTRATION_TOKEN` and configuring it statically.
|
||||
public func getRegistrationToken() async throws -> String {
|
||||
let request = try makeRequest(
|
||||
method: "POST",
|
||||
path: "/api/v1/admin/actions/runners/registration-token")
|
||||
let response = try await send(request, as: RegistrationTokenResponse.self)
|
||||
let token = response.token.trimmingCharacters(in: .whitespacesAndNewlines)
|
||||
guard !token.isEmpty else {
|
||||
throw CoreError.gitea(status: 200, message: "registration-token response carried an empty token")
|
||||
}
|
||||
return token
|
||||
}
|
||||
|
||||
/// Cheap reachability + auth probe used by `doctor`.
|
||||
///
|
||||
/// - Throws: ``CoreError/gitea(status:message:)`` when the instance is
|
||||
/// reachable but rejects the token.
|
||||
public func ping() async throws {
|
||||
// The runners list rather than /api/v1/version: version is anonymously
|
||||
// readable on most instances, so it would report "reachable" for a token
|
||||
// that is expired, wrong, or simply not an admin's — which is the exact
|
||||
// failure `doctor` exists to catch.
|
||||
_ = try await listRunners()
|
||||
}
|
||||
|
||||
// MARK: - Request plumbing
|
||||
|
||||
/// Builds an authenticated request against an API path.
|
||||
///
|
||||
/// - Parameters:
|
||||
/// - method: HTTP method.
|
||||
/// - path: API path relative to the instance root, e.g.
|
||||
/// `/api/v1/admin/actions/runners`.
|
||||
/// - query: Optional query items.
|
||||
/// - body: Optional request body; sets `Content-Type: application/json`.
|
||||
/// - Returns: A request carrying `Authorization` and `Accept` headers.
|
||||
public func makeRequest(
|
||||
method: String,
|
||||
path: String,
|
||||
query: [URLQueryItem] = [],
|
||||
body: Data? = nil
|
||||
) throws -> URLRequest {
|
||||
// Built by string-joining rather than `URL(string:relativeTo:)`, which
|
||||
// would discard any path component of `baseURL` — instances served under
|
||||
// a subpath (https://example.com/gitea) are common enough to matter.
|
||||
var base = baseURL.absoluteString
|
||||
while base.hasSuffix("/") { base.removeLast() }
|
||||
let suffix = path.hasPrefix("/") ? path : "/" + path
|
||||
|
||||
guard var components = URLComponents(string: base + suffix) else {
|
||||
throw CoreError.configInvalid("cannot form a request URL from \(base + suffix)")
|
||||
}
|
||||
if !query.isEmpty {
|
||||
components.queryItems = query
|
||||
}
|
||||
guard let url = components.url else {
|
||||
throw CoreError.configInvalid("cannot form a request URL from \(base + suffix)")
|
||||
}
|
||||
|
||||
var request = URLRequest(url: url)
|
||||
request.httpMethod = method
|
||||
// Gitea's PAT scheme. `Bearer` also works on recent versions, but
|
||||
// `token` is the documented form and works on every 1.x.
|
||||
request.setValue("token \(token)", forHTTPHeaderField: "Authorization")
|
||||
request.setValue("application/json", forHTTPHeaderField: "Accept")
|
||||
request.setValue(
|
||||
"gitea-macos-runner/\(RunnerVersion.current)", forHTTPHeaderField: "User-Agent")
|
||||
if let body {
|
||||
request.httpBody = body
|
||||
request.setValue("application/json", forHTTPHeaderField: "Content-Type")
|
||||
}
|
||||
return request
|
||||
}
|
||||
|
||||
/// Sends a request and decodes a JSON body, mapping non-2xx to
|
||||
/// ``CoreError/gitea(status:message:)``.
|
||||
public func send<T: Decodable>(_ request: URLRequest, as type: T.Type) async throws -> T {
|
||||
let (data, status) = try await transport.send(request)
|
||||
guard (200..<300).contains(status) else {
|
||||
throw CoreError.gitea(status: status, message: GiteaClient.errorMessage(from: data))
|
||||
}
|
||||
do {
|
||||
return try GiteaClient.makeDecoder().decode(T.self, from: data)
|
||||
} catch {
|
||||
throw CoreError.gitea(
|
||||
status: status,
|
||||
message: "could not decode \(T.self): \(error) — body: \(GiteaClient.excerpt(data))")
|
||||
}
|
||||
}
|
||||
|
||||
/// Sends a request that is expected to have no useful body.
|
||||
public func sendIgnoringBody(_ request: URLRequest, acceptingStatuses: Set<Int>) async throws {
|
||||
let (data, status) = try await transport.send(request)
|
||||
guard acceptingStatuses.contains(status) || (200..<300).contains(status) else {
|
||||
throw CoreError.gitea(status: status, message: GiteaClient.errorMessage(from: data))
|
||||
}
|
||||
}
|
||||
|
||||
/// Gitea's error bodies are `{"message": "...", "url": "..."}`. Prefer that
|
||||
/// message; fall back to a truncated raw body so nothing is ever reported as
|
||||
/// an empty error.
|
||||
private static func errorMessage(from data: Data) -> String {
|
||||
struct APIError: Decodable {
|
||||
let message: String?
|
||||
let errors: [String]?
|
||||
}
|
||||
if let decoded = try? JSONDecoder().decode(APIError.self, from: data) {
|
||||
if let message = decoded.message?.trimmingCharacters(in: .whitespacesAndNewlines),
|
||||
!message.isEmpty
|
||||
{
|
||||
return message
|
||||
}
|
||||
if let errors = decoded.errors, !errors.isEmpty {
|
||||
return errors.joined(separator: "; ")
|
||||
}
|
||||
}
|
||||
return excerpt(data)
|
||||
}
|
||||
|
||||
/// At most `limit` characters of a response body, for error messages.
|
||||
private static func excerpt(_ data: Data, limit: Int = 512) -> String {
|
||||
guard !data.isEmpty else { return "<empty body>" }
|
||||
let text = String(decoding: data, as: UTF8.self)
|
||||
.trimmingCharacters(in: .whitespacesAndNewlines)
|
||||
guard text.count > limit else { return text }
|
||||
return String(text.prefix(limit)) + "… (\(data.count) bytes)"
|
||||
}
|
||||
|
||||
/// A `JSONDecoder` configured for Gitea's timestamps (RFC 3339 / ISO 8601
|
||||
/// with an offset).
|
||||
///
|
||||
/// Go's `time.Time` marshals as RFC 3339 **Nano**: the fractional-seconds
|
||||
/// part is present only when non-zero, so a single strict formatter fails
|
||||
/// intermittently on real traffic. Both spellings are tried, plus a plain
|
||||
/// `YYYY-MM-DD` for good measure.
|
||||
public static func makeDecoder() -> JSONDecoder {
|
||||
let d = JSONDecoder()
|
||||
d.dateDecodingStrategy = .custom { decoder in
|
||||
let raw = try decoder.singleValueContainer().decode(String.self)
|
||||
if let date = parseTimestamp(raw) { return date }
|
||||
throw DecodingError.dataCorrupted(
|
||||
DecodingError.Context(
|
||||
codingPath: decoder.codingPath,
|
||||
debugDescription: "not an RFC 3339 timestamp: \"\(raw)\""))
|
||||
}
|
||||
return d
|
||||
}
|
||||
|
||||
/// Parses an RFC 3339 timestamp with or without fractional seconds.
|
||||
///
|
||||
/// - Parameter raw: The timestamp string.
|
||||
/// - Returns: The instant, or `nil` if it is in no recognized form.
|
||||
public static func parseTimestamp(_ raw: String) -> Date? {
|
||||
// Formatters are built per call rather than cached in a `static let`:
|
||||
// `ISO8601DateFormatter` is a non-Sendable reference type, and this is
|
||||
// called a handful of times per poll — not a hot path.
|
||||
let withFractional = ISO8601DateFormatter()
|
||||
withFractional.formatOptions = [.withInternetDateTime, .withFractionalSeconds]
|
||||
if let date = withFractional.date(from: raw) { return date }
|
||||
|
||||
let plain = ISO8601DateFormatter()
|
||||
plain.formatOptions = [.withInternetDateTime]
|
||||
if let date = plain.date(from: raw) { return date }
|
||||
|
||||
let dateOnly = ISO8601DateFormatter()
|
||||
dateOnly.formatOptions = [.withFullDate]
|
||||
return dateOnly.date(from: raw)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,330 @@
|
||||
import Foundation
|
||||
|
||||
/// A single job within a workflow run, as reported by
|
||||
/// `GET /api/v1/admin/actions/jobs` (Gitea 1.25+).
|
||||
///
|
||||
/// - Note: `status` values are the *external* strings. `queued` is the one we
|
||||
/// act on: it maps to Gitea's internal `StatusWaiting`, meaning "ready and
|
||||
/// waiting for a matching runner". The string `waiting` means something quite
|
||||
/// different — the job is **blocked** on a dependency — and must never be
|
||||
/// treated as schedulable.
|
||||
public struct WorkflowJob: Codable, Sendable, Equatable, Identifiable {
|
||||
/// Job id. Unique across the instance; the scheduler dedups on this.
|
||||
public let id: Int64
|
||||
/// The workflow run this job belongs to.
|
||||
public let runID: Int64
|
||||
/// The job's display name.
|
||||
public let name: String
|
||||
/// One of `queued`, `waiting`, `running`, `success`, `failure`, `cancelled`,
|
||||
/// `skipped`, `blocked`.
|
||||
public let status: String
|
||||
/// The job's `runs-on:` values, as **bare** label names.
|
||||
public let labels: [String]
|
||||
/// The runner that claimed the job, if any.
|
||||
public let runnerID: Int64?
|
||||
/// That runner's name, if any. Lets the reconcile loop tie a Gitea runner
|
||||
/// row back to one of our VMs.
|
||||
public let runnerName: String?
|
||||
/// When the job was created.
|
||||
public let createdAt: Date?
|
||||
/// When a runner picked it up.
|
||||
public let startedAt: Date?
|
||||
/// When it finished.
|
||||
public let completedAt: Date?
|
||||
|
||||
public init(
|
||||
id: Int64,
|
||||
runID: Int64,
|
||||
name: String,
|
||||
status: String,
|
||||
labels: [String],
|
||||
runnerID: Int64? = nil,
|
||||
runnerName: String? = nil,
|
||||
createdAt: Date? = nil,
|
||||
startedAt: Date? = nil,
|
||||
completedAt: Date? = nil
|
||||
) {
|
||||
self.id = id
|
||||
self.runID = runID
|
||||
self.name = name
|
||||
self.status = status
|
||||
self.labels = labels
|
||||
self.runnerID = runnerID
|
||||
self.runnerName = runnerName
|
||||
self.createdAt = createdAt
|
||||
self.startedAt = startedAt
|
||||
self.completedAt = completedAt
|
||||
}
|
||||
|
||||
private enum CodingKeys: String, CodingKey {
|
||||
case id
|
||||
case runID = "run_id"
|
||||
case name
|
||||
case status
|
||||
case labels
|
||||
case runnerID = "runner_id"
|
||||
case runnerName = "runner_name"
|
||||
case createdAt = "created_at"
|
||||
case startedAt = "started_at"
|
||||
case completedAt = "completed_at"
|
||||
}
|
||||
|
||||
/// Whether this job is waiting for a runner right now.
|
||||
public var isQueued: Bool { status == "queued" }
|
||||
|
||||
/// ``status`` as a case, with an ``JobStatus/unknown(_:)`` catch-all.
|
||||
public var jobStatus: JobStatus { JobStatus(rawValue: status) }
|
||||
|
||||
/// Decodes tolerantly: `runner_id` / `runner_name` carry `omitempty` in
|
||||
/// Gitea, and the timestamps are Go `time.Time` values that serialize as
|
||||
/// `0001-01-01T00:00:00Z` when unset (a queued job has no `started_at`).
|
||||
/// Those zero instants are surfaced as `nil` rather than as a year-1 date.
|
||||
public init(from decoder: Decoder) throws {
|
||||
let c = try decoder.container(keyedBy: CodingKeys.self)
|
||||
self.id = try c.decode(Int64.self, forKey: .id)
|
||||
self.runID = try c.decodeIfPresent(Int64.self, forKey: .runID) ?? 0
|
||||
self.name = try c.decodeIfPresent(String.self, forKey: .name) ?? ""
|
||||
self.status = try c.decodeIfPresent(String.self, forKey: .status) ?? ""
|
||||
self.labels = try c.decodeIfPresent([String].self, forKey: .labels) ?? []
|
||||
self.runnerID = try c.decodeIfPresent(Int64.self, forKey: .runnerID)
|
||||
self.runnerName = try c.decodeIfPresent(String.self, forKey: .runnerName)
|
||||
self.createdAt = WorkflowJob.nonZero(try c.decodeIfPresent(Date.self, forKey: .createdAt))
|
||||
self.startedAt = WorkflowJob.nonZero(try c.decodeIfPresent(Date.self, forKey: .startedAt))
|
||||
self.completedAt = WorkflowJob.nonZero(try c.decodeIfPresent(Date.self, forKey: .completedAt))
|
||||
}
|
||||
|
||||
/// Maps Go's zero `time.Time` (year 1) to `nil`.
|
||||
private static func nonZero(_ date: Date?) -> Date? {
|
||||
guard let date else { return nil }
|
||||
// 0001-01-01T00:00:00Z is ~62.1e9 seconds before the reference date.
|
||||
return date.timeIntervalSinceReferenceDate <= -62_135_596_800 ? nil : date
|
||||
}
|
||||
}
|
||||
|
||||
/// The *external* status strings Gitea reports for a workflow job.
|
||||
///
|
||||
/// Gitea maps its internal statuses onto GitHub's vocabulary in
|
||||
/// `convert.ToActionsStatus`: `StatusWaiting → "queued"`,
|
||||
/// `StatusBlocked → "waiting"`, `StatusRunning → "in_progress"`, and every
|
||||
/// terminal status → `"completed"` (the detail moves to a separate `conclusion`
|
||||
/// field). The `unknown` case exists because that mapping is Gitea's to change:
|
||||
/// a closed enum that threw on an unrecognized string would turn a new server
|
||||
/// version into a decode failure and stop the poll loop dead.
|
||||
public enum JobStatus: RawRepresentable, Sendable, Equatable, Hashable {
|
||||
/// Ready and waiting for a matching runner — Gitea's internal `StatusWaiting`.
|
||||
/// This is the only status that is schedulable.
|
||||
case queued
|
||||
/// **Blocked** on a dependency — Gitea's internal `StatusBlocked`. Despite
|
||||
/// the name, this is *not* a job waiting for a runner.
|
||||
case waiting
|
||||
/// Claimed by a runner and executing.
|
||||
case inProgress
|
||||
/// Terminal, in any of success / failure / cancelled / skipped.
|
||||
case completed
|
||||
/// A status string this build does not know about.
|
||||
case unknown(String)
|
||||
|
||||
public init(rawValue: String) {
|
||||
switch rawValue {
|
||||
case "queued": self = .queued
|
||||
case "waiting": self = .waiting
|
||||
case "in_progress": self = .inProgress
|
||||
case "completed": self = .completed
|
||||
default: self = .unknown(rawValue)
|
||||
}
|
||||
}
|
||||
|
||||
public var rawValue: String {
|
||||
switch self {
|
||||
case .queued: return "queued"
|
||||
case .waiting: return "waiting"
|
||||
case .inProgress: return "in_progress"
|
||||
case .completed: return "completed"
|
||||
case .unknown(let raw): return raw
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Envelope returned by `GET /api/v1/admin/actions/jobs`.
|
||||
///
|
||||
/// - Important: The array key has been observed as `jobs`, which is what is
|
||||
/// decoded here; some Gitea builds/OpenAPI revisions have used `workflow_jobs`
|
||||
/// for the equivalent repo-scoped endpoint. ``CodingKeys`` is written out
|
||||
/// explicitly so that adding a fallback is a one-line change, and
|
||||
/// ``jobs`` is optional so an empty response body decodes rather than throwing.
|
||||
public struct WorkflowJobsResponse: Codable, Sendable, Equatable {
|
||||
/// Total matching jobs server-side, ignoring `limit`.
|
||||
public let totalCount: Int?
|
||||
/// The page of jobs. `nil` and `[]` both mean "nothing queued".
|
||||
public let jobs: [WorkflowJob]?
|
||||
|
||||
public init(totalCount: Int?, jobs: [WorkflowJob]?) {
|
||||
self.totalCount = totalCount
|
||||
self.jobs = jobs
|
||||
}
|
||||
|
||||
private enum CodingKeys: String, CodingKey {
|
||||
case totalCount = "total_count"
|
||||
case jobs
|
||||
case workflowJobs = "workflow_jobs"
|
||||
case entries
|
||||
}
|
||||
|
||||
/// The jobs, never `nil`.
|
||||
public var items: [WorkflowJob] { jobs ?? [] }
|
||||
|
||||
/// Decodes the array under `jobs`, falling back to `workflow_jobs` and
|
||||
/// `entries`.
|
||||
///
|
||||
/// Gitea 1.25's `ActionWorkflowJobsResponse` tags its slice `json:"jobs"`
|
||||
/// (the Go field is named `Entries`), which is what the fallbacks guard
|
||||
/// against: a future rename of the tag, or a proxy that reserializes from
|
||||
/// the Go field name.
|
||||
public init(from decoder: Decoder) throws {
|
||||
let c = try decoder.container(keyedBy: CodingKeys.self)
|
||||
self.totalCount = try c.decodeIfPresent(Int.self, forKey: .totalCount)
|
||||
if let jobs = try c.decodeIfPresent([WorkflowJob].self, forKey: .jobs) {
|
||||
self.jobs = jobs
|
||||
} else if let jobs = try c.decodeIfPresent([WorkflowJob].self, forKey: .workflowJobs) {
|
||||
self.jobs = jobs
|
||||
} else {
|
||||
self.jobs = try c.decodeIfPresent([WorkflowJob].self, forKey: .entries)
|
||||
}
|
||||
}
|
||||
|
||||
public func encode(to encoder: Encoder) throws {
|
||||
var c = encoder.container(keyedBy: CodingKeys.self)
|
||||
try c.encodeIfPresent(totalCount, forKey: .totalCount)
|
||||
try c.encodeIfPresent(jobs, forKey: .jobs)
|
||||
}
|
||||
}
|
||||
|
||||
/// A registered Actions runner, from `GET /api/v1/admin/actions/runners`.
|
||||
public struct ActionRunner: Codable, Sendable, Equatable, Identifiable {
|
||||
/// Runner id, used for `DELETE /api/v1/admin/actions/runners/{id}`.
|
||||
public let id: Int64
|
||||
/// Runner name. Ours always start with the configured `namePrefix`.
|
||||
public let name: String
|
||||
/// Bare label names the runner advertises.
|
||||
public let labels: [String]
|
||||
/// Server-side liveness, e.g. `online` / `offline`.
|
||||
public let status: String?
|
||||
/// Whether the runner is currently executing a task.
|
||||
public let busy: Bool?
|
||||
/// Whether the runner registered with `--ephemeral`, i.e. the server will
|
||||
/// hand it exactly one task and then auto-deregister it (Gitea 1.24+).
|
||||
public let ephemeral: Bool?
|
||||
|
||||
public init(
|
||||
id: Int64,
|
||||
name: String,
|
||||
labels: [String],
|
||||
status: String? = nil,
|
||||
busy: Bool? = nil,
|
||||
ephemeral: Bool? = nil
|
||||
) {
|
||||
self.id = id
|
||||
self.name = name
|
||||
self.labels = labels
|
||||
self.status = status
|
||||
self.busy = busy
|
||||
self.ephemeral = ephemeral
|
||||
}
|
||||
|
||||
private enum CodingKeys: String, CodingKey {
|
||||
case id, name, labels, status, busy, ephemeral
|
||||
}
|
||||
|
||||
/// A single entry of `ActionRunner.labels` as Gitea actually serializes it:
|
||||
/// an object, not a string.
|
||||
///
|
||||
/// Verified against `modules/structs/repo_actions.go` at tag `v1.25.0`,
|
||||
/// where `ActionRunner.Labels` is `[]*ActionRunnerLabel` and
|
||||
/// `ActionRunnerLabel` is `{id int64, name string, type string}`. Only
|
||||
/// ``name`` is of any use here.
|
||||
private struct LabelObject: Decodable {
|
||||
let name: String?
|
||||
}
|
||||
|
||||
/// Decodes `labels` from either shape.
|
||||
///
|
||||
/// Gitea's job payload gives labels as plain strings, its runner payload
|
||||
/// gives them as objects. Both are accepted so that this one model keeps
|
||||
/// working if a version, a proxy, or a hand-written fixture disagrees.
|
||||
public init(from decoder: Decoder) throws {
|
||||
let c = try decoder.container(keyedBy: CodingKeys.self)
|
||||
self.id = try c.decode(Int64.self, forKey: .id)
|
||||
self.name = try c.decodeIfPresent(String.self, forKey: .name) ?? ""
|
||||
self.status = try c.decodeIfPresent(String.self, forKey: .status)
|
||||
self.busy = try c.decodeIfPresent(Bool.self, forKey: .busy)
|
||||
self.ephemeral = try c.decodeIfPresent(Bool.self, forKey: .ephemeral)
|
||||
|
||||
if let strings = try? c.decode([String].self, forKey: .labels) {
|
||||
self.labels = strings
|
||||
} else if let objects = try? c.decode([LabelObject].self, forKey: .labels) {
|
||||
self.labels = objects.compactMap(\.name)
|
||||
} else {
|
||||
// Absent or explicitly null. Gitea does emit `"labels": null` for a
|
||||
// runner registered without any.
|
||||
self.labels = []
|
||||
}
|
||||
}
|
||||
|
||||
/// `busy` treated as `false` when the server omits it.
|
||||
public var isBusy: Bool { busy ?? false }
|
||||
|
||||
/// `ephemeral` treated as `false` when the server omits it. The reconcile
|
||||
/// loop only ever deletes rows it is sure are ephemeral.
|
||||
public var isEphemeral: Bool { ephemeral ?? false }
|
||||
}
|
||||
|
||||
/// Envelope returned by `GET /api/v1/admin/actions/runners`.
|
||||
public struct RunnersResponse: Codable, Sendable, Equatable {
|
||||
public let totalCount: Int?
|
||||
public let runners: [ActionRunner]?
|
||||
|
||||
public init(totalCount: Int?, runners: [ActionRunner]?) {
|
||||
self.totalCount = totalCount
|
||||
self.runners = runners
|
||||
}
|
||||
|
||||
private enum CodingKeys: String, CodingKey {
|
||||
case totalCount = "total_count"
|
||||
case runners
|
||||
case entries
|
||||
}
|
||||
|
||||
/// The runners, never `nil`.
|
||||
public var items: [ActionRunner] { runners ?? [] }
|
||||
|
||||
/// Decodes the array under `runners` (Gitea 1.25's tag), falling back to
|
||||
/// `entries` (the Go field name).
|
||||
public init(from decoder: Decoder) throws {
|
||||
let c = try decoder.container(keyedBy: CodingKeys.self)
|
||||
self.totalCount = try c.decodeIfPresent(Int.self, forKey: .totalCount)
|
||||
if let runners = try c.decodeIfPresent([ActionRunner].self, forKey: .runners) {
|
||||
self.runners = runners
|
||||
} else {
|
||||
self.runners = try c.decodeIfPresent([ActionRunner].self, forKey: .entries)
|
||||
}
|
||||
}
|
||||
|
||||
public func encode(to encoder: Encoder) throws {
|
||||
var c = encoder.container(keyedBy: CodingKeys.self)
|
||||
try c.encodeIfPresent(totalCount, forKey: .totalCount)
|
||||
try c.encodeIfPresent(runners, forKey: .runners)
|
||||
}
|
||||
}
|
||||
|
||||
/// Response from `POST /api/v1/admin/actions/runners/registration-token`.
|
||||
///
|
||||
/// - Warning: This is scope-wide and **reusable**. Treat the returned value as
|
||||
/// "the current token for this scope", not as a freshly minted per-VM secret.
|
||||
public struct RegistrationTokenResponse: Codable, Sendable, Equatable {
|
||||
/// The registration token.
|
||||
public let token: String
|
||||
|
||||
public init(token: String) {
|
||||
self.token = token
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,84 @@
|
||||
import Foundation
|
||||
|
||||
/// The set of labels this host's runners advertise, used to decide whether a
|
||||
/// queued Gitea job is ours to pick up.
|
||||
///
|
||||
/// ## Bare names only
|
||||
///
|
||||
/// Gitea's label syntax at *registration* time is `name:schema` (for example
|
||||
/// `macos-arm64:host`), where the schema defaults to `host` when omitted. The
|
||||
/// schema is a runner-side execution hint — it tells `gitea-runner` to run the
|
||||
/// job directly on the machine instead of inside a container. **The server only
|
||||
/// ever stores and reports the bare name.** A workflow's `runs-on:` value, and
|
||||
/// therefore the `labels` array on a queued job, likewise contains bare names.
|
||||
///
|
||||
/// So: pass `macos-arm64:host` to `gitea-runner register --labels`, but match
|
||||
/// against `macos-arm64` here. Matching is case-sensitive, because Gitea's own
|
||||
/// comparison is.
|
||||
///
|
||||
/// - Warning: If the guest's `.runner`/`config.yaml` sets `runner.labels`, it
|
||||
/// silently overrides whatever `--labels` was passed at registration. The
|
||||
/// guest must therefore never ship a config file containing labels.
|
||||
public struct LabelSet: Sendable, Equatable, Hashable {
|
||||
/// The bare label names this host serves, e.g. `["macos-arm64", "macos"]`.
|
||||
public let names: Set<String>
|
||||
|
||||
/// Creates a label set from bare names.
|
||||
///
|
||||
/// Any `:schema` suffix present in `names` is stripped, so it is safe to
|
||||
/// hand this the same array that is written into the config file.
|
||||
///
|
||||
/// - Parameter names: Label names, with or without a `:schema` suffix.
|
||||
public init(_ names: [String]) {
|
||||
self.names = Set(
|
||||
names
|
||||
.map(LabelSet.bareName)
|
||||
.filter { !$0.isEmpty }
|
||||
)
|
||||
}
|
||||
|
||||
/// Whether a queued job's `labels` array can be satisfied by this host.
|
||||
///
|
||||
/// Returns `true` if and only if `jobLabels` is non-empty *and* every entry
|
||||
/// is a member of ``names``. An empty job label array is treated as "no
|
||||
/// declared requirement" and is deliberately **not** matched — a job that
|
||||
/// asks for nothing must not be scheduled onto a scarce macOS VM.
|
||||
///
|
||||
/// The job side is put through ``bareName(_:)`` too. The server normally
|
||||
/// stores bare names, so this changes nothing in the common case — but a
|
||||
/// workflow that writes `runs-on: [macos-arm64:host]` would otherwise never
|
||||
/// match anything and its job would be skipped with no log line at all.
|
||||
///
|
||||
/// - Parameter jobLabels: The `labels` array from a `WorkflowJob`.
|
||||
/// - Returns: `true` when this host should boot a VM for the job.
|
||||
public func matches(jobLabels: [String]) -> Bool {
|
||||
guard !jobLabels.isEmpty else { return false }
|
||||
let wanted = Set(jobLabels.map(LabelSet.bareName).filter { !$0.isEmpty })
|
||||
guard !wanted.isEmpty else { return false }
|
||||
return wanted.isSubset(of: names)
|
||||
}
|
||||
|
||||
/// The value to pass to `gitea-runner register --labels`, i.e. each bare
|
||||
/// name suffixed with the given schema and joined by commas.
|
||||
///
|
||||
/// - Parameter schema: The execution schema; `host` for a bare-metal guest.
|
||||
/// - Returns: For example `"macos-arm64:host,macos:host"`.
|
||||
public func registrationArgument(schema: String = "host") -> String {
|
||||
// Sorted so the argument is stable across process runs — a `Set` has no
|
||||
// inherent order, and an unstable registration argument would make the
|
||||
// guest command line (and its logs) needlessly non-reproducible.
|
||||
names.sorted()
|
||||
.map { schema.isEmpty ? $0 : "\($0):\(schema)" }
|
||||
.joined(separator: ",")
|
||||
}
|
||||
|
||||
/// Strips an optional `:schema` suffix from a single label token.
|
||||
///
|
||||
/// - Parameter label: A label such as `macos-arm64:host` or `macos-arm64`.
|
||||
/// - Returns: The bare name.
|
||||
public static func bareName(_ label: String) -> String {
|
||||
let trimmed = label.trimmingCharacters(in: .whitespaces)
|
||||
guard let colon = trimmed.firstIndex(of: ":") else { return trimmed }
|
||||
return String(trimmed[trimmed.startIndex..<colon])
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,35 @@
|
||||
import Foundation
|
||||
|
||||
/// Naming scheme for the ephemeral runners this host registers with Gitea.
|
||||
///
|
||||
/// Every booted VM registers under a **globally unique** name. That uniqueness
|
||||
/// is what makes the reconcile loop safe: when a VM dies uncleanly, Gitea keeps
|
||||
/// the runner row forever (rows are only swept at midnight, and never at all if
|
||||
/// the runner never claimed a task), so we must be able to look at a runner row
|
||||
/// and decide "this name is mine and no live VM of mine owns it" without any
|
||||
/// ambiguity. A shared or reused name would make that decision impossible.
|
||||
public enum RunnerNaming {
|
||||
/// Generates a fresh runner name.
|
||||
///
|
||||
/// - Parameter prefix: The configured prefix, e.g. `macos-vm-`.
|
||||
/// - Returns: `prefix` followed by a lowercase UUID, e.g.
|
||||
/// `macos-vm-3f1c2f8e-...`.
|
||||
public static func makeRunnerName(prefix: String) -> String {
|
||||
prefix + UUID().uuidString.lowercased()
|
||||
}
|
||||
|
||||
/// Whether a runner name reported by Gitea was minted by this host.
|
||||
///
|
||||
/// - Parameters:
|
||||
/// - name: A runner name from `GET /api/v1/admin/actions/runners`.
|
||||
/// - prefix: The configured prefix.
|
||||
/// - Returns: `true` when the reconcile loop may consider deleting it.
|
||||
public static func hasPrefix(_ name: String, prefix: String) -> Bool {
|
||||
// An empty prefix would match every runner on the instance, including
|
||||
// other people's. The reconcile loop deletes what this returns true for,
|
||||
// so refuse rather than match everything. `validated()` also rejects an
|
||||
// empty `namePrefix`; this is the second line of defence.
|
||||
guard !prefix.isEmpty else { return false }
|
||||
return name.hasPrefix(prefix)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,589 @@
|
||||
import Foundation
|
||||
import NIOCore
|
||||
import NIOPosix
|
||||
import NIOSSH
|
||||
|
||||
/// The outcome of a command run inside a guest.
|
||||
public struct SSHCommandResult: Sendable, Equatable {
|
||||
/// The remote process's exit status. `0` on success.
|
||||
public let exitCode: Int32
|
||||
/// Captured stdout, UTF-8 decoded with lossy replacement.
|
||||
public let stdout: String
|
||||
/// Captured stderr, UTF-8 decoded with lossy replacement.
|
||||
public let stderr: String
|
||||
|
||||
public init(exitCode: Int32, stdout: String, stderr: String) {
|
||||
self.exitCode = exitCode
|
||||
self.stdout = stdout
|
||||
self.stderr = stderr
|
||||
}
|
||||
|
||||
/// Whether the command exited zero.
|
||||
public var succeeded: Bool { exitCode == 0 }
|
||||
|
||||
/// Throws ``CoreError/sshFailed(_:)`` unless the command exited zero.
|
||||
///
|
||||
/// - Parameter command: Echoed into the error message for context.
|
||||
public func throwIfFailed(command: String) throws {
|
||||
guard exitCode != 0 else { return }
|
||||
let detail = stderr.isEmpty ? stdout : stderr
|
||||
let trimmed = detail.trimmingCharacters(in: .whitespacesAndNewlines)
|
||||
throw CoreError.sshFailed(
|
||||
"command failed (exit \(exitCode)): \(command)" + (trimmed.isEmpty ? "" : "\n\(trimmed)")
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
/// The seam for talking to a guest.
|
||||
///
|
||||
/// The orchestrator and the provisioner are written against this rather than
|
||||
/// against ``SSHExecutor`` so that provisioning logic can be unit-tested with a
|
||||
/// recording fake, and so a future vsock-based transport could be dropped in
|
||||
/// without touching callers.
|
||||
public protocol GuestExecutor: Sendable {
|
||||
/// Runs a shell command in the guest and waits for it to exit.
|
||||
///
|
||||
/// - Parameters:
|
||||
/// - command: A `/bin/sh`-compatible command line.
|
||||
/// - timeout: Wall-clock ceiling; exceeding it throws
|
||||
/// ``CoreError/timeout(_:)`` and closes the channel.
|
||||
/// - Returns: Exit status and captured output.
|
||||
func run(_ command: String, timeout: Duration) async throws -> SSHCommandResult
|
||||
|
||||
/// Copies a local file into the guest.
|
||||
///
|
||||
/// - Parameters:
|
||||
/// - localPath: Source path on the host.
|
||||
/// - remotePath: Destination path in the guest.
|
||||
func upload(localPath: String, remotePath: String) async throws
|
||||
|
||||
/// Writes bytes to a guest file with an explicit mode.
|
||||
///
|
||||
/// Used for secrets — notably the registration token, which is written with
|
||||
/// mode `0600` and deleted immediately after `gitea-runner register` reads
|
||||
/// it, so it never appears in a process argument list.
|
||||
///
|
||||
/// - Parameters:
|
||||
/// - data: File contents.
|
||||
/// - remotePath: Destination path in the guest.
|
||||
/// - mode: Octal mode string, e.g. `"0600"`.
|
||||
func uploadData(_ data: Data, remotePath: String, mode: String) async throws
|
||||
}
|
||||
|
||||
extension GuestExecutor {
|
||||
/// ``run(_:timeout:)`` with a two-minute default ceiling.
|
||||
public func run(_ command: String) async throws -> SSHCommandResult {
|
||||
try await run(command, timeout: .seconds(120))
|
||||
}
|
||||
|
||||
/// Runs a command and throws unless it exits zero.
|
||||
///
|
||||
/// - Returns: The successful result.
|
||||
@discardableResult
|
||||
public func runChecked(_ command: String, timeout: Duration = .seconds(120)) async throws -> SSHCommandResult {
|
||||
let result = try await run(command, timeout: timeout)
|
||||
try result.throwIfFailed(command: command)
|
||||
return result
|
||||
}
|
||||
}
|
||||
|
||||
/// SSH client over swift-nio-ssh using password authentication.
|
||||
///
|
||||
/// Password auth (rather than keys) is deliberate: the guest is a throwaway VM
|
||||
/// on a host-private NAT network whose credentials come from the same config
|
||||
/// that created it, and injecting a key would mean another provisioning step
|
||||
/// during the window before SSH is up.
|
||||
///
|
||||
/// - Important: Host keys are **not** verified. The peer is a VM this process
|
||||
/// just booted, on a link no other host shares; there is no trust-on-first-use
|
||||
/// story that would add security here, and pinning would break on every clone.
|
||||
public final class SSHExecutor: GuestExecutor, @unchecked Sendable {
|
||||
/// Guest IP, as learned from ``DHCPLeaseParser``.
|
||||
public let host: String
|
||||
/// SSH port; `22` for a stock guest with Remote Login enabled.
|
||||
public let port: Int
|
||||
/// Guest account name.
|
||||
public let username: String
|
||||
/// Guest account password.
|
||||
public let password: String
|
||||
|
||||
/// Creates an executor. No connection is made until the first command.
|
||||
///
|
||||
/// - Parameters:
|
||||
/// - host: Guest IP address.
|
||||
/// - port: SSH port. Defaults to `22`.
|
||||
/// - username: Guest account.
|
||||
/// - password: Guest password.
|
||||
public init(host: String, port: Int = 22, username: String, password: String) {
|
||||
self.host = host
|
||||
self.port = port
|
||||
self.username = username
|
||||
self.password = password
|
||||
}
|
||||
|
||||
public func run(_ command: String, timeout: Duration) async throws -> SSHCommandResult {
|
||||
do {
|
||||
return try await execute(command, stdin: nil, timeout: timeout)
|
||||
} catch let error as SSHTransportError {
|
||||
throw error.asCoreError
|
||||
}
|
||||
}
|
||||
|
||||
public func upload(localPath: String, remotePath: String) async throws {
|
||||
let url = URL(fileURLWithPath: localPath)
|
||||
guard let data = try? Data(contentsOf: url) else {
|
||||
throw CoreError.notFound("local file for upload: \(localPath)")
|
||||
}
|
||||
try await uploadData(data, remotePath: remotePath, mode: "0644")
|
||||
}
|
||||
|
||||
public func uploadData(_ data: Data, remotePath: String, mode: String) async throws {
|
||||
// Deliberately not SFTP or SCP: a stock macOS guest runs an sshd whose
|
||||
// subsystem set we do not control at this point in provisioning, and an
|
||||
// exec channel with the payload as stdin needs nothing beyond what we
|
||||
// already use for every other command.
|
||||
let quotedPath = Self.shellQuote(remotePath)
|
||||
guard mode.allSatisfy(\.isNumber), !mode.isEmpty else {
|
||||
throw CoreError.sshFailed("invalid file mode \(mode.debugDescription) for \(remotePath)")
|
||||
}
|
||||
|
||||
let command = """
|
||||
mkdir -p "$(dirname \(quotedPath))" && cat > \(quotedPath) && chmod \(mode) \(quotedPath)
|
||||
"""
|
||||
|
||||
let result: SSHCommandResult
|
||||
do {
|
||||
result = try await execute(command, stdin: data, timeout: .seconds(300))
|
||||
} catch let error as SSHTransportError {
|
||||
throw error.asCoreError
|
||||
}
|
||||
try result.throwIfFailed(command: "upload to \(remotePath)")
|
||||
}
|
||||
|
||||
/// Releases any pooled connection and event loop resources.
|
||||
///
|
||||
/// Connections are not pooled — each command opens and closes its own — and
|
||||
/// the event loop group is NIO's process-wide singleton, so there is nothing
|
||||
/// to release. Kept so callers can be written against a lifecycle that a
|
||||
/// future pooling or vsock transport may need.
|
||||
public func close() async {}
|
||||
|
||||
// MARK: - Transport
|
||||
|
||||
/// Opens a connection, runs one exec channel, and tears both down.
|
||||
///
|
||||
/// - Parameters:
|
||||
/// - command: The `/bin/sh` command line to exec.
|
||||
/// - stdin: Bytes to stream as the command's standard input. Standard
|
||||
/// input is closed (channel EOF) either way, so a command that would
|
||||
/// otherwise read from the terminal exits instead of hanging.
|
||||
/// - timeout: Wall-clock ceiling on the whole exchange.
|
||||
/// - Throws: ``SSHTransportError`` for connect/auth problems (which
|
||||
/// ``waitForSSH(host:port:username:password:timeout:pollInterval:)`` needs
|
||||
/// to tell apart), or ``CoreError/timeout(_:)`` when the ceiling elapses.
|
||||
func execute(_ command: String, stdin: Data?, timeout: Duration) async throws -> SSHCommandResult {
|
||||
let group = MultiThreadedEventLoopGroup.singleton
|
||||
let auth = SSHAuthOutcome()
|
||||
let username = self.username
|
||||
let password = self.password
|
||||
let host = self.host
|
||||
let port = self.port
|
||||
|
||||
let bootstrap = ClientBootstrap(group: group)
|
||||
.channelOption(ChannelOptions.socketOption(.tcp_nodelay), value: 1)
|
||||
.channelInitializer { channel in
|
||||
channel.eventLoop.makeCompletedFuture {
|
||||
let configuration = SSHClientConfiguration(
|
||||
userAuthDelegate: PasswordOnlyAuthDelegate(
|
||||
username: username,
|
||||
password: password,
|
||||
outcome: auth
|
||||
),
|
||||
serverAuthDelegate: AcceptAnyHostKeyDelegate()
|
||||
)
|
||||
try channel.pipeline.syncOperations.addHandler(
|
||||
NIOSSHHandler(
|
||||
role: .client(configuration),
|
||||
allocator: channel.allocator,
|
||||
inboundChildChannelInitializer: nil
|
||||
)
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
let channel: Channel
|
||||
do {
|
||||
channel = try await bootstrap.connect(host: host, port: port).get()
|
||||
} catch {
|
||||
// No TCP connection at all: sshd is not listening yet (or the guest
|
||||
// is unreachable). Recoverable — this is what waitForSSH retries on.
|
||||
throw SSHTransportError.connectFailed(host: host, port: port, underlying: error)
|
||||
}
|
||||
|
||||
let loop = channel.eventLoop
|
||||
let resultPromise = loop.makePromise(of: SSHCommandResult.self)
|
||||
let stdinBuffer = stdin.map { ByteBuffer(bytes: $0) }
|
||||
let description = command
|
||||
|
||||
let timeoutTask = loop.scheduleTask(in: .nanoseconds(Self.nanoseconds(timeout))) {
|
||||
resultPromise.fail(CoreError.timeout("ssh command on \(host): \(description)"))
|
||||
channel.close(promise: nil)
|
||||
}
|
||||
|
||||
// Completing an already-completed NIO promise is a no-op, so these
|
||||
// racing completions are safe: whichever fires first wins.
|
||||
channel.closeFuture.whenComplete { _ in
|
||||
if auth.wasRejected {
|
||||
resultPromise.fail(
|
||||
SSHTransportError.authenticationFailed(host: host, username: username)
|
||||
)
|
||||
} else {
|
||||
resultPromise.fail(
|
||||
CoreError.sshFailed("ssh connection to \(host):\(port) closed before the command finished")
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
channel.pipeline.handler(type: NIOSSHHandler.self).flatMap { sshHandler -> EventLoopFuture<Channel> in
|
||||
let childPromise = loop.makePromise(of: Channel.self)
|
||||
sshHandler.createChannel(childPromise, channelType: .session) { child, channelType in
|
||||
guard channelType == .session else {
|
||||
return child.eventLoop.makeFailedFuture(
|
||||
CoreError.sshFailed("unexpected SSH channel type \(channelType)")
|
||||
)
|
||||
}
|
||||
return child.eventLoop.makeCompletedFuture {
|
||||
try child.pipeline.syncOperations.addHandler(
|
||||
ExecChannelHandler(
|
||||
command: description,
|
||||
stdin: stdinBuffer,
|
||||
promise: resultPromise
|
||||
)
|
||||
)
|
||||
}
|
||||
// Without this the guest's EOF would close the channel before the
|
||||
// exit-status request arrives.
|
||||
.flatMap { child.setOption(ChannelOptions.allowRemoteHalfClosure, value: true) }
|
||||
}
|
||||
return childPromise.futureResult
|
||||
}.whenFailure { error in
|
||||
resultPromise.fail(error)
|
||||
channel.close(promise: nil)
|
||||
}
|
||||
|
||||
do {
|
||||
let result = try await resultPromise.futureResult.get()
|
||||
timeoutTask.cancel()
|
||||
channel.close(promise: nil)
|
||||
return result
|
||||
} catch {
|
||||
timeoutTask.cancel()
|
||||
channel.close(promise: nil)
|
||||
throw error
|
||||
}
|
||||
}
|
||||
|
||||
/// Wraps a path (or any argument) so `/bin/sh` sees it literally.
|
||||
static func shellQuote(_ value: String) -> String {
|
||||
"'" + value.replacingOccurrences(of: "'", with: "'\\''") + "'"
|
||||
}
|
||||
|
||||
private static func nanoseconds(_ duration: Duration) -> Int64 {
|
||||
let components = duration.components
|
||||
let seconds = components.seconds.multipliedReportingOverflow(by: 1_000_000_000)
|
||||
guard !seconds.overflow else { return .max }
|
||||
let sum = seconds.partialValue.addingReportingOverflow(
|
||||
Int64(components.attoseconds / 1_000_000_000)
|
||||
)
|
||||
return sum.overflow ? .max : sum.partialValue
|
||||
}
|
||||
}
|
||||
|
||||
// MARK: - Transport failures
|
||||
|
||||
/// Connection-level failures, kept distinct from ``CoreError`` so that
|
||||
/// ``waitForSSH(host:port:username:password:timeout:pollInterval:)`` can tell
|
||||
/// "sshd is not up yet" (retry) from "the password is wrong" (give up now).
|
||||
enum SSHTransportError: Error {
|
||||
/// No TCP connection could be established.
|
||||
case connectFailed(host: String, port: Int, underlying: Error)
|
||||
/// The server rejected our credentials.
|
||||
case authenticationFailed(host: String, username: String)
|
||||
|
||||
var asCoreError: CoreError {
|
||||
switch self {
|
||||
case .connectFailed(let host, let port, let underlying):
|
||||
return .sshFailed("cannot connect to \(host):\(port): \(underlying)")
|
||||
case .authenticationFailed(let host, let username):
|
||||
return .sshFailed("authentication failed for \(username)@\(host)")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Shared, thread-safe record of whether the server rejected our password.
|
||||
///
|
||||
/// The auth delegate runs on the event loop; the value is read from the async
|
||||
/// caller, hence the lock.
|
||||
final class SSHAuthOutcome: @unchecked Sendable {
|
||||
private let lock = NSLock()
|
||||
private var rejected = false
|
||||
|
||||
var wasRejected: Bool {
|
||||
lock.lock()
|
||||
defer { lock.unlock() }
|
||||
return rejected
|
||||
}
|
||||
|
||||
func markRejected() {
|
||||
lock.lock()
|
||||
defer { lock.unlock() }
|
||||
rejected = true
|
||||
}
|
||||
}
|
||||
|
||||
// MARK: - Delegates
|
||||
|
||||
/// Accepts every host key.
|
||||
///
|
||||
/// The peer is a VM this process booted seconds ago, on a host-private NAT link
|
||||
/// no other machine shares, from an image that is destroyed after one job. Its
|
||||
/// host key is freshly generated per clone, so there is nothing to pin:
|
||||
/// trust-on-first-use would accept whatever the first connection presented —
|
||||
/// exactly what this does — while a pinned key would reject every legitimate
|
||||
/// guest. See docs/DESIGN.md §6.
|
||||
final class AcceptAnyHostKeyDelegate: NIOSSHClientServerAuthenticationDelegate {
|
||||
func validateHostKey(hostKey: NIOSSHPublicKey, validationCompletePromise: EventLoopPromise<Void>) {
|
||||
validationCompletePromise.succeed(())
|
||||
}
|
||||
}
|
||||
|
||||
/// Offers the configured password once, then reports rejection.
|
||||
///
|
||||
/// NIOSSH asks again after a failed attempt; a second ask means the server
|
||||
/// refused the first, which is a provisioning bug rather than a transient
|
||||
/// condition, so it is recorded for the caller to fail fast on.
|
||||
final class PasswordOnlyAuthDelegate: NIOSSHClientUserAuthenticationDelegate {
|
||||
private let username: String
|
||||
private let password: String
|
||||
private let outcome: SSHAuthOutcome
|
||||
private var offered = false
|
||||
|
||||
init(username: String, password: String, outcome: SSHAuthOutcome) {
|
||||
self.username = username
|
||||
self.password = password
|
||||
self.outcome = outcome
|
||||
}
|
||||
|
||||
func nextAuthenticationType(
|
||||
availableMethods: NIOSSHAvailableUserAuthenticationMethods,
|
||||
nextChallengePromise: EventLoopPromise<NIOSSHUserAuthenticationOffer?>
|
||||
) {
|
||||
guard !offered, availableMethods.contains(.password) else {
|
||||
// Either the server refused our password, or it never offered
|
||||
// password auth at all. Both mean this guest will not let us in.
|
||||
outcome.markRejected()
|
||||
nextChallengePromise.succeed(nil)
|
||||
return
|
||||
}
|
||||
|
||||
offered = true
|
||||
nextChallengePromise.succeed(
|
||||
NIOSSHUserAuthenticationOffer(
|
||||
username: username,
|
||||
serviceName: "",
|
||||
offer: .password(.init(password: password))
|
||||
)
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
// MARK: - Exec channel
|
||||
|
||||
/// Drives one exec channel: sends the request, streams stdin, splits stdout from
|
||||
/// stderr, and captures the exit status.
|
||||
final class ExecChannelHandler: ChannelInboundHandler {
|
||||
typealias InboundIn = SSHChannelData
|
||||
typealias OutboundOut = SSHChannelData
|
||||
|
||||
/// SSH channel data is framed into packets; keep writes comfortably under
|
||||
/// the 128 KiB default maximum packet size.
|
||||
private static let chunkSize = 32 * 1024
|
||||
|
||||
private let command: String
|
||||
private var stdin: ByteBuffer?
|
||||
private var promise: EventLoopPromise<SSHCommandResult>?
|
||||
|
||||
private var stdout = ByteBufferAllocator().buffer(capacity: 0)
|
||||
private var stderr = ByteBufferAllocator().buffer(capacity: 0)
|
||||
private var exitCode: Int32?
|
||||
|
||||
init(command: String, stdin: ByteBuffer?, promise: EventLoopPromise<SSHCommandResult>) {
|
||||
self.command = command
|
||||
self.stdin = stdin
|
||||
self.promise = promise
|
||||
}
|
||||
|
||||
func channelActive(context: ChannelHandlerContext) {
|
||||
let request = SSHChannelRequestEvent.ExecRequest(command: command, wantReply: true)
|
||||
let sent = context.eventLoop.makePromise(of: Void.self)
|
||||
// Capture only Sendable values: the handler and its context must not
|
||||
// escape onto another thread.
|
||||
let resultPromise = promise
|
||||
let channel = context.channel
|
||||
sent.futureResult.whenFailure { error in
|
||||
resultPromise?.fail(error)
|
||||
channel.close(promise: nil)
|
||||
}
|
||||
context.triggerUserOutboundEvent(request, promise: sent)
|
||||
context.fireChannelActive()
|
||||
}
|
||||
|
||||
func userInboundEventTriggered(context: ChannelHandlerContext, event: Any) {
|
||||
switch event {
|
||||
case is ChannelSuccessEvent:
|
||||
sendStandardInput(context: context)
|
||||
|
||||
case is ChannelFailureEvent:
|
||||
fail(context: context, error: CoreError.sshFailed("guest refused to exec: \(command)"))
|
||||
|
||||
case let status as SSHChannelRequestEvent.ExitStatus:
|
||||
exitCode = Int32(truncatingIfNeeded: status.exitStatus)
|
||||
|
||||
case let signal as SSHChannelRequestEvent.ExitSignal:
|
||||
// A signalled process has no exit status; report it the way a shell
|
||||
// would, and keep the signal name in stderr so it is not lost.
|
||||
exitCode = 128
|
||||
var note = ByteBuffer(string: "\nterminated by SIG\(signal.signalName): \(signal.errorMessage)\n")
|
||||
stderr.writeBuffer(¬e)
|
||||
|
||||
default:
|
||||
context.fireUserInboundEventTriggered(event)
|
||||
}
|
||||
}
|
||||
|
||||
func channelRead(context: ChannelHandlerContext, data: NIOAny) {
|
||||
let channelData = unwrapInboundIn(data)
|
||||
guard case .byteBuffer(var bytes) = channelData.data else { return }
|
||||
|
||||
switch channelData.type {
|
||||
case .channel: stdout.writeBuffer(&bytes)
|
||||
case .stdErr: stderr.writeBuffer(&bytes)
|
||||
default: break // An extended data type we did not ask for.
|
||||
}
|
||||
}
|
||||
|
||||
func channelInactive(context: ChannelHandlerContext) {
|
||||
complete()
|
||||
context.fireChannelInactive()
|
||||
}
|
||||
|
||||
func handlerRemoved(context: ChannelHandlerContext) {
|
||||
complete()
|
||||
}
|
||||
|
||||
func errorCaught(context: ChannelHandlerContext, error: Error) {
|
||||
fail(context: context, error: error)
|
||||
}
|
||||
|
||||
private func sendStandardInput(context: ChannelHandlerContext) {
|
||||
if var payload = stdin {
|
||||
stdin = nil
|
||||
while payload.readableBytes > 0 {
|
||||
let slice = payload.readSlice(length: min(Self.chunkSize, payload.readableBytes))!
|
||||
context.write(
|
||||
wrapOutboundOut(SSHChannelData(type: .channel, data: .byteBuffer(slice))),
|
||||
promise: nil
|
||||
)
|
||||
}
|
||||
context.flush()
|
||||
}
|
||||
|
||||
// EOF either way: `cat > file` needs it to finish, and a command that
|
||||
// would otherwise block reading stdin gets an immediate end of input.
|
||||
context.close(mode: .output, promise: nil)
|
||||
}
|
||||
|
||||
private func complete() {
|
||||
guard let promise else { return }
|
||||
self.promise = nil
|
||||
|
||||
if let exitCode {
|
||||
promise.succeed(
|
||||
SSHCommandResult(
|
||||
exitCode: exitCode,
|
||||
stdout: String(buffer: stdout),
|
||||
stderr: String(buffer: stderr)
|
||||
)
|
||||
)
|
||||
} else {
|
||||
promise.fail(
|
||||
CoreError.sshFailed("guest closed the channel without an exit status: \(command)")
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
private func fail(context: ChannelHandlerContext, error: Error) {
|
||||
if let promise {
|
||||
self.promise = nil
|
||||
promise.fail(error)
|
||||
}
|
||||
context.close(promise: nil)
|
||||
}
|
||||
}
|
||||
|
||||
/// Blocks until a guest accepts an authenticated SSH session, or the deadline
|
||||
/// passes.
|
||||
///
|
||||
/// Called after a DHCP lease appears but before any provisioning: a fresh guest
|
||||
/// answers on port 22 only once `launchd` has started `sshd`, which lags the
|
||||
/// lease by tens of seconds.
|
||||
///
|
||||
/// - Parameters:
|
||||
/// - host: Guest IP.
|
||||
/// - port: SSH port. Defaults to `22`.
|
||||
/// - username: Guest account.
|
||||
/// - password: Guest password.
|
||||
/// - timeout: Overall ceiling.
|
||||
/// - pollInterval: Delay between attempts. Defaults to 2 s.
|
||||
/// - Throws: ``CoreError/timeout(_:)`` if the guest never answers.
|
||||
public func waitForSSH(
|
||||
host: String,
|
||||
port: Int = 22,
|
||||
username: String,
|
||||
password: String,
|
||||
timeout: Duration,
|
||||
pollInterval: Duration = .seconds(2)
|
||||
) async throws {
|
||||
let executor = SSHExecutor(host: host, port: port, username: username, password: password)
|
||||
let started = ContinuousClock.now
|
||||
var lastError: Error?
|
||||
|
||||
while true {
|
||||
do {
|
||||
// A real authenticated session running a trivial command, not a bare
|
||||
// TCP probe: sshd binds the port before it is ready to authenticate,
|
||||
// so a connect that succeeds proves very little.
|
||||
_ = try await executor.execute("true", stdin: nil, timeout: .seconds(20))
|
||||
return
|
||||
} catch let error as SSHTransportError {
|
||||
if case .authenticationFailed = error {
|
||||
// Wrong credentials will not become right by waiting: the guest
|
||||
// was provisioned with a different account or password, which is
|
||||
// a build failure, not a boot delay.
|
||||
throw error.asCoreError
|
||||
}
|
||||
lastError = error
|
||||
} catch {
|
||||
// Timeouts and mid-handshake closures are what a guest that is still
|
||||
// starting `sshd` looks like. Keep waiting.
|
||||
lastError = error
|
||||
}
|
||||
|
||||
guard ContinuousClock.now - started < timeout else { break }
|
||||
try await Task.sleep(for: pollInterval)
|
||||
guard ContinuousClock.now - started < timeout else { break }
|
||||
}
|
||||
|
||||
let detail = lastError.map { "; last error: \($0)" } ?? ""
|
||||
throw CoreError.timeout("ssh on \(host):\(port)\(detail)")
|
||||
}
|
||||
@@ -0,0 +1,339 @@
|
||||
import Foundation
|
||||
|
||||
/// What a VM slot is doing.
|
||||
///
|
||||
/// Slots are fixed in number (two, matching both the kernel's concurrent-VM cap
|
||||
/// and our two persistent MAC addresses) and are recycled, never created.
|
||||
public enum SlotState: Sendable, Equatable {
|
||||
/// No VM. Available to boot.
|
||||
case idle
|
||||
|
||||
/// A VM is being cloned/booted/provisioned; not yet registered with Gitea.
|
||||
/// - Parameter since: When the transition happened, for boot-timeout checks.
|
||||
case provisioning(since: Date)
|
||||
|
||||
/// A VM is up with `gitea-runner daemon` attached.
|
||||
///
|
||||
/// - Parameters:
|
||||
/// - jobHint: The queued job whose presence motivated this boot, if known.
|
||||
/// **Only a hint** — the server, not us, decides which job this runner
|
||||
/// actually claims.
|
||||
/// - since: When the VM went live, for job-timeout checks.
|
||||
case running(jobHint: Int64?, since: Date)
|
||||
|
||||
/// Whether the slot currently holds a VM (booting or live).
|
||||
public var isOccupied: Bool {
|
||||
if case .idle = self { return false }
|
||||
return true
|
||||
}
|
||||
|
||||
/// When the slot entered its current state, or `nil` when idle.
|
||||
public var since: Date? {
|
||||
switch self {
|
||||
case .idle: return nil
|
||||
case .provisioning(let t): return t
|
||||
case .running(_, let t): return t
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// One recyclable VM slot.
|
||||
public struct VMSlot: Sendable, Equatable, Identifiable {
|
||||
/// Stable index, `0..<maxVMs`. Also indexes the persistent per-slot MAC.
|
||||
public let id: Int
|
||||
/// Current state.
|
||||
public var state: SlotState
|
||||
|
||||
public init(id: Int, state: SlotState = .idle) {
|
||||
self.id = id
|
||||
self.state = state
|
||||
}
|
||||
}
|
||||
|
||||
/// The scheduler's complete observable state.
|
||||
public struct SchedulerState: Sendable, Equatable {
|
||||
/// Fixed-size slot table.
|
||||
public var slots: [VMSlot]
|
||||
|
||||
/// Job ids that have already caused a boot.
|
||||
///
|
||||
/// This is the dedup ledger. Without it, a job that stays queued for the
|
||||
/// several seconds a VM takes to come up would trigger a second boot on the
|
||||
/// next poll, and a third after that — burning the entire slot budget on one
|
||||
/// job. Entries are dropped once the job stops appearing as queued.
|
||||
public var dispatchedJobIDs: Set<Int64>
|
||||
|
||||
/// Creates a state with `count` idle slots and an empty ledger.
|
||||
public init(slotCount: Int) {
|
||||
self.slots = (0..<slotCount).map { VMSlot(id: $0) }
|
||||
self.dispatchedJobIDs = []
|
||||
}
|
||||
|
||||
public init(slots: [VMSlot], dispatchedJobIDs: Set<Int64> = []) {
|
||||
self.slots = slots
|
||||
self.dispatchedJobIDs = dispatchedJobIDs
|
||||
}
|
||||
|
||||
/// Slots not currently holding a VM.
|
||||
public var idleSlots: [VMSlot] { slots.filter { !$0.state.isOccupied } }
|
||||
|
||||
/// Slots holding a VM.
|
||||
public var occupiedSlots: [VMSlot] { slots.filter { $0.state.isOccupied } }
|
||||
}
|
||||
|
||||
/// A side effect the orchestrator should perform.
|
||||
///
|
||||
/// The planner returns these; it never performs I/O itself, which is what makes
|
||||
/// the whole scheduling policy unit-testable against a fixed `now`.
|
||||
public enum SchedulerAction: Sendable, Equatable {
|
||||
/// Clone, boot, provision, and register a VM in the given slot.
|
||||
/// - Parameters:
|
||||
/// - slot: Slot id.
|
||||
/// - jobHint: The queued job that motivated the boot.
|
||||
case bootVM(slot: Int, jobHint: Int64)
|
||||
|
||||
/// Stop and delete the VM in the given slot.
|
||||
/// - Parameters:
|
||||
/// - slot: Slot id.
|
||||
/// - reason: Human-readable cause, logged and used in tests.
|
||||
case teardownVM(slot: Int, reason: String)
|
||||
|
||||
/// Explicit no-op. Returned so a caller can distinguish "planner ran and
|
||||
/// chose to do nothing" from "planner returned an empty list".
|
||||
case none
|
||||
}
|
||||
|
||||
/// The pure scheduling state machine.
|
||||
///
|
||||
/// ## Capacity, not assignment
|
||||
///
|
||||
/// A booted VM is **capacity**, not a promise to run a specific job. We register
|
||||
/// an ephemeral runner and the *server* decides which queued job it claims —
|
||||
/// possibly not the one that triggered the boot. That is fine and in fact
|
||||
/// desirable: it means we never have to reimplement Gitea's matching rules. The
|
||||
/// `jobHint` carried through ``SchedulerAction/bootVM(slot:jobHint:)`` and
|
||||
/// ``SlotState/running(jobHint:since:)`` exists purely for logs and for the
|
||||
/// dedup ledger.
|
||||
///
|
||||
/// Because `--ephemeral` makes the server hand each runner exactly one task and
|
||||
/// then deregister it, a slot's life is: boot → register → claim one job → the
|
||||
/// `gitea-runner daemon` process exits → we tear down. There is no reuse, which
|
||||
/// is what makes the VM genuinely disposable.
|
||||
///
|
||||
/// ## Rules
|
||||
///
|
||||
/// 1. Only jobs whose labels ``LabelSet/matches(jobLabels:)`` are considered.
|
||||
/// 2. A job id already in ``SchedulerState/dispatchedJobIDs`` never boots a
|
||||
/// second VM.
|
||||
/// 3. At most `maxVMs` slots may be occupied (hard-clamped to 2 — the kernel
|
||||
/// fails a third `start()` with `VZError.virtualMachineLimitExceeded`).
|
||||
/// 4. Ledger entries for jobs no longer visible as queued are expired, so a
|
||||
/// slot freed by a completed job can be re-earned by a genuinely new job.
|
||||
/// 5. A slot in ``SlotState/provisioning(since:)`` longer than `bootTimeout`, or
|
||||
/// ``SlotState/running(jobHint:since:)`` longer than `jobTimeout`, is torn
|
||||
/// down.
|
||||
public enum SchedulerCore {
|
||||
/// Computes the next state and the actions to reach it.
|
||||
///
|
||||
/// Deterministic and side-effect free: same inputs, same outputs. `now` is
|
||||
/// injected rather than read so timeout behaviour is testable.
|
||||
///
|
||||
/// - Parameters:
|
||||
/// - state: Current state.
|
||||
/// - queuedJobs: Jobs Gitea currently reports as `queued`. Callers must
|
||||
/// not include `waiting` (blocked) jobs.
|
||||
/// - labels: This host's label set.
|
||||
/// - maxVMs: Concurrency cap; values above 2 are clamped.
|
||||
/// - now: Reference time for timeout arithmetic.
|
||||
/// - jobTimeout: Ceiling on ``SlotState/running(jobHint:since:)``.
|
||||
/// - bootTimeout: Ceiling on ``SlotState/provisioning(since:)``.
|
||||
/// - Returns: The updated state and the actions to execute, teardowns first
|
||||
/// so a freed slot can be reused within the same pass.
|
||||
public static func plan(
|
||||
state: SchedulerState,
|
||||
queuedJobs: [WorkflowJob],
|
||||
labels: LabelSet,
|
||||
maxVMs: Int,
|
||||
now: Date,
|
||||
jobTimeout: TimeInterval,
|
||||
bootTimeout: TimeInterval
|
||||
) -> (SchedulerState, [SchedulerAction]) {
|
||||
// The kernel fails a third concurrent guest, so the config never gets to
|
||||
// negotiate this. Clamped here as well as in `RunnerConfig.validated()`.
|
||||
let cap = min(max(maxVMs, 0), 2)
|
||||
|
||||
var newState = state
|
||||
var teardowns: [SchedulerAction] = []
|
||||
var boots: [SchedulerAction] = []
|
||||
|
||||
// 1. Which of the queued jobs are ours to serve, in the order Gitea
|
||||
// reported them (so the plan is a deterministic function of input).
|
||||
let matching = queuedJobs.filter { labels.matches(jobLabels: $0.labels) }
|
||||
let queuedIDs = Set(queuedJobs.map(\.id))
|
||||
|
||||
// 2. Expire the dedup ledger against reality rather than against a
|
||||
// timer: an id that is no longer queued was either claimed or
|
||||
// cancelled, and dedup only matters while a job is still waiting.
|
||||
newState.dispatchedJobIDs.formIntersection(queuedIDs)
|
||||
|
||||
// 3. Timeouts, emitted before any boot so a slot freed here can be
|
||||
// reused in this same pass.
|
||||
for index in newState.slots.indices {
|
||||
let slot = newState.slots[index]
|
||||
switch slot.state {
|
||||
case .idle:
|
||||
continue
|
||||
|
||||
case .provisioning(let since):
|
||||
let age = now.timeIntervalSince(since)
|
||||
guard age > bootTimeout else { continue }
|
||||
teardowns.append(
|
||||
.teardownVM(
|
||||
slot: slot.id,
|
||||
reason: "boot timeout: provisioning for \(Int(age))s (limit \(Int(bootTimeout))s)"
|
||||
)
|
||||
)
|
||||
newState.slots[index].state = .idle
|
||||
// Losing a boot must not permanently strand the job that
|
||||
// motivated it. `SlotState.provisioning` deliberately carries no
|
||||
// jobHint (the hint is a log/dedup detail, not an assignment), so
|
||||
// there is no specific id to drop here. Instead we release one
|
||||
// ledger entry — the lowest still-queued dispatched id, i.e. the
|
||||
// oldest such job, since Gitea's ids increase monotonically.
|
||||
// That is deterministic, releases exactly the capacity we lost,
|
||||
// and lets a replacement VM boot (possibly on this very tick).
|
||||
if let oldest = newState.dispatchedJobIDs.min() {
|
||||
newState.dispatchedJobIDs.remove(oldest)
|
||||
}
|
||||
|
||||
case .running(let jobHint, let since):
|
||||
let age = now.timeIntervalSince(since)
|
||||
guard age > jobTimeout else { continue }
|
||||
teardowns.append(
|
||||
.teardownVM(
|
||||
slot: slot.id,
|
||||
reason: "job timeout: running for \(Int(age))s (limit \(Int(jobTimeout))s)"
|
||||
)
|
||||
)
|
||||
newState.slots[index].state = .idle
|
||||
// Same reasoning as above, except here we do know the hint. It is
|
||||
// usually gone from the ledger already (a claimed job stops being
|
||||
// queued), so this is normally a no-op.
|
||||
if let jobHint { newState.dispatchedJobIDs.remove(jobHint) }
|
||||
}
|
||||
}
|
||||
|
||||
// 4. Boot capacity for jobs we have not already booted for.
|
||||
//
|
||||
// A booted VM is CAPACITY, not an assignment: the ephemeral runner we
|
||||
// register may legally claim a DIFFERENT matching job than the one
|
||||
// whose presence motivated the boot. The counting still works out —
|
||||
// one queued matching job earns one VM, and whichever job that VM
|
||||
// claims stops being queued and drops out of the ledger.
|
||||
for job in matching {
|
||||
guard !newState.dispatchedJobIDs.contains(job.id) else { continue }
|
||||
guard newState.occupiedSlots.count < cap else { break }
|
||||
guard let free = newState.slots.firstIndex(where: { !$0.state.isOccupied }) else { break }
|
||||
|
||||
boots.append(.bootVM(slot: newState.slots[free].id, jobHint: job.id))
|
||||
newState.slots[free].state = .provisioning(since: now)
|
||||
newState.dispatchedJobIDs.insert(job.id)
|
||||
}
|
||||
|
||||
// Teardowns first, boots second. An empty list is the no-op; `.none` is
|
||||
// never emitted, so callers never have to filter it out of a real plan.
|
||||
return (newState, teardowns + boots)
|
||||
}
|
||||
|
||||
/// Records that a slot began booting for a job.
|
||||
///
|
||||
/// Called by the orchestrator once it has actually started the clone/boot,
|
||||
/// so that a failed `plan` execution does not leave a phantom occupied slot.
|
||||
///
|
||||
/// - Returns: The updated state.
|
||||
public static func markProvisioning(
|
||||
state: SchedulerState,
|
||||
slot: Int,
|
||||
jobHint: Int64,
|
||||
now: Date
|
||||
) -> SchedulerState {
|
||||
var newState = state
|
||||
guard let index = newState.slots.firstIndex(where: { $0.id == slot }) else { return newState }
|
||||
newState.slots[index].state = .provisioning(since: now)
|
||||
newState.dispatchedJobIDs.insert(jobHint)
|
||||
return newState
|
||||
}
|
||||
|
||||
/// Promotes a slot from provisioning to running.
|
||||
public static func markRunning(
|
||||
state: SchedulerState,
|
||||
slot: Int,
|
||||
now: Date
|
||||
) -> SchedulerState {
|
||||
var newState = state
|
||||
guard let index = newState.slots.firstIndex(where: { $0.id == slot }) else { return newState }
|
||||
// The hint, if any, is carried over purely so logs and the job-timeout
|
||||
// teardown reason can name a job. It is never an assignment.
|
||||
let hint: Int64?
|
||||
if case .running(let existing, _) = newState.slots[index].state {
|
||||
hint = existing
|
||||
} else {
|
||||
hint = nil
|
||||
}
|
||||
newState.slots[index].state = .running(jobHint: hint, since: now)
|
||||
return newState
|
||||
}
|
||||
|
||||
/// Promotes a slot to running while recording the job that motivated its
|
||||
/// boot, which ``SlotState/provisioning(since:)`` does not carry.
|
||||
///
|
||||
/// Additive convenience over ``markRunning(state:slot:now:)``; the hint is
|
||||
/// still only ever used for logging and the job-timeout reason string.
|
||||
public static func markRunning(
|
||||
state: SchedulerState,
|
||||
slot: Int,
|
||||
jobHint: Int64?,
|
||||
now: Date
|
||||
) -> SchedulerState {
|
||||
var newState = state
|
||||
guard let index = newState.slots.firstIndex(where: { $0.id == slot }) else { return newState }
|
||||
newState.slots[index].state = .running(jobHint: jobHint, since: now)
|
||||
return newState
|
||||
}
|
||||
|
||||
/// Drops a job id from the dedup ledger.
|
||||
///
|
||||
/// The ledger's only automatic expiry is "the job stopped being queued"
|
||||
/// (``plan(state:queuedJobs:labels:maxVMs:now:jobTimeout:bootTimeout:)``,
|
||||
/// step 2), which is exactly wrong for a boot that never happened: the job
|
||||
/// is *still* queued, so its entry is retained and no further VM is ever
|
||||
/// booted for it. Every failure path — a refused boot, a clone error, a lost
|
||||
/// lease, a dead SSH channel — must call this, or the job waits out Gitea's
|
||||
/// 24 h `ABANDONED_JOB_TIMEOUT` for nothing.
|
||||
///
|
||||
/// Safe to call for an id that was never dispatched, or twice.
|
||||
///
|
||||
/// - Parameters:
|
||||
/// - state: Current state.
|
||||
/// - jobID: The job to release.
|
||||
/// - Returns: The updated state.
|
||||
public static func releaseJob(
|
||||
state: SchedulerState,
|
||||
jobID: Int64
|
||||
) -> SchedulerState {
|
||||
var newState = state
|
||||
newState.dispatchedJobIDs.remove(jobID)
|
||||
return newState
|
||||
}
|
||||
|
||||
/// Returns a slot to ``SlotState/idle`` after teardown.
|
||||
public static func markIdle(
|
||||
state: SchedulerState,
|
||||
slot: Int
|
||||
) -> SchedulerState {
|
||||
var newState = state
|
||||
guard let index = newState.slots.firstIndex(where: { $0.id == slot }) else { return newState }
|
||||
newState.slots[index].state = .idle
|
||||
return newState
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
import Foundation
|
||||
|
||||
/// Version stamp for the runner host tool itself.
|
||||
///
|
||||
/// This is *not* the version of the `gitea-runner` binary installed into the
|
||||
/// guest — that one lives in ``RunnerConfig/RunnerSection/version``.
|
||||
public enum RunnerVersion {
|
||||
/// The semantic version of this build, reported by `--version` and by
|
||||
/// `doctor`.
|
||||
public static let current = "0.1.0"
|
||||
}
|
||||
Reference in New Issue
Block a user