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
- Variables
- func Fingerprint(obs *Observation, key string) string
- func LaunchTier(r *Role, e *EscalationState) int
- func OwnerProcess(dir string) (string, int)
- func OwnerState(dir string) string
- func ProgramDir(c *Config, program string) string
- func ReadStates(ctx context.Context, c *Config, tickets []Ticket) []string
- func ReaderContainment(dir, program string) (state string, quarantined bool, diagnostic string)
- func Records(l *Ledger, intent string) bool
- func Render(s string, values map[string]string) string
- func RenderOperatorNote(ticket string, n *NoteView) string
- func RetryCooldown(cooldown, maxCooldown, n int) time.Duration
- func Settled(l *Ledger) bool
- func Summary(logDir string) string
- func ValidName(s string) bool
- type Assignment
- func Roster(c *Config, obs *Observation, busy []Busy, skip map[string]bool) []Assignment
- func RosterTiers(c *Config, obs *Observation, busy []Busy, skip map[string]bool, ...) []Assignment
- func RosterWithPressure(c *Config, obs *Observation, busy []Busy, skip map[string]bool, ...) (out, held []Assignment)
- type Attempt
- type Backoff
- type BackoffState
- type Busy
- type Config
- type Dispatcher
- type EscalationState
- type Event
- type GateCandidate
- type GateMatch
- type GateSubject
- type GateTrust
- type GateView
- type Heal
- type HeldLaunch
- type Host
- type InfraEpisode
- type InfraRetryConfig
- type Lane
- type LaunchControl
- type Ledger
- type LoopHold
- type Match
- type Member
- type NoteView
- type Observation
- type OpenRequest
- type PoolSweepConfig
- type PoolSweepQueue
- type PoolSweepRecord
- type PoolSweepRequest
- type PoolSweepResult
- type PressureBudget
- type PressureConfig
- type PressureRecord
- type PressureSample
- type PressureState
- type Proc
- type ProgressHistory
- type Queue
- type Role
- type Seen
- type Ticket
- type Tier
- type UnparkRequest
- type WorkState
- type Worker
Constants ¶
const ( ConfigProfile = "taskman-dispatch/0" StateProfile = "taskman-dispatch-state/0" EventProfile = "taskman-dispatch-event/0" MaxConfig = 256 << 10 )
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.
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.
const ( StateNone = "NONE" StateUnknown = "UNKNOWN" )
Work states that are not program-defined.
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.
const MaxLoopGenerations = 1024
MaxLoopGenerations bounds the generations one recorded loop hold names.
Variables ¶
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.
var ErrSettled = errors.New("dispatcher settled by its service control")
ErrSettled ends Run when the controlled dispatcher's control settled it.
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.
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 ¶
OwnerProcess is OwnerState with the recorded holder's PID, which is 0 unless the state is RUNNING.
func OwnerState ¶
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 ¶
ProgramDir is the dispatcher's state directory for one program.
func ReadStates ¶
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 ¶
ReaderContainment is a pure read. Absence is NOT_OBSERVED, never RELEASED.
func Records ¶
Records reports whether a saved ledger records the effect an admission intent named: a worker or a pool sweep request.
func RenderOperatorNote ¶
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 ¶
RetryCooldown is the cooldown before retry ordinal n (1-based): min(maxCooldown, cooldown*2^(n-1)), saturating.
func Settled ¶
Settled reports a saved ledger with no recorded worker and no pending pool sweep: what a draining service needs before it may stop.
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 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 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 ¶
DecodeConfig parses and validates closed configuration bytes.
func (*Config) TicketPools ¶
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 ¶
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.
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 ¶
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 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.
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 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 ¶
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 ¶
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 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 ProgressHistory ¶
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 ¶
Deescalates reports whether progress returns the role to its base model (the default).
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.
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.