dispatch

package
v0.2.1 Latest Latest
Warning

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

Go to latest
Published: Oct 3, 2026 License: Apache-2.0 Imports: 26 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 onto sessions within a concurrency cap derived from the host's measured idle cores (harness.WorkerCap). 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.

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. There is no polling.

Index

Constants

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)
)

Functions

This section is empty.

Types

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) 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) 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) List

List reports every slice in enqueue order with its priority breakdown, the derived cap 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) 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) Reprioritize

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

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 from the database, when one was granted, and starts the owner goroutine on the scope. Slices the previous process was running return to the frontier, and the restored frontier is dispatched at once.

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)
}

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 and to ask one to stop at its next safepoint. *harness.AgentSessionService satisfies it.

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 WithHostMeasures

func WithHostMeasures(measures harness.IHostMeasures) Option

WithHostMeasures replaces the host measures the concurrency cap is derived from.

func WithLogger

func WithLogger(logger *slog.Logger) Option

WithLogger receives the service's records.

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.

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 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"
)

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