Documentation
¶
Overview ¶
Package harness defines the public behavior boundary between deploy and a compiled-in agent runtime implementation.
Index ¶
- Constants
- Variables
- func BackendName(backend deployv1.HarnessBackend) string
- func WorkerCap(cores uint32, load float64) uint32
- type AgentSessionService
- func (service *AgentSessionService) Cancel(ctx context.Context, request *harnessv1.CancelAgentSessionRequest) (*harnessv1.CancelAgentSessionResponse, error)
- func (service *AgentSessionService) Get(ctx context.Context, request *harnessv1.GetAgentSessionRequest) (*harnessv1.GetAgentSessionResponse, error)
- func (service *AgentSessionService) HoldAdmission(reason string)
- func (service *AgentSessionService) List(ctx context.Context, _ *harnessv1.ListAgentSessionsRequest) (*harnessv1.ListAgentSessionsResponse, error)
- func (service *AgentSessionService) OpenTail(ctx context.Context, assignmentID string, from int) (*EventTail, error)
- func (service *AgentSessionService) Register(router gin.IRouter)
- func (service *AgentSessionService) ReleaseAdmission()
- func (service *AgentSessionService) Send(ctx context.Context, request *harnessv1.SendAgentSessionMessageRequest) (*harnessv1.SendAgentSessionMessageResponse, error)
- func (service *AgentSessionService) Start(scope *runtime.Scope) error
- func (service *AgentSessionService) Stop(ctx context.Context, _ *harnessv1.StopHarnessRequest) (*harnessv1.StopHarnessResponse, error)
- func (service *AgentSessionService) Submit(ctx context.Context, request *harnessv1.SubmitAgentSessionRequest) (*harnessv1.SubmitAgentSessionResponse, error)
- type AgentSessionServiceOption
- func WithClock(clock IClock) AgentSessionServiceOption
- func WithHostMeasures(measures IHostMeasures) AgentSessionServiceOption
- func WithHostPID(pid int) AgentSessionServiceOption
- func WithInterruptBudget(budget time.Duration) AgentSessionServiceOption
- func WithServiceLogger(logger *slog.Logger) AgentSessionServiceOption
- func WithSessionObserver(observe SessionObserver) AgentSessionServiceOption
- func WithSessionRunner(runner *session.AgentSessionRunner) AgentSessionServiceOption
- func WithStopRequest(stop func()) AgentSessionServiceOption
- type EventTail
- type FactoryFunc
- type HostMeasures
- type IClock
- type IFactory
- type IHost
- type IHostMeasures
- type IRuntime
- type Instance
- type Runner
- func (runner *Runner[Event]) Abort(ctx context.Context) error
- func (runner *Runner[Event]) Activate(replay []Event)
- func (runner *Runner[Event]) BeginStart() error
- func (runner *Runner[Event]) Close() error
- func (runner *Runner[Event]) Install(send func(ctx context.Context, prompt *deployv1.HarnessPrompt) error, ...) error
- func (runner *Runner[Event]) Publish(event Event)
- func (runner *Runner[Event]) Send(ctx context.Context, prompt *deployv1.HarnessPrompt) error
- type SessionObserver
- type SystemClock
Constants ¶
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.
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.
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 ¶
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 = errors.New("harness runtime is unavailable") )
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") )
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.
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 ¶
func (service *AgentSessionService) Cancel(ctx context.Context, request *harnessv1.CancelAgentSessionRequest) (*harnessv1.CancelAgentSessionResponse, error)
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 ¶
func (service *AgentSessionService) Get(ctx context.Context, request *harnessv1.GetAgentSessionRequest) (*harnessv1.GetAgentSessionResponse, error)
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 ¶
func (service *AgentSessionService) List(ctx context.Context, _ *harnessv1.ListAgentSessionsRequest) (*harnessv1.ListAgentSessionsResponse, error)
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 ¶
func (service *AgentSessionService) Send(ctx context.Context, request *harnessv1.SendAgentSessionMessageRequest) (*harnessv1.SendAgentSessionMessageResponse, error)
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 ¶
func (service *AgentSessionService) Stop(ctx context.Context, _ *harnessv1.StopHarnessRequest) (*harnessv1.StopHarnessResponse, error)
Stop asks the host to shut down and reports how many sessions were running.
func (*AgentSessionService) Submit ¶
func (service *AgentSessionService) Submit(ctx context.Context, request *harnessv1.SubmitAgentSessionRequest) (*harnessv1.SubmitAgentSessionResponse, error)
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 ¶
Available returns the next record already in the file, or nil at its current end; it never waits.
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) 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 ¶
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 ¶
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 ¶
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.
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) 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. |