operator

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

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

Constants

This section is empty.

Variables

View Source
var (
	ErrApprovalRecording = errors.New("approval request is still being recorded")
	ErrApprovalResolving = errors.New("approval resolution is already in progress")
	ErrApprovalExpired   = errors.New("approval expired")
)
View Source
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

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

type RunStarted struct {
	RunID     string
	SessionID string
	Title     string
	Prompt    string
	At        time.Time
}

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.

Jump to

Keyboard shortcuts

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