harness

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: 41 Imported by: 0

Documentation

Overview

Package harness defines the public behavior boundary between deploy and a compiled-in agent runtime implementation.

Index

Constants

View Source
const (

	// MergeScript is the repository's own merge path, relative to the
	// worktree; a repository without one is merged with gh's squash merge.
	MergeScript = "tools/merge-pr.sh"
	// MergeBazelDirectory is the merge path's Bazel output base under the
	// state directory, shared by every operator merge.
	MergeBazelDirectory = "merge-bazel"
)

The control plane's pull request operations: what they run and where.

View Source
const (
	// EventsPath serves one session's records as server-sent events. The
	// query parameter from skips that many records.
	EventsPath = "/api/harness/sessions/:assignment/events"

	// DefaultTailInterval is how often a tail that reached the end of the
	// file looks for records the gates, in their own processes, appended. It
	// is also the frame a transcript is streamed in: everything one look finds
	// is one patch. Of 6,912 gaps between consecutive executor records in
	// three session logs (2026-10-05), 44% were under 50 ms and 50% under
	// 100 ms, so 50 ms folds 44% of records into their predecessor's frame
	// and doubling it would gain 6 points for twice the delay. It was 250 ms,
	// which held a record a median 107 ms before a page saw it.
	DefaultTailInterval = 50 * time.Millisecond
)

The live event stream of one session: its events.jsonl, every record the session, its turn executor and its gates wrote, followed as it grows.

View Source
const (
	// ResidentFile is the series file at the root of the state directory.
	ResidentFile = "resident.jsonl"
	// ResidentWindow is how many samples the file keeps: a day at the default
	// interval.
	ResidentWindow = 24 * 60
	// DefaultResidentSampleInterval is how often the series is sampled. It is
	// a cadence, not a bound: nothing is decided on it.
	DefaultResidentSampleInterval = time.Minute
)

The resident series: what the harness and its executors hold in memory, sampled from the process table it is granted as /proc and written under the state directory for the ops view, which follows the file like the mutation series. The file is rewritten whole with the last ResidentWindow samples, so it stays small enough to read whole.

View Source
const (
	// ResumeQueueFile is the typed record of the resumes a restart holds,
	// under the state directory: why each is held and its position. It is
	// rewritten at every resume pass and lists none once every open run is
	// back.
	ResumeQueueFile = "resume.json"
	// ResumeRetryInterval is how often held resumes are tried again. The
	// kernel recomputes the load average every 5 s (LOAD_FREQ), so trying
	// more often reads the same load and admits nothing more.
	ResumeRetryInterval = 5 * time.Second
)
View Source
const (
	// BazelDiskCacheVariable, BazelOutputVariable and
	// OCamlToolchainCacheVariable name the three caches in the environment
	// of every program the harness, or a service beside it, runs Bazel in.
	BazelDiskCacheVariable      = "CANDACE_BAZEL_DISK_CACHE"
	BazelOutputVariable         = "CANDACE_BAZEL_CACHE"
	OCamlToolchainCacheVariable = "CANDACE_OCAML_TOOLCHAIN_CACHE"
	// BazelDiskCacheDirectory is the shared cache, under the state directory.
	BazelDiskCacheDirectory = "bazel-disk-cache"
	// BazelOutputDirectory is a session's output base, under its run directory.
	BazelOutputDirectory = "bazel"
	// OCamlToolchainCacheDirectory is the shared OCaml toolchain, under the state directory.
	OCamlToolchainCacheDirectory = "ocaml-toolchain"

	// DefaultInterruptBudget is how long a canceled turn may take to reach a
	// safepoint before its executor is killed.
	DefaultInterruptBudget = 60 * time.Second
	// DefaultCloseBudget bounds the close of one session's executor.
	DefaultCloseBudget = 30 * time.Second

	// QueueFile is a run's queued messages and the text of the turn running,
	// under its run directory, so a restart neither loses nor skips one.
	QueueFile = "queue.json"
	// EndedFile marks a run that was canceled or failed: a restarted harness
	// does not reopen it.
	EndedFile = "ended"
	// RedeliveredPrefix opens a turn delivered again because the harness
	// restarted while it ran.
	RedeliveredPrefix = "(Delivered again: the harness restarted while this turn ran.)\n\n"
)

The environment every session's turn executor inherits for Bazel: the shared, content-addressed disk cache under the state directory and an output base of its own under the run directory, removed when the session ends.

View Source
const AllowedModelsFile = "allowed-models.json"

AllowedModelsFile is this host's model policy, under the state directory: the models a real session may run on, each with the ruling that allows it. It is data the operator edits; every check reads it again, so an edit takes effect without a restart.

View Source
const ClawSystemInvariants = `` /* 547-byte string literal not displayed */

ClawSystemInvariants are the Core-owned operating constraints shared by built-in agent runtimes. A provider implementation may append capability- specific instructions, but must not weaken these ownership boundaries.

View Source
const ExecutorDefaultFile = "executor-default.json"

ExecutorDefaultFile records the operator's last switch of the default executor and model, under the state directory, so a restarted host keeps it.

Variables

View Source
var (
	// ErrNoPullRequest reports a ready or merge for a session that has not
	// opened a pull request.
	ErrNoPullRequest = fmt.Errorf("%w: the session has no pull request", csf.ErrConflict)
	// ErrNoLauncher reports a ready or merge on a host that granted the
	// service no process capability.
	ErrNoLauncher = errors.New("harness: this host grants no launcher for pull request operations")
)
View Source
var (
	// ErrModelNotAllowed reports a session asked to run on a model the host's
	// policy does not allow.
	ErrModelNotAllowed = fmt.Errorf("%w: the model is not allowed on this host", csf.ErrInvalidRequest)
	// ErrNoAllowedModels reports a policy that allows no model at all, so no
	// session may run.
	ErrNoAllowedModels = fmt.Errorf("%w: the host's model policy allows no model", csf.ErrInvalidRequest)
	// ErrModelPolicy reports a policy file that cannot be read as one.
	ErrModelPolicy = errors.New("harness: the model policy cannot be read")
)
View Source
var (
	// ErrRunnerClosed reports an attempted start after lifecycle shutdown began.
	ErrRunnerClosed = errors.New("harness runner is closed")
	// ErrRunnerStarted reports an attempted second start or adapter installation.
	ErrRunnerStarted = errors.New("harness runner is already started")
	// ErrRuntimeUnavailable reports an operation without an active adapter.
	ErrRuntimeUnavailable = errors.New("harness runtime is unavailable")
)
View Source
var (
	// ErrNoRunner reports a service built without the session runner.
	ErrNoRunner = errors.New("harness: a session runner is required")
	// ErrInvalidServiceOption reports a nil option or a value the service
	// cannot use.
	ErrInvalidServiceOption = errors.New("harness: invalid option")
	// ErrNotStarted reports an operation before the service was mounted and
	// started, or after it stopped.
	ErrNotStarted = errors.New("harness: the session service is not running")
	// ErrUnknownSession reports an assignment this harness holds no session
	// for.
	ErrUnknownSession = fmt.Errorf("%w: no session for this assignment", csf.ErrNotFound)
	// ErrSessionExists reports a second Submit of an assignment this harness
	// already runs.
	ErrSessionExists = fmt.Errorf("%w: the assignment already has a session", csf.ErrConflict)
	// ErrSessionFinished reports a Send to a session that has ended.
	ErrSessionFinished = fmt.Errorf("%w: the session has finished", csf.ErrConflict)
	// ErrAdmissionHeld reports a Submit refused while admission is held, as
	// housekeeping holds it while free disk is below its floor.
	ErrAdmissionHeld = fmt.Errorf("%w: admission is held", csf.ErrConflict)
	// ErrInterruptBudget is the cause a running turn is killed with when its
	// executor did not reach a safepoint within the interrupt budget.
	ErrInterruptBudget = errors.New("harness: the turn did not stop at a safepoint within the interrupt budget")
)
View Source
var DefaultModelPolicy = ModelPolicy{Allowed: []AllowedModel{{
	Model:   "claude-opus-5-5",
	Ruling:  "every real session runs claude-opus-5-5. Never Fable.",
	RuledBy: "operator",
	RuledOn: "2026-10-05",
}}}

DefaultModelPolicy is the policy a host starts from when its state directory records none: the operator's ruling of 2026-10-05.

View Source
var (
	// ErrNoDefaultModel reports a default that moves recipes to an executor
	// other than the one they are written for without naming the model.
	ErrNoDefaultModel = fmt.Errorf("%w: a default executor other than %s needs a model, spelled as it spells it", csf.ErrInvalidRequest, session.ExecutorClaudeCode)
)
View Source
var ErrNoEvents = fmt.Errorf("%w: no event log for this assignment", ErrUnknownSession)

ErrNoEvents reports a session whose event log does not exist.

View Source
var ErrResumeHeld = errors.New("harness: resume held by the launch check")

ErrResumeHeld reports a resume the launch check did not admit; it stays in the resume queue and is tried again.

Functions

func BackendName

func BackendName(backend deployv1.HarnessBackend) string

BackendName projects the generated HarnessBackend enum into its stable configuration and diagnostic spelling.

func RecordModelPolicy added in v0.3.0

func RecordModelPolicy(stateDirectory string, policy ModelPolicy) error

RecordModelPolicy replaces the model policy recorded under the state directory.

func TurnGaps added in v0.3.0

func TurnGaps(stateDirectory string) ([]time.Duration, error)

TurnGaps measures, across every run directory under stateDirectory, the gap from each turn's end to the next turn's request on the same run. A run whose log cannot be read contributes nothing.

func WorkerCap

func WorkerCap(cores uint32, load float64) uint32

WorkerCap is how many sessions cores under a one-minute load admit: the idle cores, floored at one. It is the concurrency the launch check reports and the cap the dispatch service schedules within.

Types

type AgentSessionService

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

AgentSessionService runs agent sessions concurrently in one process. Each session is one owner goroutine on a child scope of the service's scope: it opens the session, runs its turns from a queue Send fills, and stops at a safepoint when canceled or when the service stops. The table of sessions is owned by one goroutine behind a mailbox, so no field needs a lock.

func NewAgentSessionService

func NewAgentSessionService(options ...AgentSessionServiceOption) (*AgentSessionService, error)

NewAgentSessionService validates the whole option set before building the service. WithSessionRunner is required.

func (*AgentSessionService) Cancel

Cancel asks the session's owner to stop at its next safepoint. Between turns the owner stops at once; during a turn the executor is interrupted and the owner stops when the turn reports its result, or kills the executor after the interrupt budget. A finished session is reported as it is.

func (*AgentSessionService) Check added in v0.3.0

Check is the launch check a Submit made now would get, without submitting: the worker cap from the cores and the one-minute load, and the disk floor from the largest run directory measured so far. While admission is held it reports ErrAdmissionHeld with the reason instead.

func (*AgentSessionService) CheckAdmission added in v0.3.0

CheckAdmission reports the launch check the recipe would meet now, and why admission is held when it is, without admitting it.

func (*AgentSessionService) DropInboxMessage added in v0.3.0

DropInboxMessage removes a message from the queue.

func (*AgentSessionService) Get

Get reports one session.

func (*AgentSessionService) GetExecutorDefault added in v0.3.0

GetExecutorDefault reports the executor and model a recipe naming no executor runs on now.

func (*AgentSessionService) HoldAdmission

func (service *AgentSessionService) HoldAdmission(reason string)

HoldAdmission refuses every Submit, with reason, until ReleaseAdmission. Sessions already running are not affected.

func (*AgentSessionService) IdleBound added in v0.3.0

func (service *AgentSessionService) IdleBound() IdleBoundReport

IdleBound is the derived idle bound with its derivation, after Start.

func (*AgentSessionService) List

List reports every session in submission order and the process they run in.

func (*AgentSessionService) ListInbox added in v0.3.0

ListInbox reports every message in a session's inbox with its state.

func (*AgentSessionService) Merge added in v0.3.0

Merge merges the session's pull request through the repository's merge path: the worktree's tools/merge-pr.sh, which merges the head with main, runs the merge checks and squash-merges only a passing result, or gh's squash merge where the repository has no such script. It waits for the path to finish; the report is what it printed.

The merge runs on the service's own scope, not the caller's: a caller that goes away — a closed Workbench tab, an MCP client that timed out — stops waiting, and the merge carries on and records its outcome in the session's event log. It records merge_started first, so the log says a merge is in flight until the merge record that ends it.

func (*AgentSessionService) MoveInboxMessage added in v0.3.0

MoveInboxMessage reorders a message to a specific index in the queue.

func (*AgentSessionService) OpenTail

func (service *AgentSessionService) OpenTail(ctx context.Context, assignmentID string, from int) (*EventTail, error)

OpenTail opens the event log of assignmentID, skipping the first from records. A session this harness is still opening is waited for; a session it does not hold is read from disk as it is, which is how a finished run from an earlier process is replayed.

func (*AgentSessionService) Propose added in v0.3.0

Propose applies a proposed patch to a session's worktree.

func (*AgentSessionService) Ready added in v0.3.0

Ready marks the session's draft pull request ready for review.

func (*AgentSessionService) RecordRuling added in v0.3.0

RecordRuling records one operator ruling under the harness's state directory and returns the rulings in force. Every turn sent afterwards, in every session, carries them to the question gate; the Workbench, the rulings view and the rulings series read the same records.

func (*AgentSessionService) Register

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

Register mounts the event stream route on the caller's router. The generated operations are registered by the CSF service the harness is mounted behind.

func (*AgentSessionService) ReleaseAdmission

func (service *AgentSessionService) ReleaseAdmission()

ReleaseAdmission admits sessions again.

func (*AgentSessionService) Resumed added in v0.3.0

func (service *AgentSessionService) Resumed() <-chan struct{}

Resumed is closed once Start's first resume pass has run: every open run a previous process left is resumed or recorded as held in ResumeQueueFile, so a host that waits on it is ready only with them back.

func (*AgentSessionService) Send

Send queues message for the session's next turn and acknowledges it with a receipt containing the message state.

func (*AgentSessionService) SetExecutorDefault added in v0.3.0

SetExecutorDefault switches the default for every session submitted from now on and records it under the state directory. Running sessions keep the executor they opened on.

func (*AgentSessionService) Start

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

Start creates the shared Bazel disk cache, measures the existing runs for the disk floor and the idle bound, and starts the registry owner on the scope. The registry retires when the scope is canceled; every session's child scope is canceled with it and joined by it. With a process table granted, the resident series sampler runs on the scope too.

func (*AgentSessionService) Stop

Stop asks the host to shut down and reports how many sessions were running.

func (*AgentSessionService) Submit

Submit admits one recipe as a session: it is checked against the launch bounds, recorded, and opened by its owner goroutine on a child scope. It returns once the session is open, with its receipt, or with the failure that kept it from opening.

func (*AgentSessionService) TopInboxMessage added in v0.3.0

TopInboxMessage promotes a message to the front of the queue.

type AgentSessionServiceOption

type AgentSessionServiceOption func(service *AgentSessionService) error

AgentSessionServiceOption configures an AgentSessionService.

func WithClock

func WithClock(clock IClock) AgentSessionServiceOption

WithClock replaces the clock.

func WithExecutorDefault added in v0.3.0

func WithExecutorDefault(executor session.Executor, model string) AgentSessionServiceOption

WithExecutorDefault is the executor and model a recipe naming no executor runs on until the operator switches it; a switch recorded under the state directory takes precedence. The default default is Claude Code with each recipe's own model.

func WithHostMeasures

func WithHostMeasures(measures IHostMeasures) AgentSessionServiceOption

WithHostMeasures replaces the host measures the launch check reads.

func WithHostPID

func WithHostPID(pid int) AgentSessionServiceOption

WithHostPID names the process every session runs in, for List and for finding each session's executor among its children.

func WithInterruptBudget

func WithInterruptBudget(budget time.Duration) AgentSessionServiceOption

WithInterruptBudget bounds how long a canceled turn may run on after the interrupt before its executor is killed.

func WithLauncher added in v0.3.0

func WithLauncher(launcher proc.ILauncher) AgentSessionServiceOption

WithLauncher grants the process capability gh and the merge path start through. Without it Ready and Merge report ErrNoLauncher.

func WithProcessTable added in v0.3.0

func WithProcessTable(processes iofs.IFiles) AgentSessionServiceOption

WithProcessTable grants the host's process table, as /proc. With it the service closes idle executors, never one with a live background child, and samples the resident memory of the harness and its executors into the resident series. Without it executors stay open and nothing is sampled.

func WithResidentSampleInterval added in v0.3.0

func WithResidentSampleInterval(interval time.Duration) AgentSessionServiceOption

WithResidentSampleInterval sets how often the resident series is sampled.

func WithServiceLogger

func WithServiceLogger(logger *slog.Logger) AgentSessionServiceOption

WithServiceLogger receives the service's own records; sessions keep writing their events.jsonl.

func WithSessionObserver

func WithSessionObserver(observe SessionObserver) AgentSessionServiceOption

WithSessionObserver adds an observer of every session's state changes. It may be given more than once; observers are called in the order given.

func WithSessionRunner

func WithSessionRunner(runner *session.AgentSessionRunner) AgentSessionServiceOption

WithSessionRunner grants the runner that prepares and opens each session. Required.

func WithStopRequest

func WithStopRequest(stop func()) AgentSessionServiceOption

WithStopRequest grants the host's shutdown: StopHarness calls it. Without it StopHarness reports that the host cannot be stopped this way.

type AllowedModel added in v0.3.0

type AllowedModel struct {
	Model  string `json:"model"`
	Ruling string `json:"ruling"`
	// RuledBy and RuledOn are who ruled and the day, as YYYY-MM-DD.
	RuledBy string `json:"ruled_by"`
	RuledOn string `json:"ruled_on"`
}

AllowedModel is one model a real session may run on and the ruling that allows it.

type EventTail

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

EventTail reads one session's records in order and waits for more while the session is live. One goroutine reads it.

func (*EventTail) Available

func (tail *EventTail) Available() []byte

Available returns the next record already in the file, or nil at its current end; it never waits.

func (*EventTail) Close

func (tail *EventTail) Close() error

Close closes the log file.

func (*EventTail) Next

func (tail *EventTail) Next(ctx context.Context) ([]byte, error)

Next returns the next record, waiting for it while the session is live. io.EOF reports that the session has finished and every record was read.

func (*EventTail) Sequence

func (tail *EventTail) Sequence() int

Sequence is the number of records read so far: the next record's position.

type FactoryFunc

type FactoryFunc func(harnessContext *deployv1.HarnessContext, host IHost) (*Instance, error)

FactoryFunc adapts a function to IFactory.

func (FactoryFunc) New

func (factory FactoryFunc) New(
	harnessContext *deployv1.HarnessContext,
	host IHost,
) (*Instance, error)

New implements IFactory.

type HeldResume added in v0.3.0

type HeldResume struct {
	AssignmentID string `json:"assignment_id"`
	// Position is its place in the resume order, from 1.
	Position int `json:"position"`
	// Reason is the launch check's findings that held it.
	Reason string `json:"reason"`
	// HeldSince is when it was first held.
	HeldSince time.Time `json:"held_since"`
}

HeldResume is one open run a restart has not resumed yet.

type HostMeasures

type HostMeasures struct{}

HostMeasures reads this host's kernel.

func (HostMeasures) Cores

func (HostMeasures) Cores() int

Cores is the Go runtime's CPU count.

func (HostMeasures) DirectoryBytes

func (HostMeasures) DirectoryBytes(directory string) (uint64, error)

DirectoryBytes sums the regular files under directory; a directory that does not exist is empty.

func (HostMeasures) FreeBytes

func (HostMeasures) FreeBytes(path string) (uint64, error)

FreeBytes asks the filesystem holding path.

func (HostMeasures) LoadAverage

func (HostMeasures) LoadAverage() (float64, error)

LoadAverage reads the first field of /proc/loadavg.

func (HostMeasures) MemoryAvailable added in v0.3.0

func (HostMeasures) MemoryAvailable() (uint64, error)

MemoryAvailable reads MemAvailable from /proc/meminfo.

func (HostMeasures) Pressure added in v0.3.0

func (HostMeasures) Pressure(resource proc.PressureResource) (proc.Pressure, error)

Pressure reads the resource's file under /proc/pressure.

type IClock

type IClock interface {
	Now() time.Time
	// After delivers once after d, unless the returned stop ran first.
	After(d time.Duration) (<-chan time.Time, func() bool)
	// AfterFunc runs f after d, unless the returned stop ran first.
	AfterFunc(d time.Duration, f func()) (stop func() bool)
}

IClock is the time source the service measures with: the system clock in production, a controllable one in specs.

type IFactory

type IFactory interface {
	New(harnessContext *deployv1.HarnessContext, host IHost) (*Instance, error)
}

IFactory constructs one harness runtime against Core-owned host services. Core passes an owned context snapshot and keeps IHost valid until the returned runtime is closed; implementations must not mutate protobuf inputs.

type IHost

type IHost interface {
	Publish(ctx context.Context, event *deployv1.HarnessEvent) error
	FleetStatus(ctx context.Context) (*deployv1.HarnessFleetStatus, error)
	Reconcile(
		ctx context.Context,
		request *deployv1.HarnessReconcileRequest,
	) (*deployv1.ReconcileEvidence, error)
}

IHost is the Core-owned capability surface available to a harness. Core retains fleet observation, approval, fencing, and reconciliation authority. Its methods are safe for concurrent provider callbacks, snapshot protobuf inputs before returning, and honor cancellation. Reconcile may remain blocked while Core waits for operator approval.

type IHostMeasures

type IHostMeasures interface {
	// Cores is the number of CPUs the scheduler may run threads on.
	Cores() int
	// LoadAverage is the one-minute run-queue average.
	LoadAverage() (float64, error)
	// FreeBytes is the space available to this process on path's filesystem.
	FreeBytes(path string) (uint64, error)
	// DirectoryBytes is the size of every regular file under directory.
	DirectoryBytes(directory string) (uint64, error)
	// Pressure is one resource's pressure stall information.
	Pressure(resource proc.PressureResource) (proc.Pressure, error)
	// MemoryAvailable is the memory the kernel estimates a new workload can
	// use without swapping.
	MemoryAvailable() (uint64, error)
}

IHostMeasures is what the launch check reads from the host: the measures its two bounds are derived from. The host implementation reads the kernel; specs grant a double.

type IRuntime

type IRuntime interface {
	Start(ctx context.Context) (*deployv1.HarnessSession, error)
	Activate(ctx context.Context) error
	Send(ctx context.Context, prompt *deployv1.HarnessPrompt) error
	Abort(ctx context.Context) error
	Close() error
}

IRuntime owns one harness session and its turn lifecycle. Core calls Start once, then Activate once after a successful start, and serializes Send and Abort. Close may overlap an in-flight call, so implementations must synchronize provider resources, respect each call's context, and make Close idempotent, including after a failed Start or Activate.

type IdleBoundReport added in v0.3.0

type IdleBoundReport struct {
	// Bound is how long a session may sit between turns before its executor
	// is suspended.
	Bound time.Duration `json:"bound"`
	// Quantile is the share of measured gaps at or below the bound; zero when
	// the bound is the fallback.
	Quantile float64 `json:"quantile"`
	// Gaps is how many gaps were measured.
	Gaps int `json:"gaps"`
	// Quantiles are the reported quantiles of the measured gaps.
	Quantiles []Quantile `json:"quantiles"`
	// Fallback reports that fewer than two gaps were measured, so the bound
	// is the fallback the caller gave: the prompt cache lifetime.
	Fallback bool `json:"fallback"`
}

IdleBoundReport is the derived idle bound with its derivation: the gaps it was measured on, the quantile the bound sits at and the reported quantiles.

func IdleBound added in v0.3.0

func IdleBound(gaps []time.Duration, fallback time.Duration) IdleBoundReport

IdleBound derives the idle bound from gaps: the knee of the empirical distribution of log(1+seconds), the point of the cumulative curve farthest from the chord between its ends. With fewer than two gaps there is no curve, and the bound is fallback.

func (IdleBoundReport) String added in v0.3.0

func (report IdleBoundReport) String() string

String is the report in a line, for the log and the pull request body.

type Instance

type Instance struct {
	Runtime  IRuntime
	Identity *deployv1.HarnessRuntimeIdentity
}

Instance binds a runtime to the immutable identity recorded by Core.

type ModelPolicy added in v0.3.0

type ModelPolicy struct {
	Allowed []AllowedModel `json:"allowed"`
}

ModelPolicy is the allowed-models file.

func (ModelPolicy) Admit added in v0.3.0

func (policy ModelPolicy) Admit(model string) error

Admit refuses model unless the policy allows it. The refusal names the allowed models and the rulings that allow them.

type Quantile added in v0.3.0

type Quantile struct {
	P     float64       `json:"p"`
	Value time.Duration `json:"value"`
}

Quantile is one quantile of the gap distribution.

type ResidentSample added in v0.3.0

type ResidentSample struct {
	At time.Time `json:"at"`
	// HarnessRSSBytes is the harness process's own resident set.
	HarnessRSSBytes uint64 `json:"harness_rss_bytes"`
	// ExecutorsRSSBytes is the resident set of every process in the groups the
	// harness's children lead: the executors and whatever they started.
	ExecutorsRSSBytes uint64 `json:"executors_rss_bytes"`
	// OpenSessions counts the sessions that have not ended.
	OpenSessions int `json:"open_sessions"`
	// ExecutorsAlive counts the open sessions whose executor process is open.
	ExecutorsAlive int `json:"executors_alive"`
	// Resumes counts the sessions resumed after a suspend since the
	// harness started.
	Resumes int `json:"resumes"`
	// ResumeTimeToFirstTokenMs is the latest resumed turn's time to its first
	// token; zero until one has run.
	ResumeTimeToFirstTokenMs float64 `json:"resume_ttft_ms"`
	// IdleBoundSeconds and IdleBoundQuantile are the derived bound and where
	// it sits in the Gaps measured, so the derivation travels with the series.
	IdleBoundSeconds  float64 `json:"idle_bound_s"`
	IdleBoundQuantile float64 `json:"idle_bound_quantile"`
	Gaps              int     `json:"gaps"`
}

ResidentSample is one sample of the series.

func ReadResidentSeries added in v0.3.0

func ReadResidentSeries(content []byte) []ResidentSample

ReadResidentSeries decodes the series file's content, skipping a line that is not a sample. It is pure.

type ResumeQueue added in v0.3.0

type ResumeQueue struct {
	UpdatedAt time.Time    `json:"updated_at"`
	Resumed   int          `json:"resumed"`
	Held      []HeldResume `json:"held"`
}

ResumeQueue is the resume queue at one pass: how many open runs are back and which are held.

type Runner

type Runner[Event any] struct {
	// contains filtered or unexported fields
}

Runner serializes one provider session's lifecycle and event activation. Its dedicated receiver goroutine owns the installed operations, callback buffer, activation state, in-flight count, and close result.

func NewRunner

func NewRunner[Event any](publish func(event Event), eventID func(event Event) string) *Runner[Event]

NewRunner starts the lifecycle owner. publish receives replay before callbacks accepted prior to Activate; eventID supplies the deduplication key. A nil publish discards events and a nil eventID disables deduplication.

func (*Runner[Event]) Abort

func (runner *Runner[Event]) Abort(ctx context.Context) error

Abort runs native provider cancellation without blocking the lifecycle owner.

func (*Runner[Event]) Activate

func (runner *Runner[Event]) Activate(replay []Event)

Activate publishes replay before buffered callbacks, deduplicating nonempty event IDs. A second call is a FIFO barrier and otherwise has no effect.

func (*Runner[Event]) BeginStart

func (runner *Runner[Event]) BeginStart() error

BeginStart reserves the runner for one provider session. It must precede provider startup so callbacks emitted while the session opens are retained.

func (*Runner[Event]) Close

func (runner *Runner[Event]) Close() error

Close cancels accepted operations, runs provider cleanup once, and waits for both. It is idempotent and concurrent callers receive the same cleanup error.

func (*Runner[Event]) Install

func (runner *Runner[Event]) Install(
	send func(ctx context.Context, prompt *deployv1.HarnessPrompt) error,
	abort func(ctx context.Context) error,
	cleanup func() error,
) error

Install makes native provider operations available after its session opens. Nil operations remain unavailable; nil cleanup means no provider cleanup. Until Install succeeds, the caller retains cleanup ownership. Installed operations must return when their context is canceled or cleanup runs.

func (*Runner[Event]) Publish

func (runner *Runner[Event]) Publish(event Event)

Publish accepts a provider callback. Before Activate it is buffered; after Activate it is delivered synchronously by the lifecycle owner.

func (*Runner[Event]) Send

func (runner *Runner[Event]) Send(
	ctx context.Context,
	prompt *deployv1.HarnessPrompt,
) error

Send runs one native provider turn without blocking the lifecycle owner.

type SessionObserver

type SessionObserver func(state *harnessv1.AgentSessionState)

SessionObserver receives a copy of a session's state after each change to it: the phase transitions, turn counts and pull request URL. It is called on the goroutine that made the change, never on the registry, so it may enqueue work but must return promptly.

type SystemClock

type SystemClock struct{}

SystemClock is the wall clock.

func (SystemClock) After

func (SystemClock) After(d time.Duration) (<-chan time.Time, func() bool)

func (SystemClock) AfterFunc

func (SystemClock) AfterFunc(d time.Duration, f func()) func() bool

func (SystemClock) Now

func (SystemClock) Now() time.Time

Directories

Path Synopsis
Package chat is the Workbench chat for harness sessions: a gotth-live page that shows one session's transcript as its event log grows and sends the operator's messages into the open session through the harness service, in process.
Package chat is the Workbench chat for harness sessions: a gotth-live page that shows one session's transcript as its event log grows and sends the operator's messages into the open session through the harness service, in process.
Package costs is the cost model of the harness's session operations.
Package costs is the cost model of the harness's session operations.
Package endpoint is the endpoint registry: the typed record, in the harness state directory, of every operator-facing endpoint csf serve has served (the Workbench, the API and MCP), each with its owner, its users and every address it has been served at, and of every address the operator retired, with the address it moved to, the day and the operator's acknowledgement.
Package endpoint is the endpoint registry: the typed record, in the harness state directory, of every operator-facing endpoint csf serve has served (the Workbench, the API and MCP), each with its owner, its users and every address it has been served at, and of every address the operator retired, with the address it moved to, the day and the operator's acknowledgement.
Package longturns detects and manages long-running turns that may indicate stuck execution.
Package longturns detects and manages long-running turns that may indicate stuck execution.
Package mocks is a generated GoMock package.
Package mocks is a generated GoMock package.
Package opencode implements the built-in OpenCode agent runtime behind the public deploy harness seam.
Package opencode implements the built-in OpenCode agent runtime behind the public deploy harness seam.
Package progress projects each harness session into one session_progress row — the assignment, agent, phase, turn, queue depth, tool calls, commits ahead, changed files, model, last activity and last note — and folds the session's event log into the readable lines a tail follows.
Package progress projects each harness session into one session_progress row — the assignment, agent, phase, turn, queue depth, tool calls, commits ahead, changed files, model, last activity and last note — and folds the session's event log into the readable lines a tail follows.
Package provider is the inference provider every session this host launches reaches its model through.
Package provider is the inference provider every session this host launches reaches its model through.
Package relayhook provides the UserPromptSubmit hook that forwards prompts to an orchestrator session and blocks the proxy session from taking turns.
Package relayhook provides the UserPromptSubmit hook that forwards prompts to an orchestrator session and blocks the proxy session from taking turns.
Package routing maps virtual sessions onto real sessions.
Package routing maps virtual sessions onto real sessions.
Package session is the agent harness's session runner: it runs an agent's session for an assignment on a turn executor, in a git worktree of its own, with the session gates installed and every event logged under the run's trace.
Package session is the agent harness's session runner: it runs an agent's session for an assignment on a turn executor, in a git worktree of its own, with the session gates installed and every event logged under the run's trace.
mocks
Package mocks is a generated GoMock package.
Package mocks is a generated GoMock package.
Package sessiongate holds the session gates the agent harness installs into every session it runs: structural checks Claude Code calls as hooks around a tool call and at the end of a turn.
Package sessiongate holds the session gates the agent harness installs into every session it runs: structural checks Claude Code calls as hooks around a tool call and at the end of a turn.
Package upgrade replaces the csf binary with a released one and restarts the harness on it, keeping the binary it replaced for a rollback.
Package upgrade replaces the csf binary with a released one and restarts the harness on it, keeping the binary it replaced for a rollback.

Jump to

Keyboard shortcuts

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