Documentation
¶
Overview ¶
Package store owns all controller state. The allocation transaction here is the mutex that replaces flock: it cannot be bypassed or pointed at a divergent path.
Index ¶
- Variables
- type AllocateRequest
- type EnqueueRequest
- type JobFilter
- type KillOutcome
- type Recovery
- type RecoveryEvent
- type RecoveryHistory
- type Store
- func (s *Store) ActiveJobOnDevice(deviceID string) (string, error)
- func (s *Store) ActiveJobs() ([]model.Job, error)
- func (s *Store) AllLabels() (map[string][]model.Label, error)
- func (s *Store) Allocate(req AllocateRequest) (*model.Job, error)
- func (s *Store) AssignedJobsFor(workerID string) ([]model.Job, error)
- func (s *Store) AutoRecover(workerID string, proof model.RecoveryProof) ([]Recovery, error)
- func (s *Store) CancelQueued(jobID, reason string) (bool, error)
- func (s *Store) ClearDevice(id string) (bool, error)
- func (s *Store) Close() error
- func (s *Store) Devices() ([]model.Device, error)
- func (s *Store) Enqueue(req EnqueueRequest) (*model.Job, error)
- func (s *Store) HostDoc(host, deviceID string) (string, time.Time, error)
- func (s *Store) Job(id string) (*model.Job, error)
- func (s *Store) LabelSnapshot() (map[string]map[string]string, error)
- func (s *Store) LabelsFor(deviceID string) ([]model.Label, error)
- func (s *Store) LeaseDevices(ids []string) (map[string]string, error)
- func (s *Store) Leases() ([]model.Lease, error)
- func (s *Store) ListJobs(filter JobFilter) ([]model.Job, error)
- func (s *Store) MarkRunning(jobID string, at time.Time) error
- func (s *Store) MatchingDevices(sel string) ([]string, error)
- func (s *Store) Now() time.Time
- func (s *Store) QuarantineDetails(ids []string) (map[string]string, error)
- func (s *Store) QuarantineReasons(ids []string) (map[string]string, error)
- func (s *Store) QueuePosition(jobID string) (int, error)
- func (s *Store) QueuedJobs() ([]model.Job, error)
- func (s *Store) RecentJobsForDevice(deviceID string, limit int) ([]model.Job, error)
- func (s *Store) RecordHeartbeat(workerID string, at time.Time, runningJobIDs []string, options ...SweepOptions) error
- func (s *Store) RecoveryHistoryFor(deviceID string) (RecoveryHistory, error)
- func (s *Store) Release(jobID string, state model.JobState, exitCode *int, reason string) error
- func (s *Store) ReplaceLabels(deviceID, source string, labels map[string]string, at time.Time) error
- func (s *Store) RequestKill(jobID, reason string, finalizeDisconnected bool) (KillOutcome, error)
- func (s *Store) RetireDevice(id string) error
- func (s *Store) ScheduleOnce() ([]model.Job, error)
- func (s *Store) SetDeviceMaxRuntime(deviceID string, d time.Duration) error
- func (s *Store) SetDeviceState(id string, state model.DeviceState, at time.Time, detail string) error
- func (s *Store) Sweep(grace, unhealthyAfter time.Duration, holdOffUntil time.Time, ...) (SweepResult, error)
- func (s *Store) TakeKillRequests(workerID string) ([]string, error)
- func (s *Store) UpsertHostDoc(host, deviceID, body string, at time.Time) error
- func (s *Store) UpsertWorker(w model.Worker, devices []model.Device) error
- func (s *Store) WorkerHost(workerID string) (string, error)
- type SweepOptions
- type SweepResult
Constants ¶
This section is empty.
Variables ¶
var ErrDeviceBusy = errors.New("device has a live lease")
ErrDeviceBusy is returned when an operation needs a device nobody is holding and something holds it.
var ErrDeviceNotFound = errors.New("device not found")
ErrDeviceNotFound means the device ID named by a caller matches no row. It exists so an update that changed nothing can never be mistaken for one that did: see SetDeviceState.
var ErrNoDevice = errors.New("no device available")
ErrNoDevice means nothing matching the request was free. Stage 1 has no queue, so the caller is told immediately rather than parked.
var ErrNoMatchingDevice = errors.New("no device matches the selector")
ErrNoMatchingDevice means the selector matches no device that exists. We reject at submit rather than queue forever: a selector matching nothing is far more often a typo than a bet on a host registering later.
var ErrRuntimeAboveCeiling = errors.New("requested runtime exceeds the device ceiling")
ErrRuntimeAboveCeiling means the job asked for more wall clock than the device tolerates. We reject rather than clamp: a silently shortened budget produces a run whose submitter believes it had longer.
var ErrUnknownDevice = errors.New("unknown device")
ErrUnknownDevice means the caller named a device_id that no worker has ever registered. Distinct from ErrNoDevice (the device exists but is busy) and ErrRuntimeAboveCeiling (the device exists but the requested runtime is too long) so a caller — and the HTTP layer — can answer each differently.
var ErrWorkerNotFound = errors.New("worker not found")
ErrWorkerNotFound means the worker ID named by a caller matches no row.
Functions ¶
This section is empty.
Types ¶
type AllocateRequest ¶
type EnqueueRequest ¶
type EnqueueRequest struct {
DeviceID string
Selector string
Command []string
Cwd string
Env map[string]string
Submitter string
IdempotencyKey string
Priority int
MaxRuntime time.Duration
IdleTimeout time.Duration
// Kind is model.LeaseKindJob or model.LeaseKindHold; empty defaults to
// LeaseKindJob. Validating that a hold carries no command and has a
// TTL is the submit handler's job (server.handleSubmit), not this
// layer's — Enqueue only persists what it is given.
Kind string
// Reason is why a hold was taken; carried through to the job row and,
// at assignment, copied onto the lease row (see assignQueued).
Reason string
// Stdio is model.StdioLogs, StdioTTY or StdioPipe. Validating it is the
// submit handler's job, as with Kind: Enqueue only persists it.
Stdio string
}
type JobFilter ¶
JobFilter selects job history using exact matches. Limit is required and must be between 1 and 200, inclusive.
type KillOutcome ¶
type KillOutcome int
const ( KillNotCancellable KillOutcome = iota KillRequested KillFinalized )
type Recovery ¶
Recovery is one device returned to the pool by AutoRecover, and the quarantine reason it was cleared FROM. The reason travels back because the event this feeds — an operator asking "why is this device back?" — needs to name what the device was out for, and by the time the caller looks, the row no longer says (the reason is cleared with the quarantine it explains).
type RecoveryEvent ¶ added in v0.1.2
RecoveryEvent is one automatic return to the pool, and what the device was out for when it happened.
type RecoveryHistory ¶ added in v0.1.2
type RecoveryHistory struct {
// Recoveries are the returns inside Window, newest first.
Recoveries []RecoveryEvent
// Remaining is how many automatic returns are left in the window. Zero
// means the next quarantine waits for an operator.
Remaining int
Limit int
Window time.Duration
}
RecoveryHistory is what the flap guard knows about one device: the automatic returns still inside its sliding window, and how many it has left before it stops healing itself and waits for a person.
AutoRecover has counted this since the guard was written and nothing could read it, so a device out of the pool looked the same whether it had failed once or had exhausted every automatic return it gets. That is the difference between a status and a decision.
type Store ¶
type Store struct {
// contains filtered or unexported fields
}
func (*Store) ActiveJobOnDevice ¶
ActiveJobOnDevice returns the ID of the job currently holding deviceID, or "" if the device is idle. Exclusivity is the whole point of this system, so at most one job should ever qualify; the newest wins if that invariant is ever violated, since it is the one a caller is asking about.
This is a single indexed read rather than a filter over ActiveJobs precisely because its caller is a request handler: ActiveJobs re-reads each job in full, one query apiece, which is fine for a dashboard and wasteful on a path that only needs an ID.
func (*Store) ActiveJobs ¶
ActiveJobs returns jobs that are assigned or running, newest first.
func (*Store) AllLabels ¶
AllLabels returns every stored label row for every device, grouped by device ID, in one query. It exists so a caller building a view across the whole fleet (deviceViews, behind both /v1/state and /v1/devices) can show each device's labels without issuing one LabelsFor query per device on top of the Devices()/Leases() pair it already runs — the same N+1 concern the task 7 review round found in handleExplain and fixed there.
func (*Store) Allocate ¶
func (s *Store) Allocate(req AllocateRequest) (*model.Job, error)
Allocate claims a device and creates its job in ONE transaction. Either the device flips ready -> busy and the job and lease exist, or nothing happened.
It bypasses the queue entirely: there is no queued state, no priority, no reservation, and no interaction with ScheduleOnce. That made it stage 1's only allocation path; since stage 2, production code must go through Enqueue followed by ScheduleOnce instead, which is the only route that keeps "who gets a device next" consistent with what a client seeing its queue position was told. Allocate is not deleted because the store's own tests still use it to set up a device already busy without needing a scheduling pass — but a second, queue-bypassing way to hand out a device is exactly the kind of thing that turns into a race if production code ever calls it again, so: do not wire this into any handler.
func (*Store) AssignedJobsFor ¶
AssignedJobsFor returns jobs handed to a worker that it has not started yet.
func (*Store) AutoRecover ¶
AutoRecover returns this worker's quarantined devices to the pool when the proof it presented at registration answers what they were quarantined for.
It is the cheap sibling of restoreRebootedDevicesLocked (see store.go), which does the same job on the strength of a changed boot ID. The two answer the identical question — "can anything from the interrupted job still be holding this device?" — and clear the identical set of causes (rebootClearableReasons); they differ only in what constitutes the proof. A reboot proves it by destroying every process on the machine, which costs minutes and a full restart of everything the host runs. A container restart proves it by destroying the PID namespace the jobs lived in, and a host worker proves it by looking. See model.RecoveryProof.
What is deliberately NOT here:
- fault. A self-reported hardware problem is not answered by "no processes are running": the probe that reported it tested something this proof does not. It is excluded by rebootClearableReasons, the same list and the same reasoning the reboot path uses.
- an unrecorded cause (an empty quarantine_reason). Guessing "probably a lost worker" about a device quarantined before that column existed is the one direction that can hand out bad hardware.
- a device with a live lease. Nothing may contradict the lease table — the rule ClearDevice enforces against an operator is enforced here against a proof.
It runs as its own transaction, after UpsertWorker rather than inside it, and that ordering is required: UpsertWorker's reap pass is what quarantines a device whose worker re-registered with a job in flight, so a recovery pass folded into it would either run before the quarantine it exists to answer or have to be threaded through the middle of a function that is already the most safety-critical one in this package. Between the two calls the device is simply still quarantined, which nothing schedules onto, so the window is not observable.
func (*Store) CancelQueued ¶
CancelQueued removes a job that has not started. It reports false when the job is already assigned or running — that is rc kill's job, not this one, because a running job owns a device and a live process.
func (*Store) ClearDevice ¶
ClearDevice is the explicit operator acknowledgement that a device is free. It must not be able to contradict the lease table: a device with a live lease stays unhealthy no matter what the operator asserts. The bool return tells the caller whether the device was actually cleared, so a live-lease refusal is never reported as success.
func (*Store) Enqueue ¶
func (s *Store) Enqueue(req EnqueueRequest) (*model.Job, error)
Enqueue records a job in state queued. It never assigns: ScheduleOnce does that, so there is exactly one place where a device changes hands.
func (*Store) HostDoc ¶
HostDoc returns one usage-sheet document and when it was last written. A (host, deviceID) pair with no row yields an empty body and the zero time, not an error — most hosts, and most devices, will never have written one.
func (*Store) LabelSnapshot ¶
LabelSnapshot returns the effective labels the scheduler matches against: device ID -> key -> value, with a detected value winning over a declared one for the same key.
func (*Store) LabelsFor ¶
LabelsFor returns every stored row for a device, both sources, so a caller can show the conflict rather than hide it.
func (*Store) LeaseDevices ¶
LeaseDevices maps each id in ids to the device its lease was on. The ids are the ones SweepResult.LeasesExpired reports, which is a job ID when the expired lease had a job and the lease's own ID when it did not (a hold taken with no job behind it), so both columns are matched and both are keyed in the result. An id that matches no lease is absent.
It answers for released leases too — deliberately. By the time a sweep returns, the leases it expired are already released, but the row keeps its device_id, which is the only remaining link between an expired lease and the hardware it just took out of the pool.
Same discipline as QuarantineReasons: one query, drained to exhaustion, run after Sweep's transaction has committed, never inside one. Ordered by acquisition so that if a job somehow ever held two leases, the newest is the one that wins the key.
func (*Store) ListJobs ¶
ListJobs returns matching jobs newest first. Jobs submitted in the same second are ordered by descending ID so repeated reads are deterministic.
func (*Store) MarkRunning ¶
MarkRunning records that the worker has actually spawned the process.
func (*Store) MatchingDevices ¶
MatchingDevices returns the device IDs whose effective labels satisfy the selector, sorted by ID so scheduling is deterministic.
func (*Store) Now ¶
Now returns the controller's current time as seen through its Clock, so tests using a fake clock can stamp rows consistently with the store.
func (*Store) QuarantineDetails ¶
QuarantineDetails maps each id to the operator-facing explanation recorded when it was quarantined — a verify probe's stderr, a failed acquire hook's message — as opposed to QuarantineReasons' machine-readable category. A device with no explanation on file (quarantined before this column existed, or by a path that records none, such as a sweep) maps to "".
Same discipline as QuarantineReasons: one query for all the ids, drained to exhaustion before anything else runs, because MaxOpenConns(1) turns an overlapping query into a deadlock rather than a slowdown.
func (*Store) QuarantineReasons ¶
QuarantineReasons returns the recorded quarantine reason of each device in ids, keyed by device ID. Only an id with no device row at all is absent from the map; a device that exists but has no reason recorded maps to the empty string, which a caller must read as "cause unknown" — a row written before the column existed carries the same empty string as one written today with nothing to say, and neither is evidence of anything.
This exists as a separate, post-commit read rather than as extra fields on SweepResult because Sweep is not to be restructured for it, and because the read genuinely must happen outside Sweep's transaction: the pool is capped at one connection (see store.Open), so a query issued while Sweep's own cursors are open would deadlock. The consequence is that the reason returned here is the CURRENT one, which is a moment later than the one the sweep wrote — another path (a worker self-reporting a fault, a re-registration) may have overwritten it in between. That is the honest answer for a notification anyway: it describes why the device is out of the pool now, not why it was a few milliseconds ago.
func (*Store) QueuePosition ¶
QueuePosition is 1-based. For a pinned job it counts jobs ahead of it waiting for the same device. A selector job's candidate set can change between scheduling passes (labels change, devices come and go), so it counts every queued job ahead of it in scheduling order instead — an upper bound on its wait, not an exact slot. 0 means the job is not queued.
func (*Store) QueuedJobs ¶
QueuedJobs returns queued jobs in scheduling order: priority DESC, then oldest first. The tie-break is SQLite's implicit rowid rather than queued_at or id: queued_at has one-second resolution (two jobs enqueued in the same second, or under a frozen test clock, tie on it), and id is a random UUID unrelated to submission order. rowid is strictly monotonic with insertion order, so it is the only column that reliably preserves FIFO when priority and queued_at both tie.
func (*Store) RecentJobsForDevice ¶
RecentJobsForDevice returns up to limit jobs that have run (or are running) on this device, most recent submission first — the history `rc describe` shows so an agent can see what a box has actually been doing, not just what it's doing right now. A device with no history, or one this controller has never heard of, yields an empty slice rather than an error: a freshly registered device legitimately has nothing to show.
func (*Store) RecordHeartbeat ¶
func (s *Store) RecordHeartbeat(workerID string, at time.Time, runningJobIDs []string, options ...SweepOptions) error
RecordHeartbeat refreshes a worker, restores its unknown devices, and renews the lease of every job the worker reports it is actually supervising. A device whose lease is still live returns to busy, NOT to ready: the job it was demoted with is still running on it. Promoting a leased device to ready would offer an occupied GPU to the next claimant. Devices marked unhealthy stay out until explicitly cleared.
Lease renewal here is what makes expiry (Sweep) safe: without it, any job running longer than its lease TTL would be killed and its device quarantined while still healthy.
runningJobIDs is the crucial input, and the reason this is not simply "renew everything this worker owns". A job can be recorded assigned/running on a worker that never actually received it — handleAssignments commits the running transition before writing the response, so a lost response (a controller restart mid-write, a proxy, a decode error at the worker) leaves the controller believing a job is running that the worker has no process for and will never report on. Renewing on the strength of the worker merely being alive kept that job's lease alive forever: the device stayed busy, its holder never changed, and lease expiry — the backstop the design promises — could never fire, leaving no way to recover the hardware short of restarting the worker and clearing the device by hand.
So renewal follows reality: only the leases of jobs the worker names are pushed forward. A job the worker is not running stops being renewed, its lease lapses, and Sweep reclaims it — marking the job lost and quarantining the device, which is exactly the intended behaviour for a job nobody can account for.
func (*Store) RecoveryHistoryFor ¶ added in v0.1.2
func (s *Store) RecoveryHistoryFor(deviceID string) (RecoveryHistory, error)
RecoveryHistoryFor reads the guard's own counter for one device.
It filters by the window rather than trusting the table to hold only recent rows: AutoRecover prunes lazily, and only for the devices it is considering on that pass, so a device that flapped last week and has been quiet since still has its rows until something else touches it. Reading them as current would report a device as out of automatic returns it has in fact had back for days.
func (*Store) ReplaceLabels ¶
func (s *Store) ReplaceLabels(deviceID, source string, labels map[string]string, at time.Time) error
ReplaceLabels makes the stored labels for one device and one source exactly the given set. Replacing rather than merging is deliberate: a probe that stops reporting a key means the fact is gone (a card was swapped, a driver downgraded), and a stale label is worse than a missing one.
func (*Store) RequestKill ¶
func (s *Store) RequestKill(jobID, reason string, finalizeDisconnected bool) (KillOutcome, error)
RequestKill flags an active job for termination. Ordinarily the worker sees the flag on its next poll and reports the terminal result asynchronously. When requested, a job whose device is disconnected is finalized here while retaining the flag so the original worker can still kill a surviving process if it reconnects.
func (*Store) RetireDevice ¶
RetireDevice removes a device from the fleet: the card was pulled, or the box was decommissioned. It is the counterpart to a worker registering one, and it exists because the alternative was editing the database by hand.
What it does NOT delete is job history. Jobs carry device_id as plain text with no foreign key precisely so the record of what ran where outlives the hardware — "which box ran this, and how did it go" is asked most often about a box that is no longer there.
It refuses while anything holds the device. Deleting the row under a live lease would strand a running job on hardware the controller no longer believes exists, and the lease is the one guarantee this system makes.
A caveat that belongs to the operator, not to this function: if the device's worker is still running and still declares the device, the next registration recreates it. Remove it from the worker's config first, or stop the worker. This is deliberately not enforced here — the ordinary case is dropping a device from worker.yaml on a host that is otherwise still in service, and refusing whenever the worker was alive would block exactly that.
func (*Store) ScheduleOnce ¶
ScheduleOnce makes one scheduling pass: for each queued job in priority then FIFO order, assign it if its device is free, otherwise reserve that device so nothing behind it can take it first. Returns the jobs assigned.
func (*Store) SetDeviceMaxRuntime ¶
SetDeviceMaxRuntime records the ceiling a host declares for one of its devices. Called during registration. The column is seconds by design, so a sub-second duration truncates — that's intentional, not a bug: runtime ceilings are not meant to be enforced to sub-second precision.
func (*Store) SetDeviceState ¶
func (s *Store) SetDeviceState(id string, state model.DeviceState, at time.Time, detail string) error
SetDeviceState is the host's own report about one of its devices. Marking a device unhealthy through this path is a self-reported FAULT — the worker (or a verify probe standing in for it) saying the hardware itself is not fit to hand out — which is recorded as such: a fault outlives a reboot, unlike a quarantine that merely reflects a process nobody can account for. Any other state clears the reason along with the quarantine it explained.
A device ID that matches no row yields ErrDeviceNotFound rather than a silent success. An UPDATE that changes zero rows is not an error to SQL, but it is very much one here: the caller believes it has just quarantined a device, and a worker that logs a successful fault report having changed nothing is the worst possible outcome — the device stays schedulable and nobody is looking for it. SetDeviceState moves a device and, when that move is a quarantine, records both WHY in the machine sense (quarantine_reason, a category the reaper matches on to decide what a reboot may clear) and why in the operator's sense (detail — a verify probe's stderr, a failed hook's message). The detail is free text from a worker and is never matched on; it exists to be read by whoever decides whether to clear the device.
func (*Store) Sweep ¶
func (s *Store) Sweep(grace, unhealthyAfter time.Duration, holdOffUntil time.Time, options ...SweepOptions) (SweepResult, error)
Sweep demotes devices whose worker has stopped reporting. A device is never promoted to ready by this path: silence is not evidence that it is free.
holdOffUntil is the instant before which this sweep must not conclude anything DESTRUCTIVE — it expires no lease and writes off no worker until the clock has passed it. A zero time means "judge now", which is what every caller with no startup to wait out passes.
It exists because both of those conclusions are drawn from a stored timestamp — leases.expires_at and workers.last_heartbeat_at — and a stored timestamp keeps running while this process is not. On 2026-08-18 the controller was restarted to pick up a new image, was down for a few seconds, and its first sweep found a lease whose deadline had passed during exactly that gap. It expired the lease, quarantined the device and marked a job that was running perfectly well as lost. The holder had been alive throughout; it had simply had nobody to renew against. A worker that was silent across the same gap is the same mistake one table over, and worse in proportion to the outage: a controller down longer than unhealthyAfter would write off every device in the fleet on its first tick and lose every job on it.
So for a bounded window after the controller starts (see internal/cli/serve.go, which is where the window is chosen), the sweep declines to judge and gives the holders a chance to renew. It is NOT an amnesty and it does not touch any deadline: a holder that is genuinely gone renews nothing, and the very next sweep after the window closes reaches the identical verdict. Demotion to unknown is deliberately left running through the window — it is not a quarantine, nothing schedules against it, the next heartbeat undoes it, and "we have not heard from this worker since we started" is the honest thing to say about a worker we have not heard from since we started.
func (*Store) TakeKillRequests ¶
TakeKillRequests lists the jobs on a worker whose kill flag is due for delivery, and stamps them as delivered in the same transaction. It is a take, not a read: the stamp is what bounds re-delivery (see killRedeliverInterval), so a caller that reads without stamping would reintroduce the hot loop.
func (*Store) UpsertHostDoc ¶
UpsertHostDoc stores (or replaces) one usage-sheet document. deviceID == "" is the host-wide sheet (applies to every device on that host in `rc describe`); any other value is one device's own sheet, keyed by its full "host:name" ID. A second call for the same (host, deviceID) replaces the body outright — the row is a live snapshot of whatever the worker's disk currently holds, not an append-only log, so a deleted sheet file legitimately overwrites a previous body with the empty string.
func (*Store) UpsertWorker ¶
UpsertWorker registers a worker and its declared devices. A worker that registers is a fresh process announcing it has no running jobs — nothing a brand-new process could be supervising survives a restart — so registration must reconcile whatever the previous process left in flight before it touches device state at all. Without that reconciliation a restart either strands a job "running" forever with its device stuck busy (nothing else keys off worker identity to notice), or, once the reaper has already demoted the device, falsifies it as ready while an orphaned process from the dead worker may still be pinning it. Both are exactly what "never hand out a device we cannot prove is free" forbids, so every device backing a reaped in-flight job comes back unhealthy — never ready, never left busy — and only an explicit clear (or a verify probe standing in for one) puts it back in the pool.
func (*Store) WorkerHost ¶
WorkerHost returns the host a registered worker ID belongs to, so a caller can verify a request naming that worker ID actually agrees about which host it is — see handlePushLabels for why: without this check, a worker with a typo'd or stale `host:` in worker.yaml could push labels that silently land nowhere (a device ID for a host that doesn't exist) while reporting success forever, or overwrite another host's labels outright if the typo happens to collide with a real one.
An ID that matches no row yields ErrWorkerNotFound rather than an empty string, so "unknown worker" is never silently treated as "empty host".
type SweepOptions ¶
type SweepOptions struct {
RetainDisconnectedJobs bool
}