dispatch

package
v1.0.0-rc.2 Latest Latest
Warning

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

Go to latest
Published: Oct 6, 2026 License: AGPL-3.0, AGPL-3.0-or-later Imports: 28 Imported by: 0

Documentation

Overview

Package dispatch is the CAL-V0-052..058 continuous dispatcher: a deterministic roster over native queue state that launches, supervises and heals independent host worker processes. It owns only its own state directory; every queue change goes through the existing lease transactions.

Index

Constants

View Source
const (
	ConfigProfile = "taskman-dispatch/0"
	StateProfile  = "taskman-dispatch-state/0"
	EventProfile  = "taskman-dispatch-event/0"
	MaxConfig     = 256 << 10
)
View Source
const (
	SessionProgress       = "progress"
	SessionInfrastructure = "infrastructure"
	SessionHeld           = "held"
	SessionNoProgress     = "no-progress"
)

Session classes of an ended worker (ESC-V0-008), in precedence order: checked progress, a current typed decision, scope or blocked hold, the worker's own infrastructure request under the ESC-V0-007 policy, and ordinary no-progress. A hold outranks infrastructure, so a session that raised both spends no retry and leaves no infrastructure hold behind the owner's answer.

View Source
const (
	InfraWaiting         = "WAITING"
	InfraReserved        = "RESERVED"
	InfraRunning         = "RUNNING"
	InfraIdle            = "IDLE"
	InfraRecovered       = "RECOVERED"
	InfraExhausted       = "EXHAUSTED"
	InfraDisabled        = "DISABLED"
	InfraNativeExhausted = "NATIVE_EXHAUSTED"
	InfraUnknown         = "UNKNOWN"
)

Infrastructure retry episode states (ESC-V0-007). WAITING has a reserved retry ordinal and cooldown deadline; RESERVED has also written its launch identity before spawn; RUNNING is that launch's recorded worker; IDLE keeps the debt with no retry pending; RECOVERED followed checked work progress. EXHAUSTED, DISABLED, NATIVE_EXHAUSTED and UNKNOWN hold the ticket.

View Source
const (
	StateNone    = "NONE"
	StateUnknown = "UNKNOWN"
)

Work states that are not program-defined.

View Source
const (
	GatePass        = "PASS"
	GateReturn      = "RETURN"
	GateResubmitted = "RESUBMITTED"
	GateNone        = "NONE"
)

Gate predicate states. NONE is a gate with no record; the others need a CURRENT head, so a STALE or UNKNOWN gate matches no predicate.

View Source
const MaxLoopGenerations = 1024

MaxLoopGenerations bounds the generations one recorded loop hold names.

Variables

View Source
var ErrReaderQuiescence = errors.New("QUIESCENCE_UNPROVED: reader containment or quarantine evidence UNKNOWN")

ErrReaderQuiescence means reader containment or its restart quarantine is unproved. It outranks ordinary cancellation and never permits another tick.

View Source
var ErrSettled = errors.New("dispatcher settled by its service control")

ErrSettled ends Run when the controlled dispatcher's control settled it.

View Source
var EventKinds = []string{"started", "stopped", "adopted", "launched", "launch-failed", "finished", "killing", "killed", "handoff", "handoff-refused", "reaped", "state", "claim", "release", "lane", "cooldown", "parked", "unparked", "alert", "needs-owner", "throttled", "escalated"}

EventKinds is the closed CAL-V0-058 event vocabulary.

View Source
var Placeholders = []string{"{program}", "{role}", "{slot}", "{worker}", "{holder}", "{ticket}", "{ticketLocal}", "{state}", "{pool}", "{member}", "{workRoot}", "{prompt}", "{model}", "{nextStage}", "{operatorNote}"}

Placeholders are the only substitutions in argv, env and prompts. {operatorNote} renders untrusted operator prose, so only a role prompt may use it (ON-V0-011).

Functions

func Fingerprint

func Fingerprint(obs *Observation, key string) string

Fingerprint is the CAL-V0-057 progress identity of one work key: ticket status, revision and work state plus every attempt that carries durable work, with an optional CAL-V0-064 digest from checked progress admission. An empty claim followed by a handoff leaves it unchanged.

func LaunchTier

func LaunchTier(r *Role, e *EscalationState) int

LaunchTier is the tier a role launches at on a ticket with this ladder record: the streak's tier, or the role's last tier when it is higher and the role does not deescalate.

func OwnerProcess

func OwnerProcess(dir string) (string, int)

OwnerProcess is OwnerState with the recorded holder's PID, which is 0 unless the state is RUNNING.

func OwnerState

func OwnerState(dir string) string

OwnerState is RUNNING when the recorded lock holder still runs with its recorded identity, NOT_RUNNING when no holder is recorded or it is gone, and UNKNOWN when the record cannot be read.

func ProgramDir

func ProgramDir(c *Config, program string) string

ProgramDir is the dispatcher's state directory for one program.

func ReadStates

func ReadStates(ctx context.Context, c *Config, tickets []Ticket) []string

ReadStates is CAL-V0-053's program-defined work state. It fills every ticket's State with a value, NONE, or UNKNOWN, and returns one alert per failed read. A missing reader leaves every state NONE.

func ReaderContainment

func ReaderContainment(dir, program string) (state string, quarantined bool, diagnostic string)

ReaderContainment is a pure read. Absence is NOT_OBSERVED, never RELEASED.

func Records

func Records(l *Ledger, intent string) bool

Records reports whether a saved ledger records the effect an admission intent named: a worker or a pool sweep request.

func Render

func Render(s string, values map[string]string) string

Render substitutes placeholders in one pass; values never re-expand.

func RenderOperatorNote

func RenderOperatorNote(ticket string, n *NoteView) string

RenderOperatorNote is the {operatorNote} role-prompt substitution: empty for a never-noted ticket, so a prompt renders byte-identically to one without the placeholder. Otherwise it states the note's state, revision, provenance and launch-time freshness, and that the claim result's operatorNote member, pinned by the worker's own admission, supersedes it. The note text is substituted once and never re-expanded.

func RetryCooldown

func RetryCooldown(cooldown, maxCooldown, n int) time.Duration

RetryCooldown is the cooldown before retry ordinal n (1-based): min(maxCooldown, cooldown*2^(n-1)), saturating.

func Settled

func Settled(l *Ledger) bool

Settled reports a saved ledger with no recorded worker and no pending pool sweep: what a draining service needs before it may stop.

func Summary

func Summary(logDir string) string

Summary is the sanitized last part of a worker's stdout.

func ValidName

func ValidName(s string) bool

ValidName is the bound on program and role names; worker IDs and holder labels are built from them and must stay inside a 64-byte label.

Types

type Assignment

type Assignment struct {
	Role, Key, Ticket, Local, State, Pool, Member string
	NextStage                                     string
	Slot, Tier                                    int
	// OperatorNote is the ticket's observed note, rendered only through the
	// {operatorNote} role-prompt placeholder (ON-V0-011).
	OperatorNote *NoteView
}

Assignment is one roster decision: a role slot bound to one work key at one CAL-V0-057 escalation tier (0 is the base model).

func Roster

func Roster(c *Config, obs *Observation, busy []Busy, skip map[string]bool) []Assignment

Roster is the CAL-V0-054 pure roster: the same configuration, observation, running set and skip set always produce the same assignments.

func RosterTiers

func RosterTiers(c *Config, obs *Observation, busy []Busy, skip map[string]bool, tierOf func(role, key string) int) []Assignment

RosterTiers is Roster with each candidate's escalation tier from the pure tierOf (nil means every tier is 0). A candidate whose tier is at its tier cap waits; it never falls back to a lower tier.

func RosterWithPressure

func RosterWithPressure(c *Config, obs *Observation, busy []Busy, skip map[string]bool, budget *PressureBudget) (out, held []Assignment)

RosterWithPressure is Roster with an optional CAL-V0-068 pressure budget. The budget is consulted only after every static fence admits a candidate and before the candidate consumes a slot or reserves its key, so a held candidate reserves nothing. held lists each held work key once, in roster order, while the role, tier and global caps (charged with launches and earlier holds) would have admitted it, unless a later exempt candidate launched the same key.

type Attempt

type Attempt struct {
	ID, Ticket, Phase, Stage, Generation, Holder, Candidate string
	LeaseExpires                                            time.Time
	Live                                                    bool
	Gates, Reviews                                          int

	// Pool and Member name the attempt's pool allocation; empty without one.
	Pool, Member string
}

Attempt is the dispatcher's view of one attempt.

type Backoff

type Backoff struct {
	CooldownSeconds int `json:"cooldownSeconds"`
	ParkAfter       int `json:"parkAfter"`
}

type BackoffState

type BackoffState struct {
	NoProgress      int       `json:"noProgress"`
	CooldownUntil   time.Time `json:"cooldownUntil"`
	Parked          bool      `json:"parked"`
	Fingerprint     string    `json:"fingerprint"`
	BaseFingerprint string    `json:"baseFingerprint,omitempty"`
	ProgressDigest  string    `json:"progressDigest,omitempty"`
}

Backoff is the CAL-V0-057 per-key no-progress record.

type Busy

type Busy struct {
	Role, Key  string
	Slot, Tier int
}

Busy is a running worker's claim on a role slot, a work key and a tier.

type Config

type Config struct {
	PoolSweep        *PoolSweepConfig `json:"poolSweep,omitempty"`
	Profile          string           `json:"profile"`
	StateDir         string           `json:"stateDir"`
	WorkRoot         string           `json:"workRoot"`
	TickSeconds      int              `json:"tickSeconds"`
	GlobalCap        int              `json:"globalCap"`
	KillGraceSeconds int              `json:"killGraceSeconds"`
	Hosts            map[string]Host  `json:"hosts"`
	WorkState        *WorkState       `json:"workState,omitempty"`
	Roles            []Role           `json:"roles"`
	Pinned           []string         `json:"pinned,omitempty"`
	Backoff          Backoff          `json:"backoff"`
	Heal             Heal             `json:"heal"`
	// Pressure is the optional CAL-V0-068 host-pressure launch throttle.
	Pressure *PressureConfig `json:"pressure,omitempty"`
	// InfrastructureRetry is the optional ESC-V0-007 infrastructure retry
	// policy. Without it an infrastructure session is ordinary no-progress.
	InfrastructureRetry *InfraRetryConfig `json:"infrastructureRetry,omitempty"`
}

Config is the closed taskman-dispatch/0 operator configuration.

func DecodeConfig

func DecodeConfig(raw []byte) (*Config, error)

DecodeConfig parses and validates closed configuration bytes.

func (*Config) Escalates

func (c *Config) Escalates() bool

Escalates reports whether any role has a ladder.

func (*Config) TicketPools

func (c *Config) TicketPools() []string

TicketPools is the sorted set of pools the ticket roles' workers claim: the only pools whose tickets the dispatcher's plan may select (CAL-V0-097). It is empty, never nil, when no role names a pool.

type Dispatcher

type Dispatcher struct {
	Program string
	Config  *Config
	Queue   Queue
	Out     io.Writer
	Now     func() time.Time
	// contains filtered or unexported fields
}

Dispatcher is one program's continuous dispatcher. Only one runs per state directory; the lock is held for its lifetime.

func Open

func Open(program string, c *Config, q Queue, out io.Writer) (*Dispatcher, error)

Open locks the program's state directory, loads the ledger and adopts every worker recorded by a previous dispatcher.

func OpenControlled

func OpenControlled(program string, c *Config, q Queue, out io.Writer, control LaunchControl) (*Dispatcher, error)

OpenControlled is Open for a dispatcher whose every new worker launch and pool sweep native call first passes control. It takes the lifetime lock before the fence; supervision, healing and accounting of recorded workers are never fenced.

func (*Dispatcher) Close

func (d *Dispatcher) Close() error

Close records the stop and releases the lock. Workers keep running and are adopted by the next dispatcher.

func (*Dispatcher) LastEvent

func (d *Dispatcher) LastEvent() uint64

LastEvent is the last event sequence number.

func (*Dispatcher) Run

func (d *Dispatcher) Run(ctx context.Context, ticks int) (result error)

Run ticks until ctx ends or ticks reach the bound (0 is unbounded). A failed tick is an alert, not an exit, so the dispatcher keeps supervising.

func (*Dispatcher) Running

func (d *Dispatcher) Running() int

Running is the number of supervised workers.

func (*Dispatcher) Tick

func (d *Dispatcher) Tick(ctx context.Context) error

Tick is one CAL-V0-053..058 pass: observe, supervise, heal, re-observe, account for finished workers, launch the roster, and emit changes. Idle, wall and orphan enforcement runs even when the store is unreadable; ended workers then stay recorded and are accounted on the next readable tick.

type EscalationState

type EscalationState struct {
	Streak int            `json:"streak"`
	Tiers  map[string]int `json:"tiers,omitempty"`
}

EscalationState is the CAL-V0-057 ladder record of one ticket: the trailing count of finished sessions without progress, and the tier each escalating role last launched at.

type Event

type Event struct {
	Profile string            `json:"profile"`
	Seq     uint64            `json:"seq"`
	At      string            `json:"at"`
	Program string            `json:"program"`
	Kind    string            `json:"kind"`
	Ticket  string            `json:"ticket,omitempty"`
	Role    string            `json:"role,omitempty"`
	Worker  string            `json:"worker,omitempty"`
	Message string            `json:"message"`
	Detail  map[string]string `json:"detail,omitempty"`
}

Event is one taskman-dispatch-event/0 line. Message is plain language.

func ReadEvents

func ReadEvents(dir string, n int) ([]Event, error)

ReadEvents returns the last n events of the current log.

type GateCandidate

type GateCandidate struct{ Kind, TreeOID, Sha256, Bytes string }

GateCandidate is the reviewed candidate: Kind TREE with TreeOID, or Kind EVIDENCE with Sha256 and Bytes.

type GateMatch

type GateMatch struct {
	Gate   string   `json:"gate"`
	States []string `json:"states"`
}

GateMatch requires one gate's routing state (Ticket.GateState) to be one of States: PASS, RETURN, RESUBMITTED or NONE.

type GateSubject

type GateSubject struct {
	AttemptID, Generation, ReceiptSeq, ReceiptSha256, AttemptSha256 string
}

GateSubject is the author's reviewed submission.

type GateTrust

type GateTrust struct{ ActorAuthentication, Independence, Source string }

GateTrust is the head event's trust: Source LEASE_BOUND or OPERATOR_ATTESTED; actor authentication and independence are always NOT_OBSERVED (ERG-V0-001).

type GateView

type GateView struct {
	Verdict, Status      string
	Generation, Revision string
	Resubmitted          bool
	Head                 string
	// Candidate, Subject, Trust and EvidenceSha256 complete the ERG-V0-009
	// field set. Each is read from the head event, so the progress
	// fingerprint, which already carries Head, is unchanged. EvidenceSha256
	// is the head event's digest; all four are empty without a head.
	Candidate      GateCandidate
	Subject        GateSubject
	Trust          GateTrust
	EvidenceSha256 string
}

GateView is one gate of the native ERG-V0-009 external review state, read from typed queue events, never from a workState program. Status is CURRENT, STALE or UNKNOWN. Verdict is PASS, RETURN, or empty for a current resubmission awaiting review.

type Heal

type Heal struct {
	Handoff bool `json:"handoff"`
	Reap    bool `json:"reap"`
}

type HeldLaunch

type HeldLaunch struct {
	Role   string `json:"role"`
	Key    string `json:"key"`
	Ticket string `json:"ticket,omitempty"`
}

HeldLaunch is one roster candidate held by the pressure budget.

type Host

type Host struct {
	Argv          []string          `json:"argv"`
	Env           map[string]string `json:"env,omitempty"`
	IdleIgnore    []string          `json:"idleIgnore,omitempty"`
	ActivityPaths []string          `json:"activityPaths,omitempty"`
}

Host is one worker runtime. argv[0] is an absolute executable; every argv and env value may use the placeholders in Placeholders.

type InfraEpisode

type InfraEpisode struct {
	AcceptanceRevision string    `json:"acceptanceRevision"`
	State              string    `json:"state"`
	Sessions           int       `json:"sessions"`
	Charged            int       `json:"charged"`
	Limit              int       `json:"limit"`
	CooldownUntil      time.Time `json:"cooldownUntil"`
	Launch             string    `json:"launch,omitempty"`
}

InfraEpisode is one ticket's ESC-V0-007 infrastructure retry episode at one acceptance revision, shared across roles and request IDs. Sessions counts ended infrastructure sessions, each once; Charged counts reserved retry ordinals; Limit is the maxRetries the episode started with, which a configuration reload can only narrow. Nothing here resets automatically.

func (*InfraEpisode) Holds

func (e *InfraEpisode) Holds(now time.Time) bool

Holds reports whether the episode keeps its ticket from launching at now.

type InfraRetryConfig

type InfraRetryConfig struct {
	MaxRetries         *int `json:"maxRetries,omitempty"`
	CooldownSeconds    *int `json:"cooldownSeconds,omitempty"`
	MaxCooldownSeconds *int `json:"maxCooldownSeconds,omitempty"`
}

InfraRetryConfig bounds automatic retries of a ticket's infrastructure sessions per acceptance revision. An absent member takes its default.

func (*InfraRetryConfig) Limits

func (r *InfraRetryConfig) Limits() (maxRetries, cooldown, maxCooldown int)

Limits returns maxRetries, cooldownSeconds and maxCooldownSeconds with the enabled defaults 3, 30 and 300.

type Lane

type Lane struct {
	Pool   string   `json:"pool"`
	States []string `json:"states,omitempty"`
}

Lane selects pool members in one native pool state (QUARANTINED by default).

type LaunchControl

type LaunchControl interface {
	// Admit takes the fence and admits one new effect, or refuses it. The
	// intent names the effect: a worker id, or a pool sweep request id. On
	// admission the dispatcher calls release once; recorded reports that
	// the effect's outcome is durable: it is in the saved ledger, or
	// nothing started.
	Admit(intent string) (release func(recorded bool), err error)
	// Boundary runs between ticks, never during one. recorded reports that
	// the tick's final ledger save succeeded, so every launch this
	// dispatcher made is in the saved ledger; settled that the tick also
	// completed with no recorded worker and no pending or in-flight pool
	// sweep. A true result ends Run with ErrSettled.
	Boundary(recorded, settled bool) bool
}

LaunchControl links a controlled dispatcher to its service control (SERVICE500-003).

type Ledger

type Ledger struct {
	PoolSweeps map[string]*PoolSweepRecord `json:"poolSweeps,omitempty"`
	Profile    string                      `json:"profile"`
	Program    string                      `json:"program"`
	LaunchSeq  uint64                      `json:"launchSeq"`
	EventSeq   uint64                      `json:"eventSeq"`
	Workers    []*Worker                   `json:"workers"`
	Backoff    map[string]*BackoffState    `json:"backoff"`
	Seen       *Seen                       `json:"seen,omitempty"`
	Progress   map[string]*ProgressHistory `json:"progress,omitempty"`
	Pressure   *PressureRecord             `json:"pressure,omitempty"`
	// Escalation is present only while the configuration has a ladder.
	Escalation map[string]*EscalationState `json:"escalation,omitempty"`
	// InfraRetry is present only once an ESC-V0-007 episode exists.
	InfraRetry map[string]*InfraEpisode `json:"infraRetry,omitempty"`
}

Ledger is the dispatcher's private taskman-dispatch-state/0 file. It is never an input to the native store.

func LoadLedger

func LoadLedger(dir, program string) (*Ledger, error)

LoadLedger reads the ledger; a missing ledger is a fresh one.

type LoopHold

type LoopHold struct {
	Signal             string   `json:"signal"`
	AcceptanceRevision string   `json:"acceptanceRevision"`
	Generations        []string `json:"generations"`
	// Pending marks, in the ledger only, an episode whose needs-owner event
	// is not yet appended to the event log, so the next tick and a restart
	// retry it (CAL-V0-103).
	Pending bool `json:"pending,omitempty"`
}

LoopHold is a CAL-V0-102 hold as the dispatcher records it: the signal, the acceptance revision it is bound to and the counted generations, oldest first. The ticket, signal, acceptance revision and newest generation name the episode.

type Match

type Match struct {
	Labels        []string `json:"labels,omitempty"`
	Kinds         []string `json:"kinds,omitempty"`
	IDGlob        string   `json:"idGlob,omitempty"`
	States        []string `json:"states,omitempty"`
	ExcludeStates []string `json:"excludeStates,omitempty"`
	Statuses      []string `json:"statuses,omitempty"`
	PlanSelected  bool     `json:"planSelected,omitempty"`
	// Pool is the pool the role's workers claim with `--pool {pool}`. A
	// ticket that requires a pool matches only a role naming that pool, and
	// a role naming a pool matches only its tickets (CAL-V0-029, CAL-V0-097).
	Pool string `json:"pool,omitempty"`
	// Gates is a conjunction of ERG-V0-009 native external review predicates.
	Gates []GateMatch `json:"gates,omitempty"`
}

Match is a conjunction; an empty list matches anything. Statuses defaults to OPEN. States and ExcludeStates name work states (NONE when the reader found none); UNKNOWN never matches a rule that names states. IDGlob matches the local ticket ID. A ticket with any live attempt is never a candidate; an expired lease becomes free only after heal reaps it.

type Member

type Member struct {
	Pool, Member, State, Holder, Attempt string
	Queue, Allocation, Definition        string
	SafeReuse, Owned                     bool
}

Member is one native pool member.

type NoteView

type NoteView struct {
	State, Revision, Head string
	Text                  string
	RecordedAt            string
	ActorID, ActorRole    string
	Code                  string
}

NoteView is a ticket's operator note as the native observation read it (ON-V0-011): CURRENT with its text, CLEARED, or UNAVAILABLE with the code of the resolution failure. A never-noted ticket has no NoteView.

type Observation

type Observation struct {
	Tickets  []Ticket
	Attempts []Attempt
	Members  []Member
}

Observation is one authoritative read of the native store.

type OpenRequest

type OpenRequest struct {
	RequestID  string `json:"requestId"`
	Kind       string `json:"kind"`
	RecordedAt string `json:"recordedAt"`
}

OpenRequest is one current OPEN escalation request: its kind and the RecordedAt of its audited OPEN event.

type PoolSweepConfig

type PoolSweepConfig struct {
	TimeoutSeconds  int `json:"timeoutSeconds"`
	IntervalSeconds int `json:"intervalSeconds"`
}

PoolSweepConfig opts into one owned native operation at a time.

type PoolSweepQueue

type PoolSweepQueue interface {
	PoolSweepActor() (id, role string, err error)
	PoolSweep(context.Context, PoolSweepRequest) (PoolSweepResult, error)
}

PoolSweepQueue is an explicit opt-in boundary; legacy Queue implementations never gain command authority from a configuration field.

type PoolSweepRecord

type PoolSweepRecord struct {
	PoolSweepRequest
	Phase    string          `json:"phase"`
	Started  time.Time       `json:"started"`
	Observed time.Time       `json:"observed"`
	Result   PoolSweepResult `json:"result"`
	Reason   string          `json:"reason,omitempty"`
}

type PoolSweepRequest

type PoolSweepRequest struct {
	WorkRoot       string `json:"workRoot"`
	Program        string `json:"program"`
	Queue          string `json:"queue"`
	Pool           string `json:"pool"`
	Member         string `json:"member"`
	Allocation     string `json:"allocation"`
	Definition     string `json:"definition"`
	RequestID      string `json:"requestId"`
	Actor          string `json:"actor"`
	ActorRole      string `json:"actorRole"`
	ConfigDigest   string `json:"configDigest"`
	TimeoutSeconds int    `json:"timeoutSeconds"`
}

PoolSweepRequest retains the immutable admission facts across restarts. The definition is a first-admission observation; allocation is a native selector.

type PoolSweepResult

type PoolSweepResult struct {
	Pending    bool   `json:"pending"`
	Evidence   string `json:"evidence,omitempty"`
	Receipt    string `json:"receipt,omitempty"`
	ReceiptSeq string `json:"receiptSeq,omitempty"`
	Outcome    string `json:"outcome,omitempty"`
}

type PressureBudget

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

PressureBudget counts only non-exempt concurrency. Roster integration must call Accept before consuming static slots or reserving a candidate's key.

func NewPressureBudget

func NewPressureBudget(c *PressureConfig, state PressureState, running []Assignment, pinned []string) (*PressureBudget, error)

NewPressureBudget observes running assignments without changing them. A nil configuration disables pressure and preserves ordinary admission behavior.

func (*PressureBudget) Accept

func (b *PressureBudget) Accept(a Assignment) bool

Accept reserves pressure budget only; callers retain every static/native admission fence and must not reserve a rejected candidate's key or role slot.

func (*PressureBudget) Exempt

func (b *PressureBudget) Exempt(a Assignment) bool

Exempt is based only on explicit role/ticket identities and admitted pins.

type PressureConfig

type PressureConfig struct {
	LoadPerCPUHigh     float64        `json:"loadPerCpuHigh"`
	LoadPerCPUCritical float64        `json:"loadPerCpuCritical"`
	SwapHigh           float64        `json:"swapHigh"`
	SwapCritical       float64        `json:"swapCritical"`
	CalmLoadPerCPU     float64        `json:"calmLoadPerCpu"`
	CalmSwap           float64        `json:"calmSwap"`
	TicksToChange      int            `json:"ticksToChange"`
	LevelCaps          map[string]int `json:"levelCaps"`
	ExemptRoles        []string       `json:"exemptRoles,omitempty"`
	ExemptTickets      []string       `json:"exemptTickets,omitempty"`
}

PressureConfig is optional launch admission policy. It does not authorize stopping workers or bypassing the ordinary roster and native lease checks.

func (PressureConfig) Validate

func (c PressureConfig) Validate() error

Validate checks local bounds. Config validation additionally requires each exempt role to name a configured role; exempt tickets are matched against a candidate's native ID or local name and are not resolved in advance.

type PressureRecord

type PressureRecord struct {
	State  PressureState  `json:"state"`
	Sample PressureSample `json:"sample"`
	Held   []HeldLaunch   `json:"held"`
}

PressureRecord is the CAL-V0-068 derived pressure state: hysteresis, the newest bounded sample and the work keys the last roster held. It exists only while pressure is configured and is never a native-store input.

type PressureSample

type PressureSample struct {
	SampledAt      time.Time `json:"sampledAt"`
	Source         string    `json:"source"`
	LoadAverage    float64   `json:"loadAverage"`
	CPUs           int       `json:"cpus"`
	SwapTotalBytes uint64    `json:"swapTotalBytes"`
	SwapUsedBytes  uint64    `json:"swapUsedBytes"`
	LoadKnown      bool      `json:"loadKnown"`
	CPUKnown       bool      `json:"cpuKnown"`
	SwapKnown      bool      `json:"swapKnown"`
	Problems       []string  `json:"problems,omitempty"`
}

PressureSample reports independent known/unknown observations. Swap bytes represent allocated swap utilization, not compression or OS memory pressure.

func SamplePressure

func SamplePressure(ctx context.Context) PressureSample

SamplePressure reads fixed local OS metrics; unsupported inputs stay UNKNOWN.

func (PressureSample) LoadPerCPU

func (s PressureSample) LoadPerCPU() (float64, bool)

LoadPerCPU normalizes the first load average by positive host-visible CPUs.

func (PressureSample) SwapFraction

func (s PressureSample) SwapFraction() (float64, bool)

SwapFraction treats explicitly observed zero total/used as no allocated swap (fraction zero), avoiding a fabricated value for unavailable metrics.

type PressureState

type PressureState struct {
	Level        int  `json:"level"`
	PendingLevel int  `json:"pendingLevel"`
	PendingTicks int  `json:"pendingTicks"`
	Unknown      bool `json:"unknown"`
}

PressureState retains hysteresis across ticks. UNKNOWN cancels pending dwell, while retaining the current level; it never manufactures calm evidence.

func StepPressure

func StepPressure(c PressureConfig, state PressureState, sample PressureSample) (PressureState, error)

StepPressure is pure: thresholds use one-minute load per positive host CPU, either metric can raise pressure, and both must be calm before release.

type Proc

type Proc struct {
	PID      int    `json:"pid"`
	Identity string `json:"identity"`
}

Proc is one identity-verified process.

type ProgressHistory

type ProgressHistory struct {
	Current string   `json:"current"`
	Seen    []string `json:"seen"`
}

ProgressHistory is lifetime replay protection, including deleted keys.

type Queue

type Queue interface {
	Observe(ctx context.Context) (*Observation, error)
	Release(ctx context.Context, a Attempt, evidence, requestID string) error
	Reap(ctx context.Context, a Attempt, requestID string) error
}

Queue is the native store boundary. Release and Reap go through the existing fenced lease transactions; the dispatcher never writes the store any other way.

type Role

type Role struct {
	Name        string `json:"name"`
	Host        string `json:"host"`
	Cap         int    `json:"cap"`
	Priority    int    `json:"priority"`
	Match       *Match `json:"match,omitempty"`
	Lane        *Lane  `json:"lane,omitempty"`
	Prompt      string `json:"prompt"`
	IdleSeconds int    `json:"idleSeconds"`
	WallSeconds int    `json:"wallSeconds"`
	// Model, Escalate and DeescalateOnProgress are the CAL-V0-057
	// escalation ladder. The role's host must render {model}.
	Model                string `json:"model,omitempty"`
	Escalate             []Tier `json:"escalate,omitempty"`
	DeescalateOnProgress *bool  `json:"deescalateOnProgress,omitempty"`
}

Role is one roster rule. Exactly one of Match (ticket work) and Lane (pool member work) is present.

func (*Role) Deescalates

func (r *Role) Deescalates() bool

Deescalates reports whether progress returns the role to its base model (the default).

func (*Role) ModelAt

func (r *Role) ModelAt(tier int) string

ModelAt is the model the role launches at a tier.

func (*Role) TierFor

func (r *Role) TierFor(streak int) int

TierFor is the ladder tier for a no-progress streak: 0 is the base model.

type Seen

type Seen struct {
	Tickets map[string]string `json:"tickets"`
	Claims  map[string]string `json:"claims"`
	Lanes   map[string]string `json:"lanes"`
	// Escalations keeps each ESC-V0-006 held ticket's request IDs, apart from
	// its plan reason, so status shows a hold behind another blocker.
	Escalations map[string][]string `json:"escalations,omitempty"`
	// Loops keeps each CAL-V0-102 held ticket's loop hold, so status shows
	// it and diff raises one blocked escalation event per episode.
	Loops map[string]LoopHold `json:"loops,omitempty"`
	// Requests keeps each ticket's current OPEN escalation requests, and
	// RequestsUnknown the sorted tickets whose escalation material could not
	// be validated, so status shows kinds and ages without reading the
	// native store (ESC-V0-009).
	Requests        map[string][]OpenRequest `json:"requests,omitempty"`
	RequestsUnknown []string                 `json:"requestsUnknown,omitempty"`
}

Seen is the previous observation, kept to emit change events.

type Ticket

type Ticket struct {
	ID, Local, Status, Priority, Kind, Revision string
	Order                                       uint64
	Labels                                      []string
	Plan, PlanReason                            string
	State                                       string
	// RequiresPool is the pool a claim of the ticket must name, or empty.
	RequiresPool                  string
	ProgressToken, ProgressDigest string
	// Gates is the native ERG-V0-009 external review state by gate ID. A
	// workState program cannot supply it; absent means no record.
	Gates map[string]GateView
	// GatesObserved is set only by a native observation that read every
	// gate of the ticket. While false every gate is UNKNOWN, so no gate
	// predicate (NONE included) matches an unobserved ticket.
	GatesObserved bool
	// EscalationPending is the native ESC-V0-006 derived hold: the sorted
	// request IDs of current OPEN decision, scope or blocked questions. A
	// workState program cannot supply it. A held ticket is never rostered.
	EscalationPending []string
	// Loop is the native CAL-V0-102 LOOP_DETECTED hold, nil when none. A
	// workState program cannot supply it. A held ticket is never rostered.
	Loop *LoopHold
	// AcceptanceRevision is the native acceptance revision that scopes an
	// ESC-V0-007 infrastructure retry episode.
	AcceptanceRevision string
	// Infrastructure lists the holders (dispatcher worker IDs) of current
	// OPEN or ANSWERED infrastructure requests at AcceptanceRevision, so an
	// ended session is classified by its own typed request (ESC-V0-008).
	Infrastructure []string
	// OpenRequests lists the ticket's current OPEN escalation requests of
	// every kind, with the original OPEN time, sorted by request ID, so
	// status can show kinds and ages (ESC-V0-009). A workState program
	// cannot supply it.
	OpenRequests []OpenRequest
	// EscalationUnknown is set when the ticket's escalation material could
	// not be validated: its revision binding and infrastructure observation
	// are UNKNOWN, never progress. A workState program cannot clear it.
	EscalationUnknown bool
	// NextStage is the advisory CAL-V0-084 recorded hand-off target:
	// implement, review, integrate, STALE, or NONE when none is observed.
	// It comes from the native queue, never from a workState program.
	NextStage string
	// OperatorNote is the ticket's ON-V0-011 note as read at observation,
	// nil when the ticket was never noted. Only a native observation sets
	// it; a workState program cannot supply it.
	OperatorNote *NoteView
}

Ticket is the dispatcher's read-only view of one native ticket.

func (Ticket) GateState

func (t Ticket) GateState(gate string) string

GateState is the routing state of one gate: NONE, PASS, RETURN or RESUBMITTED, or the non-CURRENT status (STALE or UNKNOWN) otherwise.

type Tier

type Tier struct {
	After int    `json:"after"`
	Model string `json:"model"`
	Cap   int    `json:"cap,omitempty"`
	Notes string `json:"notes,omitempty"`
}

Tier is one escalation step: from After consecutive no-progress sessions on a ticket, the role launches Model. Cap 0 leaves only the role cap.

type UnparkRequest

type UnparkRequest struct {
	Unpark string `json:"unpark"`
}

UnparkRequest is one queued operator request file.

type WorkState

type WorkState struct {
	Kind string   `json:"kind"`
	Path string   `json:"path,omitempty"`
	Key  string   `json:"key,omitempty"`
	Argv []string `json:"argv,omitempty"`
}

WorkState reads a program-defined per-ticket state. Kind status-line reads the first "key: value" line of a per-ticket file; kind command runs argv once per tick and expects one JSON object of ticket ID to state.

type Worker

type Worker struct {
	ID              string    `json:"id"`
	Role            string    `json:"role"`
	Host            string    `json:"host"`
	Slot            int       `json:"slot"`
	Key             string    `json:"key"`
	Ticket          string    `json:"ticket,omitempty"`
	Pool            string    `json:"pool,omitempty"`
	Member          string    `json:"member,omitempty"`
	PID             int       `json:"pid"`
	LeaderIdentity  string    `json:"leaderIdentity"`
	Members         []Proc    `json:"members"`
	Started         time.Time `json:"started"`
	LastActive      time.Time `json:"lastActive"`
	LogBytes        int64     `json:"logBytes"`
	ActivityPaths   []string  `json:"activityPaths,omitempty"`
	ActivityMtime   time.Time `json:"activityMtime"`
	State           string    `json:"state"` // RUNNING or KILLING
	KillReason      string    `json:"killReason,omitempty"`
	KillDeadline    time.Time `json:"killDeadline,omitempty"`
	Fingerprint     string    `json:"fingerprint"`
	BaseFingerprint string    `json:"baseFingerprint,omitempty"`
	ProgressDigest  string    `json:"progressDigest,omitempty"`
	Tier            int       `json:"tier,omitempty"`
	Model           string    `json:"model,omitempty"`
}

Worker is one launched host process and its supervised tree.

Jump to

Keyboard shortcuts

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