runqueue

package
v0.182.6 Latest Latest
Warning

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

Go to latest
Published: Oct 6, 2026 License: Apache-2.0 Imports: 18 Imported by: 0

Documentation

Overview

Package runqueue coordinates CPU-heavy commands across WB processes.

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func Budget

func Budget() int

Budget leaves one logical CPU available for the harness and operating system. Even a single-core machine retains one execution slot. It sizes the original small-machine spec table and is reported informationally elsewhere; the adaptive heavy-job share (see heavyShare) is computed from numCPU directly, not from this budget.

func EffectiveGOFLAGS added in v0.148.1

func EffectiveGOFLAGS() string

EffectiveGOFLAGS returns the GOFLAGS a plain `go` invocation would use: the process environment's own GOFLAGS if set, else the persisted default from `go env -w GOFLAGS=...` (review finding, PR #628, M6 — the process environment alone misses a persisted default, so a caller who pinned -race via `go env -w` rather than an exported variable was silently overridden). Best-effort: an error running `go env` (no toolchain on PATH, e.g.) degrades to "", matching the pre-existing behavior of reading only the process environment.

func GovernGOMAXPROCS added in v0.148.1

func GovernGOMAXPROCS(existing string, units int) string

GovernGoFlags returns the GOFLAGS value a governed "go" command should run with: the caller's own existing value is preserved, with -p=<units> appended so package parallelism follows the CPU allocation instead of defaulting to the whole machine — unless the caller already pinned -p, in which case the caller's flag wins untouched. argv that does not invoke the go tool, or units <= 0, returns existing unchanged. GovernGOMAXPROCS returns the GOMAXPROCS value a governed command's child should run with: the caller's own already-set value, preserved untouched (review finding, PR #628, S4: the spec says WB "sets GOMAXPROCS... from the allocation," but a caller that pinned its own value first must still win, the same way an explicit -p does for GOFLAGS), or units when the caller set nothing.

func GovernGoFlags added in v0.148.1

func GovernGoFlags(argv []string, existing string, units int) string

func LookupEnv added in v0.148.1

func LookupEnv(environment []string, key string) string

LookupEnv returns the value of key in environment (an os.Environ()-shaped slice), or "" if key is not present. Exported so every governed- environment builder (cmd/wb/run.go, cmd/wb/worker.go, internal/daemon/service.go) reads the caller's own environment consistently before deciding what to preserve.

func QueueDirForTest added in v0.164.10

func QueueDirForTest(projectsRoot string) string

QueueDirForTest returns exactly the directory queueRoot resolves for projectsRoot right now: the real, unoverridden path, or the active override's target when projectsRoot matches its key. It exists only so a test in another package — cmd/wb's TestMain, which installs the binary-wide isolation this seam exists for — can assert that its own isolation is actually in effect, rather than trusting it silently (PR #736 review finding B1, round 3). Production code must never call this.

func SetNumCPUForTest added in v0.148.1

func SetNumCPUForTest(n int) (restore func())

SetNumCPUForTest overrides the machine's logical CPU count every admission formula in this package reads, returning a restore func. It exists only so tests outside this package (cmd/wb, internal/daemon) can drive a specific N deterministically, the same way this package's own tests assign numCPU directly; production code must never call it.

func SetQueueRootForTest added in v0.164.10

func SetQueueRootForTest(fromProjectsRoot, dir string) (restore func())

SetQueueRootForTest overrides the directory queueRoot resolves to, but only for the exact fromProjectsRoot given — every other projectsRoot is unaffected. It returns a restore func. See queueRootOverrideFrom/ queueRootOverrideTo for this override's sequential-only contract. Production code must never call this; it exists for TestMain and for tests the same way SetNumCPUForTest exists for tests that need a specific NumCPU.

func Summary

func Summary(args []string) string

Summary names a program and its first non-flag verb without retaining full argv.

func Units

func Units(argv []string, budget int) int

Units is a stateless CPU-share estimate for argv: the small-machine table's exact number on a machine with numCPU < 8, or the share a large machine's job would get if admitted alone (k=1) otherwise. It exists for informational sizing (e.g. a durable operation record's initial estimate) where no live queue state is available; the real, state-dependent number for a heavy job on a large machine is decided by Admit at admission time and can differ from this estimate.

Types

type Admission added in v0.148.1

type Admission struct {
	// Lease must be released (Release is always safe, including on a nil
	// Lease or one that never held anything) once the command finishes.
	Lease *Lease
	// Units is the CPU share (GOMAXPROCS / Go -p) the command should run
	// with. Zero for KindNone.
	Units int
	// Waited is how long this call spent waiting before admission.
	Waited time.Duration
}

Admission is the result of requesting to run one governed command.

func Admit added in v0.148.1

func Admit(ctx context.Context, projectsRoot string, argv []string, self Participant, ticket *Ticket) (Admission, error)

Admit requests to run argv under the admission policy designed for sneat-dev/wb#621 (lead design, formalizing the founder's messages relayed during implementation):

  • KindNone is a no-op: Units 0, nothing to release, no wait.
  • KindFocused is admitted immediately with its fixed share (focusedShare, or the small-machine table's flat 1) and never takes a heavy slot or waits behind one.
  • On a small machine (numCPU < 8), a heavy kind (broad or coverage/race) uses the original budget-sum Acquire pool with the exact original spec-table weight (smallMachineUnits) — no adaptive sharing.
  • On a large machine (numCPU >= 8), a heavy kind is subject to admitHeavy's FIFO, k-tiered, 150%-capped share.

ticket, when non-nil, must already be registered via RegisterForAdmission against the same argv; when nil, Admit manages its own internally for a caller (worker/daemon executors) that does not render progress lines. Review finding (PR #628, S2): on a small machine every governed kind, focused included, must go through the original budget-sum Acquire and announcement exactly as on origin/main — the small-machine table has no "instant, unqueued" concept. So smallMachine() is checked first; only on a large machine does a focused job additionally get the new instant, never-queued path.

func AdmitExplicit added in v0.148.1

func AdmitExplicit(ctx context.Context, projectsRoot string, units int, self Participant) (Admission, error)

AdmitExplicit is Admit for a caller that declares its own CPU need directly instead of via argv classification — the daemon's trusted raw execution fallback ("wb v0.105.0 raw-execution policy remains available only through `wb daemon operation submit`") is the one caller today. units acquires from the plain budget-sum pool regardless of what argv actually is (that pool, not the adaptive heavy-job one, is the right fit for an operator-declared want, since bypassing classification means WB cannot tell a heavy job from a focused one); self is announced for `wb run --queue` visibility exactly as Admit's other branches do.

type Announcement added in v0.120.0

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

Announcement is the live handle returned by Lease.Announce. Heartbeat keeps the holder records fresh while the lease is held — call it from the same ~10s loop `wb run` already runs while a command executes, mirroring how a waiting Ticket is refreshed — so readHolders does not age the slots out as stale while the holder is legitimately still running. Cleanup removes the holder records once the lease is released; callers should defer it alongside Lease.Release. Both methods are safe to call on a nil Announcement (e.g. the zero-unit case where Announce never wrote anything).

func (*Announcement) Cleanup added in v0.120.0

func (announcement *Announcement) Cleanup()

Cleanup removes this announcement's holder records. Safe to call more than once and on a nil Announcement.

func (*Announcement) Heartbeat added in v0.120.0

func (announcement *Announcement) Heartbeat()

Heartbeat refreshes every announced holder record's UpdatedAt so readHolders keeps treating this lease as live. Best-effort, like Announce: a failed refresh just risks the holder aging past staleAfter and being reaped as if the process had died, which only affects visibility.

type Holder added in v0.120.0

type Holder struct {
	Participant
	Units     int       `json:"units,omitempty"`
	StartedAt time.Time `json:"started_at"`
	UpdatedAt time.Time `json:"updated_at"`
}

Holder is a Participant currently holding one or more CPU lease slots. Units is the CPU share this holder was fixed at, at admission time; the legacy budget-sum pool leaves it zero (its QueueEntry.Units is instead derived by counting held slot files — see groupRunningHolders), while a heavy holder on a large machine always stores its real, computed share here, since it holds exactly one bookkeeping record regardless of how many CPUs that share represents.

type Kind added in v0.148.1

type Kind int

Kind classifies a governed command for admission purposes.

const (
	// KindNone is not CPU-governed at all: no admission, no share.
	KindNone Kind = iota
	// KindFocused is a single-package Go test/vet or a light lint
	// (golangci-lint, staticcheck, pytest, vitest, jest, mocha, or an
	// nx/npm/pnpm/yarn/bun/npx test or lint script). Always admitted
	// immediately; never waits behind a heavy job and never takes a heavy
	// slot.
	KindFocused
	// KindBroad is a broad-scope Go/Node test or build (or an Angular/Nx
	// production build, or cargo test/build/check/clippy).
	KindBroad
	// KindRaceOrCover is any -race or -cover* Go run.
	KindRaceOrCover
)

func Classify added in v0.148.1

func Classify(argv []string) Kind

Classify reports what kind of governed work argv is. KindNone means argv is not CPU-governed at all.

type Lease

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

Lease holds units machine-wide until Release. Slot files live below the projects root so harnesses already permitted to write repositories can join the same budget without requiring access to the user's home directory.

func Acquire

func Acquire(ctx context.Context, projectsRoot string, units, budget int) (*Lease, time.Duration, error)

Acquire waits for units from one projects-root budget-sum lease pool (the small-machine table, and any other caller — e.g. internal/repositoryevents — sharing plain unit-weighted capacity). Each attempt either acquires every requested slot or releases all partial locks before waiting, preventing two multi-unit commands from deadlocking one another. There is no backfill or fairness ordering here: a flock race is fine for this pool because every caller today requests either 1 unit or a small, budget-bounded count — see Admit and admitHeavy for the FIFO, share-computing pool a heavy job on a large machine uses instead.

func (*Lease) Announce added in v0.120.0

func (lease *Lease) Announce(self Participant) *Announcement

Announce records this Lease's slots as held by self, for State/Snapshot and `wb run --queue` visibility. Best-effort: a failure to write a holder file just means that slot stays anonymous in State (readHolders skips slots with no holder file), never a hard error.

func (*Lease) Release

func (lease *Lease) Release()

type Participant added in v0.120.0

type Participant struct {
	PID      int    `json:"pid"`
	Summary  string `json:"summary"`
	Worktree string `json:"worktree,omitempty"`
}

Participant identifies who is waiting for or holding a CPU lease slot, for human-readable queue visibility (`wb run`'s queued/heartbeat/admitted lines and `wb run --queue`). Summary is a short, already-public label (the program name and its verb, e.g. "go test") — never full command arguments, paths, or flags — matching runlog's privacy-safe-telemetry contract even though these records are transient rather than durable.

type QueueEntry added in v0.120.0

type QueueEntry struct {
	PID      int           `json:"pid"`
	Summary  string        `json:"summary"`
	Worktree string        `json:"worktree,omitempty"`
	Units    int           `json:"units,omitempty"`
	Age      time.Duration `json:"age_ns"`
}

QueueEntry is one running or waiting governed command, for `wb run --queue`. Units is the number of CPU-budget slots a running holder currently occupies; it is omitted (zero) for a waiting entry, which has not been granted any slots yet.

type QueueListing added in v0.120.0

type QueueListing struct {
	Budget  int          `json:"budget"`
	Running []QueueEntry `json:"running"`
	Waiting []QueueEntry `json:"waiting"`
	HeavyK  int          `json:"heavy_k,omitempty"`
}

QueueListing is the inspectable state of the CPU lease queue for `wb run --queue`: who currently holds a slot, and who is waiting, oldest first, across both the legacy budget-sum pool and the large-machine heavy pool. HeavyK is the current k (sneat-dev/wb#621): every heavy job alive, running or waiting, right now — the same live count a heavy job arriving this instant would be sized against.

func ListQueue added in v0.120.0

func ListQueue(projectsRoot string, budget int) QueueListing

ListQueue reports every currently announced holder and registered waiter, legacy and heavy alike. It is read-only and safe to call from a separate `wb run --queue` invocation while other WB processes hold or wait for slots.

readHolders returns one record per held slot (Announce writes one holder file per unit in the lease), so a multi-unit legacy holder would otherwise print once per unit it holds. groupRunningHolders collapses those into one QueueEntry per distinct holder (PID + StartedAt identifies one Announce call, i.e. one lease) and sums the collapsed records' Units. A heavy holder needs no such collapsing — it is already exactly one record — but shares the same helper for a uniform QueueEntry shape.

type State added in v0.120.0

type State struct {
	// Position is this ticket's 1-based rank among current waiters, oldest
	// first. Zero when the ticket is not registered (units <= 0) or the
	// caller asked for Peek rather than a specific Ticket's Snapshot.
	Position int
	// Total is the number of tickets currently registered as waiting.
	Total int
	// Holders are the Participants currently holding CPU lease slots,
	// oldest first. It can be shorter than the busy-slot count when a
	// holder never announced (e.g. an older WB binary, or a caller other
	// than `wb run` sharing the same budget).
	Holders []Holder
}

State is a point-in-time snapshot of the CPU lease queue relevant to one waiter: its position among registered waiters, how many waiters are registered in total, and who currently holds slots.

func Peek added in v0.120.0

func Peek(projectsRoot string, budget int) State

Peek reports legacy-pool queue state without registering a waiter. No production caller remains: `wb run --queue` lists through ListQueue. It is kept for cmd/wb tests that observe the pool.

type Ticket added in v0.120.0

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

Ticket is one caller's registered wait, used for queue visibility and, in the heavy namespace only, for strict FIFO ordering (see admitHeavy). The legacy namespace's Ticket never gates admission itself.

Review finding (PR #628, S5): admitHeavy's own loop and a caller's external progress ticker (cmd/wb/run.go) can both call Heartbeat/Forget on the same *Ticket concurrently. mu guards every read and write of path so Forget can never race a concurrent Heartbeat into "recreating" a ticket (Heartbeat always re-checks path under the same lock Forget clears it under, so once cleared it stays cleared).

func Register added in v0.120.0

func Register(projectsRoot string, self Participant) *Ticket

Register records a legacy-pool waiter so State/Snapshot can report queue position and depth while units > 0. Callers must call Forget once they stop waiting, admitted or not — typically via defer immediately after Register. Registration is best-effort: a failure to write the ticket file degrades to an invisible waiter (Snapshot reports Position 0) rather than blocking or failing the caller, since visibility must never become a new way for `wb run` to hang or refuse work.

func RegisterForAdmission added in v0.148.1

func RegisterForAdmission(projectsRoot string, argv []string, self Participant) *Ticket

RegisterForAdmission registers argv's waiting ticket in the namespace its Kind uses — nil for KindNone always, and for KindFocused only on a large machine (numCPU >= 8), where a focused job is admitted instantly and never waits — so a caller that wants to render its own queued/heartbeat progress lines (only `wb run` does today) can call Ticket.Snapshot while Admit is in flight. The caller must call Forget on whatever is returned exactly once, regardless of outcome; Forget on a nil Ticket is a safe no-op.

func RegisterHeavy added in v0.148.1

func RegisterHeavy(projectsRoot string, self Participant) *Ticket

RegisterHeavy registers a heavy-job waiter in the FIFO queue a large machine (numCPU >= 8) uses instead of the small-machine budget-sum pool. See Register for the general contract (Forget once, best-effort).

func (*Ticket) Forget added in v0.120.0

func (ticket *Ticket) Forget()

Forget removes the waiter's ticket. Safe to call more than once and on a Ticket whose registration never succeeded.

func (*Ticket) Heartbeat added in v0.120.0

func (ticket *Ticket) Heartbeat()

Heartbeat refreshes the ticket's UpdatedAt so readTickets keeps treating it as live while its caller is still waiting. Callers already poll queue state on an interval (e.g. `wb run`'s queued-command heartbeat every ~10s); call Heartbeat from that same loop rather than adding a new one. Safe to call on a nil Ticket or one whose registration never succeeded, and best-effort like Register: a failed refresh just means the ticket may age past staleAfter and be reaped as if its process had died, which only affects visibility, never admission.

func (*Ticket) Snapshot added in v0.120.0

func (ticket *Ticket) Snapshot(budget int) State

Snapshot reports this ticket's current position and the queue depth, plus the Participants currently holding CPU lease slots in this ticket's own namespace. Safe to call repeatedly (e.g. from a heartbeat); it reflects other WB processes' registrations and slot holdings on disk at the moment of the call.

Jump to

Keyboard shortcuts

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