Documentation
¶
Overview ¶
Package operator owns deploy's agent turn: policy, approvals, and run state.
Controller is the boundary between an agent harness and the rest of Core. It owns session and run identity, the fencing rules that decide whether an incoming message steers the active run or starts a new one, the approval gate in front of every reconciliation, and the operator-visible timeline. An agent runtime supplies only the turn lifecycle, behind a deliberately unexported interface; applying desired state belongs to the exported IReconciler interface. Controller passes owned values across both boundaries and retains nothing a caller handed it.
Callers may rely on ErrSessionConflict, ErrRunConflict, and ErrRunNotActive distinguishing a stale session from a stale run from a run that cannot be steered, on ApprovalQueue resolving each request exactly once — by operator decision, timeout, or abort — and on ApprovalActorState naming the actor keys written into persisted approval history. Nothing here is durable: persistence and receipts belong to store, installed by control.
Index ¶
- Variables
- type Approval
- type ApprovalActorKeys
- type ApprovalDecision
- type ApprovalNotPendingError
- type ApprovalQueue
- func (q *ApprovalQueue) ExpireRun(runID, actor string) error
- func (q *ApprovalQueue) Get(id string) (Approval, bool)
- func (q *ApprovalQueue) Pending() []Approval
- func (q *ApprovalQueue) Request(ctx context.Context, request ApprovalRequest) (ApprovalResolution, error)
- func (q *ApprovalQueue) Resolve(id string, decision ApprovalDecision, actor string) (ApprovalResolution, error)
- type ApprovalRequest
- type ApprovalResolution
- type Controller
- func (c *Controller) Abort(ctx context.Context) error
- func (c *Controller) AbortRun(ctx context.Context, expectedSessionID, expectedRunID string) error
- func (c *Controller) ApprovalQueue() *ApprovalQueue
- func (c *Controller) Close() error
- func (c *Controller) HarnessBackend() deployv1.HarnessBackend
- func (c *Controller) HarnessBackendName() string
- func (c *Controller) HarnessIdentity() *deployv1.HarnessRuntimeIdentity
- func (c *Controller) Run() RunState
- func (c *Controller) Send(ctx context.Context, prompt string, delivery deployv1.HarnessDelivery) (string, error)
- func (c *Controller) SendToSession(ctx context.Context, expectedSessionID, expectedRunID, prompt string, ...) (string, error)
- func (c *Controller) SessionID() string
- func (c *Controller) Start(ctx context.Context) error
- func (c *Controller) Status() string
- func (c *Controller) Steer(ctx context.Context, expectedSessionID, expectedRunID, prompt string, ...) error
- func (c *Controller) Subscribe() (<-chan struct{}, func())
- type IReconciler
- type RunStarted
- type RunState
- type TimelineEntry
Constants ¶
This section is empty.
Variables ¶
var ( ErrApprovalRecording = errors.New("approval request is still being recorded") ErrApprovalResolving = errors.New("approval resolution is already in progress") ErrApprovalExpired = errors.New("approval expired") )
var ( // ErrSessionConflict reports that an action targeted a different harness // session than the one currently owned by this Controller. ErrSessionConflict = errors.New("agent session changed") // ErrRunConflict reports that an action targeted an older or otherwise // different run in the current session. ErrRunConflict = errors.New("agent run changed") // ErrRunNotActive reports that a targeted run cannot currently be steered. ErrRunNotActive = errors.New("agent run is not active") )
Functions ¶
This section is empty.
Types ¶
type Approval ¶
type Approval struct {
ID string `json:"id"`
RunID string `json:"run_id,omitempty"`
ToolCallID string `json:"tool_call_id,omitempty"`
Kind string `json:"kind"`
Title string `json:"title"`
Detail string `json:"detail"`
Risk string `json:"risk"`
PayloadSHA256 string `json:"payload_sha256"`
RequiresFleetQuorum bool `json:"requires_fleet_quorum,omitempty"`
RequestedAt time.Time `json:"requested_at"`
ExpiresAt time.Time `json:"expires_at"`
}
Approval is the operator-visible form of a pending request, carrying the payload digest rather than the payload itself.
type ApprovalActorKeys ¶
type ApprovalActorKeys struct {
Operator string
Runtime string
Timeout string
Abort string
WebOperator string
}
ApprovalActorKeys names the canonical actors recorded in approval history.
func ApprovalActorState ¶
func ApprovalActorState() ApprovalActorKeys
ApprovalActorState returns the actor keys used in persisted approval history.
type ApprovalDecision ¶
type ApprovalDecision string
ApprovalDecision is the terminal outcome recorded for an approval.
const ( // DecisionApprove and DecisionReject are operator choices; DecisionExpired is // the fail-closed outcome the queue applies when neither arrives in time. DecisionApprove ApprovalDecision = "approve" DecisionReject ApprovalDecision = "reject" DecisionExpired ApprovalDecision = "expired" )
type ApprovalNotPendingError ¶
type ApprovalNotPendingError struct {
ID string
}
ApprovalNotPendingError identifies an approval that cannot be resolved by the in-memory queue. The durable approval record owns its terminal history.
func (*ApprovalNotPendingError) Error ¶
func (err *ApprovalNotPendingError) Error() string
Error implements error.
type ApprovalQueue ¶
type ApprovalQueue struct {
OnRequested func(approval Approval) error
OnResolved func(resolution ApprovalResolution) error
// contains filtered or unexported fields
}
ApprovalQueue serializes pending approvals for one Controller. It is safe for concurrent use. OnRequested and OnResolved are durability hooks the containing core installs; a hook returning an error fails the approval closed rather than admitting work that was not recorded.
func NewApprovalQueue ¶
func NewApprovalQueue(timeout time.Duration) *ApprovalQueue
NewApprovalQueue constructs an empty queue whose requests expire after timeout.
func (*ApprovalQueue) ExpireRun ¶
func (q *ApprovalQueue) ExpireRun(runID, actor string) error
ExpireRun fails closed every pending permission associated with one run. It is used after aborting an agent turn so no background approval waiter or actionable UI card survives the turn that requested it.
func (*ApprovalQueue) Get ¶
func (q *ApprovalQueue) Get(id string) (Approval, bool)
Get returns one pending approval without claiming or resolving it.
func (*ApprovalQueue) Pending ¶
func (q *ApprovalQueue) Pending() []Approval
Pending lists the approvals awaiting a decision, oldest request first.
func (*ApprovalQueue) Request ¶
func (q *ApprovalQueue) Request(ctx context.Context, request ApprovalRequest) (ApprovalResolution, error)
Request blocks until the approval is resolved, expires, or ctx is done, and resolves each request exactly once.
func (*ApprovalQueue) Resolve ¶
func (q *ApprovalQueue) Resolve(id string, decision ApprovalDecision, actor string) (ApprovalResolution, error)
Resolve applies an operator decision to one pending approval. An empty actor defaults to the operator key. It reports ApprovalNotPendingError for an unknown approval and ErrApprovalRecording or ErrApprovalResolving while another writer holds the same one.
type ApprovalRequest ¶
type ApprovalRequest struct {
RunID string
ToolCallID string
Kind string
Title string
Detail string
Risk string
Payload any
RequiresFleetQuorum bool
}
ApprovalRequest is what a harness or tool asks the operator to authorize. Payload is hashed, not stored, so the queue never retains tool arguments.
type ApprovalResolution ¶
type ApprovalResolution struct {
Approval Approval
Decision ApprovalDecision
Actor string
At time.Time
}
ApprovalResolution is one approval's terminal outcome together with the actor that caused it; Actor is one of the keys ApprovalActorState names.
type Controller ¶
type Controller struct {
OnRunStarted func(event RunStarted) error
OnRunStatus func(runID, status string, at time.Time)
OnApprovalRequested func(approval Approval) error
OnApprovalResolved func(resolution ApprovalResolution) error
// contains filtered or unexported fields
}
Controller owns one agent harness session: run admission and fencing, approvals, the operator timeline, and change notification. The exported On* fields are durability hooks the containing core installs before Start; a hook returning an error fails the corresponding operation closed.
func NewController ¶
func NewController(cfg *deployv1.CoreConfig, fleetClient *fleet.Client, reconciler IReconciler) (*Controller, error)
NewController constructs the operator lifecycle for one configured harness and optional reconciler.
func NewControllerWithHarness ¶
func NewControllerWithHarness( cfg *deployv1.CoreConfig, fleetClient *fleet.Client, reconciler IReconciler, factory harnesssdk.IFactory, ) (*Controller, error)
NewControllerWithHarness constructs the operator lifecycle with an optional compiled-in harness. A nil factory preserves configuration-selected built-ins.
func (*Controller) Abort ¶
func (c *Controller) Abort(ctx context.Context) error
Abort stops the active run, whichever it is. Use AbortRun to stop only the run a caller has actually rendered.
func (*Controller) AbortRun ¶
func (c *Controller) AbortRun(ctx context.Context, expectedSessionID, expectedRunID string) error
AbortRun stops only the exact execution rendered by the caller.
func (*Controller) ApprovalQueue ¶
func (c *Controller) ApprovalQueue() *ApprovalQueue
ApprovalQueue exposes the queue so the transport can list and resolve pending approvals.
func (*Controller) Close ¶
func (c *Controller) Close() error
Close cancels the lifecycle, drops approved reconciliations, and closes the harness. It is safe to call on a Controller that never started.
func (*Controller) HarnessBackend ¶
func (c *Controller) HarnessBackend() deployv1.HarnessBackend
HarnessBackend returns the canonical configured adapter enum.
func (*Controller) HarnessBackendName ¶
func (c *Controller) HarnessBackendName() string
HarnessBackendName returns the configured adapter's environment spelling.
func (*Controller) HarnessIdentity ¶
func (c *Controller) HarnessIdentity() *deployv1.HarnessRuntimeIdentity
HarnessIdentity returns the immutable backend and model artifact selected when this controller was constructed.
func (*Controller) Run ¶
func (c *Controller) Run() RunState
Run returns a snapshot of the current run and its timeline.
func (*Controller) Send ¶
func (c *Controller) Send( ctx context.Context, prompt string, delivery deployv1.HarnessDelivery, ) (string, error)
Send admits one operator prompt, starting a new run or joining the active one according to delivery. Concurrent calls are serialized.
func (*Controller) SendToSession ¶
func (c *Controller) SendToSession( ctx context.Context, expectedSessionID, expectedRunID, prompt string, delivery deployv1.HarnessDelivery, ) (string, error)
SendToSession atomically steers the active run or starts an ordinary new run after an idle turn. Both paths are fenced to the session and run rendered by the caller, so a stale browser cannot affect newer work.
func (*Controller) SessionID ¶
func (c *Controller) SessionID() string
SessionID is the current harness session identifier, empty before Start. Callers fence their actions against it.
func (*Controller) Start ¶
func (c *Controller) Start(ctx context.Context) error
Start opens the harness session and binds this Controller's lifecycle to ctx. A harness that fails to start or activate is closed before the error is returned, so a failed Start leaves nothing running.
func (*Controller) Status ¶
func (c *Controller) Status() string
Status is the controller phase rendered for the operator.
func (*Controller) Steer ¶
func (c *Controller) Steer( ctx context.Context, expectedSessionID, expectedRunID, prompt string, delivery deployv1.HarnessDelivery, ) error
Steer delivers another message to the exact active execution without creating a run, resetting its transcript, or invoking durable run-start callbacks.
func (*Controller) Subscribe ¶
func (c *Controller) Subscribe() (<-chan struct{}, func())
Subscribe returns a coalescing change signal and the function that unsubscribes it. The channel carries no payload: a reader re-reads state through Run, Status, or the snapshot projection.
type IReconciler ¶
type IReconciler interface {
Prepare(
ctx context.Context,
intent *deployv1.ReconcileIntent,
) (*deployv1.ReconcileRevision, error)
ReconcileApproved(
ctx context.Context,
intent *deployv1.ReconcileIntent,
revision *deployv1.ReconcileRevision,
) (*deployv1.ReconcileEvidence, error)
}
IReconciler resolves and applies canonical desired state. Implementations may retain neither input: Controller passes owned snapshots and snapshots both returned messages before exposing or storing them.
type RunStarted ¶
RunStarted is the durability hook's view of a newly admitted run, delivered before the run produces any output.
type RunState ¶
type RunState struct {
ID string `json:"id"`
SessionID string `json:"session_id"`
Title string `json:"title"`
Prompt string `json:"prompt,omitempty"`
Status string `json:"status"`
StartedAt time.Time `json:"started_at"`
CanAbort bool `json:"can_abort"`
Entries []TimelineEntry `json:"entries"`
}
RunState is the complete operator-visible state of the current run, including its timeline. It is a value copy: mutating it does not affect the Controller.
type TimelineEntry ¶
type TimelineEntry struct {
ID string `json:"id"`
Kind string `json:"kind"`
Role string `json:"role,omitempty"`
Name string `json:"name,omitempty"`
Text string `json:"text,omitempty"`
Detail string `json:"detail,omitempty"`
Status string `json:"status,omitempty"`
At time.Time `json:"at"`
}
TimelineEntry is one operator-visible item in a run: a message, a tool call, or a status change. Entries are append-only within a run.