harness

package
v0.2.8 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: 30 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 (
	// 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.
	DefaultTailInterval = 250 * 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 (

	// 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. The OCaml toolchain cache is also shared across sessions and built once.

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.

Variables

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 ErrNoEvents = fmt.Errorf("%w: no event log for this assignment", ErrUnknownSession)

ErrNoEvents reports a session whose event log does not exist.

Functions

func BackendName

func BackendName(backend deployv1.HarnessBackend) string

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

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

Get reports one session.

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

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

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

Send queues message for the session's next turn and acknowledges it with the turn's identifier.

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

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.

type AgentSessionServiceOption

type AgentSessionServiceOption func(service *AgentSessionService) error

AgentSessionServiceOption configures an AgentSessionService.

func WithClock

func WithClock(clock IClock) AgentSessionServiceOption

WithClock replaces the clock.

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.

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

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

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 Instance

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

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

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

Jump to

Keyboard shortcuts

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