worker

package
v0.1.2 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Aug 29, 2026 License: Apache-2.0 Imports: 26 Imported by: 0

Documentation

Overview

Package worker runs jobs on a device host and reports back to the controller.

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func BootID

func BootID() string

BootID returns a value that changes when the machine reboots and is stable while it stays up. The controller uses a change as proof that nothing from before can still be holding a device; an empty value is treated as no proof, which quarantines rather than frees.

Types

type Config

type Config struct {
	ControllerURL     string         `yaml:"controller_url"`
	Token             string         `yaml:"token"`
	Host              string         `yaml:"host"`
	Devices           []DeviceConfig `yaml:"devices"`
	HeartbeatInterval time.Duration  `yaml:"heartbeat_interval"`
	PollWait          time.Duration  `yaml:"poll_wait"`
	// Hooks holds host-level defaults for the lifecycle hook timeout and
	// release linger, so a multi-GPU box need not repeat them on every
	// device. Both are optional; unset fields fall back to
	// defaultHookTimeout / defaultReleaseLinger.
	Hooks HooksConfig `yaml:"hooks"`
	// ProbeDir is where this worker looks for drop-in probe executables,
	// run in name order on top of the built-ins to gather device labels.
	// Optional; falls back to defaultProbeDir.
	ProbeDir string `yaml:"probe_dir"`
	// ProbeInterval is how often a full probe pass re-runs while the worker
	// is up, so a label picks up a change (a card swap, a driver upgrade)
	// without a restart. Optional; falls back to defaultProbeInterval.
	ProbeInterval time.Duration `yaml:"probe_interval"`
	// ProbeTimeout bounds any single probe — a built-in nvidia-smi call or
	// one drop-in executable — before it is killed as a process group and
	// skipped. Optional; falls back to defaultProbeTimeout.
	ProbeTimeout time.Duration `yaml:"probe_timeout"`
	// VerifyDir is where this worker looks for drop-in verify script
	// executables, run in name order once a job's process tree is
	// confirmed gone and before the terminal report frees the device.
	// Optional; falls back to defaultVerifyDir.
	VerifyDir string `yaml:"verify_dir"`
	// VerifyTimeout bounds any single verify script before it is killed as
	// a process group and its exit is treated as a failure. Optional;
	// falls back to defaultVerifyTimeout.
	VerifyTimeout time.Duration `yaml:"verify_timeout"`
	// VerifyPassBudget bounds the WHOLE verify pass — every script in
	// VerifyDir together — not each script's own VerifyTimeout summed.
	// Optional; falls back to defaultVerifyPassBudget.
	VerifyPassBudget time.Duration `yaml:"verify_pass_budget"`
	// SheetDir is where this worker looks for its usage-sheet documentation:
	// <SheetDir>/host.md and <SheetDir>/host.d/<device>.md. Optional; falls
	// back to defaultSheetDir.
	SheetDir string `yaml:"sheet_dir"`
	// RequireManualClear stops this worker from claiming anything at
	// registration about processes left behind by an interrupted job (see
	// recovery.go), so its quarantined devices keep the behaviour they have
	// always had: they come back when an admin runs `rc clear`, or when a
	// proven reboot answers them, and not otherwise.
	//
	// It only ever points one way. There is deliberately NO setting that
	// forces auto-recovery, because that is the only direction in which a
	// misconfigured switch hands out a device with a live process on it: a
	// mistake here costs a manual clear, and nothing worse. That asymmetry is
	// also why it is ORed rather than overridden by the environment (see
	// LoadConfig) — every source can make this worker more cautious, and none
	// can make it less.
	RequireManualClear bool `yaml:"require_manual_clear"`
}

func LoadConfig

func LoadConfig(path string) (Config, error)

LoadConfig reads /etc/rc/worker.yaml (or another path) and applies defaults.

type DeviceConfig

type DeviceConfig struct {
	Name       string        `yaml:"name"`
	MaxRuntime time.Duration `yaml:"max_runtime"`
	// OnAcquire and OnRelease are paths to scripts run (via the same
	// process-group supervision jobs get) when this device transitions
	// worker-side between free and held. Both are optional and
	// independent — either, neither, or both may be set.
	OnAcquire string `yaml:"on_acquire"`
	OnRelease string `yaml:"on_release"`
	// HookTimeout and ReleaseLinger override the host-level Hooks defaults
	// for this device only. Zero means "not overridden": LoadConfig
	// resolves each to the host default (or the built-in default, if the
	// host declared none) once loading is complete.
	HookTimeout   time.Duration `yaml:"timeout"`
	ReleaseLinger time.Duration `yaml:"release_linger"`
	// Labels are this device's declared facts — operator-asserted, as
	// opposed to what a probe detects — keyed the same flat way a probe's
	// output is: a bare key ("rack") or, though redundant for a per-device
	// list, a "<device-name>.<label>" key is accepted wherever labels are
	// merged. Optional.
	Labels map[string]string `yaml:"labels"`
}

DeviceConfig is one device this host offers. It accepts either a bare name (stage 1 style) or an object with a runtime ceiling and optional lease lifecycle hooks.

func (*DeviceConfig) UnmarshalYAML

func (d *DeviceConfig) UnmarshalYAML(value *yaml.Node) error

type HooksConfig

type HooksConfig struct {
	Timeout       time.Duration `yaml:"timeout"`
	ReleaseLinger time.Duration `yaml:"release_linger"`
}

HooksConfig is the host-level default for lease lifecycle hooks. A per-device value (DeviceConfig.HookTimeout / ReleaseLinger) overrides it.

type JobSpec

type JobSpec struct {
	Command      []string
	Cwd          string
	Env          map[string]string
	GraceCeiling time.Duration // SIGTERM -> SIGKILL window; default 10s
	MaxRuntime   time.Duration // total wall-clock ceiling; 0 means no limit
	IdleTimeout  time.Duration // max gap with no stdout/stderr output; 0 means no limit
	// Stderr, if non-nil, receives the process's stderr separately from the
	// sink passed to Run (which then carries only stdout). Nil (the
	// default) preserves today's behavior for jobs and hooks: stdout and
	// stderr both land in sink, combined, so a live `rc run` sees one
	// interleaved stream. A caller that needs to parse stdout on its own —
	// a probe emitting a single JSON object, where a stray stderr line
	// would otherwise corrupt the parse — sets this to keep the streams
	// apart.
	Stderr io.Writer
	// Stdin, if non-nil, is the process's standard input. Nil (the default)
	// leaves it as /dev/null, which is what every unattended job has always
	// had: a job nobody is watching that blocks reading a terminal nobody is
	// typing at would sit there holding the GPU until a watchdog fired.
	//
	// It is wired through an os.Pipe this package creates rather than handed
	// to os/exec directly, and that is not incidental — see Run.
	Stdin io.Reader
	// SweepJobID, when non-empty, turns on the RC_JOB_ID sweep: once the
	// process group (or PTY session) has been torn down, every remaining
	// process whose environment carries exactly this job id is signalled out
	// of existence too. That is what catches a child which called setsid(2)
	// and so left the group and the session the kill was aimed at. See
	// sweepJobSurvivors in sweep.go for the mechanism and its residual.
	//
	// It is an explicit opt-in and not read out of Env["RC_JOB_ID"], even
	// though execute() sets both to the same value, because the callers must
	// differ: a JOB gets swept, while lifecycle hooks (hooks.go) and verify
	// probes (verify.go) run with the same RC_JOB_ID in their environment and
	// must never sweep — a hook or probe hunting for processes by job id is
	// a mechanism nobody asked for, aimed at a job that is either about to
	// start or already accounted for.
	SweepJobID string
}

type ProbeResult

type ProbeResult struct {
	Host   map[string]string
	Device map[string]map[string]string

	// Unconfirmed names the devices whose OWN device-scoped facts this pass
	// could not confirm, because a source understood to speak for that
	// device failed: nvidia-smi found on PATH but erroring, timing out, or
	// emitting something that doesn't parse; nvidia-smi not found on PATH
	// THIS pass despite having been found on this worker process's FIRST
	// pass (that one snapshot, w.nvidiaSmiSeenAtStartup, is the whole
	// history kept — see gatherLabels for the "seen at startup" rule this
	// depends on; a host that has never once seen nvidia-smi does NOT mark
	// anything unconfirmed for that reason, or a device's labels could never
	// be cleared on a GPU-less host); or a drop-in probe erroring, being
	// skipped for budget, or failing to stat.
	//
	// This is what lets a caller (see labelsPayload) tell "this device's
	// probes ran and confirmed there is nothing to report" — which
	// ReplaceLabels must be allowed to turn into a clear, per Task 3's own
	// design — apart from "something that can produce THIS device's facts
	// broke this pass, so an empty result here proves nothing about the
	// device itself".
	//
	// It is per DEVICE, not one flag for the whole pass. The pass-wide
	// `Failed` bool this replaces argued that a single flag was NECESSARY,
	// because gatherLabels "has no way to know in advance which device a
	// probe would have named had it succeeded". That premise is right, and
	// it is still why the attribution below is deliberately conservative —
	// but the conclusion drawn from it was wrong, and Stage 3's final
	// review measured the damage:
	//
	//   - HOST-scoped facts apply to every device by construction, so one
	//     source failing cannot make ANOTHER source's host facts doubtful.
	//     Under the pass-wide flag it did exactly that: with one unrelated
	//     drop-in broken, disk_free_bytes and mem_total_bytes froze at
	//     their last good values on every device of the host — facts
	//     builtinLabels had gathered perfectly well on that same pass —
	//     because a device omitted from the payload never gets
	//     ReplaceLabels called on it at all (see labelsPayload here, and
	//     applyDeviceFacts server-side). A full box kept advertising
	//     disk_free_bytes=72G, `--select 'disk_free_bytes>=50G'` routed a
	//     job onto it, and the job died on write.
	//   - Which devices a source can speak for is not, in fact, wholly
	//     unknowable. nvidia-smi's keys are "gpu<N>.<label>" by
	//     construction, so it can only ever name a device this host
	//     declares under that exact form — see nvidiaScopedDevices. A
	//     drop-in script is opaque, but its LAST SUCCESSFUL run in THIS
	//     process is evidence about it: a script that named gpu0 the last
	//     time it worked is understood to speak for gpu0, and for nothing
	//     else, until it succeeds again and says otherwise — see
	//     Worker.probeSourceDevices.
	//
	// What has NOT changed is what being unconfirmed then means: preserve,
	// do not clear. A device named here is omitted from the payload
	// entirely rather than sent as an empty map, exactly as before, because
	// wiping a device's facts on a guess is fleet-wide while a stale label
	// is one device.
	//
	// A source this worker process has never once seen succeed speaks for
	// no device, so its failure marks nothing unconfirmed. That is the
	// deliberate cost of the fix: a drop-in that reported gpu0's facts for
	// a PREVIOUS process and is broken by the time this one starts will
	// have those facts cleared on the first pass, instead of frozen
	// forever. It is the same trade gatherLabels already makes for
	// nvidia-smi's "seen at startup" rule (see its comment there, and the
	// residual it accepts): a process restart is a far narrower window than
	// every probe interval on every host for the rest of a worker's life,
	// and registration is already the moment the controller reconciles a
	// worker's whole world. nvidia-smi itself is exempt — its scope is
	// static, so it needs no such history (nvidiaScopedDevices again).
	Unconfirmed map[string]bool
}

ProbeResult is what one gatherLabels pass produced: host-wide facts that apply to every device on this host, facts scoped to one device by name, and the devices whose own facts this pass could not confirm. All three maps are always non-nil, even when nothing was gathered, so a caller never has to nil-check before ranging or indexing.

type Result

type Result struct {
	ExitCode int
	Killed   bool
	Reason   string
	Err      error
	// Sweep records what the RC_JOB_ID sweep found and did, when
	// JobSpec.SweepJobID asked for one. It is a report, NOT an outcome:
	// nothing about it ever touches Killed, Reason or ExitCode, because a
	// job that exited 0 and left a detached process behind is a successful
	// job. See SweepReport in sweep.go for what happened the last time a
	// straggler sweep was allowed to relabel a job.
	Sweep SweepReport
}

func Run

func Run(ctx context.Context, spec JobSpec, sink io.Writer) Result

Run spawns the command in its OWN process group and merges stdout and stderr into sink. On cancellation the whole group is signalled, so children that outlive their parent — the ones still holding device memory — die too.

type SweepReport

type SweepReport struct {
	// Scanned reports whether /proc could be walked at all. "Found nothing"
	// and "could not look" are opposite answers and are never conflated
	// here — the same distinction liveProcExists draws for recovery.
	Scanned bool
	// Found is every live process that carried this job's exact RC_JOB_ID
	// after the job's process group had been torn down. Empty is the
	// overwhelmingly common case: the group kill already did the work.
	Found []int
	// Remaining is what was still alive after SIGTERM, the grace window and
	// SIGKILL. Non-empty means the kernel is holding a process we cannot
	// remove (uninterruptible sleep in a driver call, most realistically),
	// and the device it is sitting on is still dirty.
	Remaining []int
	// Blind reports that at least one process could not be inspected —
	// EACCES on the /proc/<pid>/environ of a process this worker does not
	// own, i.e. a worker not running as root. A process we could not look at
	// is not evidence of a survivor AND not evidence of a clean box; it is
	// only evidence that the answer here is incomplete.
	Blind bool
}

SweepReport is what the RC_JOB_ID sweep saw and did, and it exists as its own field on Result for one reason, which is the most important sentence in this file:

REAPING A SURVIVOR MUST NEVER RELABEL A SUCCESSFUL JOB AS KILLED.

On 2026-08-17 exactly that mistake shipped. A straggler sweep set Killed = true on the clean-exit path, and because cmd.Wait() returns the instant the last holder of the inherited stdout pipe becomes a zombie, 40 out of 40 ordinary successful jobs came back reporting themselves killed (fixed in 337e23d, "a zombie is not a straggler"). The blast radius was not cosmetic: worker.go maps Killed onto model.JobKilled, hooks.go treats a killed hook as a failed hook — so an on_acquire hook that backgrounds anything refuses the lease — and verify.go treats a killed probe as a failed probe, which quarantines the device.

A job that exited 0 and happened to leave a detached process behind is a SUCCESSFUL job that left a mess. The mess is reported here and logged by execute(); the job's own outcome is untouched. Nothing in this file writes to Result.Killed, Result.Reason or Result.ExitCode, and nothing in it should ever be made to.

func (SweepReport) Clean

func (s SweepReport) Clean() bool

Clean reports whether the sweep positively established that nothing of this job is left running: it looked, it could see everything it looked at, and either found nothing or removed everything it found. An empty Found on its own is NOT that claim — a worker that is not root cannot read most of the process table, so "I found nothing" and "there is nothing" are different sentences and only this one says the second.

type VerifyResult

type VerifyResult struct {
	// OK is true when every verify script exited zero, or there were none
	// to run (including a VerifyDir that does not exist at all — the
	// feature is simply off on a host that ships no scripts).
	OK bool
	// Reason is empty when OK. Otherwise it always starts with
	// verifyReasonPrefix, followed by the first failing script's name and
	// the tail of its stderr — everything an operator needs to see WHY a
	// device was quarantined, and everything the controller needs to
	// recognise this fault as verify-sourced (it keys the verify_failed
	// event off that prefix).
	//
	// Where an operator reads it: the worker's log, the controller's log,
	// and the verify_failed webhook event. NOT `rc ps` or `rc devices` —
	// the device row stores only the fixed quarantine reason `fault`, and a
	// verify failure leaves the job itself succeeded, so no job failure
	// report carries this either. Keep it short for a log line and an event
	// payload, not for a table cell.
	Reason string
}

VerifyResult is the outcome of one runVerify pass over VerifyDir.

type Worker

type Worker struct {
	// contains filtered or unexported fields
}

func New

func New(cfg Config) *Worker

func (*Worker) Start

func (w *Worker) Start(ctx context.Context) error

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL