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 RecordModelPolicy(stateDirectory string, policy ModelPolicy) error
- func TurnGaps(stateDirectory string) ([]time.Duration, error)
- 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) Check(ctx context.Context) (*harnessv1.LaunchCheck, error)
- func (service *AgentSessionService) CheckAdmission(ctx context.Context, request *harnessv1.CheckAgentSessionAdmissionRequest) (*harnessv1.CheckAgentSessionAdmissionResponse, error)
- func (service *AgentSessionService) DropInboxMessage(ctx context.Context, request *harnessv1.DropInboxMessageRequest) (*harnessv1.DropInboxMessageResponse, error)
- func (service *AgentSessionService) Get(ctx context.Context, request *harnessv1.GetAgentSessionRequest) (*harnessv1.GetAgentSessionResponse, error)
- func (service *AgentSessionService) GetExecutorDefault(ctx context.Context, _ *harnessv1.GetAgentExecutorDefaultRequest) (*harnessv1.GetAgentExecutorDefaultResponse, error)
- func (service *AgentSessionService) HoldAdmission(reason string)
- func (service *AgentSessionService) IdleBound() IdleBoundReport
- func (service *AgentSessionService) List(ctx context.Context, _ *harnessv1.ListAgentSessionsRequest) (*harnessv1.ListAgentSessionsResponse, error)
- func (service *AgentSessionService) ListInbox(ctx context.Context, request *harnessv1.ListInboxRequest) (*harnessv1.ListInboxResponse, error)
- func (service *AgentSessionService) Merge(ctx context.Context, request *harnessv1.MergeAgentSessionPullRequestRequest) (*harnessv1.MergeAgentSessionPullRequestResponse, error)
- func (service *AgentSessionService) MoveInboxMessage(ctx context.Context, request *harnessv1.MoveInboxMessageRequest) (*harnessv1.MoveInboxMessageResponse, error)
- func (service *AgentSessionService) OpenTail(ctx context.Context, assignmentID string, from int) (*EventTail, error)
- func (service *AgentSessionService) Propose(ctx context.Context, request *harnessv1.ProposeProposalRequest) (*harnessv1.ProposeProposalResponse, error)
- func (service *AgentSessionService) Ready(ctx context.Context, request *harnessv1.ReadyAgentSessionPullRequestRequest) (*harnessv1.ReadyAgentSessionPullRequestResponse, error)
- func (service *AgentSessionService) RecordRuling(ctx context.Context, request *harnessv1.RecordRulingRequest) (*harnessv1.RecordRulingResponse, error)
- func (service *AgentSessionService) Register(router gin.IRouter)
- func (service *AgentSessionService) ReleaseAdmission()
- func (service *AgentSessionService) Resumed() <-chan struct{}
- func (service *AgentSessionService) Send(ctx context.Context, request *harnessv1.SendAgentSessionMessageRequest) (*harnessv1.SendAgentSessionMessageResponse, error)
- func (service *AgentSessionService) SetExecutorDefault(ctx context.Context, request *harnessv1.SetAgentExecutorDefaultRequest) (*harnessv1.SetAgentExecutorDefaultResponse, 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)
- func (service *AgentSessionService) TopInboxMessage(ctx context.Context, request *harnessv1.TopInboxMessageRequest) (*harnessv1.TopInboxMessageResponse, error)
- type AgentSessionServiceOption
- func WithClock(clock IClock) AgentSessionServiceOption
- func WithExecutorDefault(executor session.Executor, model string) AgentSessionServiceOption
- func WithHostMeasures(measures IHostMeasures) AgentSessionServiceOption
- func WithHostPID(pid int) AgentSessionServiceOption
- func WithInterruptBudget(budget time.Duration) AgentSessionServiceOption
- func WithLauncher(launcher proc.ILauncher) AgentSessionServiceOption
- func WithProcessTable(processes iofs.IFiles) AgentSessionServiceOption
- func WithResidentSampleInterval(interval 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 AllowedModel
- type EventTail
- type FactoryFunc
- type HeldResume
- type HostMeasures
- func (HostMeasures) Cores() int
- func (HostMeasures) DirectoryBytes(directory string) (uint64, error)
- func (HostMeasures) FreeBytes(path string) (uint64, error)
- func (HostMeasures) LoadAverage() (float64, error)
- func (HostMeasures) MemoryAvailable() (uint64, error)
- func (HostMeasures) Pressure(resource proc.PressureResource) (proc.Pressure, error)
- type IClock
- type IFactory
- type IHost
- type IHostMeasures
- type IRuntime
- type IdleBoundReport
- type Instance
- type ModelPolicy
- type Quantile
- type ResidentSample
- type ResumeQueue
- 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 ( // 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.
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.
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.
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 )
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.
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.
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.
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 ¶
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") )
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") )
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 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.
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) )
var ErrNoEvents = fmt.Errorf("%w: no event log for this assignment", ErrUnknownSession)
ErrNoEvents reports a session whose event log does not exist.
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.
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) Check ¶ added in v0.3.0
func (service *AgentSessionService) Check(ctx context.Context) (*harnessv1.LaunchCheck, error)
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
func (service *AgentSessionService) CheckAdmission(ctx context.Context, request *harnessv1.CheckAgentSessionAdmissionRequest) (*harnessv1.CheckAgentSessionAdmissionResponse, error)
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
func (service *AgentSessionService) DropInboxMessage(ctx context.Context, request *harnessv1.DropInboxMessageRequest) (*harnessv1.DropInboxMessageResponse, error)
DropInboxMessage removes a message from the queue.
func (*AgentSessionService) Get ¶
func (service *AgentSessionService) Get(ctx context.Context, request *harnessv1.GetAgentSessionRequest) (*harnessv1.GetAgentSessionResponse, error)
Get reports one session.
func (*AgentSessionService) GetExecutorDefault ¶ added in v0.3.0
func (service *AgentSessionService) GetExecutorDefault(ctx context.Context, _ *harnessv1.GetAgentExecutorDefaultRequest) (*harnessv1.GetAgentExecutorDefaultResponse, error)
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 ¶
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) ListInbox ¶ added in v0.3.0
func (service *AgentSessionService) ListInbox(ctx context.Context, request *harnessv1.ListInboxRequest) (*harnessv1.ListInboxResponse, error)
ListInbox reports every message in a session's inbox with its state.
func (*AgentSessionService) Merge ¶ added in v0.3.0
func (service *AgentSessionService) Merge(ctx context.Context, request *harnessv1.MergeAgentSessionPullRequestRequest) (*harnessv1.MergeAgentSessionPullRequestResponse, error)
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
func (service *AgentSessionService) MoveInboxMessage(ctx context.Context, request *harnessv1.MoveInboxMessageRequest) (*harnessv1.MoveInboxMessageResponse, error)
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
func (service *AgentSessionService) Propose(ctx context.Context, request *harnessv1.ProposeProposalRequest) (*harnessv1.ProposeProposalResponse, error)
Propose applies a proposed patch to a session's worktree.
func (*AgentSessionService) Ready ¶ added in v0.3.0
func (service *AgentSessionService) Ready(ctx context.Context, request *harnessv1.ReadyAgentSessionPullRequestRequest) (*harnessv1.ReadyAgentSessionPullRequestResponse, error)
Ready marks the session's draft pull request ready for review.
func (*AgentSessionService) RecordRuling ¶ added in v0.3.0
func (service *AgentSessionService) RecordRuling(ctx context.Context, request *harnessv1.RecordRulingRequest) (*harnessv1.RecordRulingResponse, error)
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 ¶
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 a receipt containing the message state.
func (*AgentSessionService) SetExecutorDefault ¶ added in v0.3.0
func (service *AgentSessionService) SetExecutorDefault(ctx context.Context, request *harnessv1.SetAgentExecutorDefaultRequest) (*harnessv1.SetAgentExecutorDefaultResponse, error)
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 ¶
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.
func (*AgentSessionService) TopInboxMessage ¶ added in v0.3.0
func (service *AgentSessionService) TopInboxMessage(ctx context.Context, request *harnessv1.TopInboxMessageRequest) (*harnessv1.TopInboxMessageResponse, error)
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 ¶
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 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) 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 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 ¶
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
Source Files
¶
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. |