dispatch

package
v0.4.0 Latest Latest
Warning

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

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

Documentation

Overview

Package dispatch is the intent-ranked work queue of one harness process: a slice graph, its frontier, and the list scheduling that dispatches the frontier onto agent sessions.

Contract

Work is a slice graph, not a heap. A slice is one unit of work with the recipe its session runs. Its depends_on edges form a DAG and are the only serial relation: a slice enters the frontier once every slice it depends on has merged. Its contends edges are undirected mutual exclusion, declared or computed from overlapping touch-sets (path globs and named hotspots): two contending slices never run at once. No edge means parallel.

Dispatch is list scheduling (pkg/listsched) of the frontier, less the held slices, onto sessions within an admission the slice dispatcher derives. Rank is lexicographic: the upward critical-path length of the slice in the depends_on graph (HEFT upward rank, unit weights until durations have been measured), then the strongest urgency of the intents attached to it, then enqueue order. No weight is hand-set; the learned term is zero until data exists. A blocked higher-ranked slice does not hold back a lower-ranked one that may run (bucketed, not strictly serial, order).

An intent is a typed record: statement, scope, urgency, optional deadline and the intent it supersedes, with term slots. Route sends an operator message to the slice whose touch-set covers the intent's terms: into its session when it runs (steer), onto the slice when it does not (queue), as a new slice when none covers the terms, or back as already done when the intent compiles to no change against main. An URGENT intent preempts: when its slice cannot start because the cap is full or a contending slice runs, the lowest-ranked running session that ranks below it is asked to cancel at its next safepoint; once that session reports CANCELED the slice it ran is checkpointed and re-enqueued.

The slice dispatcher

The dispatcher is a pass, DispatchService.Dispatch, that the binary runs as a cron trigger every DispatcherInterval and after every merge the merge train makes, and once on start. A pass marks merged every slice whose recorded pull request has merged (WithMergeCheck), measures each limit (Limit) and launches the frontier — less the held slices and the slices whose provider's rate limit admits no launches — up to the admission: the dispatched sessions running plus the fewest launches any unscoped limit allows. The limits are the harness's launch check (the worker cap of idle cores and the disk floor; a held admission admits none), and those the binary adds: the mining loop's daily budget (DailyBudgetLimit) and the rate-limit headroom of the provider its newest rate_limit_event came from, read from the session logs (RateLimit). A rate limit is scoped to that provider, so one provider's event never gates another's launches: it drops its own provider's slices from the frontier rather than lowering the admission. Each limit records its derivation, and a limit that cannot be measured admits no launch. Between passes, enqueues, merges and session ends launch within the latest admission at once.

Pause and resume stop and restart every launch; hold and release keep one slice from being launched. Each control carries its reason, is written to csf_dispatch_controls and replayed on start. AddSlice, the four controls and the snapshot are typed operations, served as MCP tools (DispatchService.Tools) and HTTP routes (DispatchService.Register); every pass and control publishes the Snapshot to the binary's sink, which writes SnapshotFile.

Ownership

The graph is an in-process store owned by one goroutine, which pops commands and session events from one io/inproc.Queue. Every change is written through to csfpg when a database was granted, and the whole graph is read back on Start, so the queue survives a restart: slices that were running in the previous process return to the frontier. The service owns its goroutines through its runtime scope; it opens no listener and no pool.

Sessions are reached through ISessionHost, which services/harness's AgentSessionService satisfies, and session completion arrives through ObserveSession, which the host wires as a harness session observer. The pass reads what has no event: the limits and the pull requests' state.

Index

Constants

View Source
const (
	AddSliceTool         = "AddSlice"
	HoldSliceTool        = "HoldSlice"
	ReleaseSliceTool     = "ReleaseSlice"
	PauseDispatcherTool  = "PauseDispatcher"
	ResumeDispatcherTool = "ResumeDispatcher"
	DispatcherStateTool  = "DispatcherSnapshot"

	SnapshotPath = "/api/dispatch"
	AddPath      = "/api/dispatch/slices"
	HoldPath     = "/api/dispatch/slices/hold"
	ReleasePath  = "/api/dispatch/slices/release"
	PausePath    = "/api/dispatch/pause"
	ResumePath   = "/api/dispatch/resume"

	// ProvenanceSliceAdd is the provenance source of a slice added through
	// AddSlice.
	ProvenanceSliceAdd = "slice add"
)

The slice dispatcher's typed operations: each is one MCP tool on the CSF service and one HTTP route on the host's router, and the csf slice verb is a client of the routes.

View Source
const (
	// TriggerDispatcher is the slice dispatcher's cron trigger: one pass
	// every DispatcherInterval.
	TriggerDispatcher = "dispatch.dispatcher"
	// DispatcherInterval is the cadence of the dispatcher's passes. A pass
	// costs one harness check, one ledger read, the newest session logs and
	// one pull request read per recorded pull request; a minute keeps a
	// freed machine idle for at most that long between merges the pass
	// detects itself.
	DispatcherInterval = time.Minute
	// SnapshotFile is where the binary writes every published snapshot,
	// under the harness state directory, for the Workbench and the ops view.
	SnapshotFile = "dispatch.json"
)
View Source
const ScopeOntology = "ontology"

ScopeOntology is the intent scope whose term slots name ontology terms: an intent in this scope asks that its terms be defined, so it is already done when every term exists in the ontology source.

Variables

View Source
var (
	// ErrNoSessions reports a service built without a session host.
	ErrNoSessions = errors.New("dispatch: a session host is required")
	// ErrInvalidOption reports a nil option or a value the service cannot use.
	ErrInvalidOption = errors.New("dispatch: invalid option")
	// ErrNotStarted reports an operation before the service was mounted and
	// started, or after it stopped.
	ErrNotStarted = errors.New("dispatch: the service is not running")
	// ErrInvalidSlice reports a slice the service cannot enqueue.
	ErrInvalidSlice = fmt.Errorf("%w: invalid slice", csf.ErrInvalidRequest)
	// ErrInvalidIntent reports an intent the service cannot use.
	ErrInvalidIntent = fmt.Errorf("%w: invalid intent", csf.ErrInvalidRequest)
	// ErrSliceExists reports a second Enqueue of a slice identifier.
	ErrSliceExists = fmt.Errorf("%w: the slice is already enqueued", csf.ErrConflict)
	// ErrUnknownSlice reports an edge end or a slice this service holds no
	// node for.
	ErrUnknownSlice = fmt.Errorf("%w: no such slice", csf.ErrNotFound)
	// ErrUnknownIntent reports a Reprioritize of an intent never declared.
	ErrUnknownIntent = fmt.Errorf("%w: no such intent", csf.ErrNotFound)
	// ErrCycle reports depends_on edges that would close a cycle.
	ErrCycle = fmt.Errorf("%w: depends_on would close a cycle", csf.ErrInvalidRequest)
	// ErrSliceFinished reports a change to a merged or canceled slice.
	ErrSliceFinished = fmt.Errorf("%w: the slice has finished", csf.ErrConflict)
	// ErrInvalidControl reports a control without a reason, or a hold or
	// release naming no slice.
	ErrInvalidControl = fmt.Errorf("%w: invalid control", csf.ErrInvalidRequest)
)

Stages is every stage in the order a slice moves through them.

Functions

This section is empty.

Types

type AddSliceInput added in v0.4.0

type AddSliceInput struct {
	SliceID   string                    `json:"slice_id" jsonschema:"the slice identifier: letters, digits, dot, dash or underscore, at most 80"`
	Title     string                    `json:"title,omitempty" jsonschema:"defaults to the recipe's pull request title"`
	TicketURL string                    `json:"ticket_url" jsonschema:"the ticket the slice delivers"`
	Recipe    *pb.AgentAssignmentRecipe `json:"recipe" jsonschema:"the agent.json recipe with its brief as the task"`
	DependsOn []string                  `json:"depends_on,omitempty" jsonschema:"slices that must merge before this one starts"`
	Contends  []string                  `json:"contends,omitempty" jsonschema:"slices this one never runs beside"`
	Hotspots  []string                  `json:"hotspots,omitempty" jsonschema:"named hotspots it touches; slices sharing one contend"`
	Paths     []string                  `json:"paths,omitempty" jsonschema:"path globs it touches; slices overlapping one contend"`
	// Source is the provenance the queue shows for the slice: who added it
	// and why it ranks where it does. Empty is ProvenanceSliceAdd.
	Source string `json:"source,omitempty" jsonschema:"who added the slice and why, shown on the queue; defaults to slice add"`
	// Hold, when set, holds the slice in the same step that enqueues it, with
	// this reason, so the dispatcher cannot launch it before it is released.
	Hold string `` /* 126-byte string literal not displayed */
}

AddSliceInput is one slice to add to the graph: the ticket it delivers, the recipe its session runs, its edges and its touch-set.

type Control added in v0.4.0

type Control struct {
	Action  ControlAction
	SliceID string
	Reason  string
}

Control is one recorded control: the action, the slice a hold or release names, and the reason it was given.

type ControlAction added in v0.4.0

type ControlAction string

ControlAction is one control of the slice dispatcher; its spelling is the action column of csf_dispatch_controls.

const (
	ControlPause   ControlAction = "pause"
	ControlResume  ControlAction = "resume"
	ControlHold    ControlAction = "hold"
	ControlRelease ControlAction = "release"
)

The four controls.

type DailyBudgetReader added in v0.4.0

type DailyBudgetReader func(ctx context.Context) (ouroboros.DailyBudget, error)

DailyBudgetReader reads the mining loop's daily budget: (*ouroboros.Loop).DailyBudget.

type DispatchService

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

DispatchService is the slice graph and its dispatch onto harness sessions. One goroutine owns the graph; every operation and every session event is a command it pops from one queue.

func NewDispatchService

func NewDispatchService(options ...Option) (*DispatchService, error)

NewDispatchService validates the whole option set before building the service. WithSessions is required.

func (*DispatchService) AddSlice added in v0.4.0

func (service *DispatchService) AddSlice(ctx context.Context, input AddSliceInput) (SliceView, error)

AddSlice enqueues one slice in the graph, holds it in the same step when the input asks, and publishes the snapshot; the dispatcher launches it once it is ready, released and admitted.

func (*DispatchService) CurrentSnapshot added in v0.4.0

func (service *DispatchService) CurrentSnapshot(ctx context.Context, _ SnapshotInput) (Snapshot, error)

CurrentSnapshot reports the dispatcher's state now, under the latest pass's limits.

func (*DispatchService) DeclareIntent

DeclareIntent records an intent and attaches it to the slices its scope or terms route it to. An URGENT intent preempts a lower-ranked running session standing between its slice and the cap.

func (*DispatchService) Dispatch added in v0.4.0

func (service *DispatchService) Dispatch(ctx context.Context) (Snapshot, error)

Dispatch is one pass of the slice dispatcher. It marks merged every slice whose recorded pull request merged, measures every limit, and launches the ready frontier in rank order up to the admission they derive: the minimum of the limits' launches over the dispatched sessions running. Between passes, merges and session ends launch within that admission at once. It returns the snapshot it published.

func (*DispatchService) Enqueue

Enqueue admits one slice with its edges, touch-set and provenance. Every slice its edges name must be enqueued already; a depends_on relation that would close a cycle is rejected.

func (*DispatchService) Frontier

Frontier reports the slices whose predecessors have merged, in dispatch order with their priority breakdown.

func (*DispatchService) HeldBranch added in v0.4.0

func (service *DispatchService) HeldBranch(ctx context.Context, branch string) (string, bool, error)

HeldBranch reports whether the slice whose branch is branch is held, and the reason it is: a merge of a pull request from that branch is refused while the hold stands. It is the guard the GitHub tools ask before merging, so a held slice's branch cannot reach main.

func (*DispatchService) Hold added in v0.4.0

func (service *DispatchService) Hold(ctx context.Context, input SliceControlInput) (Snapshot, error)

Hold keeps one slice from being launched until Release; a held slice that runs keeps running.

func (*DispatchService) List

List reports every slice in enqueue order with its priority breakdown, the admission the latest dispatcher pass derived and how many sessions run under it.

func (*DispatchService) MarkMerged

MarkMerged records that a slice's pull request merged: its session, if one still runs, is asked to stop, and every slice that depended on it may enter the frontier.

func (*DispatchService) MarkPullRequestMerged added in v0.4.0

func (service *DispatchService) MarkPullRequestMerged(ctx context.Context, url string) ([]string, error)

MarkPullRequestMerged marks merged every unfinished slice that recorded the pull request at url: a merge GitHub reported, so the frontier moves without waiting for the next pass to ask. It returns the slices it marked, none when no slice recorded that pull request.

func (*DispatchService) ObserveSession

func (service *DispatchService) ObserveSession(state *harnessv1.AgentSessionState)

ObserveSession is the harness session observer: the host wires it with harness.WithSessionObserver. Session events reach the graph in order.

func (*DispatchService) Pause added in v0.4.0

func (service *DispatchService) Pause(ctx context.Context, input DispatcherControlInput) (Snapshot, error)

Pause stops every launch until Resume; running sessions are not touched.

func (*DispatchService) ReadySlices added in v0.4.0

func (service *DispatchService) ReadySlices(ctx context.Context, limit int) ([]ReadySlice, error)

ReadySlices is at most limit slices that may run away from this host now, in dispatch order: on the frontier, not held, contending with nothing that runs here and with none of the others picked. The pause and the admission govern this host's launches, not another machine's, so they do not apply.

func (*DispatchService) Register added in v0.4.0

func (service *DispatchService) Register(router gin.IRouter)

Register mounts every operation's HTTP route on the caller's router.

func (*DispatchService) Release added in v0.4.0

func (service *DispatchService) Release(ctx context.Context, input SliceControlInput) (Snapshot, error)

Release lets a held slice be launched again.

func (*DispatchService) Reprioritize

Reprioritize re-declares an existing intent: a changed urgency, deadline or supersedes takes effect on the slices it attaches to.

func (*DispatchService) Resume added in v0.4.0

func (service *DispatchService) Resume(ctx context.Context, input DispatcherControlInput) (Snapshot, error)

Resume lets the dispatcher launch again, at once within the latest admission.

func (*DispatchService) Route

Route sends an operator message to the slice whose touch-set covers the intent's terms: into its running session (steer), onto the queued slice (queue), as a new slice when none covers the terms (new), or back as already done.

func (*DispatchService) Start

func (service *DispatchService) Start(scope *runtime.Scope) error

Start restores the graph and the controls from the database, when one was granted, starts the owner goroutine on the scope and runs the first dispatcher pass. Slices the previous process was running return to the frontier and are dispatched by that pass.

func (*DispatchService) Tools added in v0.4.0

func (service *DispatchService) Tools() []csf.Option

Tools is every operation as an MCP tool for the CSF service: csf.New(append(options, dispatcher.Tools()...)...).

func (*DispatchService) Trigger added in v0.4.0

func (service *DispatchService) Trigger() cronservice.Option

Trigger declares the dispatcher's pass for the cron service.

type DispatcherControlInput added in v0.4.0

type DispatcherControlInput struct {
	Reason string `json:"reason" jsonschema:"why, recorded with the control"`
}

DispatcherControlInput is a pause or a resume of the slice dispatcher.

type Flow added in v0.4.0

type Flow struct {
	Slices   map[Stage]int
	Oldest   map[Stage]time.Duration
	LeadTime *time.Duration
}

Flow is the snapshot's slices as a flow: how many stand in each stage, how long the oldest of them has stood there, and the mean lead time from enqueue to merge of the slices merged within the window before the snapshot. LeadTime is nil when none merged in the window.

type IDispatchDatabase

type IDispatchDatabase interface {
	IDispatchQueries
	Transact(ctx context.Context, work func(queries IDispatchQueries) error) error
}

IDispatchDatabase is the service's database: the queries outside a transaction, and Transact for work that must commit or roll back as one unit.

type IDispatchQueries

type IDispatchQueries interface {
	InsertSlice(ctx context.Context, arg csfpg.InsertSliceParams) (csfpg.CsfSlice, error)
	UpdateSliceState(ctx context.Context, arg csfpg.UpdateSliceStateParams) (csfpg.CsfSlice, error)
	ListSlices(ctx context.Context) ([]csfpg.CsfSlice, error)
	InsertSliceEdge(ctx context.Context, arg csfpg.InsertSliceEdgeParams) error
	ListSliceEdges(ctx context.Context) ([]csfpg.CsfSliceEdge, error)
	InsertIntent(ctx context.Context, arg csfpg.InsertIntentParams) (csfpg.CsfIntent, error)
	ListIntents(ctx context.Context) ([]csfpg.CsfIntent, error)
	InsertDispatchControl(ctx context.Context, arg csfpg.InsertDispatchControlParams) (csfpg.CsfDispatchControl, error)
	ListDispatchControls(ctx context.Context) ([]csfpg.CsfDispatchControl, error)
}

IDispatchQueries is the part of csfpg's generated queries the service runs. *csfpg.Queries satisfies it; the SQL is owned by ipc/db/csfpg/dispatch.sql.

type ISessionHost

ISessionHost is what the service needs of the agent harness: to open a session for a slice, to steer a running one, to ask one to stop at its next safepoint, and the launch check the dispatcher's admission starts from. *harness.AgentSessionService satisfies it.

type Limit added in v0.4.0

type Limit struct {
	Name LimitName `json:"name"`
	// Provider is the provider this limit applies to, empty when it bounds
	// every launch. A rate limit is scoped to the provider its newest
	// rate_limit_event came from, so it bounds only that provider's launches
	// and never another's.
	Provider string `json:"provider,omitempty"`
	// Bounded is false when the measure sets no bound now.
	Bounded bool `json:"bounded"`
	// Launches is how many more dispatched sessions the measure admits now;
	// it is read only when Bounded.
	Launches int `json:"launches"`
	// Remaining is what is left of the measured resource, in Unit.
	Remaining float64 `json:"remaining"`
	Unit      string  `json:"unit"`
	// Derivation is how the bound follows from the data.
	Derivation string `json:"derivation"`
}

Limit is one bound on the slice dispatcher's launches, measured from data at one pass, with the derivation that produced it in words.

type LimitName added in v0.4.0

type LimitName string

LimitName says which measure a Limit was derived from.

const (
	// LimitHarness is the harness's launch check: the idle cores (cores
	// minus the one-minute load) and the disk floor.
	LimitHarness LimitName = "harness"
	// LimitBudget is the mining loop's daily budget.
	LimitBudget LimitName = "budget"
	// LimitRate is the provider's rate-limit headroom, read from the latest
	// rate_limit_event any session logged.
	LimitRate LimitName = "rate"
)

The measures the slice dispatcher's admission is the minimum of.

type LimitSource added in v0.4.0

type LimitSource func(ctx context.Context, running int) Limit

LimitSource measures one limit at a pass; running is how many dispatched sessions hold a session now. A source that cannot measure returns a bounded limit of no launches saying why: the dispatcher fails closed.

func DailyBudgetLimit added in v0.4.0

func DailyBudgetLimit(read DailyBudgetReader) LimitSource

DailyBudgetLimit bounds launches by the mining loop's daily budget: what is left of the day's cap after the spend and the running fixers' reservation buys that many sessions at what one is expected to cost — the ledger's mean once it has measured enough finished sessions, the measured baseline before. The ledger reserves only for its fixers, so each running dispatched session reserves one more of those, and the launches are what remains.

func RateLimit added in v0.4.0

func RateLimit(state stdfs.FS, clock harness.IClock, activeProvider string) LimitSource

RateLimit bounds launches by the active provider's rate-limit headroom, read from the latest rate_limit_event the runs launched through that provider logged under state (the harness state directory, one run directory per session). A run is attributed to the provider its event log's harness_provider_applied record names, and a run with none ran on the executor's own provider; only the active provider's records count, so a window spent on a provider the host has switched away from does not hold launches. The headroom is one less the highest utilization of a window that has not reset. Launches stop while that utilization is at or above rateWarnUtilization, the lowest threshold the provider has warned at, and while the record is rejected and its window has not reset. Otherwise, or with no record at all, the limit sets no bound.

activeProvider is the active provider's base URL, or empty for the executor's own provider, which the host records no router entry for.

type MergeCheck added in v0.4.0

type MergeCheck func(ctx context.Context, url string) (bool, error)

MergeCheck reports whether the pull request at url has merged: the dispatcher asks it about every slice that recorded one.

type Option

type Option func(service *DispatchService) error

Option configures a DispatchService.

func WithClock

func WithClock(clock harness.IClock) Option

WithClock replaces the clock.

func WithDatabase

func WithDatabase(database IDispatchDatabase) Option

WithDatabase grants the csfpg-backed persistence every change is written through to and the graph is restored from. Without it the graph lives only in this process.

func WithLimit added in v0.4.0

func WithLimit(source LimitSource) Option

WithLimit adds one measure the dispatcher's admission is the minimum of, beside the harness's launch check, which is always one: DailyBudgetLimit and RateLimit are the binary's.

func WithLogger

func WithLogger(logger *slog.Logger) Option

WithLogger receives the service's records.

func WithMergeCheck added in v0.4.0

func WithMergeCheck(check MergeCheck) Option

WithMergeCheck grants the check a dispatcher pass asks about every recorded pull request, so a slice whose pull request merged anywhere — by the merge train, its own session or the operator — releases what depends on it. Without it only MarkSliceMerged records a merge.

func WithOntologySource

func WithOntologySource(path string) Option

WithOntologySource names the architecture.csf whose terms Route checks an ontology-scoped intent against.

func WithSessions

func WithSessions(sessions ISessionHost) Option

WithSessions grants the agent harness the frontier is dispatched onto. Required.

func WithSnapshotSink added in v0.4.0

func WithSnapshotSink(sink SnapshotSink) Option

WithSnapshotSink receives every published snapshot.

type PostgresDispatchDatabase

type PostgresDispatchDatabase struct {
	*csfpg.Queries
	// contains filtered or unexported fields
}

PostgresDispatchDatabase is the IDispatchDatabase over the PostgreSQL capability. It borrows the capability and never closes it.

func NewPostgresDispatchDatabase

func NewPostgresDispatchDatabase(database csfpg.IDB) (*PostgresDispatchDatabase, error)

NewPostgresDispatchDatabase returns the dispatch database over a pool the binary opened through ipc/db/csfpg.

func (*PostgresDispatchDatabase) Transact

func (database *PostgresDispatchDatabase) Transact(ctx context.Context, work func(queries IDispatchQueries) error) error

Transact runs work in one transaction.

type ReadySlice added in v0.4.0

type ReadySlice struct {
	SliceID   string
	TicketURL string
	Recipe    *pb.AgentAssignmentRecipe
}

ReadySlice is one slice that may run elsewhere now, with the recipe its session runs.

type Relation

type Relation string

Relation is an edge kind of the slice graph; its spelling is the relation column of csf_slice_edges.

const (
	// RelationDependsOn is the one serial relation: directed and acyclic,
	// from the slice that must merge first to the slice that waits.
	RelationDependsOn Relation = "depends_on"
	// RelationContends is mutual exclusion: undirected; the two never run at
	// once.
	RelationContends Relation = "contends"
)

type SliceControlInput added in v0.4.0

type SliceControlInput struct {
	SliceID string `json:"slice_id" jsonschema:"the slice to hold or release"`
	Reason  string `json:"reason" jsonschema:"why, recorded with the control"`
}

SliceControlInput is a hold or a release of one slice.

type SliceView added in v0.4.0

type SliceView struct {
	SliceID        string `json:"slice_id"`
	Title          string `json:"title"`
	TicketURL      string `json:"ticket_url,omitempty"`
	State          string `json:"state"`
	Rank           uint32 `json:"rank,omitempty"`
	CriticalPath   uint32 `json:"critical_path"`
	Urgency        string `json:"urgency"`
	Attempts       uint32 `json:"attempts"`
	AssignmentID   string `json:"assignment_id,omitempty"`
	PullRequestURL string `json:"pull_request_url,omitempty"`
	Held           string `json:"held,omitempty"`
	// Source is the slice's provenance: who added it and why.
	Source string `json:"source,omitempty"`
	// Ready is true for a slice in the frontier: every slice it depends on
	// has merged.
	Ready bool `json:"ready"`
	// Waiting says why a queued slice does not run yet.
	Waiting string `json:"waiting,omitempty"`
	// DependsOn is the slices that must merge before this one; Contends the
	// slices it never runs beside.
	DependsOn []string `json:"depends_on,omitempty"`
	Contends  []string `json:"contends,omitempty"`
	// CreatedAt is when the slice was enqueued and UpdatedAt when its state
	// last moved: time in stage, and for a merged slice its lead time.
	CreatedAt time.Time `json:"created_at,omitzero"`
	UpdatedAt time.Time `json:"updated_at,omitzero"`
}

SliceView is one slice as the snapshot shows it.

type Snapshot added in v0.4.0

type Snapshot struct {
	At          time.Time `json:"at"`
	Paused      bool      `json:"paused"`
	PauseReason string    `json:"pause_reason,omitempty"`
	// Capacity is how many dispatched sessions the limits admit at once.
	Capacity   int         `json:"capacity"`
	Limits     []Limit     `json:"limits"`
	LimitsAt   time.Time   `json:"limits_at"`
	NextPassAt time.Time   `json:"next_pass_at"`
	Queue      []SliceView `json:"queue"`
	Running    []SliceView `json:"running"`
	Held       []SliceView `json:"held"`
	// Finished is every merged, failed or canceled slice, in enqueue order:
	// the rest of the slice graph, so a reader can draw all of it.
	Finished []SliceView `json:"finished"`
	// NextLaunch is the queued slice the dispatcher launches next, absent
	// when none is ready.
	NextLaunch *SliceView `json:"next_launch,omitempty"`
}

Snapshot is what the slice dispatcher shows: whether it is paused, the admission and the limits it was derived from, the queue in dispatch order, what runs, what is held, and the next launch.

func (Snapshot) Flow added in v0.4.0

func (snapshot Snapshot) Flow(window time.Duration) Flow

Flow measures the snapshot at its own instant.

func (Snapshot) Staged added in v0.4.0

func (snapshot Snapshot) Staged() []StagedSlice

Staged is every slice of the snapshot with its stage: the queue in dispatch order, then what runs, then what finished. A held running slice is listed once, as running.

type SnapshotInput added in v0.4.0

type SnapshotInput struct{}

SnapshotInput asks for the dispatcher's snapshot; it has no fields.

type SnapshotSink added in v0.4.0

type SnapshotSink func(snapshot Snapshot) error

SnapshotSink receives every snapshot the dispatcher publishes, on the service's own goroutine: the binary writes it to the state directory.

type Stage added in v0.4.0

type Stage string

Stage is where a slice stands on its way from the queue to main, as the snapshot shows it: the three ways a queued slice waits, running, and the three ways a slice finishes.

const (
	// StageReady is a queued slice on the frontier that nothing running
	// contends with: it launches when the admission allows.
	StageReady Stage = "ready"
	// StageBlocked is a queued slice with a depends_on predecessor not yet
	// merged, or one held back.
	StageBlocked Stage = "blocked"
	// StageContended is a queued slice on the frontier that a running slice
	// contends with.
	StageContended Stage = "contended"
	StageRunning   Stage = "running"
	StageMerged    Stage = "merged"
	StageFailed    Stage = "failed"
	StageCanceled  Stage = "canceled"
)

type StagedSlice added in v0.4.0

type StagedSlice struct {
	Stage Stage
	Slice SliceView
}

StagedSlice is one slice of the snapshot with its stage.

Directories

Path Synopsis
Package mocks is a generated GoMock package.
Package mocks is a generated GoMock package.

Jump to

Keyboard shortcuts

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