340 lines
14 KiB
Swift
340 lines
14 KiB
Swift
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
|
|
}
|
|
}
|