Merge nucleic/upbeat-opal-viper-cc8l into main
build / build (push) Canceled after 0s

This commit is contained in:
2026-08-07 05:11:01 -07:00
parent ee19744496
commit 6b0c01b74b
8 changed files with 395 additions and 28 deletions
+89 -6
View File
@@ -50,7 +50,7 @@ public struct LiveVM: Sendable {
///
/// `ensureFreeSpace` → `cloneImage(named:slotMAC:)` → ``VMInstance/start(options:)``
/// → poll `/var/db/dhcpd_leases` for the slot MAC until `bootTimeout` →
/// ``waitForSSH(host:port:username:password:timeout:pollInterval:)`` → write the
/// ``waitForSSH(host:port:username:password:timeout:pollInterval:reportInterval:onAttemptFailure:)`` → write the
/// registration token into a guest file with mode `0600` → over SSH:
///
/// ```sh
@@ -105,6 +105,15 @@ public actor Orchestrator {
/// The supervising task per slot: clone → boot → register → run → teardown.
private var slotTasks: [Int: Task<Void, Never>] = [:]
/// Why a slot's supervising task was cancelled, left for that task to find.
///
/// `Task.cancel()` carries no payload and `CancellationError` no detail, so
/// a lifecycle that catches one knows only *that* it was stopped. Everything
/// worth reading — "boot timeout: provisioning for 312s (limit 300s)" — is
/// known only to the canceller. Without this hand-off the operator sees
/// `reason=cancelled`, which names the mechanism and hides the cause.
private var slotCancelReasons: [Int: String] = [:]
/// The "the guest stopped on its own" watcher per slot.
private var deathWatchTasks: [Int: Task<Void, Never>] = [:]
@@ -252,6 +261,7 @@ public actor Orchestrator {
// let them run their own teardown; whatever they miss we clean up below.
let tasks = slotTasks
slotTasks.removeAll()
for (slot, _) in tasks { slotCancelReasons[slot] = "shutting down" }
for (_, task) in tasks { task.cancel() }
for (_, task) in tasks { await task.value }
@@ -307,6 +317,9 @@ public actor Orchestrator {
// are about to kill, so it would otherwise clear that entry only
// after the boot had already been refused.
if let task = slotTasks.removeValue(forKey: slot) {
// Before the cancel, not after: the task may reach its
// `catch` the instant it is cancelled.
slotCancelReasons[slot] = reason
task.cancel()
// Its own teardown runs to completion here, which also means
// it cannot race a successor booted later in this pass.
@@ -404,6 +417,9 @@ public actor Orchestrator {
let generation = (slotGeneration[slot] ?? 0) + 1
slotGeneration[slot] = generation
// Taken alongside the `Date()` handed to the planner, so the lifecycle
// can measure against the same deadline the planner will enforce.
let provisioningStarted = ContinuousClock.now
state = SchedulerCore.markProvisioning(state: state, slot: slot, jobHint: jobHint, now: Date())
logger.info(
@@ -417,7 +433,13 @@ public actor Orchestrator {
slotTasks[slot] = Task { [weak self] in
guard let self else { return }
await self.runSlotLifecycle(slot: slot, jobHint: jobHint, runnerName: runnerName, generation: generation)
await self.runSlotLifecycle(
slot: slot,
jobHint: jobHint,
runnerName: runnerName,
generation: generation,
provisioningStarted: provisioningStarted
)
}
}
@@ -425,11 +447,40 @@ public actor Orchestrator {
///
/// Every failure path funnels into the same teardown, because a slot that is
/// neither live nor idle is a slot leaked for the process's lifetime.
private func runSlotLifecycle(slot: Int, jobHint: Int64, runnerName: String, generation: Int) async {
///
/// - Parameter provisioningStarted: When the *planner's* boot deadline began
/// ticking — earlier than this function's own first instruction. Used to
/// budget the readiness waits against the deadline that will actually be
/// enforced rather than against a fresh copy of it.
private func runSlotLifecycle(
slot: Int,
jobHint: Int64,
runnerName: String,
generation: Int,
provisioningStarted: ContinuousClock.Instant
) async {
let bootTimeout = Duration.seconds(max(30, config.scheduler.bootTimeoutSeconds))
let jobTimeout = Duration.seconds(max(60, config.scheduler.jobTimeoutMinutes * 60))
var teardownReason = "job finished"
// Both readiness waits below are already supervised by the planner's
// boot deadline, which started ticking at `provisioningStarted` — before
// the clone, the boot and the DHCP lease had spent any of it. Handing
// either wait the full `bootTimeout` puts its deadline strictly *after*
// the planner's, so the planner always wins the race: this task is
// cancelled mid-wait and the specific error the wait was about to throw
// ("ssh on 192.168.65.233:22 after 41 attempts; last error: …") is
// discarded in favour of a bare cancellation. Budget from what is left
// and the wait gets to speak first.
func remainingBootBudget() -> Duration {
let spent = ContinuousClock.now - provisioningStarted
// Landing a little before the planner, so its next tick finds the
// slot already failing for a stated reason. The floor keeps an
// already-overrun budget from collapsing to zero attempts, which
// would trade one useless message for another.
return max(.seconds(15), bootTimeout - spent - .seconds(5))
}
do {
let mac = try store.macAddress(forSlot: slot, slotCount: slotCount)
// Whatever lease this MAC already holds belongs to the *previous*
@@ -454,15 +505,40 @@ public actor Orchestrator {
await self?.vmStoppedUnexpectedly(slot: slot, generation: generation, reason: reason)
}
let ip = try await waitForLease(mac: mac, timeout: bootTimeout, replacing: priorLease)
let ip = try await waitForLease(
mac: mac,
timeout: remainingBootBudget(),
replacing: priorLease
)
live[slot]?.ipAddress = ip
logger.info("guest leased address", metadata: ["slot": .stringConvertible(slot), "ip": .string(ip)])
// A `Logger` is a value type, so the callback below gets its own
// copy and never touches the actor — which is what lets it be a
// plain synchronous closure called from inside the poll loop.
let log = logger
try await waitForSSH(
host: ip,
username: config.guest.username,
password: config.guest.password,
timeout: bootTimeout
timeout: remainingBootBudget(),
onAttemptFailure: { attempt in
// A guest sharing a host with other Virtualization guests
// can take minutes to start `sshd`. Without this the wait is
// indistinguishable from a hang: the log goes quiet between
// "guest leased address" and teardown, which is exactly the
// window an operator most wants to see into.
log.info(
"waiting for guest ssh",
metadata: [
"slot": .stringConvertible(slot),
"host": .string(ip),
"attempt": .stringConvertible(attempt.attempt),
"elapsed": .string("\(attempt.elapsed.components.seconds)s"),
"error": .string(attempt.error),
]
)
}
)
let token = try await registrationToken()
@@ -511,7 +587,9 @@ public actor Orchestrator {
)
}
} catch is CancellationError {
teardownReason = "cancelled"
// Whoever cancelled us knows why; `CancellationError` does not.
// Falling back to "cancelled" only when nobody left a note.
teardownReason = slotCancelReasons[slot] ?? "cancelled"
} catch {
teardownReason = "\(error)"
logger.error(
@@ -539,6 +617,10 @@ public actor Orchestrator {
state = SchedulerCore.releaseJob(state: state, jobID: jobHint)
}
// A note left for a cancel that arrived after the lifecycle had already
// finished on its own would otherwise be read by the *next* occupant of
// this slot, mislabelling its teardown.
slotCancelReasons[slot] = nil
slotTasks[slot] = nil
}
@@ -551,6 +633,7 @@ public actor Orchestrator {
"guest stopped unexpectedly",
metadata: ["slot": .stringConvertible(slot), "reason": .string("\(reason)")]
)
slotCancelReasons[slot] = "guest stopped: \(reason)"
slotTasks[slot]?.cancel()
await teardownSlot(slot, reason: "guest stopped: \(reason)")
}