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.. /// Creates a state with `count` idle slots and an empty ledger. public init(slotCount: Int) { self.slots = (0.. = []) { 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 } }