runtime

package
v0.1.59 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Sep 14, 2026 License: Apache-2.0 Imports: 47 Imported by: 0

Documentation

Overview

Package runtime —— runtime 模块占位。

阶段 2 由对应 subagent 填充:handler.go / service.go / dto.go 等。 见 docs/13-phase1-prd.md 和 docs/14-week1-day-by-day.md。

Index

Constants

View Source
const (
	DefaultRuntimeHTTPRequestsPerSecond      = 100
	DefaultRuntimeHTTPBurst                  = 200
	DefaultRuntimeWebSocketMessagesPerSecond = 200
	DefaultRuntimeWebSocketMessageBurst      = 400
	DefaultRuntimeWebSocketsPerIdentity      = 16
)
View Source
const (
	BrowserObserverMinFrameIntervalMS = 100
	BrowserObserverMaxFrameIntervalMS = 5000
	BrowserObserverMaxFrameBytes      = 1 << 20
)
View Source
const (
	// RuntimeSchemaVersion is the schema contract compiled into this Core
	// release. A migration that changes the current runtime schema contract must
	// update both this version and RuntimeSchemaChecksum.
	RuntimeSchemaVersion       int32 = 80
	RuntimeSchemaMigrationName       = "080_runtime_attempt_transport_evidence"
	// RuntimeSchemaChecksum is SHA-256 over the canonical current schema
	// contract tuple:
	// 80:080_runtime_attempt_transport_evidence:<contract id>:<contract digest>.
	RuntimeSchemaChecksum = "48d87c8033a33eaa62ab4f47d1b09ed259ebde2a36d8918f2e03c856037af55c"

	RuntimeClusterModeNormal          RuntimeClusterMode = "normal"
	RuntimeClusterModeDraining        RuntimeClusterMode = "draining"
	RuntimeClusterModeHardMaintenance RuntimeClusterMode = "hard_maintenance"

	RuntimeClusterNewRun     RuntimeClusterOperation = "new_run"
	RuntimeClusterNewSession RuntimeClusterOperation = "new_session"
	RuntimeClusterClaim      RuntimeClusterOperation = "claim"

	// RuntimeClusterMemberLiveWindow is the canonical liveness boundary for
	// readiness and cutover retirement. A heartbeat exactly at the cutoff is live.
	RuntimeClusterMemberLiveWindow = 15 * time.Second
)
View Source
const (
	RuntimeProtocolVersion = 2
	RuntimeContractID      = "openlinker.runtime.v2"
	RuntimeContractDigest  = "4be9b2fe09eeedf0e37119075134064be88f93b301c502cdfa21a6cb978c6481"
)
View Source
const (
	RunEffectTypeAgentWebhook     = "run.agent_webhook"
	RunEffectTypeTaskCallback     = "run.task_callback"
	RunEffectTypeDefaultDelivery  = "run.default_delivery"
	RunEffectTypeParentCompletion = "run.parent_completion"
)
View Source
const (
	// MaxRuntimeMessageBytes is the contract's limit for every non-artifact
	// HTTP body and WebSocket message. The limit applies to the complete wire
	// message, not only to its business payload.
	MaxRuntimeMessageBytes int64 = 4 * 1024 * 1024

	RuntimeOfferTTLSeconds       = 30
	RuntimeLeaseTTLSeconds       = 60
	RuntimeHelloTimeoutSeconds   = 5
	RuntimeMaxPullWaitSeconds    = 30
	RuntimeMaximumNodeCapacity   = 1024
	RuntimeMaximumResumeAttempts = 1024
)
View Source
const (
	RuntimeWSCloseAuthenticationFailed   = 4401
	RuntimeWSCloseClientUpgradeRequired  = 4406
	RuntimeWSCloseSessionConflict        = 4409
	RuntimeWSCloseRequiredFeatureMissing = 4412
	RuntimeWSCloseProtocolError          = 1002
	RuntimeWSCloseInternalError          = 1011
)
View Source
const (
	RuntimeAttachmentIDHeader   = "OpenLinker-Runtime-Attachment"
	RuntimeNodeIDHeader         = "OpenLinker-Runtime-Node"
	RuntimeFallbackReasonHeader = "OpenLinker-Runtime-Fallback-Reason"
)
View Source
const (
	RuntimeTransportWebSocket RuntimeTransport = "websocket"
	RuntimeTransportLongPoll  RuntimeTransport = "long_poll"
	RuntimeTransportUnknown   RuntimeTransport = "unknown"

	RuntimeTransportReasonExplicit             RuntimeTransportReason = "explicit"
	RuntimeTransportReasonWebSocketUnavailable RuntimeTransportReason = "websocket_unavailable"
	RuntimeTransportReasonPolicyForced         RuntimeTransportReason = "policy_forced"
	RuntimeTransportReasonRecovery             RuntimeTransportReason = "recovery"

	RuntimeHeartbeatInterval = 20 * time.Second
	RuntimeSessionStaleAfter = 45 * time.Second
	RuntimePresenceTTL       = 60 * time.Second

	RuntimeTransportPolicyVersion        = 1
	RuntimeRetryMinimum                  = 250 * time.Millisecond
	RuntimeRetryMaximum                  = 15 * time.Second
	RuntimeWebSocketProbeInterval        = 15 * time.Second
	RuntimeWebSocketProbeTimeout         = 10 * time.Second
	RuntimeDefaultTransportPolicy string = "auto"

	// RuntimeTransportForbiddenSignal and RuntimePolicyChangedSignal are
	// canonical, wire-compatible message values carried with the existing
	// FORBIDDEN error code. They deliberately are not new stable error enum
	// members: current and previous SDKs can distinguish the recovery action
	// without changing the current contract bytes or digest.
	RuntimeTransportForbiddenSignal = "RUNTIME_TRANSPORT_FORBIDDEN"
	RuntimePolicyChangedSignal      = "RUNTIME_POLICY_CHANGED"
)
View Source
const BrowserObservationFeature = "browser_authenticated_observation.v1"

BrowserObservationFeature is the optional Runtime capability that gates this whole surface.

View Source
const MaxRuntimeEventPayloadBytes = 4 * 1024 * 1024

MaxRuntimeEventPayloadBytes is the Runtime max_non_artifact_message_bytes limit. It is applied to RFC 8785 canonical payload bytes, not to a transport-dependent encoded envelope.

View Source
const MaxRuntimeSignalConnections = 64
View Source
const (
	// RunFingerprintSchemaV1 makes changes to the semantic fingerprint envelope
	// explicit. Changing this value is a breaking idempotency-contract change.
	RunFingerprintSchemaV1 = "openlinker.run-create-fingerprint.v1"
)
View Source
const RuntimeBrowserExecutionProfileFeature = "browser_execution_profile.v1"
View Source
const RuntimeBrowserFullInteractionFeature = "browser_full_interaction.v1"
View Source
const RuntimeCredentialProjectionTTL = 15 * time.Minute
View Source
const (
	RuntimeCredentialRevocationWakeTopic = "runtime.credential.revoke"
)
View Source
const (
	RuntimeDelegatedRunReadFeature = "delegated_run_read.v1"
)
View Source
const (
	// RuntimeResultFingerprintVersion makes the immutable Result fingerprint
	// domain explicit. Changing it is a breaking runtime protocol change.
	RuntimeResultFingerprintVersion = "openlinker.runtime-result.v1"
)

Variables

View Source
var (
	ErrRuntimeCancellationNotFound = errors.New("Runtime Run is not owned by the requester")
	ErrRuntimeCancellationRunEnded = errors.New("Runtime Run is already terminal")
	ErrRuntimeCancellationInvalid  = errors.New("invalid runtime cancellation request")
)
View Source
var (
	ErrExternalExecutionLaunchFenceRejected = errors.New("external execution launch fence rejected")
	ErrWorkflowChildLaunchFenceRejected     = errors.New("workflow child launch fence rejected")
)
View Source
var (
	ErrInvalidRuntimeInvocation = errors.New("invalid runtime invocation capability")
	ErrExpiredRuntimeInvocation = errors.New("runtime invocation capability expired")
)
View Source
var (
	ErrRuntimeReconcilerNotConfigured = errors.New("Runtime deadline reconciler is not configured")
	ErrRuntimeReconcileBatchInvalid   = errors.New("Runtime deadline reconcile batch must be between 1 and 1000")
)
View Source
var (
	ErrRuntimeSignalBusClosed      = errors.New("runtime signal bus is closed")
	ErrRuntimeSignalBusUnavailable = errors.New("runtime signal bus is unavailable")
	ErrRuntimeSignalInvalid        = errors.New("runtime signal is invalid")
)
View Source
var (
	// ErrInvalidRuntimeEvent is returned before touching PostgreSQL when the
	// request cannot satisfy the Runtime event contract.
	ErrInvalidRuntimeEvent = errors.New("invalid runtime event")
)
View Source
var ErrObservationAlreadyActive = errors.New("browser observation is already active for this Run")

ErrObservationAlreadyActive is returned when a Run is already being observed. The Runtime lease is singular, so a second observation cannot exist.

View Source
var ErrObservationBusy = errors.New("browser observation capacity is exhausted on this Core instance")

ErrObservationBusy is returned when this instance is already at its concurrent-observation ceiling.

View Source
var ErrObservationChannelUnavailable = errors.New("browser observation channel is unavailable on this Core instance")

ErrObservationChannelUnavailable is returned when this Core process does not hold the Worker's WebSocket. Frames and wakeups live in process memory, so an observation started here could never deliver anything. Failing closed with a distinguishable error keeps that from presenting as "connected but no frames ever arrive", which is the shape a multi-instance deployment would otherwise produce.

View Source
var ErrObservationForbidden = errors.New("browser observation is not permitted for this caller")

ErrObservationForbidden is returned when the caller may not observe this Run.

View Source
var ErrObservationInactive = errors.New("browser observation is not active for this Run")

ErrObservationInactive is returned when a Run has no live observation to read from. It is separate from a failure so a viewer polling an observation that just ended gets a plain answer instead of an internal error.

View Source
var ErrObservationNotConfirmed = errors.New("browser observation was not confirmed by the Worker")

ErrObservationNotConfirmed is returned when the Worker never confirmed a start. The lease is torn down before it surfaces, so the Run is left observable rather than pinned by an observation that never began.

View Source
var ErrObservationUnsupported = errors.New("browser observation is unsupported by this Runtime")

ErrObservationUnsupported is returned when the Runtime holding this Run never declared the observation feature. Reporting it distinctly keeps an old Runtime from looking like a broken one.

View Source
var ErrObservationViewerCapacity = errors.New("browser observation viewer capacity is exhausted for this Run")

ErrObservationViewerCapacity is returned when one Run already has the maximum number of concurrent frame long polls. It is distinct from ErrObservationBusy: the observation itself is healthy, and closing another viewer for this Run frees capacity.

View Source
var ErrRuntimeCredentialSessionScopeChanged = errors.New("runtime credential session scope changed")

ErrRuntimeCredentialSessionScopeChanged asks the API caller to retry after bounded lock-order retries could not capture a concurrently created Session.

View Source
var ErrRuntimeDispatchWakeReconcilerNotConfigured = errors.New("Runtime dispatch wake reconciler is not configured")
View Source
var (
	ErrRuntimeSessionReaperNotConfigured = errors.New("runtime session reaper is not configured")
)

Functions

func BuildRuntimeInvocationProof added in v0.1.56

func BuildRuntimeInvocationProof(token string, request RuntimeInvocationProofRequest) (string, error)

BuildRuntimeInvocationProof is used by a Node holding the short-lived invocation token. The token itself is never put into the proof payload.

func CanonicalizeRFC8785 added in v0.1.56

func CanonicalizeRFC8785(value any) ([]byte, error)

CanonicalizeRFC8785 serializes an already parsed I-JSON value using JSON Canonicalization Scheme rules. It accepts nil, booleans, strings, finite IEEE-754 numbers, string-keyed maps, arrays/slices, pointers, interfaces, and json.Number. Structs and custom marshalers are rejected so hidden toJSON-like behavior cannot change an idempotency fingerprint.

func DecodeRuntimeBody added in v0.1.56

func DecodeRuntimeBody[P any](reader io.Reader) (P, error)

DecodeRuntimeBody decodes one strict HTTP Runtime request body.

func DecodeRuntimeMessagePayload added in v0.1.56

func DecodeRuntimeMessagePayload[P any](envelope RuntimeEnvelope, expected RuntimeMessageType) (P, error)

DecodeRuntimeMessagePayload strictly decodes the routed envelope payload.

func FingerprintRunCreation added in v0.1.56

func FingerprintRunCreation(input RunFingerprintInput) ([sha256.Size]byte, error)

FingerprintRunCreation computes SHA-256 over the RFC 8785 canonical JSON representation of the normalized Run creation semantics.

func HashIdempotencyKey added in v0.1.56

func HashIdempotencyKey(key string) ([sha256.Size]byte, error)

HashIdempotencyKey validates the wire value and returns the only form that may be persisted. It does not trim or otherwise normalize the key.

func IsRuntimeEventError added in v0.1.56

func IsRuntimeEventError(err error, code RuntimeEventErrorCode) bool

IsRuntimeEventError reports whether err has the requested stable code.

func IsRuntimeLeaseError added in v0.1.56

func IsRuntimeLeaseError(err error, code RuntimeLeaseErrorCode) bool

func IsRuntimeResultError added in v0.1.56

func IsRuntimeResultError(err error, code RuntimeResultErrorCode) bool

IsRuntimeResultError reports whether err carries the requested stable code.

func IsRuntimeSessionError added in v0.1.56

func IsRuntimeSessionError(err error, code RuntimeSessionErrorCode) bool

func MarshalRuntimeSignal added in v0.1.56

func MarshalRuntimeSignal(signal RuntimeSignal) ([]byte, error)

func RequireRuntimeClusterOperation added in v0.1.56

func RequireRuntimeClusterOperation(
	ctx context.Context,
	querier interface {
		QueryRow(context.Context, string, ...any) pgx.Row
	},
	operation RuntimeClusterOperation,
) error

RequireRuntimeClusterOperation takes a row lock that conflicts with a control-mode update. This makes the mode transition the linearization point for new Run inserts, new Session inserts, and claims.

func RevokeAgentCredential added in v0.1.56

func RevokeAgentCredential(
	ctx context.Context,
	pool *pgxpool.Pool,
	creatorID uuid.UUID,
	credentialID uuid.UUID,
) (bool, error)

RevokeAgentCredential closes every Session owned by one Agent Token before revoking it. The transaction follows the global principal lock order: Session -> Node -> Token -> Attachment. Runtime Nodes are only locked for serialization; they remain available to Sessions using other credentials.

func RuntimeEventFingerprint added in v0.1.56

func RuntimeEventFingerprint(request RuntimeEventRequest) ([]byte, error)

RuntimeEventFingerprint computes the Core-owned payload fingerprint. The stable client event ID and transport envelope are deliberately excluded so callers cannot choose or spoof the digest domain.

func RuntimeHTTPStatus added in v0.1.56

func RuntimeHTTPStatus(code RuntimeErrorCode) int

RuntimeHTTPStatus maps every stable runtime error to a deterministic status.

func RuntimeRequiredFeatures added in v0.1.56

func RuntimeRequiredFeatures() []string

RuntimeRequiredFeatures returns a copy so callers cannot mutate the handshake requirements for the running process.

func RuntimeResultFingerprint added in v0.1.56

func RuntimeResultFingerprint(request RuntimeResultRequest) ([]byte, error)

RuntimeResultFingerprint computes the Core-owned RFC 8785 SHA-256 digest. ResultID, Attempt identity, and the transport envelope are excluded; those identities are compared separately under the Run lock.

func RuntimeWebSocketCloseCode added in v0.1.56

func RuntimeWebSocketCloseCode(code RuntimeErrorCode) (int, bool)

RuntimeWebSocketCloseCode returns a close code only for connection-fatal errors. Durable Run errors are sent as runtime.error and keep the session alive, so their second return value is false.

func StartRunEffectWorker added in v0.1.56

func StartRunEffectWorker(
	ctx context.Context,
	svc *Service,
	cfg RunEffectWorkerConfig,
)

func StartRunEffectWorkerWithWake added in v0.1.56

func StartRunEffectWorkerWithWake(
	ctx context.Context,
	svc *Service,
	cfg RunEffectWorkerConfig,
	source eventwake.TopicSource,
)

StartRunEffectWorkerWithWake preserves the legacy worker as the degraded path. While LISTEN is healthy, a topic notification, the earliest durable due timestamp, or the 60-second reconciliation pass is the only reason to query/claim Effects.

func StartRuntimeCredentialRevocationWake added in v0.1.56

func StartRuntimeCredentialRevocationWake(
	ctx context.Context,
	source eventwake.TopicSource,
	hub *RuntimeWakeHub,
)

StartRuntimeCredentialRevocationWake turns the credential-specific transactional NOTIFY into a bounded database-fallback hint. The notification carries no credential material and never makes an authorization decision.

func StartRuntimeDispatchWakeReconciler added in v0.1.56

func StartRuntimeDispatchWakeReconciler(
	ctx context.Context,
	reconciler *RuntimeDispatchWakeReconciler,
	cfg RuntimeDispatchWakeReconcilerConfig,
)

func StartRuntimeMaintenanceWorker added in v0.1.56

func StartRuntimeMaintenanceWorker(
	ctx context.Context,
	reconciler runtimeDeadlineReconcileWorker,
	cancellations runtimeCancellationReapWorker,
	sessions runtimeSessionReapWorker,
	cfg RuntimeMaintenanceWorkerConfig,
)

StartRuntimeMaintenanceWorker runs an immediate pass and then continues until shutdown. PostgreSQL remains the truth source; failures are logged and retried on the next tick instead of terminating the API process.

func StartRuntimeMaintenanceWorkerWithWake added in v0.1.56

func StartRuntimeMaintenanceWorkerWithWake(
	ctx context.Context,
	reconciler runtimeDeadlineReconcileWorker,
	cancellations runtimeCancellationReapWorker,
	sessions runtimeSessionReapWorker,
	cfg RuntimeMaintenanceWorkerConfig,
	source eventwake.TopicSource,
)

StartRuntimeMaintenanceWorkerWithWake keeps the legacy worker available for tests and degraded deployments, while replacing healthy PostgreSQL deadline and cancellation polling with transactional Run wakes, exact due timers and one low-frequency reconciliation. Session expiry retains its existing cadence because healthy Redis lease checks do not touch PostgreSQL and its externally visible offline convergence must not be extended.

func StartRuntimeSignalOutboxWorker added in v0.1.56

func StartRuntimeSignalOutboxWorker(
	ctx context.Context,
	worker *RuntimeSignalOutboxWorker,
	cfg RuntimeSignalOutboxWorkerConfig,
)

func StartRuntimeSignalOutboxWorkerWithWake added in v0.1.56

func StartRuntimeSignalOutboxWorkerWithWake(
	ctx context.Context,
	worker *RuntimeSignalOutboxWorker,
	cfg RuntimeSignalOutboxWorkerConfig,
	source eventwake.TopicSource,
)

StartRuntimeSignalOutboxWorkerWithWake keeps the legacy polling entry point intact and cuts over only while the advisory LISTEN source is healthy. A disconnect immediately returns to the legacy interval; durable claims and database-clock retry timestamps remain authoritative in both modes.

func StartRuntimeSignalSubscriber added in v0.1.56

func StartRuntimeSignalSubscriber(
	ctx context.Context,
	bus RuntimeSignalBus,
	instanceID uuid.UUID,
	hub *RuntimeWakeHub,
	service *Service,
)

StartRuntimeSignalSubscriber supervises the blocking subscription and reconnects with bounded backoff. It deliberately does not expose a readiness bit: the existing signal-bus Health check fails HA readiness, while PostgreSQL polling and reconciliation continue to converge state.

func ValidateRuntimeEnvelope added in v0.1.56

func ValidateRuntimeEnvelope(envelope RuntimeEnvelope) error

ValidateRuntimeEnvelope enforces the version, contract, identity, message set, payload-object, and reply-presence rules in core-runtime.json.

func ValidateRuntimePayload added in v0.1.56

func ValidateRuntimePayload(payload any) error

ValidateRuntimePayload applies the semantic constraints that encoding/json cannot express. Required-field presence and unknown fields are handled by the strict decoders above.

func ValidateRuntimeReplyCorrelation added in v0.1.56

func ValidateRuntimeReplyCorrelation(request, reply RuntimeEnvelope) error

ValidateRuntimeReplyCorrelation proves that a business ACK/error belongs to a concrete request. Socket write order is not accepted as correlation.

func ValidateRuntimeSignal added in v0.1.56

func ValidateRuntimeSignal(signal RuntimeSignal) error

func VerifyRuntimeInvocationProof added in v0.1.56

func VerifyRuntimeInvocationProof(token, proof string, request RuntimeInvocationProofRequest) error

Types

type AgentA2AContext

type AgentA2AContext struct {
	CurrentRunID        string   `json:"current_run_id"`
	MessageID           string   `json:"message_id,omitempty"`
	Protocol            string   `json:"protocol,omitempty"`
	Method              string   `json:"method,omitempty"`
	ParentRunID         string   `json:"parent_run_id,omitempty"`
	CallerAgentID       string   `json:"caller_agent_id,omitempty"`
	ProtocolContextID   string   `json:"protocol_context_id,omitempty"`
	ProtocolTaskID      string   `json:"protocol_task_id,omitempty"`
	RootContextID       string   `json:"root_context_id,omitempty"`
	ParentContextID     string   `json:"parent_context_id,omitempty"`
	ParentTaskID        string   `json:"parent_task_id,omitempty"`
	TraceID             string   `json:"trace_id,omitempty"`
	ReferenceTaskIDs    []string `json:"reference_task_ids,omitempty"`
	AcceptedOutputModes []string `json:"accepted_output_modes,omitempty"`
	Extensions          []string `json:"extensions,omitempty"`
}

AgentA2AContext tells an Agent how to delegate from its current run without making a human copy/paste a parent run id from the UI.

type AgentError

type AgentError struct {
	Code    string `json:"code"`
	Message string `json:"message"`
}

AgentError 创作者侧错误(业务级 4xx 等场景)。

type AgentEvent

type AgentEvent struct {
	EventType string                 `json:"event_type"`
	Payload   map[string]interface{} `json:"payload"`
}

AgentEvent 是 Agent endpoint 可选返回的运行中事件。

Phase 2 先允许少量 OpenLinker-native 事件;后续可映射到 A2A Message / Part / Artifact。

type AgentRequest

type AgentRequest struct {
	Input         map[string]interface{} `json:"input"`
	Metadata      map[string]interface{} `json:"metadata,omitempty"`
	RunID         string                 `json:"run_id"`
	ParentRunID   string                 `json:"parent_run_id,omitempty"`
	CallerAgentID string                 `json:"caller_agent_id,omitempty"`
	A2A           *AgentA2AContext       `json:"a2a,omitempty"`
	Conversation  *ConversationContext   `json:"conversation,omitempty"`
}

AgentRequest 平台 → 创作者 endpoint 的请求体。

RunID 让创作者侧可对账 / 排查;Metadata 是经过第三方出站投影的公开上下文。

type AgentResponse

type AgentResponse struct {
	Output   map[string]interface{} `json:"output"`
	Events   []AgentEvent           `json:"events,omitempty"`
	CostUSD  *float64               `json:"cost_usd,omitempty"`
	Metadata map[string]interface{} `json:"metadata,omitempty"`
	Error    *AgentError            `json:"error,omitempty"`
}

AgentResponse 创作者 endpoint → 平台的响应体。

Output 业务结果(成功时必填);Error 业务错误(失败时必填); CostUSD 创作者透明告知本次调用真实成本(可选,目前平台不做核对,仅记录用)。

type AttemptIdentity added in v0.1.56

type AttemptIdentity struct {
	RunID            uuid.UUID `json:"run_id" runtime:"required"`
	AttemptID        uuid.UUID `json:"attempt_id" runtime:"required"`
	LeaseID          uuid.UUID `json:"lease_id" runtime:"required"`
	FencingToken     int64     `json:"fencing_token" runtime:"required"`
	NodeID           uuid.UUID `json:"node_id" runtime:"required"`
	AgentID          uuid.UUID `json:"agent_id" runtime:"required"`
	WorkerID         string    `json:"worker_id" runtime:"required"`
	RuntimeSessionID uuid.UUID `json:"runtime_session_id" runtime:"required"`
}

AttemptIdentity is the exact runtime wire identity. Unlike the transport-neutral RuntimeAttemptIdentity, every Node/session field is mandatory here.

func (AttemptIdentity) RuntimeIdentity added in v0.1.56

func (i AttemptIdentity) RuntimeIdentity() RuntimeAttemptIdentity

type AuthenticatedRuntimePrincipal added in v0.1.56

type AuthenticatedRuntimePrincipal struct {
	AgentID      uuid.UUID             `json:"agent_id"`
	CredentialID uuid.UUID             `json:"credential_id"`
	Device       RuntimeDeviceIdentity `json:"device"`
}

AuthenticatedRuntimePrincipal combines the independently authenticated Agent Token and Node device identities. The request body may repeat these identifiers for protocol clarity, but it never establishes them.

type BrowserHumanControl added in v0.1.56

type BrowserHumanControl struct {
	// contains filtered or unexported fields
}

BrowserHumanControl persists only low-frequency ownership transitions. Frames, input and their wakeups remain bounded process memory.

func NewBrowserHumanControl added in v0.1.56

func NewBrowserHumanControl(pool *pgxpool.Pool) *BrowserHumanControl

func (*BrowserHumanControl) BindCommandSender added in v0.1.56

func (control *BrowserHumanControl) BindCommandSender(
	sender browserViewerCommandSender,
)

func (*BrowserHumanControl) Claim added in v0.1.56

func (control *BrowserHumanControl) Claim(
	ctx context.Context,
	userID uuid.UUID,
	runID uuid.UUID,
) (BrowserHumanControlState, error)

func (*BrowserHumanControl) Input added in v0.1.56

func (control *BrowserHumanControl) Input(
	ctx context.Context,
	userID uuid.UUID,
	runID uuid.UUID,
	input BrowserViewerInputPayload,
) error

func (*BrowserHumanControl) PauseFromEvent added in v0.1.56

func (control *BrowserHumanControl) PauseFromEvent(
	ctx context.Context,
	identity RuntimeAttemptIdentity,
	payload map[string]any,
) error

func (*BrowserHumanControl) PublishFrame added in v0.1.56

func (control *BrowserHumanControl) PublishFrame(
	frame BrowserViewerFramePayload,
) error

func (*BrowserHumanControl) Release added in v0.1.56

func (control *BrowserHumanControl) Release(
	ctx context.Context,
	userID uuid.UUID,
	runID uuid.UUID,
) (BrowserHumanControlState, error)

func (*BrowserHumanControl) Resume added in v0.1.56

func (control *BrowserHumanControl) Resume(
	ctx context.Context,
	userID uuid.UUID,
	runID uuid.UUID,
) (BrowserHumanControlState, error)

func (*BrowserHumanControl) RunGC added in v0.1.56

func (control *BrowserHumanControl) RunGC(ctx context.Context)

func (*BrowserHumanControl) State added in v0.1.56

func (control *BrowserHumanControl) State(
	ctx context.Context,
	userID uuid.UUID,
	runID uuid.UUID,
) (BrowserHumanControlState, error)

func (*BrowserHumanControl) WaitFrame added in v0.1.56

func (control *BrowserHumanControl) WaitFrame(
	ctx context.Context,
	userID uuid.UUID,
	runID uuid.UUID,
	after uint64,
) (*BrowserViewerFramePayload, error)

type BrowserHumanControlState added in v0.1.56

type BrowserHumanControlState struct {
	RunID            uuid.UUID  `json:"run_id"`
	AttemptID        uuid.UUID  `json:"attempt_id"`
	RuntimeSessionID uuid.UUID  `json:"runtime_session_id"`
	BrowserSessionID uuid.UUID  `json:"browser_session_id"`
	SessionEpoch     uint64     `json:"session_epoch"`
	AttachmentID     uuid.UUID  `json:"attachment_id"`
	ControlEpoch     uint64     `json:"control_epoch"`
	Controller       string     `json:"controller"`
	State            string     `json:"state"`
	PauseReason      string     `json:"pause_reason"`
	PauseExpiresAt   time.Time  `json:"pause_expires_at"`
	HumanExpiresAt   *time.Time `json:"human_expires_at,omitempty"`
	ClaimedAt        *time.Time `json:"claimed_at,omitempty"`
	HumanDurationMS  int64      `json:"human_duration_ms"`
	UpdatedAt        time.Time  `json:"updated_at"`

	UserID       uuid.UUID `json:"-"`
	AgentID      uuid.UUID `json:"-"`
	LeaseID      uuid.UUID `json:"-"`
	FencingToken int64     `json:"-"`
	NodeID       uuid.UUID `json:"-"`
	WorkerID     string    `json:"-"`
	RunDeadline  time.Time `json:"-"`
}

type BrowserObservation added in v0.1.56

type BrowserObservation struct {
	// contains filtered or unexported fields
}

func NewBrowserObservation added in v0.1.56

func NewBrowserObservation(
	pool *pgxpool.Pool,
	now func() time.Time,
	instance uuid.UUID,
	quota int,
) *BrowserObservation

func (*BrowserObservation) AuthorizeOwner added in v0.1.56

func (observation *BrowserObservation) AuthorizeOwner(
	ctx context.Context,
	runID uuid.UUID,
	callerUserID uuid.UUID,
) error

AuthorizeOwner rejects a caller who does not own the Run. Every observation endpoint calls it, including the read-only ones: state reveals that a Run is being watched and the frame endpoint returns the page itself, so neither may rely on start having been authorized earlier.

func (*BrowserObservation) BindCommandSender added in v0.1.56

func (observation *BrowserObservation) BindCommandSender(
	sender browserObserverCommandSender,
)

func (*BrowserObservation) BindInstance added in v0.1.56

func (observation *BrowserObservation) BindInstance(instance uuid.UUID)

BindInstance records which Core process owns observations started here. It is separate from construction because the process identity is only known once cluster membership is configured.

func (*BrowserObservation) CloseSessionObservations added in v0.1.56

func (observation *BrowserObservation) CloseSessionObservations(
	ctx context.Context,
	runtimeSessionID uuid.UUID,
	endReason string,
) error

CloseSessionObservations ends every observation bound to a Runtime Session. A disconnected Worker cannot answer a stop command and will drop its lease on its own TTL, so the audit has to be closed from this side or it stays active and its unique index blocks the Run from ever being observed again.

func (*BrowserObservation) HandleEvent added in v0.1.56

HandleEvent consumes one Worker event. It returns the ack the Worker is waiting on: the window is a single unacknowledged event, so failing to ack stops the stream rather than dropping one frame.

func (*BrowserObservation) ProjectFromEvent added in v0.1.56

func (observation *BrowserObservation) ProjectFromEvent(
	ctx context.Context,
	identity RuntimeAttemptIdentity,
	payload map[string]any,
	eventSeq int64,
) error

ProjectFromEvent records the live Browser identity of a running Run.

Observation cannot read browser_run_controls: that row is written only when a challenge pauses the Run for takeover, so a normally executing Run has none -- and "watch while the Agent keeps working" is the entire point of this surface. The ready lifecycle event carries the same identity fields the pause event does, so the observable identity is projected from there instead.

func (*BrowserObservation) ReconcileAbandoned added in v0.1.56

func (observation *BrowserObservation) ReconcileAbandoned(ctx context.Context)

ReconcileAbandoned ends observations no viewer is reading any more. The browser sends a stop when it can, but a crashed tab or a lost network sends nothing, so the absence of polling is what Core actually goes on. Without this such an observation holds its lease for the whole TTL and the Run cannot be observed by anyone else in the meantime.

func (*BrowserObservation) ReconcileExpired added in v0.1.56

func (observation *BrowserObservation) ReconcileExpired(ctx context.Context) error

ReconcileExpired closes observations whose lease has passed while they were still marked active. A Core that exits mid-observation cannot run its own teardown, so without this the record would stay active forever and an auditor could not tell "still watching" from "Core crashed".

func (*BrowserObservation) ReconcileForeignExpired added in v0.1.56

func (observation *BrowserObservation) ReconcileForeignExpired(ctx context.Context) error

ReconcileForeignExpired closes expired audits left behind by an instance that is no longer running. It is separate from ReconcileExpired because there is no in-process state to release, and because an audit must not be closed out from under an instance that is still serving it.

func (*BrowserObservation) RecordFrames added in v0.1.56

func (observation *BrowserObservation) RecordFrames(
	ctx context.Context,
	leaseID uuid.UUID,
	count int64,
) error

RecordFrames persists a lower bound for the frame count while an observation is running. It is deliberately not called per frame: the point is that a Core which exits mid-observation leaves a bound behind, not that the running count is exact.

func (*BrowserObservation) ResolveIdentity added in v0.1.56

func (observation *BrowserObservation) ResolveIdentity(
	ctx context.Context,
	runID uuid.UUID,
	callerUserID uuid.UUID,
	isAdmin bool,
) (BrowserObserverIdentity, error)

ResolveIdentity finds the Attempt to observe and authorizes the caller.

An owner may only observe their own Run. An admin may observe any Run, but arrives through a separate route with its own permission and a recorded reason, so the two never collapse into one check.

func (*BrowserObservation) RunGC added in v0.1.56

func (observation *BrowserObservation) RunGC(ctx context.Context)

RunGC reconciles on startup and then periodically, so a crash that happened while this process was down is repaired as soon as one comes back.

func (*BrowserObservation) Start added in v0.1.56

func (observation *BrowserObservation) Start(
	ctx context.Context,
	runID uuid.UUID,
	observerUserID uuid.UUID,
	isAdmin bool,
	reason string,
	identity BrowserObserverIdentity,
) (BrowserObservationState, error)

Start opens an observation. The channel check happens before the audit record is written: creating an active record for an observation that can never deliver a frame would leave a dangling row that reconciliation later has to explain away.

func (*BrowserObservation) State added in v0.1.56

func (observation *BrowserObservation) State(
	ctx context.Context,
	runID uuid.UUID,
) (BrowserObservationState, error)

State reports the live observation for a Run, if any.

func (*BrowserObservation) Stop added in v0.1.56

func (observation *BrowserObservation) Stop(
	ctx context.Context,
	runID uuid.UUID,
	endReason string,
) error

Stop ends an observation and closes its audit record. It is safe to call for an observation that has already ended, because every teardown path -- explicit stop, Run terminal, disconnect, TTL -- has to be able to call it.

func (*BrowserObservation) WaitFrame added in v0.1.56

func (observation *BrowserObservation) WaitFrame(
	ctx context.Context,
	runID uuid.UUID,
	after int64,
) (*BrowserObservationFrame, error)

WaitFrame serves one long poll. A nil frame with no error means the poll timed out with nothing new, which is a normal empty response rather than a failure.

type BrowserObservationFrame added in v0.1.56

type BrowserObservationFrame struct {
	FrameSeq   int64     `json:"frame_seq"`
	CapturedAt time.Time `json:"captured_at"`
	MIMEType   string    `json:"mime_type"`
	Data       []byte    `json:"data"`
	Width      int       `json:"width"`
	Height     int       `json:"height"`
}

BrowserObservationFrame is the only shape the browser ever receives. It deliberately carries no Runtime, Node, Attachment or Session identity: those are internal invariants between the Worker and Core, and a viewer has no use for them.

type BrowserObservationState added in v0.1.56

type BrowserObservationState struct {
	RunID              uuid.UUID  `json:"run_id"`
	Active             bool       `json:"active"`
	LeaseID            uuid.UUID  `json:"lease_id,omitempty"`
	LeaseExpiresAt     *time.Time `json:"lease_expires_at,omitempty"`
	FrameCount         int64      `json:"frame_count"`
	FrameCountComplete bool       `json:"frame_count_complete"`
}

type BrowserObserverAction added in v0.1.56

type BrowserObserverAction string
const (
	BrowserObserverStart BrowserObserverAction = "start"
	BrowserObserverStop  BrowserObserverAction = "stop"
)

type BrowserObserverCommandPayload added in v0.1.56

type BrowserObserverCommandPayload struct {
	AttemptIdentity      AttemptIdentity       `json:"attempt_identity"`
	SessionEpoch         int64                 `json:"session_epoch"`
	BrowserSessionSHA256 string                `json:"browser_session_sha256"`
	AttachmentSHA256     string                `json:"browser_attachment_sha256"`
	CommandID            uuid.UUID             `json:"command_id"`
	Action               BrowserObserverAction `json:"action"`
	LeaseID              uuid.UUID             `json:"lease_id"`
	LeaseExpiresAt       time.Time             `json:"lease_expires_at"`
	DeadlineAt           time.Time             `json:"deadline_at"`
	FrameIntervalMS      int                   `json:"frame_interval_ms"`
}

func (BrowserObserverCommandPayload) Validate added in v0.1.56

func (payload BrowserObserverCommandPayload) Validate() error

type BrowserObserverEventAckPayload added in v0.1.56

type BrowserObserverEventAckPayload struct {
	AttemptIdentity      AttemptIdentity `json:"attempt_identity"`
	SessionEpoch         int64           `json:"session_epoch"`
	BrowserSessionSHA256 string          `json:"browser_session_sha256"`
	AttachmentSHA256     string          `json:"browser_attachment_sha256"`
	LeaseID              uuid.UUID       `json:"lease_id"`
	EventSeq             int64           `json:"event_seq"`
}

func (BrowserObserverEventAckPayload) Validate added in v0.1.56

func (payload BrowserObserverEventAckPayload) Validate() error

type BrowserObserverEventKind added in v0.1.56

type BrowserObserverEventKind string
const (
	BrowserObserverStarted BrowserObserverEventKind = "started"
	BrowserObserverFrame   BrowserObserverEventKind = "frame"
	BrowserObserverStopped BrowserObserverEventKind = "stopped"
	BrowserObserverError   BrowserObserverEventKind = "error"
)

type BrowserObserverEventPayload added in v0.1.56

type BrowserObserverEventPayload struct {
	AttemptIdentity      AttemptIdentity              `json:"attempt_identity"`
	SessionEpoch         int64                        `json:"session_epoch"`
	BrowserSessionSHA256 string                       `json:"browser_session_sha256"`
	AttachmentSHA256     string                       `json:"browser_attachment_sha256"`
	CommandID            uuid.UUID                    `json:"command_id"`
	LeaseID              uuid.UUID                    `json:"lease_id"`
	EventSeq             int64                        `json:"event_seq"`
	Kind                 BrowserObserverEventKind     `json:"kind"`
	CapturedAt           *time.Time                   `json:"captured_at,omitempty"`
	Frame                *BrowserObserverFramePayload `json:"frame,omitempty"`
	ErrorCode            string                       `json:"error_code,omitempty"`
}

func (BrowserObserverEventPayload) Validate added in v0.1.56

func (payload BrowserObserverEventPayload) Validate() error

type BrowserObserverFramePayload added in v0.1.56

type BrowserObserverFramePayload struct {
	MIMEType string `json:"mime_type"`
	Data     []byte `json:"data"`
	Width    int    `json:"width"`
	Height   int    `json:"height"`
}

type BrowserObserverIdentity added in v0.1.56

type BrowserObserverIdentity struct {
	RunID                uuid.UUID `json:"run_id"`
	AttemptID            uuid.UUID `json:"attempt_id"`
	RuntimeLeaseID       uuid.UUID `json:"runtime_lease_id"`
	FencingToken         int64     `json:"fencing_token"`
	NodeID               uuid.UUID `json:"node_id"`
	AgentID              uuid.UUID `json:"agent_id"`
	WorkerID             string    `json:"worker_id"`
	SessionEpoch         int64     `json:"session_epoch"`
	BrowserSessionSHA256 string    `json:"browser_session_sha256"`
	AttachmentSHA256     string    `json:"browser_attachment_sha256"`
	RuntimeSessionID     uuid.UUID `json:"runtime_session_id"`
}

BrowserObserverIdentity carries hashed Browser identity. The ready lifecycle event publishes browser_session_sha256 and browser_attachment_sha256 and never the underlying UUIDs, so Core cannot send raw IDs here without being told something it is deliberately not told. The Worker rehashes its own identity to verify, which is no weaker.

func (BrowserObserverIdentity) RuntimeIdentity added in v0.1.56

func (identity BrowserObserverIdentity) RuntimeIdentity() AttemptIdentity

type BrowserViewerAction added in v0.1.56

type BrowserViewerAction string
const (
	BrowserViewerActionClaim     BrowserViewerAction = "claim"
	BrowserViewerActionRelease   BrowserViewerAction = "release"
	BrowserViewerActionResume    BrowserViewerAction = "resume"
	BrowserViewerActionTerminate BrowserViewerAction = "terminate"
	BrowserViewerActionInput     BrowserViewerAction = "input"
)

type BrowserViewerCommand added in v0.1.56

type BrowserViewerCommand = RuntimeCommand[BrowserViewerCommandPayload]

type BrowserViewerCommandMessage added in v0.1.56

type BrowserViewerCommandMessage = RuntimeTypedEnvelope[BrowserViewerCommandPayload]

type BrowserViewerCommandPayload added in v0.1.56

type BrowserViewerCommandPayload struct {
	AttemptIdentity      AttemptIdentity            `json:"attempt_identity" runtime:"required"`
	Action               BrowserViewerAction        `json:"action" runtime:"required"`
	BrowserSessionID     uuid.UUID                  `json:"browser_session_id" runtime:"required"`
	SessionEpoch         uint64                     `json:"session_epoch" runtime:"required"`
	AttachmentID         uuid.UUID                  `json:"attachment_id" runtime:"required"`
	PreviousControlEpoch uint64                     `json:"previous_control_epoch" runtime:"required"`
	ControlEpoch         uint64                     `json:"control_epoch" runtime:"required"`
	Input                *BrowserViewerInputPayload `json:"input,omitempty"`
	DeadlineAt           time.Time                  `json:"deadline_at" runtime:"required"`
}

type BrowserViewerFrameAckMessage added in v0.1.56

type BrowserViewerFrameAckMessage = RuntimeTypedEnvelope[BrowserViewerFrameAckPayload]

type BrowserViewerFrameAckPayload added in v0.1.56

type BrowserViewerFrameAckPayload struct {
	AttemptIdentity AttemptIdentity `json:"attempt_identity" runtime:"required"`
	ControlEpoch    uint64          `json:"control_epoch" runtime:"required"`
	FrameSeq        uint64          `json:"frame_seq" runtime:"required"`
}

type BrowserViewerFrameMessage added in v0.1.56

type BrowserViewerFrameMessage = RuntimeTypedEnvelope[BrowserViewerFramePayload]

type BrowserViewerFramePayload added in v0.1.56

type BrowserViewerFramePayload struct {
	AttemptIdentity  AttemptIdentity `json:"attempt_identity" runtime:"required"`
	BrowserSessionID uuid.UUID       `json:"browser_session_id" runtime:"required"`
	SessionEpoch     uint64          `json:"session_epoch" runtime:"required"`
	AttachmentID     uuid.UUID       `json:"attachment_id" runtime:"required"`
	ControlEpoch     uint64          `json:"control_epoch" runtime:"required"`
	FrameSeq         uint64          `json:"frame_seq" runtime:"required"`
	MIMEType         string          `json:"mime_type" runtime:"required"`
	Data             []byte          `json:"data" runtime:"required"`
	Width            int             `json:"width" runtime:"required"`
	Height           int             `json:"height" runtime:"required"`
}

type BrowserViewerInputPayload added in v0.1.56

type BrowserViewerInputPayload struct {
	Kind           string  `json:"kind" runtime:"required"`
	PointerAction  string  `json:"pointer_action,omitempty"`
	KeyboardAction string  `json:"keyboard_action,omitempty"`
	X              *int    `json:"x,omitempty"`
	Y              *int    `json:"y,omitempty"`
	Button         string  `json:"button,omitempty"`
	ClickCount     int     `json:"click_count,omitempty"`
	Key            string  `json:"key,omitempty"`
	Text           string  `json:"text,omitempty"`
	DeltaX         float64 `json:"delta_x,omitempty"`
	DeltaY         float64 `json:"delta_y,omitempty"`
}

type CallAgentRequest added in v0.1.56

type CallAgentRequest struct {
	TargetAgentID uuid.UUID      `json:"target_agent_id" runtime:"required"`
	Input         map[string]any `json:"input" runtime:"required"`
	Metadata      map[string]any `json:"metadata,omitempty"`
	Reason        string         `json:"reason,omitempty"`
}

type CancelCommand added in v0.1.56

type CancelCommand = RuntimeCommand[RunCancelPayload]

type ConversationContext

type ConversationContext struct {
	ID                   string                `json:"id"`
	SessionKey           string                `json:"session_key"`
	ProtocolContextID    string                `json:"protocol_context_id,omitempty"`
	RootContextID        string                `json:"root_context_id,omitempty"`
	CurrentRunID         string                `json:"current_run_id"`
	CurrentProtocolTask  string                `json:"current_protocol_task_id,omitempty"`
	HistoryBeforeCurrent []ConversationMessage `json:"history_before_current,omitempty"`
	Truncated            bool                  `json:"truncated"`
	Source               string                `json:"source"`
}

ConversationContext is the Core-owned multi-run message context for a protocol/root context. Callers send the current message; Core packages prior persisted messages so runtimes do not trust client-supplied history.

type ConversationMessage

type ConversationMessage struct {
	RunID         string                 `json:"run_id"`
	EventSequence *int32                 `json:"event_sequence,omitempty"`
	Role          string                 `json:"role"`
	Content       string                 `json:"content"`
	Payload       map[string]interface{} `json:"payload,omitempty"`
	CreatedAt     string                 `json:"created_at,omitempty"`
}

type ConversationRunItem added in v0.1.56

type ConversationRunItem struct {
	RunID                    string     `json:"run_id"`
	ParentRunID              string     `json:"parent_run_id,omitempty"`
	ConversationOrdinal      *int32     `json:"conversation_ordinal,omitempty"`
	Status                   string     `json:"status"`
	BrowserInteractionPolicy string     `json:"browser_interaction_policy,omitempty"`
	StartedAt                time.Time  `json:"started_at"`
	FinishedAt               *time.Time `json:"finished_at,omitempty"`
}

type ConversationRunListResponse added in v0.1.56

type ConversationRunListResponse struct {
	AnchorRunID                string                `json:"anchor_run_id"`
	ConversationIdentitySHA256 string                `json:"conversation_identity_sha256,omitempty"`
	Linear                     bool                  `json:"linear"`
	Revision                   string                `json:"revision"`
	Items                      []ConversationRunItem `json:"items"`
}

type DBRuntimeNodeCredentialVerifier added in v0.1.56

type DBRuntimeNodeCredentialVerifier struct {
	// contains filtered or unexported fields
}

DBRuntimeNodeCredentialVerifier resolves the enrolled runtime_nodes record by certificate serial plus SPKI SHA-256 thumbprint. Certificate validity is compared with PostgreSQL clock_timestamp(), keeping credential decisions on the same clock as Session and lease state.

func NewDBRuntimeNodeCredentialVerifier added in v0.1.56

func NewDBRuntimeNodeCredentialVerifier(pool *pgxpool.Pool) *DBRuntimeNodeCredentialVerifier

func (*DBRuntimeNodeCredentialVerifier) VerifyRuntimeNodeCredential added in v0.1.56

func (v *DBRuntimeNodeCredentialVerifier) VerifyRuntimeNodeCredential(
	ctx context.Context,
	presented RuntimePresentedCertificate,
) (RuntimeDeviceIdentity, error)

type DecodedPendingCommand added in v0.1.56

type DecodedPendingCommand struct {
	Type     RuntimeMessageType
	Cancel   *RunCancelPayload
	Drain    *RuntimeDrainPayload
	Revoke   *RunLeaseRevokedPayload
	Viewer   *BrowserViewerCommandPayload
	Observer *BrowserObserverCommandPayload
}

func DecodePendingCommand added in v0.1.56

func DecodePendingCommand(command PendingCommand) (DecodedPendingCommand, error)

DecodePendingCommand strictly decodes the command union and rejects unknown fields inside its raw payload.

type DelegatedRunReadRequest added in v0.1.59

type DelegatedRunReadRequest struct {
	RunID uuid.UUID `json:"run_id"`
}

type DelegatedRunView added in v0.1.59

type DelegatedRunView struct {
	RunSummary
	Output       json.RawMessage `json:"output,omitempty"`
	ErrorCode    string          `json:"error_code,omitempty"`
	ErrorMessage string          `json:"error_message,omitempty"`
}

DelegatedRunView intentionally excludes input, credentials, transport evidence and owner metadata. Only the direct caller may inspect this result.

type Delegation

type Delegation struct {
	ParentRunID   uuid.UUID
	CallerAgentID uuid.UUID
	Reason        string
}

Delegation describes an Agent-mediated child run executed within an active parent run.

type DeliveryRunEffectHandler added in v0.1.56

type DeliveryRunEffectHandler interface {
	AttemptDefaultDeliveryEffect(context.Context, db.RunEffectOutbox) RunEffectAttemptResult
	ResetDefaultDeliveryEffect(context.Context, db.RunEffectOutbox) error
}

DeliveryRunEffectHandler owns automatic delivery-target Effects. Manual user-triggered deliveries remain on delivery.Service's legacy queue.

type DrainCommand added in v0.1.56

type DrainCommand = RuntimeCommand[RuntimeDrainPayload]

type EventRange added in v0.1.56

type EventRange struct {
	Start int64 `json:"start"`
	End   int64 `json:"end"`
}

EventRange is an inclusive range of missing client event sequences.

type EventStore added in v0.1.56

type EventStore struct {
	// contains filtered or unexported fields
}

EventStore appends runtime execution events using PostgreSQL as the sole linearization point.

func NewEventStore added in v0.1.56

func NewEventStore(pool *pgxpool.Pool) *EventStore

func NewRuntimeEventStore added in v0.1.56

func NewRuntimeEventStore(pool *pgxpool.Pool) *EventStore

NewRuntimeEventStore is an explicit alias for call sites where EventStore would otherwise be ambiguous with domain or delivery event stores.

func (*EventStore) Append added in v0.1.56

Append is a READ COMMITTED transaction. The Run row is the per-Run linearization lock. A stored client_event_id is resolved before any active lease/deadline validation so a legal replay remains ACK-able after expiry, cancellation, completion, failover, or fence rotation.

func (*EventStore) AppendEvent added in v0.1.56

func (s *EventStore) AppendEvent(
	ctx context.Context,
	principal RuntimeEventPrincipal,
	identity RuntimeAttemptIdentity,
	request RuntimeEventRequest,
) (RuntimeEventAck, error)

AppendEvent is a naming-compatible wrapper for callers that make the event nature explicit at the call site.

func (*EventStore) MissingClientEventRanges added in v0.1.56

func (s *EventStore) MissingClientEventRanges(
	ctx context.Context,
	runID uuid.UUID,
	attemptID uuid.UUID,
	finalClientEventSeq int64,
) ([]EventRange, error)

MissingClientEventRanges returns inclusive gaps in 1..finalClientEventSeq. It scans only persisted sequences with implicit 0/N+1 sentinels and lag; it never expands the range with generate_series, so work is proportional to events actually retained rather than the claimed final sequence.

func (*EventStore) RequireCompleteClientEvents added in v0.1.56

func (s *EventStore) RequireCompleteClientEvents(
	ctx context.Context,
	runID uuid.UUID,
	attemptID uuid.UUID,
	finalClientEventSeq int64,
) error

RequireCompleteClientEvents turns exact missing ranges into the stable EVENTS_MISSING error consumed by the Result finalizer.

type ExternalExecutionLaunchFence added in v0.1.56

type ExternalExecutionLaunchFence struct {
	CallerServiceID   string
	ExternalRequestID uuid.UUID
	ActorUserID       uuid.UUID
	LaunchToken       uuid.UUID
}

type Handler

type Handler struct {
	// contains filtered or unexported fields
}

Handler 调用执行 HTTP 入口。

func NewHandler

func NewHandler(svc runtimeService, cfg ...*config.Config) *Handler

NewHandler 构造 Handler。cfg 可选(测试可省略)。

func (*Handler) ActivateRuntimeNode added in v0.1.56

func (h *Handler) ActivateRuntimeNode(c echo.Context) error

func (*Handler) CancelRun added in v0.1.41

func (h *Handler) CancelRun(c echo.Context) error

CancelRun cancels an owned, cancellable run. The concrete Service already implements this method; the narrow assertion keeps existing handler fakes source-compatible.

func (*Handler) ClaimBrowserControl added in v0.1.56

func (h *Handler) ClaimBrowserControl(c echo.Context) error

func (*Handler) DrainRuntimeNode added in v0.1.56

func (h *Handler) DrainRuntimeNode(c echo.Context) error

func (*Handler) GetAdminBrowserObservation added in v0.1.56

func (h *Handler) GetAdminBrowserObservation(c echo.Context) error

func (*Handler) GetAdminBrowserObservationFrame added in v0.1.56

func (h *Handler) GetAdminBrowserObservationFrame(c echo.Context) error

func (*Handler) GetBrowserControl added in v0.1.56

func (h *Handler) GetBrowserControl(c echo.Context) error

func (*Handler) GetBrowserControlFrame added in v0.1.56

func (h *Handler) GetBrowserControlFrame(c echo.Context) error

func (*Handler) GetBrowserObservation added in v0.1.56

func (h *Handler) GetBrowserObservation(c echo.Context) error

func (*Handler) GetBrowserObservationFrame added in v0.1.56

func (h *Handler) GetBrowserObservationFrame(c echo.Context) error

func (*Handler) GetConversationRuns added in v0.1.56

func (h *Handler) GetConversationRuns(c echo.Context) error

func (*Handler) GetRun

func (h *Handler) GetRun(c echo.Context) error

GetRun 查询单条调用详情(仅 owner)。

func (*Handler) GetRunArtifacts

func (h *Handler) GetRunArtifacts(c echo.Context) error

GetRunArtifacts 查询 run 持久化产物。只返回给 run owner。

func (*Handler) GetRunEvents

func (h *Handler) GetRunEvents(c echo.Context) error

GetRunEvents 查询 run 事件流。SSE 接口后续会复用同一 service 方法。

func (*Handler) GetRunMessages

func (h *Handler) GetRunMessages(c echo.Context) error

GetRunMessages 查询 run 的稳定消息回放。只返回给 run owner。

func (*Handler) ListRuntimeDeadLetters added in v0.1.56

func (h *Handler) ListRuntimeDeadLetters(c echo.Context) error

func (*Handler) ListRuntimeNodes added in v0.1.56

func (h *Handler) ListRuntimeNodes(c echo.Context) error

func (*Handler) PostRun

func (h *Handler) PostRun(c echo.Context) error

PostRun 调用 Agent。

Endpoint 连接模式会同步等待 Agent 返回;其他运行模式由各自的调度路径处理。 失败 / 超时 / 取消 → status='failed' or 'timeout' or 'canceled',已退款。

func (*Handler) PostRunAsync

func (h *Handler) PostRunAsync(c echo.Context) error

PostRunAsync 启动异步调用,立即返回 run_id,调用结果通过 GET /runs/:id 或 SSE 查询。

func (*Handler) RegisterAdmin added in v0.1.56

func (h *Handler) RegisterAdmin(api *echo.Group, jwtMw, adminMw echo.MiddlewareFunc)

func (*Handler) RegisterAgentRuntime

func (h *Handler) RegisterAgentRuntime(api *echo.Group)

RegisterAgentRuntime mounts the canonical Agent Runtime transport. The public API listener blocks this entire route family; only the dedicated mTLS listener dispatches it. Version compatibility is negotiated in the runtime handshake instead of being exposed in the URL.

func (*Handler) RegisterAgentRuntimeAttachOnly added in v0.1.56

func (h *Handler) RegisterAgentRuntimeAttachOnly(api *echo.Group)

RegisterAgentRuntimeAttachOnly mounts the cutover-only Session lifecycle. Normal execution routes remain owned by RegisterAgentRuntime.

func (*Handler) RegisterObservation added in v0.1.56

func (h *Handler) RegisterObservation(api *echo.Group, jwtMw echo.MiddlewareFunc)

RegisterAdmin mounts read-only runtime operational inventory. Core API owns the concrete JWT/admin middleware wiring so this package stays independent from admin policy. RegisterObservation mounts read-only observation on JWT only. It is not folded into RegisterProtected because that group runs hybrid middleware: observation is a person watching a live screen, so it must bind to a short-lived session rather than a long-lived token that can be scripted.

func (*Handler) RegisterProtected

func (h *Handler) RegisterProtected(api *echo.Group, runMw, queryMw echo.MiddlewareFunc)

RegisterProtected 注册需要鉴权的端点,分别接收 /run 与 /runs/:id 的 middleware。

POST /run            同步调用 Agent   —— runMw(JWT + User Token 混合)
POST /runs           异步启动调用     —— runMw(JWT + User Token 混合)
GET  /runs/:id       单条调用详情     —— queryMw(可按部署选择 JWT-only 或 hybrid)
GET  /runs/:id/events 调用事件流      —— queryMw(轮询)
GET  /runs/:id/artifacts 运行产物      —— queryMw
GET  /runs/:id/messages 运行消息回放    —— queryMw
GET  /runs/:id/stream 调用事件 SSE    —— queryMw
POST /runs/:id/cancel 取消运行         —— queryMw
POST /runs/:id/replay 回放死信运行      —— runMw

GET /runs 列表由 dashboard 模块(subagent-6a)提供,本模块不挂。

调用方若两条路由想共用同一个 middleware,传入相同实例即可。

func (*Handler) ReleaseBrowserControl added in v0.1.56

func (h *Handler) ReleaseBrowserControl(c echo.Context) error

func (*Handler) ReplayRun added in v0.1.56

func (h *Handler) ReplayRun(c echo.Context) error

ReplayRun creates a new Run from one owned dead-letter Run. Agent policy and availability are re-evaluated by the normal creation path.

func (*Handler) ResumeBrowserControl added in v0.1.56

func (h *Handler) ResumeBrowserControl(c echo.Context) error

func (*Handler) RevokeRuntimeNode added in v0.1.56

func (h *Handler) RevokeRuntimeNode(c echo.Context) error

func (*Handler) RuntimeController added in v0.1.56

func (h *Handler) RuntimeController() *RuntimeHTTPController

RuntimeController exposes the transport lifecycle owner so the process can drain hijacked WebSocket connections before shutting down its HTTP servers.

func (*Handler) SendBrowserControlInput added in v0.1.56

func (h *Handler) SendBrowserControlInput(c echo.Context) error

func (*Handler) SetRunUpdateSource added in v0.1.56

func (h *Handler) SetRunUpdateSource(source RunUpdateSource)

func (*Handler) SetRuntimeDependencies added in v0.1.56

func (h *Handler) SetRuntimeDependencies(dependencies RuntimeHTTPDependencies)

SetRuntimeDependencies completes explicit production wiring. When the token validator is omitted, Handler's runtime service supplies authentication.

func (*Handler) SetWorkerObserver added in v0.1.56

func (h *Handler) SetWorkerObserver(observer WorkerObserver)

SetWorkerObserver installs payload-free test instrumentation. It does not change response, retry, timeout, or query behavior.

func (*Handler) StartAdminBrowserObservation added in v0.1.56

func (h *Handler) StartAdminBrowserObservation(c echo.Context) error

func (*Handler) StartBrowserObservation added in v0.1.56

func (h *Handler) StartBrowserObservation(c echo.Context) error

func (*Handler) StopAdminBrowserObservation added in v0.1.56

func (h *Handler) StopAdminBrowserObservation(c echo.Context) error

func (*Handler) StopBrowserObservation added in v0.1.56

func (h *Handler) StopBrowserObservation(c echo.Context) error

func (*Handler) StreamRunEvents

func (h *Handler) StreamRunEvents(c echo.Context) error

StreamRunEvents 以 SSE 输出 run events。

已结束的 run 会回放事件后关闭;运行中的 run 会轮询等待新事件直到终态或客户端断开。

type IdempotencyError added in v0.1.56

type IdempotencyError struct {
	Class IdempotencyErrorClass
	// contains filtered or unexported fields
}

IdempotencyError reports a safe error class without echoing request data.

func (*IdempotencyError) Error added in v0.1.56

func (e *IdempotencyError) Error() string

func (*IdempotencyError) Unwrap added in v0.1.56

func (e *IdempotencyError) Unwrap() error

type IdempotencyErrorClass added in v0.1.56

type IdempotencyErrorClass string

IdempotencyErrorClass is a stable, transport-independent classification. Handlers map these classes to their protocol-specific validation errors; the raw Idempotency-Key is deliberately never retained in an error.

const (
	IdempotencyErrorKeyRequired  IdempotencyErrorClass = "IDEMPOTENCY_KEY_REQUIRED"
	IdempotencyErrorKeyInvalid   IdempotencyErrorClass = "IDEMPOTENCY_KEY_INVALID"
	IdempotencyErrorInputInvalid IdempotencyErrorClass = "IDEMPOTENCY_INPUT_NOT_IJSON"
	IdempotencyErrorKeyReused    IdempotencyErrorClass = "IDEMPOTENCY_KEY_REUSED"
)

func IdempotencyErrorClassOf added in v0.1.56

func IdempotencyErrorClassOf(err error) (IdempotencyErrorClass, bool)

IdempotencyErrorClassOf returns the stable class for service/handler mapping.

type LocalSignalBus added in v0.1.56

type LocalSignalBus struct {
	// contains filtered or unexported fields
}

LocalSignalBus is a single-process implementation for development, tests, and explicitly single-instance deployments. It does not claim HA delivery.

func NewLocalSignalBus added in v0.1.56

func NewLocalSignalBus(instanceID uuid.UUID) *LocalSignalBus

func (*LocalSignalBus) Close added in v0.1.56

func (b *LocalSignalBus) Close() error

func (*LocalSignalBus) Health added in v0.1.56

func (b *LocalSignalBus) Health(ctx context.Context) error

func (*LocalSignalBus) Publish added in v0.1.56

func (b *LocalSignalBus) Publish(ctx context.Context, signal RuntimeSignal) error

func (*LocalSignalBus) Subscribe added in v0.1.56

func (b *LocalSignalBus) Subscribe(ctx context.Context, handler RuntimeSignalHandler) error

type MTLSRuntimeDeviceAuthenticator added in v0.1.56

type MTLSRuntimeDeviceAuthenticator struct {
	// contains filtered or unexported fields
}

MTLSRuntimeDeviceAuthenticator authenticates only a TLS peer certificate. Forwarded certificate headers are deliberately ignored.

func NewMTLSRuntimeDeviceAuthenticator added in v0.1.56

func NewMTLSRuntimeDeviceAuthenticator(verifier RuntimeNodeCredentialVerifier) *MTLSRuntimeDeviceAuthenticator

func (*MTLSRuntimeDeviceAuthenticator) AuthenticateHTTP added in v0.1.56

AuthenticateHTTP authenticates the verified leaf certificate attached by net/http and resolves its durable Node credential. VerifiedChains is required so RequestClientCert-only server configurations cannot silently turn a presented, unverified certificate into a trusted identity.

type PendingCommand added in v0.1.56

type PendingCommand struct {
	Type    RuntimeMessageType `json:"type" runtime:"required"`
	Payload json.RawMessage    `json:"payload" runtime:"required"`
}

PendingCommand is a discriminated union. Payload is decoded strictly by DecodePendingCommand after Type has selected the only legal schema.

type RedisRuntimePresenceStore added in v0.1.56

type RedisRuntimePresenceStore struct {
	// contains filtered or unexported fields
}

func NewRedisRuntimePresenceStore added in v0.1.56

func NewRedisRuntimePresenceStore(client redis.UniversalClient, prefix string) (*RedisRuntimePresenceStore, error)

func (*RedisRuntimePresenceStore) ListByAgent added in v0.1.56

func (s *RedisRuntimePresenceStore) ListByAgent(ctx context.Context, agentID uuid.UUID) ([]RuntimePresence, error)

func (*RedisRuntimePresenceStore) Refresh added in v0.1.56

func (*RedisRuntimePresenceStore) Remove added in v0.1.56

type RedisRuntimeSessionLeaseStore added in v0.1.56

type RedisRuntimeSessionLeaseStore struct {
	// contains filtered or unexported fields
}

func NewRedisRuntimeSessionLeaseStore added in v0.1.56

func NewRedisRuntimeSessionLeaseStore(
	client redis.UniversalClient,
	leasePrefix string,
	presencePrefix string,
) (*RedisRuntimeSessionLeaseStore, error)

func (*RedisRuntimeSessionLeaseStore) Forget added in v0.1.56

func (s *RedisRuntimeSessionLeaseStore) Forget(
	ctx context.Context,
	runtimeSessionID uuid.UUID,
) error

func (*RedisRuntimeSessionLeaseStore) ListExpired added in v0.1.56

func (s *RedisRuntimeSessionLeaseStore) ListExpired(
	ctx context.Context,
	limit int,
) ([]uuid.UUID, error)

func (*RedisRuntimeSessionLeaseStore) Lookup added in v0.1.56

func (s *RedisRuntimeSessionLeaseStore) Lookup(
	ctx context.Context,
	runtimeSessionID uuid.UUID,
) (RuntimeSessionLease, bool, error)

func (*RedisRuntimeSessionLeaseStore) RefreshBatch added in v0.1.56

func (s *RedisRuntimeSessionLeaseStore) RefreshBatch(
	ctx context.Context,
	records []RuntimeSessionLeaseRecord,
	leaseTTL time.Duration,
	presenceTTL time.Duration,
) error

func (*RedisRuntimeSessionLeaseStore) Remove added in v0.1.56

func (*RedisRuntimeSessionLeaseStore) ScheduleCheck added in v0.1.56

func (s *RedisRuntimeSessionLeaseStore) ScheduleCheck(
	ctx context.Context,
	runtimeSessionID uuid.UUID,
	after time.Duration,
) error

type RedisSignalBus added in v0.1.56

type RedisSignalBus struct {
	// contains filtered or unexported fields
}

RedisSignalBus owns its Pub/Sub subscriptions but not the supplied Redis client. Construction performs no network I/O: temporary Redis loss must not prevent Core from starting its PostgreSQL reconciliation workers. Health is checked dynamically by readiness.

func NewRedisSignalBus added in v0.1.56

func NewRedisSignalBus(client redis.UniversalClient, cfg RedisSignalBusConfig) (*RedisSignalBus, error)

func (*RedisSignalBus) Check added in v0.1.56

func (*RedisSignalBus) Close added in v0.1.56

func (b *RedisSignalBus) Close() error

func (*RedisSignalBus) Health added in v0.1.56

func (b *RedisSignalBus) Health(ctx context.Context) error

func (*RedisSignalBus) MarkActive added in v0.1.56

func (b *RedisSignalBus) MarkActive(
	ctx context.Context,
	registrations []RuntimeConnectionRegistration,
) error

func (*RedisSignalBus) Publish added in v0.1.56

func (b *RedisSignalBus) Publish(ctx context.Context, signal RuntimeSignal) error

func (*RedisSignalBus) RuntimeCredentialProjectionStore added in v0.1.56

func (b *RedisSignalBus) RuntimeCredentialProjectionStore() (RuntimeCredentialProjectionStore, error)

func (*RedisSignalBus) RuntimePresenceStore added in v0.1.56

func (b *RedisSignalBus) RuntimePresenceStore() (RuntimePresenceStore, error)

func (*RedisSignalBus) RuntimeSessionLeaseStore added in v0.1.56

func (b *RedisSignalBus) RuntimeSessionLeaseStore() (RuntimeSessionLeaseStore, error)

func (*RedisSignalBus) Subscribe added in v0.1.56

func (b *RedisSignalBus) Subscribe(ctx context.Context, handler RuntimeSignalHandler) error

type RedisSignalBusConfig added in v0.1.56

type RedisSignalBusConfig struct {
	Channel    string
	InstanceID uuid.UUID
}

type ResultClassificationInput added in v0.1.56

type ResultClassificationInput struct {
	ExecutorType        string
	EndpointIdempotency bool
	Request             RuntimeResultRequest
}

ResultClassificationInput is the immutable server policy snapshot supplied to ResultClassifier. The caller-provided retryable hint is never sufficient for Core-owned HTTP/MCP attempts without endpoint idempotency evidence.

type ResultClassifier added in v0.1.56

type ResultClassifier interface {
	ClassifyResult(ResultClassificationInput) RuntimeResultClassification
}

ResultClassifier owns retry policy. Implementations must return one of the three persisted Attempt classifications.

type ResultClassifierFunc added in v0.1.56

type ResultClassifierFunc func(ResultClassificationInput) RuntimeResultClassification

ResultClassifierFunc adapts a function for focused policy tests.

func (ResultClassifierFunc) ClassifyResult added in v0.1.56

type ResultFinalizer added in v0.1.56

type ResultFinalizer struct {
	// contains filtered or unexported fields
}

ResultFinalizer is the sole atomic write path for an accepted Runtime Result. Public transports are intentionally wired in Task 6.

func NewResultFinalizer added in v0.1.56

func NewResultFinalizer(pool *pgxpool.Pool, classifier ResultClassifier, retryPlanner ResultRetryPlanner) *ResultFinalizer

NewResultFinalizer builds an isolated Result finalizer. Nil policies select conservative server defaults.

func (*ResultFinalizer) Finalize added in v0.1.56

Finalize validates, fingerprints, and atomically persists a Result. The implementation uses PostgreSQL's Run row as the per-Run linearization point.

type ResultRetryPlanner added in v0.1.56

type ResultRetryPlanner interface {
	NextRetryDelay(attemptNo int32) time.Duration
}

ResultRetryPlanner persists a single delay for an accepted retryable Result. Task 7 replaces the deterministic placeholder with the approved jittered policy without changing Finalizer's atomic boundary.

type ResultRetryPlannerFunc added in v0.1.56

type ResultRetryPlannerFunc func(int32) time.Duration

ResultRetryPlannerFunc adapts a function for focused tests.

func (ResultRetryPlannerFunc) NextRetryDelay added in v0.1.56

func (fn ResultRetryPlannerFunc) NextRetryDelay(attemptNo int32) time.Duration

type ResumeAttempt added in v0.1.56

type ResumeAttempt struct {
	AttemptIdentity          AttemptIdentity `json:"attempt_identity" runtime:"required"`
	LastAckedClientEventSeq  int64           `json:"last_acked_client_event_seq" runtime:"required"`
	PendingClientEventRanges []EventRange    `json:"pending_client_event_ranges" runtime:"required"`
	PendingResultID          *uuid.UUID      `json:"pending_result_id,omitempty"`
	FinalClientEventSeq      *int64          `json:"final_client_event_seq,omitempty"`
}

type RevokeCommand added in v0.1.56

type RevokeCommand = RuntimeCommand[RunLeaseRevokedPayload]

type RevokeRuntimeNodeRequest added in v0.1.56

type RevokeRuntimeNodeRequest struct {
	Reason string `json:"reason" validate:"required,min=1,max=500"`
}

type RunA2AContextRequest

type RunA2AContextRequest struct {
	MessageID           string                 `json:"message_id,omitempty"`
	ProtocolContextID   string                 `json:"protocol_context_id,omitempty"`
	ProtocolTaskID      string                 `json:"protocol_task_id,omitempty"`
	RootContextID       string                 `json:"root_context_id,omitempty"`
	ParentContextID     string                 `json:"parent_context_id,omitempty"`
	ParentTaskID        string                 `json:"parent_task_id,omitempty"`
	ParentRunID         string                 `json:"parent_run_id,omitempty"`
	CallerAgentID       string                 `json:"caller_agent_id,omitempty"`
	TargetAgentID       string                 `json:"target_agent_id,omitempty"`
	TraceID             string                 `json:"trace_id,omitempty"`
	ReferenceTaskIDs    []string               `json:"reference_task_ids,omitempty"`
	Source              string                 `json:"source,omitempty"`
	AcceptedOutputModes []string               `json:"accepted_output_modes,omitempty"`
	Extensions          []string               `json:"extensions,omitempty"`
	Visibility          string                 `json:"visibility,omitempty"`
	Options             map[string]interface{} `json:"options,omitempty"`
}

type RunA2AContextResponse

type RunA2AContextResponse struct {
	ProtocolContextID string   `json:"protocol_context_id"`
	ProtocolTaskID    string   `json:"protocol_task_id"`
	RootContextID     string   `json:"root_context_id"`
	ParentContextID   string   `json:"parent_context_id,omitempty"`
	ParentTaskID      string   `json:"parent_task_id,omitempty"`
	ParentRunID       string   `json:"parent_run_id,omitempty"`
	CallerAgentID     string   `json:"caller_agent_id,omitempty"`
	TargetAgentID     string   `json:"target_agent_id,omitempty"`
	TraceID           string   `json:"trace_id,omitempty"`
	ReferenceTaskIDs  []string `json:"reference_task_ids,omitempty"`
	Source            string   `json:"source,omitempty"`
}

type RunArtifactResponse

type RunArtifactResponse struct {
	ID               string                 `json:"id"`
	RunID            string                 `json:"run_id"`
	ArtifactType     string                 `json:"artifact_type"`
	Title            string                 `json:"title"`
	Content          map[string]interface{} `json:"content"`
	Visibility       string                 `json:"visibility"`
	SourceArtifactID string                 `json:"source_artifact_id,omitempty"`
	MimeType         string                 `json:"mime_type,omitempty"`
	FileURI          string                 `json:"file_uri,omitempty"`
	FileName         string                 `json:"file_name,omitempty"`
	FileSHA256       string                 `json:"file_sha256,omitempty"`
	FileSizeBytes    *int64                 `json:"file_size_bytes,omitempty"`
	CreatedAt        time.Time              `json:"created_at"`
}

RunArtifactResponse is a persisted, owner-readable artifact produced by a run.

type RunAssignedMessage added in v0.1.56

type RunAssignedMessage = RuntimeTypedEnvelope[RunAssignedPayload]

type RunAssignedPayload added in v0.1.56

type RunAssignedPayload struct {
	AttemptIdentity      AttemptIdentity `json:"attempt_identity" runtime:"required"`
	OfferNo              int64           `json:"offer_no" runtime:"required"`
	OfferExpiresAt       time.Time       `json:"offer_expires_at" runtime:"required"`
	AttemptDeadlineAt    time.Time       `json:"attempt_deadline_at" runtime:"required"`
	RunDeadlineAt        time.Time       `json:"run_deadline_at" runtime:"required"`
	Input                map[string]any  `json:"input" runtime:"required"`
	Metadata             map[string]any  `json:"metadata,omitempty"`
	NodeEnvelope         string          `json:"node_envelope" runtime:"required"`
	AgentInvocationToken string          `json:"agent_invocation_token" runtime:"required"`
}

type RunAssignmentAckMessage added in v0.1.56

type RunAssignmentAckMessage = RuntimeTypedEnvelope[RunAssignmentAckPayload]

type RunAssignmentAckPayload added in v0.1.56

type RunAssignmentAckPayload struct {
	AttemptIdentity AttemptIdentity `json:"attempt_identity" runtime:"required"`
}

type RunAssignmentConfirmedMessage added in v0.1.56

type RunAssignmentConfirmedMessage = RuntimeTypedEnvelope[RunAssignmentConfirmedPayload]

type RunAssignmentConfirmedPayload added in v0.1.56

type RunAssignmentConfirmedPayload struct {
	AttemptIdentity AttemptIdentity `json:"attempt_identity" runtime:"required"`
	AttemptNo       int64           `json:"attempt_no" runtime:"required"`
	LeaseExpiresAt  time.Time       `json:"lease_expires_at" runtime:"required"`
}

type RunAssignmentRejectMessage added in v0.1.56

type RunAssignmentRejectMessage = RuntimeTypedEnvelope[RunAssignmentRejectPayload]

type RunAssignmentRejectPayload added in v0.1.56

type RunAssignmentRejectPayload struct {
	AttemptIdentity AttemptIdentity               `json:"attempt_identity" runtime:"required"`
	ReasonCode      RuntimeAssignmentRejectReason `json:"reason_code" runtime:"required"`
	Capacity        int64                         `json:"capacity" runtime:"required"`
	Inflight        int64                         `json:"inflight" runtime:"required"`
}

type RunAssignmentRejectedMessage added in v0.1.56

type RunAssignmentRejectedMessage = RuntimeTypedEnvelope[RunAssignmentRejectedPayload]

type RunAssignmentRejectedPayload added in v0.1.56

type RunAssignmentRejectedPayload struct {
	AttemptIdentity AttemptIdentity                `json:"attempt_identity" runtime:"required"`
	Outcome         RuntimeAssignmentRejectOutcome `json:"outcome" runtime:"required"`
	DispatchState   RuntimeDispatchState           `json:"dispatch_state" runtime:"required"`
}

type RunCancelAckMessage added in v0.1.56

type RunCancelAckMessage = RuntimeTypedEnvelope[RunCancelAckPayload]

type RunCancelAckPayload added in v0.1.56

type RunCancelAckPayload struct {
	CancellationID  uuid.UUID          `json:"cancellation_id" runtime:"required"`
	AttemptIdentity AttemptIdentity    `json:"attempt_identity" runtime:"required"`
	CancelState     RuntimeCancelState `json:"cancel_state" runtime:"required"`
	ErrorCode       string             `json:"error_code,omitempty"`
}

type RunCancelMessage added in v0.1.56

type RunCancelMessage = RuntimeTypedEnvelope[RunCancelPayload]

type RunCancelPayload added in v0.1.56

type RunCancelPayload struct {
	CancellationID  uuid.UUID       `json:"cancellation_id" runtime:"required"`
	AttemptIdentity AttemptIdentity `json:"attempt_identity" runtime:"required"`
	ReasonCode      string          `json:"reason_code" runtime:"required"`
	DeadlineAt      time.Time       `json:"deadline_at" runtime:"required"`
}

type RunCancellationEvidence added in v0.1.56

type RunCancellationEvidence struct {
	CancellationID uuid.UUID
	State          string
	ErrorCode      string
	RequestedAt    time.Time
	FinishedAt     *time.Time
}

type RunCancellationState added in v0.1.56

type RunCancellationState struct {
	CancellationID uuid.UUID          `json:"cancellation_id" runtime:"required"`
	CancelState    RuntimeCancelState `json:"cancel_state" runtime:"required"`
	UpdatedAt      time.Time          `json:"updated_at" runtime:"required"`
	ErrorCode      string             `json:"error_code,omitempty"`
}

type RunCreationIdentity added in v0.1.56

type RunCreationIdentity struct {
	IdempotencyKeyHash  []byte
	CreationFingerprint []byte
}

RunCreationIdentity is the non-sensitive identity needed to recover a Run after creation committed but its caller crashed before recording the Run ID.

type RunEffectAttemptResult added in v0.1.56

type RunEffectAttemptResult struct {
	Succeeded   bool
	Retryable   bool
	ErrorCode   string
	SafeMessage string
	Err         error
}

RunEffectAttemptResult is returned by one downstream delivery attempt. A handler must perform at most one external HTTP request for a claimed Effect. ErrorCode and SafeMessage are persisted, so they must never contain target URLs, credentials, request/response payloads, or caller-controlled text.

func PermanentRunEffectFailure added in v0.1.56

func PermanentRunEffectFailure(code, safeMessage string, err error) RunEffectAttemptResult

func RetryableRunEffectAttempt added in v0.1.56

func RetryableRunEffectAttempt(code, safeMessage string, err error) RunEffectAttemptResult

func RunEffectAttemptSucceeded added in v0.1.56

func RunEffectAttemptSucceeded() RunEffectAttemptResult

type RunEffectWorker added in v0.1.56

type RunEffectWorker struct {
	// contains filtered or unexported fields
}

func NewRunEffectWorker added in v0.1.56

func NewRunEffectWorker(
	queries runEffectStore,
	webhook WebhookRunEffectHandler,
	delivery DeliveryRunEffectHandler,
) *RunEffectWorker

func (*RunEffectWorker) ProcessOnce added in v0.1.56

func (w *RunEffectWorker) ProcessOnce(
	ctx context.Context,
	cfg RunEffectWorkerConfig,
) (int, error)

func (*RunEffectWorker) Replay added in v0.1.56

func (w *RunEffectWorker) Replay(
	ctx context.Context,
	effectID uuid.UUID,
	actorType string,
	actorID *uuid.UUID,
	reason string,
) (*db.RunEffectOutbox, error)

func (*RunEffectWorker) SetHandlers added in v0.1.56

func (w *RunEffectWorker) SetHandlers(
	webhook WebhookRunEffectHandler,
	delivery DeliveryRunEffectHandler,
)

type RunEffectWorkerConfig added in v0.1.56

type RunEffectWorkerConfig struct {
	Interval      time.Duration
	LeaseDuration time.Duration
	BatchSize     int32
	Observer      WorkerObserver
}

type RunErrorPayload added in v0.1.56

type RunErrorPayload struct {
	ErrorCode     string `json:"error_code" runtime:"required"`
	Message       string `json:"message" runtime:"required"`
	RetryableHint bool   `json:"retryable_hint,omitempty"`
}

type RunEventAckMessage added in v0.1.56

type RunEventAckMessage = RuntimeTypedEnvelope[RunEventAckPayload]

type RunEventAckPayload added in v0.1.56

type RunEventAckPayload struct {
	ClientEventID  uuid.UUID `json:"client_event_id" runtime:"required"`
	ClientEventSeq int64     `json:"client_event_seq" runtime:"required"`
	Sequence       int64     `json:"sequence" runtime:"required"`
	Replayed       bool      `json:"replayed" runtime:"required"`
}

type RunEventMessage added in v0.1.56

type RunEventMessage = RuntimeTypedEnvelope[RunEventPayload]

type RunEventPageMeta added in v0.1.56

type RunEventPageMeta struct {
	RequestedAfterSequence    int32  `json:"requested_after_sequence"`
	EffectiveAfterSequence    int32  `json:"effective_after_sequence"`
	RetainedThroughSequence   int32  `json:"retained_through_sequence"`
	EarliestAvailableSequence *int32 `json:"earliest_available_sequence"`
	LatestAvailableSequence   *int32 `json:"latest_available_sequence"`
	RetentionGap              bool   `json:"retention_gap"`
	Terminal                  bool   `json:"terminal"`
	StreamComplete            bool   `json:"stream_complete"`
}

RunEventPageMeta makes an incomplete event history explicit. Sequence zero is the cursor before the first event; unavailable bounds are encoded as null.

type RunEventPageResponse added in v0.1.56

type RunEventPageResponse struct {
	Items []RunEventResponse `json:"items"`
	Meta  RunEventPageMeta   `json:"meta"`
}

RunEventPageResponse is an owner-readable event page plus the durable retention boundary used to interpret that page.

type RunEventPayload added in v0.1.56

type RunEventPayload struct {
	AttemptIdentity AttemptIdentity `json:"attempt_identity" runtime:"required"`
	ClientEventID   uuid.UUID       `json:"client_event_id" runtime:"required"`
	ClientEventSeq  int64           `json:"client_event_seq" runtime:"required"`
	EventType       string          `json:"event_type" runtime:"required"`
	Payload         map[string]any  `json:"payload" runtime:"required"`
}

func (RunEventPayload) StoreRequest added in v0.1.56

func (p RunEventPayload) StoreRequest() RuntimeEventRequest

type RunEventResponse

type RunEventResponse struct {
	EventID     string                 `json:"event_id"`
	RunID       string                 `json:"run_id"`
	ParentRunID string                 `json:"parent_run_id,omitempty"`
	Sequence    int32                  `json:"sequence"`
	EventType   string                 `json:"event_type"`
	Payload     map[string]interface{} `json:"payload"`
	CreatedAt   time.Time              `json:"created_at"`
}

RunEventResponse GET /api/v1/runs/:id/events 响应项。

sequence 在单个 run 内单调递增,后续 SSE 可用作 Last-Event-ID。

type RunEvidenceSummary

type RunEvidenceSummary struct {
	Status            string  `json:"status"`
	CoverageStatus    string  `json:"coverage_status"`
	MatchedSkillCount int     `json:"matched_skill_count"`
	MissingSkillCount int     `json:"missing_skill_count"`
	UsedMCPToolCount  int     `json:"used_mcp_tool_count"`
	ArtifactCount     int     `json:"artifact_count"`
	MessageCount      int     `json:"message_count"`
	DeliveryStatus    *string `json:"delivery_status,omitempty"`
	PublicSafe        bool    `json:"public_safe"`
	EvidenceURL       string  `json:"evidence_url"`
}

RunEvidenceSummary gives UI and external clients a compact view of why a run is trustworthy without forcing them to stitch together multiple endpoints.

type RunFingerprintA2A added in v0.1.56

type RunFingerprintA2A struct {
	MessageID           string
	ProtocolContextID   string
	ProtocolTaskID      string
	RootContextID       string
	ParentContextID     string
	ParentTaskID        string
	ReferenceTaskIDs    []string
	Source              string
	Visibility          string
	AcceptedOutputModes []string
	Extensions          []string
	Options             map[string]any

	TraceID string
}

RunFingerprintA2A contains A2A context, references, and output semantics. TraceID is accepted solely to make exclusion explicit; it is transport correlation data and is not serialized into the fingerprint.

type RunFingerprintDelegation added in v0.1.56

type RunFingerprintDelegation struct {
	Reason  string
	Mode    string
	Depth   int
	Options map[string]any
}

RunFingerprintDelegation contains execution-visible delegation options. ParentRunID and CallerAgentID live on RunFingerprintInput so they cannot be accidentally omitted when Delegation has no extra options.

type RunFingerprintInput added in v0.1.56

type RunFingerprintInput struct {
	Target        RunFingerprintTarget
	Input         map[string]any
	Metadata      map[string]any
	Callback      any
	Delivery      any
	Source        RunFingerprintSource
	ParentRunID   string
	CallerAgentID string
	Delegation    *RunFingerprintDelegation
	A2A           *RunFingerprintA2A
	Visibility    string
	Options       any

	Transport RunFingerprintTransport
}

RunFingerprintInput is the normalized semantic snapshot used by every Run creation entrypoint. Callback, Delivery, and Options must contain JSON-domain values after protocol aliases and defaults have been collapsed. Transport is intentionally excluded by FingerprintRunCreation.

type RunFingerprintSource added in v0.1.56

type RunFingerprintSource struct {
	Protocol string
	Method   string
}

RunFingerprintSource identifies the entry protocol and operation after aliases have been normalized (for example, rest/create-run or a2a/message-send).

type RunFingerprintTarget added in v0.1.56

type RunFingerprintTarget struct {
	AgentID   string
	ReleaseID string
}

RunFingerprintTarget identifies the immutable execution target. ReleaseID is optional until Agent Releases are introduced, but is already part of the fingerprint contract so the cutover does not need another envelope shape.

type RunFingerprintTransport added in v0.1.56

type RunFingerprintTransport struct {
	RequestID      string
	TraceID        string
	TraceParent    string
	RetryCount     int
	IdempotencyKey string
}

RunFingerprintTransport records intentionally excluded delivery metadata. Keeping it separate prevents recursive name-based filtering from deleting a legitimate request_id, trace_id, retry_count, or idempotency_key inside the user-controlled input or metadata objects.

type RunLeaseRenewMessage added in v0.1.56

type RunLeaseRenewMessage = RuntimeTypedEnvelope[RunLeaseRenewPayload]

type RunLeaseRenewPayload added in v0.1.56

type RunLeaseRenewPayload struct {
	AttemptIdentity    AttemptIdentity `json:"attempt_identity" runtime:"required"`
	LastClientEventSeq int64           `json:"last_client_event_seq" runtime:"required"`
	Capacity           int64           `json:"capacity" runtime:"required"`
	Inflight           int64           `json:"inflight" runtime:"required"`
}

type RunLeaseRenewedMessage added in v0.1.56

type RunLeaseRenewedMessage = RuntimeTypedEnvelope[RunLeaseRenewedPayload]

type RunLeaseRenewedPayload added in v0.1.56

type RunLeaseRenewedPayload struct {
	AttemptIdentity AttemptIdentity `json:"attempt_identity" runtime:"required"`
	LeaseExpiresAt  time.Time       `json:"lease_expires_at" runtime:"required"`
	PendingCommand  *PendingCommand `json:"pending_command,omitempty" runtime:"nullable"`
}

type RunLeaseRevokedMessage added in v0.1.56

type RunLeaseRevokedMessage = RuntimeTypedEnvelope[RunLeaseRevokedPayload]

type RunLeaseRevokedPayload added in v0.1.56

type RunLeaseRevokedPayload struct {
	AttemptIdentity AttemptIdentity      `json:"attempt_identity" runtime:"required"`
	ReasonCode      string               `json:"reason_code" runtime:"required"`
	DispatchState   RuntimeDispatchState `json:"dispatch_state" runtime:"required"`
	RunStatus       RuntimeRunStatus     `json:"run_status" runtime:"required"`
}

type RunMessageResponse

type RunMessageResponse struct {
	ID            string                 `json:"id"`
	RunID         string                 `json:"run_id"`
	EventSequence *int32                 `json:"event_sequence,omitempty"`
	Role          string                 `json:"role"`
	Content       string                 `json:"content"`
	Payload       map[string]interface{} `json:"payload"`
	CreatedAt     time.Time              `json:"created_at"`
}

RunMessageResponse is a stable replay record derived from user input and agent message events.

type RunNextAction

type RunNextAction struct {
	Type            string                 `json:"type"`
	Label           string                 `json:"label"`
	Hint            string                 `json:"hint"`
	Href            string                 `json:"href,omitempty"`
	Method          string                 `json:"method,omitempty"`
	RequiresHuman   bool                   `json:"requires_human,omitempty"`
	ResourceType    string                 `json:"resource_type,omitempty"`
	ResourceID      string                 `json:"resource_id,omitempty"`
	Source          string                 `json:"source,omitempty"`
	AdditionalProps map[string]interface{} `json:"additional_props,omitempty"`
}

RunNextAction is a machine-readable hint for the UI or external clients.

Type is stable enough for clients to branch on. Hint is human-facing.

type RunRequest

type RunRequest struct {
	AgentID        string                 `json:"agent_id" validate:"required,uuid"`
	Input          map[string]interface{} `json:"input" validate:"required"`
	Metadata       map[string]interface{} `json:"metadata,omitempty"`
	A2AContext     *RunA2AContextRequest  `json:"a2a_context,omitempty"`
	IdempotencyKey string                 `json:"-"`

	// CreationProtocol and CreationMethod are normalized by each entrypoint.
	// They are execution semantics, not user-controlled JSON fields, and are
	// included in the idempotency fingerprint.
	CreationProtocol string         `json:"-"`
	CreationMethod   string         `json:"-"`
	CreationOptions  map[string]any `json:"-"`
	// TaskCallback is OpenLinker's canonical callback field. PushNotification,
	// PushNotificationAlias, and PushNotificationConfig are A2A compatibility
	// aliases; taskCallbackConfigFromRunRequest chooses the first non-empty
	// value in this declaration order.
	TaskCallback           *TaskCallbackConfig `json:"task_callback,omitempty"`
	PushNotification       *TaskCallbackConfig `json:"push_notification,omitempty"`
	PushNotificationAlias  *TaskCallbackConfig `json:"pushNotification,omitempty"`
	PushNotificationConfig *TaskCallbackConfig `json:"pushNotificationConfig,omitempty"`
}

RunRequest POST /api/v1/run 请求体。

AgentID 在 service 层 uuid.Parse 校验。 Input 必填,为创作者 endpoint 接收的入参(透传)。 Metadata 可选。普通调用方字段会转发给第三方 endpoint;平台身份、内部编排 和私有 authority 字段会在出站边界被剔除,并由 Core 加入不可逆的 principal_scope_id。queued Runtime 保留完整的私有 authority。

type RunRequirementEvidenceResponse

type RunRequirementEvidenceResponse struct {
	RunID            string    `json:"run_id"`
	TaskID           string    `json:"task_id"`
	AgentID          string    `json:"agent_id"`
	RequiredSkillIDs []string  `json:"required_skill_ids"`
	RequiredMCPTools []string  `json:"required_mcp_tools"`
	AgentSkillIDs    []string  `json:"agent_skill_ids"`
	MatchedSkillIDs  []string  `json:"matched_skill_ids"`
	MissingSkillIDs  []string  `json:"missing_skill_ids"`
	UsedMCPTools     []string  `json:"used_mcp_tools"`
	MissingMCPTools  []string  `json:"missing_mcp_tools"`
	CoverageStatus   string    `json:"coverage_status"`
	EvidenceSource   string    `json:"evidence_source"`
	CreatedAt        time.Time `json:"created_at"`
}

RunRequirementEvidenceResponse proves which task Skill/MCP requirements were snapshotted onto a run and how the selected Agent covered them.

type RunResponse

type RunResponse struct {
	CanReplay           bool                   `json:"can_replay"`
	RunID               string                 `json:"run_id"`
	AgentID             string                 `json:"agent_id,omitempty"`
	AgentSlug           string                 `json:"agent_slug,omitempty"`
	AgentName           string                 `json:"agent_name,omitempty"`
	AgentConnectionMode string                 `json:"agent_connection_mode,omitempty"`
	Status              string                 `json:"status"`
	Input               map[string]interface{} `json:"input,omitempty"`
	Output              map[string]interface{} `json:"output,omitempty"`
	ErrorCode           string                 `json:"error_code,omitempty"`
	// ErrorMsg keeps the historical Go field name while preserving the public
	// JSON contract as error_message.
	ErrorMsg                           string                          `json:"error_message,omitempty"`
	CostCents                          int32                           `json:"cost_cents"`
	DurationMs                         int32                           `json:"duration_ms"`
	StartedAt                          time.Time                       `json:"started_at"`
	FinishedAt                         *time.Time                      `json:"finished_at,omitempty"`
	Source                             string                          `json:"source,omitempty"`
	RuntimeContractID                  string                          `json:"runtime_contract_id"`
	RuntimeTransport                   string                          `json:"runtime_transport,omitempty"`
	RuntimeTransportReason             string                          `json:"runtime_transport_reason,omitempty"`
	RuntimeTransportChangedAt          *time.Time                      `json:"runtime_transport_changed_at,omitempty"`
	BrowserInteractionPolicy           string                          `json:"browser_interaction_policy,omitempty"`
	BrowserInteractionPolicyGeneration int64                           `json:"browser_interaction_policy_generation,omitempty"`
	BrowserMutationOrigins             []string                        `json:"browser_mutation_origins,omitempty"`
	BrowserMutationOriginsSHA256       string                          `json:"browser_mutation_origins_sha256,omitempty"`
	BrowserContractID                  string                          `json:"browser_contract_id,omitempty"`
	DispatchState                      string                          `json:"dispatch_state"`
	AttemptCount                       int32                           `json:"attempt_count"`
	MaxAttempts                        int32                           `json:"max_attempts"`
	NextAttemptAt                      *time.Time                      `json:"next_attempt_at,omitempty"`
	LatestAttemptID                    string                          `json:"latest_attempt_id,omitempty"`
	ActiveAttemptID                    string                          `json:"active_attempt_id,omitempty"`
	CancelState                        string                          `json:"cancel_state,omitempty"`
	CancelRequestedAt                  *time.Time                      `json:"cancel_requested_at,omitempty"`
	CancelAcknowledgedAt               *time.Time                      `json:"cancel_acknowledged_at,omitempty"`
	CancelReason                       string                          `json:"cancel_reason,omitempty"`
	DeadLetteredAt                     *time.Time                      `json:"dead_lettered_at,omitempty"`
	ReplayOfRunID                      string                          `json:"replay_of_run_id,omitempty"`
	ParentRunID                        string                          `json:"parent_run_id,omitempty"`
	CallerAgentID                      string                          `json:"caller_agent_id,omitempty"`
	BillingMode                        string                          `json:"billing_mode,omitempty"`
	A2AContext                         *RunA2AContextResponse          `json:"a2a_context,omitempty"`
	TaskCallback                       *RunTaskCallbackResponse        `json:"task_callback,omitempty"`
	RequirementEvidence                *RunRequirementEvidenceResponse `json:"requirement_evidence,omitempty"`
	EvidenceSummary                    *RunEvidenceSummary             `json:"evidence_summary,omitempty"`
	NextAction                         *RunNextAction                  `json:"next_action,omitempty"`
	Replayed                           bool                            `json:"replayed"`
}

RunResponse POST /api/v1/run 同步响应体,或 POST /api/v1/runs 异步启动响应体。

Status: 'running' / 'success' / 'failed' / 'timeout'。 失败 / 超时 时 Output 为空、ErrorCode + ErrorMsg 必填,CostCents=0(已退款)。 Source: 'web' / 'mcp' / 'api',由 handler 从鉴权方式推导。

type RunResultAckMessage added in v0.1.56

type RunResultAckMessage = RuntimeTypedEnvelope[RunResultAckPayload]

type RunResultAckPayload added in v0.1.56

type RunResultAckPayload struct {
	ResultID       uuid.UUID                   `json:"result_id" runtime:"required"`
	Classification RuntimeResultClassification `json:"classification" runtime:"required"`
	RunStatus      RuntimeRunStatus            `json:"run_status" runtime:"required"`
	DispatchState  RuntimeDispatchState        `json:"dispatch_state" runtime:"required"`
	Replayed       bool                        `json:"replayed" runtime:"required"`
	NextAttemptAt  *time.Time                  `json:"next_attempt_at,omitempty"`
}

type RunResultMessage added in v0.1.56

type RunResultMessage = RuntimeTypedEnvelope[RunResultPayload]

type RunResultPayload added in v0.1.56

type RunResultPayload struct {
	AttemptIdentity     AttemptIdentity  `json:"attempt_identity" runtime:"required"`
	ResultID            uuid.UUID        `json:"result_id" runtime:"required"`
	Status              string           `json:"status" runtime:"required"`
	Output              map[string]any   `json:"output,omitempty"`
	Error               *RunErrorPayload `json:"error,omitempty"`
	DurationMS          int64            `json:"duration_ms" runtime:"required"`
	FinalClientEventSeq int64            `json:"final_client_event_seq" runtime:"required"`
}

func (RunResultPayload) FinalizerRequest added in v0.1.56

func (p RunResultPayload) FinalizerRequest() (RuntimeResultRequest, error)

FinalizerRequest converts a validated wire Result into the transport-neutral Finalizer request without trusting a lossy integer conversion.

type RunResumeAcceptedMessage added in v0.1.56

type RunResumeAcceptedMessage = RuntimeTypedEnvelope[RunResumeAcceptedPayload]

type RunResumeAcceptedPayload added in v0.1.56

type RunResumeAcceptedPayload struct {
	AttemptIdentity AttemptIdentity       `json:"attempt_identity" runtime:"required"`
	Decision        RuntimeResumeDecision `json:"decision" runtime:"required"`
	LeaseExpiresAt  *time.Time            `json:"lease_expires_at,omitempty"`
	AllowedActions  []RuntimeResumeAction `json:"allowed_actions" runtime:"required"`
}

type RunSummary added in v0.1.56

type RunSummary struct {
	RunID         uuid.UUID            `json:"run_id" runtime:"required"`
	Status        RuntimeRunStatus     `json:"status" runtime:"required"`
	DispatchState RuntimeDispatchState `json:"dispatch_state" runtime:"required"`
}

type RunTaskCallbackResponse

type RunTaskCallbackResponse struct {
	ID                  string   `json:"id"`
	RunID               string   `json:"run_id"`
	TargetURL           string   `json:"target_url"`
	EventTypes          []string `json:"event_types"`
	AuthScheme          string   `json:"auth_scheme,omitempty"`
	Status              string   `json:"status"`
	ConsecutiveFailures int32    `json:"consecutive_failures"`
	Secret              string   `json:"secret,omitempty"`
	CreatedAt           string   `json:"created_at"`
	UpdatedAt           string   `json:"updated_at"`
}

RunTaskCallbackResponse describes a caller-owned task callback created while starting or delegating a run. Secret is only populated on creation.

type RunUpdateHub added in v0.1.56

type RunUpdateHub struct {
	// contains filtered or unexported fields
}

RunUpdateHub adapts the generic event-wake infrastructure to Run IDs. It carries no Run fields or event payloads; every waiter must re-read PostgreSQL after a wake.

func NewRunUpdateHub added in v0.1.56

func NewRunUpdateHub(infrastructure *eventwake.Infrastructure) *RunUpdateHub

func (*RunUpdateHub) Healthy added in v0.1.56

func (h *RunUpdateHub) Healthy() bool

func (*RunUpdateHub) SubscribeRun added in v0.1.56

func (h *RunUpdateHub) SubscribeRun(runID uuid.UUID) (RunUpdateSubscription, error)

type RunUpdateSource added in v0.1.56

type RunUpdateSource interface {
	SubscribeRun(uuid.UUID) (RunUpdateSubscription, error)
	Healthy() bool
}

type RunUpdateSubscription added in v0.1.56

type RunUpdateSubscription interface {
	Wait(context.Context) error
	Close()
}

type RuntimeAdmissionIdentity added in v0.1.56

type RuntimeAdmissionIdentity struct {
	AgentID uuid.UUID
	NodeID  uuid.UUID
}

RuntimeAdmissionIdentity contains only durable, non-secret identifiers. The limiter hashes this value before it is used as an internal key; Agent Token plaintext and certificate material are never accepted by this API.

type RuntimeAdmissionLimitConfig added in v0.1.56

type RuntimeAdmissionLimitConfig struct {
	HTTPRequestsPerSecond      int
	HTTPBurst                  int
	WebSocketMessagesPerSecond int
	WebSocketMessageBurst      int
	MaxWebSocketsPerIdentity   int
}

type RuntimeAdmissionLimiter added in v0.1.56

type RuntimeAdmissionLimiter interface {
	AllowHTTP(RuntimeAdmissionIdentity) bool
	AcquireWebSocket(RuntimeAdmissionIdentity) (release func(), allowed bool)
	AllowWebSocketMessage(RuntimeAdmissionIdentity) bool
}

RuntimeAdmissionLimiter applies authenticated, per-principal admission to Runtime HTTP requests and WebSocket traffic. Implementations must not derive keys from request source IP because the raw mTLS edge intentionally preserves TLS bytes, not the downstream socket address.

func NewRuntimeAdmissionLimiter added in v0.1.56

func NewRuntimeAdmissionLimiter(config RuntimeAdmissionLimitConfig) RuntimeAdmissionLimiter

type RuntimeAssignmentRejectOutcome added in v0.1.56

type RuntimeAssignmentRejectOutcome string
const (
	RuntimeOfferRejected RuntimeAssignmentRejectOutcome = "offer_rejected"
	RuntimeLeaseRevoked  RuntimeAssignmentRejectOutcome = "lease_revoked"
)

type RuntimeAssignmentRejectReason added in v0.1.56

type RuntimeAssignmentRejectReason string
const (
	RuntimeRejectNodeAtCapacity         RuntimeAssignmentRejectReason = "NODE_AT_CAPACITY"
	RuntimeRejectNodeDraining           RuntimeAssignmentRejectReason = "NODE_DRAINING"
	RuntimeRejectClientUpgradeRequired  RuntimeAssignmentRejectReason = "RUNTIME_CLIENT_UPGRADE_REQUIRED"
	RuntimeRejectRequiredFeatureMissing RuntimeAssignmentRejectReason = "RUNTIME_REQUIRED_FEATURE_MISSING"
)

type RuntimeAttemptIdentity added in v0.1.56

type RuntimeAttemptIdentity struct {
	RunID            uuid.UUID  `json:"run_id"`
	AttemptID        uuid.UUID  `json:"attempt_id"`
	LeaseID          uuid.UUID  `json:"lease_id"`
	FencingToken     int64      `json:"fencing_token"`
	NodeID           *uuid.UUID `json:"node_id,omitempty"`
	AgentID          uuid.UUID  `json:"agent_id"`
	WorkerID         *string    `json:"worker_id,omitempty"`
	RuntimeSessionID *uuid.UUID `json:"runtime_session_id,omitempty"`
}

RuntimeAttemptIdentity is the immutable identity from the Runtime wire contract. AttemptNo is intentionally absent: Core derives it from the locked run_attempts row after assignment confirmation.

NodeID, RuntimeSessionID, and WorkerID are required for runtime attempts and absent for Core-owned HTTP/MCP attempts.

type RuntimeAuthenticationMode added in v0.1.56

type RuntimeAuthenticationMode string
const (
	RuntimeAuthenticationMTLS      RuntimeAuthenticationMode = "mtls"
	RuntimeAuthenticationTokenOnly RuntimeAuthenticationMode = "token_only"
)

type RuntimeCancelState added in v0.1.56

type RuntimeCancelState string
const (
	RuntimeCancelRequested   RuntimeCancelState = "requested"
	RuntimeCancelDelivered   RuntimeCancelState = "delivered"
	RuntimeCancelStopping    RuntimeCancelState = "stopping"
	RuntimeCancelStopped     RuntimeCancelState = "stopped"
	RuntimeCancelUnsupported RuntimeCancelState = "unsupported"
	RuntimeCancelFailed      RuntimeCancelState = "failed"
	RuntimeCancelUnconfirmed RuntimeCancelState = "unconfirmed"
)

type RuntimeCancellationCoordinator added in v0.1.56

type RuntimeCancellationCoordinator struct {
	// contains filtered or unexported fields
}

RuntimeCancellationCoordinator linearizes owner cancellation, durable command delivery, and Node stop acknowledgement through PostgreSQL.

func NewRuntimeCancellationCoordinator added in v0.1.56

func NewRuntimeCancellationCoordinator(pool *pgxpool.Pool) *RuntimeCancellationCoordinator

func (*RuntimeCancellationCoordinator) AckCancel added in v0.1.56

AckCancel advances stop evidence. Only a terminal stop ACK ends the target Attempt and releases its capacity, guarded by the Attempt-owned CAS.

func (*RuntimeCancellationCoordinator) AcknowledgeCoreStopped added in v0.1.56

func (c *RuntimeCancellationCoordinator) AcknowledgeCoreStopped(
	ctx context.Context,
	coreInstanceID uuid.UUID,
	identity RuntimeAttemptIdentity,
) (RunCancellationState, error)

AcknowledgeCoreStopped is called only after the Core-owned HTTP/MCP call has returned. It records stopped evidence and ends the immutable Attempt in one Run -> Attempt -> Cancellation transaction; the owner request path never claims that an in-flight goroutine has already stopped.

func (*RuntimeCancellationCoordinator) CancelOwnedRun added in v0.1.56

func (c *RuntimeCancellationCoordinator) CancelOwnedRun(
	ctx context.Context,
	requesterID, runID uuid.UUID,
	reason string,
) (RuntimeCancellationResult, error)

CancelOwnedRun atomically creates cancellation evidence and the complete public canceled terminal fact. Any active Attempt remains unfinished until its executor has actually stopped (or the deadline reaper records an unconfirmed stop after a Core/Node crash).

func (*RuntimeCancellationCoordinator) NextCommand added in v0.1.56

NextCommand returns one at-least-once durable cancellation command for the authenticated Session. A nil command with no error is a normal empty poll. Contention is distinct from an empty queue so transports retain the wake.

func (*RuntimeCancellationCoordinator) PollCommands added in v0.1.56

func (*RuntimeCancellationCoordinator) ReapExpiredCancellation added in v0.1.56

func (c *RuntimeCancellationCoordinator) ReapExpiredCancellation(
	ctx context.Context,
) (*RunCancellationState, error)

ReapExpiredCancellation converts one overdue cancellation into durable unconfirmed stop evidence, ends its fenced Attempt, and releases capacity in the same transaction. A nil result means no cancellation is currently due.

func (*RuntimeCancellationCoordinator) ReapExpiredCancellations added in v0.1.56

func (c *RuntimeCancellationCoordinator) ReapExpiredCancellations(
	ctx context.Context,
	limit int,
) (int, error)

ReapExpiredCancellations drains at most limit overdue cancellations. Each winner commits independently so one long batch never holds unrelated locks.

type RuntimeCancellationResult added in v0.1.56

type RuntimeCancellationResult struct {
	Run          db.Run
	Cancellation db.RunCancellation
	Replayed     bool
}

RuntimeCancellationResult is the durable owner-facing result. Replayed is true when the Run was already canceled by the same immutable cancellation evidence; no terminal artifact or command signal is emitted twice.

type RuntimeClaimRequest added in v0.1.56

type RuntimeClaimRequest struct {
	RuntimeSessionID uuid.UUID `json:"runtime_session_id" runtime:"required"`
	Capacity         int64     `json:"capacity" runtime:"required"`
	Inflight         int64     `json:"inflight" runtime:"required"`
}

type RuntimeClusterControlSnapshot added in v0.1.56

type RuntimeClusterControlSnapshot struct {
	Mode             RuntimeClusterMode `json:"mode"`
	ExpectedReplicas int32              `json:"expected_replicas"`
}

type RuntimeClusterCoordinator added in v0.1.56

type RuntimeClusterCoordinator struct {
	// contains filtered or unexported fields
}

func NewRuntimeClusterCoordinator added in v0.1.56

func NewRuntimeClusterCoordinator(
	pool *pgxpool.Pool,
	signalBus RuntimeSignalBus,
	identity RuntimeClusterIdentity,
	requireSignal bool,
) (*RuntimeClusterCoordinator, error)

NewRuntimeClusterCoordinator creates a coordinator with the protocol's fixed five-second heartbeat and fifteen-second live-member window.

func (*RuntimeClusterCoordinator) Close added in v0.1.56

Close stops future heartbeats before deleting this process's membership. Stale rows left by SIGKILL naturally stop counting after the live window.

func (*RuntimeClusterCoordinator) Readiness added in v0.1.56

func (*RuntimeClusterCoordinator) Start added in v0.1.56

func (c *RuntimeClusterCoordinator) Start(parent context.Context)

Start registers immediately, then refreshes membership every five seconds. A transient dependency or database failure is reported through readiness and retried; it never stops the database reconciliation workers.

type RuntimeClusterIdentity added in v0.1.56

type RuntimeClusterIdentity struct {
	InstanceID            uuid.UUID `json:"instance_id"`
	ReleaseVersion        string    `json:"release_version"`
	ReleaseCommit         string    `json:"release_commit"`
	SchemaVersion         int32     `json:"schema_version"`
	SchemaChecksum        string    `json:"schema_checksum"`
	RuntimeContractID     string    `json:"runtime_contract_id"`
	RuntimeContractDigest string    `json:"runtime_contract_digest"`
}

RuntimeClusterIdentity is immutable process evidence stored on every heartbeat. Release fields deliberately come from deployment metadata rather than a mutable database setting.

type RuntimeClusterMemberSnapshot added in v0.1.56

type RuntimeClusterMemberSnapshot struct {
	RuntimeClusterIdentity
	HeartbeatAt time.Time `json:"heartbeat_at"`
	Draining    bool      `json:"draining"`
	Ready       bool      `json:"ready"`
}

type RuntimeClusterMode added in v0.1.56

type RuntimeClusterMode string

type RuntimeClusterOperation added in v0.1.56

type RuntimeClusterOperation string

type RuntimeClusterReadiness added in v0.1.56

type RuntimeClusterReadiness struct {
	Ready             bool               `json:"ready"`
	Status            string             `json:"status"`
	Reasons           []string           `json:"reasons,omitempty"`
	Mode              RuntimeClusterMode `json:"mode,omitempty"`
	ExpectedReplicas  int32              `json:"expected_replicas,omitempty"`
	LiveReplicas      int                `json:"live_replicas"`
	InstanceID        uuid.UUID          `json:"instance_id"`
	ReleaseVersion    string             `json:"release_version"`
	ReleaseCommit     string             `json:"release_commit"`
	SchemaVersion     int32              `json:"schema_version"`
	SchemaChecksum    string             `json:"schema_checksum"`
	RuntimeContractID string             `json:"runtime_contract_id"`
	DatabaseTime      *time.Time         `json:"database_time,omitempty"`
}

RuntimeClusterReadiness is safe for an unauthenticated health endpoint: it contains version/contract evidence and stable reason codes, never secrets, Run payloads, credentials, or infrastructure error strings.

func (RuntimeClusterReadiness) HTTPStatus added in v0.1.56

func (r RuntimeClusterReadiness) HTTPStatus() int

type RuntimeClusterRepository added in v0.1.56

type RuntimeClusterRepository interface {
	UpsertMember(context.Context, RuntimeClusterIdentity, bool, bool) error
	Snapshot(context.Context, time.Duration) (RuntimeClusterSnapshot, error)
	CloseMember(context.Context, uuid.UUID) error
}

type RuntimeClusterSnapshot added in v0.1.56

type RuntimeClusterSnapshot struct {
	DatabaseTime  time.Time                      `json:"database_time"`
	Control       RuntimeClusterControlSnapshot  `json:"control"`
	CurrentSchema RuntimeSchemaContractSnapshot  `json:"current_schema"`
	LiveMembers   []RuntimeClusterMemberSnapshot `json:"live_members"`
}

type RuntimeCommand added in v0.1.56

type RuntimeCommand[P any] struct {
	Type    RuntimeMessageType `json:"type" runtime:"required"`
	Payload P                  `json:"payload" runtime:"required"`
}

type RuntimeCommandsResponse added in v0.1.56

type RuntimeCommandsResponse struct {
	Commands     []PendingCommand `json:"commands" runtime:"required"`
	DatabaseTime time.Time        `json:"database_time" runtime:"required"`
}

type RuntimeConnectionIdentity added in v0.1.56

type RuntimeConnectionIdentity struct {
	RuntimeSessionID uuid.UUID `json:"runtime_session_id"`
	SessionEpoch     int64     `json:"session_epoch"`
	AttachmentID     uuid.UUID `json:"attachment_id"`
}

RuntimeConnectionIdentity is the immutable attachment-generation fence for one live Runtime transport. A lifecycle signal must match every field before it can affect a connection.

type RuntimeConnectionRegistration added in v0.1.56

type RuntimeConnectionRegistration struct {
	Identity     RuntimeConnectionIdentity
	CredentialID uuid.UUID
}

type RuntimeCredentialConnectionValidator added in v0.1.56

type RuntimeCredentialConnectionValidator interface {
	Validate(
		context.Context,
		[]RuntimeConnectionRegistration,
	) ([]RuntimeCredentialValidationResult, error)
}

func NewPostgresRuntimeCredentialConnectionValidator added in v0.1.56

func NewPostgresRuntimeCredentialConnectionValidator(
	pool *pgxpool.Pool,
	coreInstanceID uuid.UUID,
) RuntimeCredentialConnectionValidator

type RuntimeCredentialProjectionResult added in v0.1.56

type RuntimeCredentialProjectionResult struct {
	Registration RuntimeConnectionRegistration
	State        RuntimeCredentialProjectionState
}

type RuntimeCredentialProjectionState added in v0.1.56

type RuntimeCredentialProjectionState string
const (
	RuntimeCredentialProjectionMissing   RuntimeCredentialProjectionState = "missing"
	RuntimeCredentialProjectionActive    RuntimeCredentialProjectionState = "active"
	RuntimeCredentialProjectionRevoked   RuntimeCredentialProjectionState = "revoked"
	RuntimeCredentialProjectionMalformed RuntimeCredentialProjectionState = "malformed"
)

type RuntimeCredentialProjectionStore added in v0.1.56

type RuntimeCredentialProjectionStore interface {
	Check(context.Context, []RuntimeConnectionRegistration) ([]RuntimeCredentialProjectionResult, error)
	MarkActive(context.Context, []RuntimeConnectionRegistration) error
}

RuntimeCredentialProjectionStore is a bounded Redis cache. PostgreSQL remains authoritative: missing or malformed values must be revalidated there, never interpreted as active.

type RuntimeCredentialProjectionStoreProvider added in v0.1.56

type RuntimeCredentialProjectionStoreProvider interface {
	RuntimeCredentialProjectionStore() (RuntimeCredentialProjectionStore, error)
}

type RuntimeCredentialReconciler added in v0.1.56

type RuntimeCredentialReconciler struct {
	// contains filtered or unexported fields
}

func (*RuntimeCredentialReconciler) Run added in v0.1.56

type RuntimeCredentialReconcilerConfig added in v0.1.56

type RuntimeCredentialReconcilerConfig struct {
	Interval         time.Duration
	AuditInterval    time.Duration
	QueryTimeout     time.Duration
	ProjectionSettle time.Duration
	BatchSize        int
	Observer         WorkerObserver
	// contains filtered or unexported fields
}

type RuntimeCredentialValidationResult added in v0.1.56

type RuntimeCredentialValidationResult struct {
	Registration RuntimeConnectionRegistration
	Valid        bool
}

type RuntimeDeadLetterListItem added in v0.1.56

type RuntimeDeadLetterListItem struct {
	DeadLetterID     string     `json:"dead_letter_id"`
	RunID            string     `json:"run_id"`
	AgentID          string     `json:"agent_id"`
	AgentSlug        string     `json:"agent_slug"`
	AgentName        string     `json:"agent_name"`
	Status           string     `json:"status"`
	DispatchState    string     `json:"dispatch_state"`
	AttemptCount     int32      `json:"attempt_count"`
	MaxAttempts      int32      `json:"max_attempts"`
	FinalAttemptID   string     `json:"final_attempt_id,omitempty"`
	FinalAttemptNo   int32      `json:"final_attempt_no"`
	ErrorCode        string     `json:"error_code,omitempty"`
	ErrorMessage     string     `json:"error_message,omitempty"`
	ErrorDetail      string     `json:"error_detail_redacted,omitempty"`
	ReasonCode       string     `json:"reason_code"`
	Reason           string     `json:"reason_redacted,omitempty"`
	DeadLetteredAt   *time.Time `json:"dead_lettered_at,omitempty"`
	CreatedAt        time.Time  `json:"created_at"`
	ReplayOfRunID    string     `json:"replay_of_run_id,omitempty"`
	ReplayedAsRunIDs []string   `json:"replayed_as_run_ids"`
}

type RuntimeDeadLetterListResponse added in v0.1.56

type RuntimeDeadLetterListResponse struct {
	Items  []RuntimeDeadLetterListItem `json:"items"`
	Total  int32                       `json:"total"`
	Limit  int32                       `json:"limit"`
	Offset int32                       `json:"offset"`
}

RuntimeDeadLetterListResponse is the admin-only, input-free DLQ inventory. It intentionally exposes only redacted execution evidence and replay lineage.

type RuntimeDeadlineReconciler added in v0.1.56

type RuntimeDeadlineReconciler struct {
	// contains filtered or unexported fields
}

RuntimeDeadlineReconciler is a standalone worker boundary. It is not wired to coreapi here: callers choose scheduling, cadence, and shutdown.

func NewRuntimeDeadlineReconciler added in v0.1.56

func NewRuntimeDeadlineReconciler(
	pool *pgxpool.Pool,
	retryPlanner ResultRetryPlanner,
) *RuntimeDeadlineReconciler

NewRuntimeDeadlineReconciler builds the lease/deadline worker. A nil retry planner uses the same bounded retry policy as ResultFinalizer.

func (*RuntimeDeadlineReconciler) ReconcileBatch added in v0.1.56

func (r *RuntimeDeadlineReconciler) ReconcileBatch(
	ctx context.Context,
	limit int,
) (RuntimeReconcileBatchResult, error)

ReconcileBatch discovers at most limit due Runs with PostgreSQL's clock and commits each winner independently. A failed candidate rolls back its entire Attempt/capacity/Run/artifact transaction and stops the batch.

type RuntimeDelegationAPI added in v0.1.56

type RuntimeDelegationAPI interface {
	CallAgent(context.Context, RuntimeDelegationAuthorization) (RunSummary, error)
}

type RuntimeDelegationAuthorization added in v0.1.56

type RuntimeDelegationAuthorization struct {
	Device            RuntimeDeviceIdentity
	InvocationContext string
	InvocationToken   string
	InvocationProof   string
	IdempotencyKey    string
	ProofRequest      RuntimeInvocationProofRequest
}

RuntimeDelegationAuthorization contains transport-authenticated evidence. ProofRequest.Body is the exact byte sequence read from HTTP and later decoded as Payload; callers must not marshal the body a second time.

type RuntimeDelegationService added in v0.1.56

type RuntimeDelegationService struct {
	// contains filtered or unexported fields
}

RuntimeDelegationService creates a child Run only while the signed parent Attempt is accepted, current, uncanceled, and unexpired. PostgreSQL locks and time are authoritative; the preliminary capability read is only used to assemble a candidate request before the creation transaction.

func NewRuntimeDelegationService added in v0.1.56

func NewRuntimeDelegationService(
	pool *pgxpool.Pool,
	runtimeService *Service,
	verifier RuntimeInvocationVerifier,
) *RuntimeDelegationService

func (*RuntimeDelegationService) CallAgent added in v0.1.56

func (*RuntimeDelegationService) ReadDelegatedRun added in v0.1.59

ReadDelegatedRun is a separate, explicitly negotiated authority from the original call-agent audience. It never falls back to owner-facing Run APIs. Authorization and the read share a transaction: a terminal/canceled/replaced parent Attempt cannot race the read after its authority has been checked.

func (*RuntimeDelegationService) ResolveInvocationDevice added in v0.1.56

func (s *RuntimeDelegationService) ResolveInvocationDevice(
	ctx context.Context,
	invocationToken string,
) (RuntimeDeviceIdentity, error)

ResolveInvocationDevice verifies the assignment-scoped Bearer capability against PostgreSQL time, then resolves the same durable token-to-key binding used by ordinary Runtime requests. It is used only when mTLS is explicitly disabled; request bodies and headers never establish Node identity.

type RuntimeDeviceAuthenticator added in v0.1.56

type RuntimeDeviceAuthenticator interface {
	AuthenticateHTTP(context.Context, *http.Request) (RuntimeDeviceIdentity, error)
}

RuntimeDeviceAuthenticator authenticates the independently enrolled Node device. An implementation must use the verified TLS peer certificate, never forwarded certificate or identity headers.

type RuntimeDeviceIdentity added in v0.1.56

type RuntimeDeviceIdentity struct {
	NodeID             uuid.UUID                 `json:"node_id"`
	AuthenticationMode RuntimeAuthenticationMode `json:"-"`
	// CertificateSerial is the stable enrolled credential serial used by
	// durable Sessions. PresentedCertificateSerial is the current rotating leaf
	// serial and is never accepted from request data.
	CertificateSerial            string `json:"certificate_serial"`
	PresentedCertificateSerial   string `json:"-"`
	CertificateFingerprintSHA256 string `json:"certificate_fingerprint_sha256"`
	PublicKeyThumbprintSHA256    string `json:"public_key_thumbprint_sha256"`
}

RuntimeDeviceIdentity is the immutable Node identity established by either a verified mTLS peer certificate or a durable token-only credential binding. The Node selector may come from a strict header in token-only mode, but the durable identity values are always resolved or deterministically derived by the server and never trusted from a Runtime hello body.

type RuntimeDispatchState added in v0.1.56

type RuntimeDispatchState string
const (
	RuntimeDispatchPending    RuntimeDispatchState = "pending"
	RuntimeDispatchOffered    RuntimeDispatchState = "offered"
	RuntimeDispatchExecuting  RuntimeDispatchState = "executing"
	RuntimeDispatchRetryWait  RuntimeDispatchState = "retry_wait"
	RuntimeDispatchTerminal   RuntimeDispatchState = "terminal"
	RuntimeDispatchDeadLetter RuntimeDispatchState = "dead_letter"
)

type RuntimeDispatchWakeReconcileResult added in v0.1.56

type RuntimeDispatchWakeReconcileResult struct {
	Scanned int  `json:"scanned"`
	Woken   int  `json:"woken"`
	Wrapped bool `json:"wrapped"`
	// Counts refer to Agent queues, not Run assignments. AgentIDs is bounded
	// independently of the scan batch; payloads and credentials never appear.
	Queued    int      `json:"queued"`
	Coalesced int      `json:"coalesced"`
	AgentIDs  []string `json:"agent_ids,omitempty"`
}

type RuntimeDispatchWakeReconciler added in v0.1.56

type RuntimeDispatchWakeReconciler struct {
	// contains filtered or unexported fields
}

RuntimeDispatchWakeReconciler is the durable safety net for lossy wake hints. It performs one process-level, bounded PostgreSQL discovery pass and wakes only Agents with claimable Runtime work and a registered local waiter. WebSocket and Pull connections remain free of periodic database polling.

func NewRuntimeDispatchWakeReconciler added in v0.1.56

func NewRuntimeDispatchWakeReconciler(pool *pgxpool.Pool, hub *RuntimeWakeHub) *RuntimeDispatchWakeReconciler

func (*RuntimeDispatchWakeReconciler) ReconcileOnce added in v0.1.56

ReconcileOnce advances a process-local UUID cursor so a large offline backlog cannot permanently starve later Agent IDs. A cursor wrap may issue a second bounded query; the normal empty/partial-page path issues exactly one.

type RuntimeDispatchWakeReconcilerConfig added in v0.1.56

type RuntimeDispatchWakeReconcilerConfig struct {
	Interval  time.Duration
	BatchSize int
	Observer  WorkerObserver
}

type RuntimeDrainMessage added in v0.1.56

type RuntimeDrainMessage = RuntimeTypedEnvelope[RuntimeDrainPayload]

type RuntimeDrainPayload added in v0.1.56

type RuntimeDrainPayload struct {
	DeadlineAt time.Time `json:"deadline_at" runtime:"required"`
	ReasonCode string    `json:"reason_code" runtime:"required"`
	Capacity   int64     `json:"capacity" runtime:"required"`
	Inflight   int64     `json:"inflight" runtime:"required"`
}

type RuntimeEnvelope added in v0.1.56

type RuntimeEnvelope struct {
	RuntimeEnvelopeFields
	Payload json.RawMessage `json:"payload" runtime:"required"`
}

RuntimeEnvelope is the routing form of a WebSocket message. Payload remains raw until Type has selected the one permitted concrete payload schema.

func DecodeRuntimeEnvelope added in v0.1.56

func DecodeRuntimeEnvelope(reader io.Reader) (RuntimeEnvelope, error)

DecodeRuntimeEnvelope decodes one complete WebSocket message. It validates only the common envelope; callers then use DecodeRuntimeMessagePayload after dispatching on Type so unknown payload fields cannot slip through.

func ParseRuntimeEnvelope added in v0.1.56

func ParseRuntimeEnvelope(frame []byte) (RuntimeEnvelope, error)

ParseRuntimeEnvelope is the byte-slice form used by WebSocket readers.

type RuntimeEnvelopeFields added in v0.1.56

type RuntimeEnvelopeFields struct {
	ProtocolVersion   int                `json:"protocol_version" runtime:"required"`
	RuntimeContractID string             `json:"runtime_contract_id" runtime:"required"`
	MessageID         uuid.UUID          `json:"message_id" runtime:"required"`
	ReplyToMessageID  *uuid.UUID         `json:"reply_to_message_id,omitempty"`
	Type              RuntimeMessageType `json:"type" runtime:"required"`
	SentAt            time.Time          `json:"sent_at" runtime:"required"`
}

RuntimeEnvelopeFields are common to every Runtime WebSocket message. Payload is deliberately not part of this struct so concrete messages keep a strongly typed payload while the router can use RuntimeEnvelope below.

type RuntimeError added in v0.1.56

type RuntimeError struct {
	Error RuntimeErrorBody `json:"error" runtime:"required"`
}

RuntimeError is the stable HTTP error envelope from the canonical contract.

type RuntimeErrorBody added in v0.1.56

type RuntimeErrorBody struct {
	Code                 RuntimeErrorCode     `json:"code" runtime:"required"`
	Message              string               `json:"message" runtime:"required"`
	Retryable            bool                 `json:"retryable,omitempty"`
	MissingEventRanges   []EventRange         `json:"missing_event_ranges,omitempty"`
	CurrentRunStatus     RuntimeRunStatus     `json:"current_run_status,omitempty"`
	CurrentDispatchState RuntimeDispatchState `json:"current_dispatch_state,omitempty"`
}

type RuntimeErrorCode added in v0.1.56

type RuntimeErrorCode string

RuntimeErrorCode is the stable error vocabulary shared by HTTP and WebSocket Runtime transports.

const (
	RuntimeErrorBadRequest             RuntimeErrorCode = "BAD_REQUEST"
	RuntimeErrorUnauthorized           RuntimeErrorCode = "UNAUTHORIZED"
	RuntimeErrorForbidden              RuntimeErrorCode = "FORBIDDEN"
	RuntimeErrorPermissionDenied       RuntimeErrorCode = "PERMISSION_DENIED"
	RuntimeErrorNotFound               RuntimeErrorCode = "NOT_FOUND"
	RuntimeErrorConflict               RuntimeErrorCode = "CONFLICT"
	RuntimeErrorValidationFailed       RuntimeErrorCode = "VALIDATION_FAILED"
	RuntimeErrorRateLimited            RuntimeErrorCode = "RATE_LIMITED"
	RuntimeErrorInternal               RuntimeErrorCode = "INTERNAL_ERROR"
	RuntimeErrorServiceUnavailable     RuntimeErrorCode = "SERVICE_UNAVAILABLE"
	RuntimeErrorIdempotencyKeyReused   RuntimeErrorCode = "IDEMPOTENCY_KEY_REUSED"
	RuntimeErrorRunAlreadyTerminal     RuntimeErrorCode = "RUN_ALREADY_TERMINAL"
	RuntimeErrorStaleLease             RuntimeErrorCode = "STALE_LEASE"
	RuntimeErrorLeaseExpired           RuntimeErrorCode = "LEASE_EXPIRED"
	RuntimeErrorLeaseIdentityMismatch  RuntimeErrorCode = "LEASE_IDENTITY_MISMATCH"
	RuntimeErrorResultIDConflict       RuntimeErrorCode = "RESULT_ID_CONFLICT"
	RuntimeErrorEventIDConflict        RuntimeErrorCode = "EVENT_ID_CONFLICT"
	RuntimeErrorNodeAtCapacity         RuntimeErrorCode = "NODE_AT_CAPACITY"
	RuntimeErrorClientUpgradeRequired  RuntimeErrorCode = "RUNTIME_CLIENT_UPGRADE_REQUIRED"
	RuntimeErrorRequiredFeatureMissing RuntimeErrorCode = "RUNTIME_REQUIRED_FEATURE_MISSING"
	RuntimeErrorRunCancelRequested     RuntimeErrorCode = "RUN_CANCEL_REQUESTED"
	RuntimeErrorRunCancelUnconfirmed   RuntimeErrorCode = "RUN_CANCEL_UNCONFIRMED"
	RuntimeErrorRetryExhausted         RuntimeErrorCode = "RUNTIME_RETRY_EXHAUSTED"
	RuntimeErrorDispatchTimeout        RuntimeErrorCode = "RUNTIME_DISPATCH_TIMEOUT"
	RuntimeErrorRunDeadlineExceeded    RuntimeErrorCode = "RUN_DEADLINE_EXCEEDED"
	RuntimeErrorEventsMissing          RuntimeErrorCode = "EVENTS_MISSING"
	RuntimeErrorReplayInputUnavailable RuntimeErrorCode = "REPLAY_INPUT_UNAVAILABLE"
	RuntimeErrorEndpointResultUnknown  RuntimeErrorCode = "ENDPOINT_RESULT_UNKNOWN"
	RuntimeErrorSessionConflict        RuntimeErrorCode = "RUNTIME_SESSION_CONFLICT"
	RuntimeErrorSpoolCorrupt           RuntimeErrorCode = "RUNTIME_SPOOL_CORRUPT"
)

type RuntimeErrorMessage added in v0.1.56

type RuntimeErrorMessage = RuntimeTypedEnvelope[RuntimeErrorBody]

type RuntimeEventAck added in v0.1.56

type RuntimeEventAck struct {
	ClientEventID  uuid.UUID `json:"client_event_id"`
	ClientEventSeq int64     `json:"client_event_seq"`
	Sequence       int32     `json:"sequence"`
	Replayed       bool      `json:"replayed"`
	Inserted       bool      `json:"-"`
	EventID        uuid.UUID `json:"-"`
	CreatedAt      time.Time `json:"-"`
}

RuntimeEventAck matches RunEventAckPayload. Inserted is internal-only: an upper layer may publish side effects exclusively when it is true.

type RuntimeEventError added in v0.1.56

type RuntimeEventError struct {
	Code          RuntimeEventErrorCode `json:"code"`
	MissingRanges []EventRange          `json:"missing_ranges,omitempty"`
	// contains filtered or unexported fields
}

RuntimeEventError carries a stable reason code and, for EVENTS_MISSING, the exact client sequence ranges that must be uploaded before retrying Result.

func (*RuntimeEventError) Error added in v0.1.56

func (e *RuntimeEventError) Error() string

func (*RuntimeEventError) Unwrap added in v0.1.56

func (e *RuntimeEventError) Unwrap() error

type RuntimeEventErrorCode added in v0.1.56

type RuntimeEventErrorCode string

Runtime event errors are transport-neutral. HTTP, WebSocket, gRPC, and MCP adapters map these codes to their own error envelopes without parsing text.

const (
	RuntimeEventErrorIDConflict            RuntimeEventErrorCode = "EVENT_ID_CONFLICT"
	RuntimeEventErrorStaleLease            RuntimeEventErrorCode = "STALE_LEASE"
	RuntimeEventErrorLeaseExpired          RuntimeEventErrorCode = "LEASE_EXPIRED"
	RuntimeEventErrorLeaseIdentityMismatch RuntimeEventErrorCode = "LEASE_IDENTITY_MISMATCH"
	RuntimeEventErrorRunAlreadyTerminal    RuntimeEventErrorCode = "RUN_ALREADY_TERMINAL"
	RuntimeEventErrorEventsMissing         RuntimeEventErrorCode = "EVENTS_MISSING"
)

type RuntimeEventPrincipal added in v0.1.56

type RuntimeEventPrincipal struct {
	AgentID                         uuid.UUID
	RuntimeContractDigest           string
	CredentialID                    *uuid.UUID
	NodeID                          *uuid.UUID
	WorkerID                        *string
	RuntimeSessionID                *uuid.UUID
	CoreInstanceID                  *uuid.UUID
	AttachmentID                    *uuid.UUID
	DeviceCertificateSerial         *string
	DevicePublicKeyThumbprintSHA256 *string
}

RuntimeEventPrincipal is the already-authenticated execution principal. EventStore never accepts a token or transport assertion in place of it.

type RuntimeEventProjector added in v0.1.56

type RuntimeEventProjector interface {
	AppendRuntimeEvent(context.Context, RuntimeEventPrincipal, RuntimeAttemptIdentity, RuntimeEventRequest) (RuntimeEventAck, error)
}

RuntimeEventProjector is the only execution-event entrypoint exposed to transport adapters. Persistence and message/artifact/callback projections must stay behind this boundary so WebSocket and Pull cannot diverge.

type RuntimeEventRequest added in v0.1.56

type RuntimeEventRequest struct {
	ClientEventID  uuid.UUID      `json:"client_event_id"`
	ClientEventSeq int64          `json:"client_event_seq"`
	EventType      string         `json:"event_type"`
	Payload        map[string]any `json:"payload"`
}

RuntimeEventRequest is the transport-neutral RunEventPayload body after its AttemptIdentity has been separated by the caller.

type RuntimeHTTPController added in v0.1.56

type RuntimeHTTPController struct {
	// contains filtered or unexported fields
}

RuntimeHTTPController is the strict HTTP transport adapter for the durable Runtime state machine.

func NewRuntimeHTTPController added in v0.1.56

func NewRuntimeHTTPController(dependencies RuntimeHTTPDependencies) *RuntimeHTTPController

func (*RuntimeHTTPController) AckAssignment added in v0.1.56

func (h *RuntimeHTTPController) AckAssignment(c echo.Context) error

func (*RuntimeHTTPController) AckCancel added in v0.1.56

func (h *RuntimeHTTPController) AckCancel(c echo.Context) error

func (*RuntimeHTTPController) AppendEvent added in v0.1.56

func (h *RuntimeHTTPController) AppendEvent(c echo.Context) error

func (*RuntimeHTTPController) AuthenticateAgentRequest added in v0.1.56

AuthenticateAgentRequest authenticates the Agent Token and device mTLS identity used by server-owned compatibility adapters. It deliberately exposes only the already validated Runtime principal; adapters must derive all user and Agent ownership data from that principal on Core.

func (*RuntimeHTTPController) CallAgent added in v0.1.56

func (h *RuntimeHTTPController) CallAgent(c echo.Context) error

CallAgent authenticates the Node and both assignment-scoped signed capabilities before the exact request bytes are decoded as business input. The invocation capability, not a long-lived Agent Token, is the Bearer credential for this endpoint. Token-only transport resolves the same durable Node/key binding from that signed capability.

func (*RuntimeHTTPController) ClaimRun added in v0.1.56

func (h *RuntimeHTTPController) ClaimRun(c echo.Context) error

func (*RuntimeHTTPController) CloseSession added in v0.1.56

func (h *RuntimeHTTPController) CloseSession(c echo.Context) error

func (*RuntimeHTTPController) CreateSession added in v0.1.56

func (h *RuntimeHTTPController) CreateSession(c echo.Context) error

func (*RuntimeHTTPController) DrainSession added in v0.1.56

func (h *RuntimeHTTPController) DrainSession(c echo.Context) error

DrainSession is the Pull half of the server-authoritative drain handshake. The payload cannot establish identity or capacity: Core authenticates the path Session and attachment, atomically commits draining/capacity=0, then returns only the database receipt.

func (*RuntimeHTTPController) FinalizeResult added in v0.1.56

func (h *RuntimeHTTPController) FinalizeResult(c echo.Context) error

func (*RuntimeHTTPController) HeartbeatSession added in v0.1.56

func (h *RuntimeHTTPController) HeartbeatSession(c echo.Context) error

func (*RuntimeHTTPController) PollCommands added in v0.1.56

func (h *RuntimeHTTPController) PollCommands(c echo.Context) error

PollCommands long-polls commands for one explicit Session. A token and device may legitimately own several worker Sessions, so resolving a command poll from authentication alone would deliver work to an arbitrary process.

func (*RuntimeHTTPController) ReadDelegatedRun added in v0.1.59

func (h *RuntimeHTTPController) ReadDelegatedRun(c echo.Context) error

func (*RuntimeHTTPController) Register added in v0.1.56

func (h *RuntimeHTTPController) Register(api *echo.Group)

func (*RuntimeHTTPController) RegisterAttachOnly added in v0.1.56

func (h *RuntimeHTTPController) RegisterAttachOnly(api *echo.Group)

RegisterAttachOnly mounts only the Runtime Session lifecycle required to prove SDK connectivity during a release cutover. Execution, command, resume, result, event, and delegation routes deliberately remain absent.

func (*RuntimeHTTPController) RejectAssignment added in v0.1.56

func (h *RuntimeHTTPController) RejectAssignment(c echo.Context) error

func (*RuntimeHTTPController) RenewLease added in v0.1.56

func (h *RuntimeHTTPController) RenewLease(c echo.Context) error

func (*RuntimeHTTPController) ResumeRuns added in v0.1.56

func (h *RuntimeHTTPController) ResumeRuns(c echo.Context) error

ResumeRuns authorizes recovery against the currently authenticated target Session. Attempt identities in the payload continue to name their immutable source Sessions; the body therefore cannot be resolved through an Attempt's source runtime_session_id.

func (*RuntimeHTTPController) SendBrowserObserverCommand added in v0.1.56

func (h *RuntimeHTTPController) SendBrowserObserverCommand(
	runtimeSessionID uuid.UUID,
	payload BrowserObserverCommandPayload,
) error

SendBrowserObserverCommand routes an observation command to the Worker held by this process. Frames and wakeups live in process memory, so a command that cannot be delivered locally must fail closed here rather than appear to start an observation no frame will ever reach.

func (*RuntimeHTTPController) SendBrowserViewerCommand added in v0.1.56

func (h *RuntimeHTTPController) SendBrowserViewerCommand(
	runtimeSessionID uuid.UUID,
	payload BrowserViewerCommandPayload,
) error

func (*RuntimeHTTPController) Shutdown added in v0.1.56

func (h *RuntimeHTTPController) Shutdown(ctx context.Context) error

Shutdown rejects new Runtime WebSockets, interrupts every hijacked connection, and waits until each handler has completed durable cleanup.

func (*RuntimeHTTPController) WebSocket added in v0.1.56

func (h *RuntimeHTTPController) WebSocket(c echo.Context) error

WebSocket authenticates both Agent Token and Node certificate before the HTTP connection is upgraded. No unauthenticated peer can consume a socket or create/attach a durable Session.

type RuntimeHTTPDependencies added in v0.1.56

type RuntimeHTTPDependencies struct {
	TokenValidator       RuntimeTokenValidator
	DeviceAuthenticator  RuntimeDeviceAuthenticator
	PrincipalBinder      RuntimePrincipalBinder
	TokenOnlyTransport   bool
	TransportPolicy      RuntimeTransportPolicyProvider
	Sessions             RuntimeSessionAPI
	Leases               RuntimeLeaseAPI
	EventProjector       RuntimeEventProjector
	Finalizer            RuntimeResultFinalizer
	Resume               RuntimeResumeAPI
	Delegation           RuntimeDelegationAPI
	Cancellations        RuntimeCancellationAPI
	WakeHub              *RuntimeWakeHub
	Presence             RuntimePresenceStore
	SessionLeases        *RuntimeSessionLeaseManager
	AdmissionLimiter     RuntimeAdmissionLimiter
	Observer             WorkerObserver
	BrowserControl       *BrowserHumanControl
	BrowserObservation   *BrowserObservation
	CoreInstanceID       uuid.UUID
	WebSocketConcurrency RuntimeWebSocketConcurrencyConfig
	// AttachOnly is a release-cutover safety mode. It permits authenticated
	// Session lifecycle traffic, but never claims Runs or accepts execution
	// events/results before the normal Core producer boundary is crossed.
	AttachOnly bool
}

RuntimeHTTPDependencies are deliberately narrow so the HTTP adapter has no database access and cannot reconstruct trusted principals from JSON.

type RuntimeHelloMessage added in v0.1.56

type RuntimeHelloMessage = RuntimeTypedEnvelope[RuntimeHelloPayload]

type RuntimeHelloPayload added in v0.1.56

type RuntimeHelloPayload struct {
	NodeID           uuid.UUID `json:"node_id" runtime:"required"`
	AgentID          uuid.UUID `json:"agent_id" runtime:"required"`
	WorkerID         string    `json:"worker_id" runtime:"required"`
	RuntimeSessionID uuid.UUID `json:"runtime_session_id" runtime:"required"`
	SessionEpoch     int64     `json:"session_epoch" runtime:"required"`
	NodeVersion      string    `json:"node_version" runtime:"required"`
	Capacity         int64     `json:"capacity" runtime:"required"`
	Features         []string  `json:"features" runtime:"required"`
	ContractDigest   string    `json:"contract_digest" runtime:"required"`
}

type RuntimeInvocationCapability added in v0.1.56

type RuntimeInvocationCapability struct {
	// Audience is immutable Session-negotiated authority. Empty preserves the
	// original create-only capability and its canonical bytes.
	Audience         string
	RunID            uuid.UUID
	AttemptID        uuid.UUID
	LeaseID          uuid.UUID
	FencingToken     int64
	AgentID          uuid.UUID
	CredentialID     uuid.UUID
	NodeID           uuid.UUID
	WorkerID         string
	RuntimeSessionID uuid.UUID
	InputSHA256      [sha256.Size]byte
	IssuedAt         time.Time
	ExpiresAt        time.Time
}

RuntimeInvocationCapability is the immutable authority delegated to one accepted runtime offer. It intentionally contains only an input digest, not user input or a long-lived Agent token.

type RuntimeInvocationCapabilityIssuer added in v0.1.56

type RuntimeInvocationCapabilityIssuer interface {
	Issue(RuntimeInvocationCapability) (nodeEnvelope, invocationToken string, err error)
}

RuntimeInvocationCapabilityIssuer issues the two assignment-scoped capabilities. Implementations must be deterministic for the same immutable Attempt evidence so a repeated claim returns byte-identical authority.

type RuntimeInvocationProofRequest added in v0.1.56

type RuntimeInvocationProofRequest struct {
	Method         string
	Path           string
	IdempotencyKey string
	Context        string
	Body           []byte
}

RuntimeInvocationProofRequest binds a delegated call to its exact HTTP request. Body is hashed before signing, so no request payload is embedded in a header.

func RuntimeInvocationProofRequestFromHTTP added in v0.1.56

func RuntimeInvocationProofRequestFromHTTP(r *http.Request, body []byte) RuntimeInvocationProofRequest

RuntimeInvocationProofRequestFromHTTP preserves the exact escaped path and body bytes used for proof verification. The caller must restore Body if a downstream JSON decoder still needs it.

type RuntimeInvocationSigner added in v0.1.56

type RuntimeInvocationSigner struct {
	// contains filtered or unexported fields
}

RuntimeInvocationSigner issues and verifies stateless short-lived runtime capabilities. The configured secret is domain-derived before use so these signatures cannot be confused with user JWTs or another HMAC protocol.

func NewRuntimeInvocationSigner added in v0.1.56

func NewRuntimeInvocationSigner(secret string) (*RuntimeInvocationSigner, error)

func NewRuntimeInvocationSignerKeyring added in v0.1.56

func NewRuntimeInvocationSignerKeyring(activeKeyID string, secrets map[string]string) (*RuntimeInvocationSigner, error)

NewRuntimeInvocationSignerKeyring creates a signer that issues with one active key while continuing to verify explicitly configured predecessor keys. Outstanding Attempt capabilities therefore survive a deliberate key rotation only for as long as operators retain the old key in this bounded ring.

func NewRuntimeInvocationSignerWithPrevious added in v0.1.56

func NewRuntimeInvocationSignerWithPrevious(
	activeKeyID, activeSecret, previousKeyID, previousSecret string,
) (*RuntimeInvocationSigner, error)

NewRuntimeInvocationSignerWithPrevious is the configuration-oriented keyring constructor. Both predecessor values are optional together; partial or same-ID rotation configuration fails closed.

func (*RuntimeInvocationSigner) Issue added in v0.1.56

func (s *RuntimeInvocationSigner) Issue(capability RuntimeInvocationCapability) (nodeEnvelope, invocationToken string, err error)

Issue returns a signed Node context and bearer capability. For the same immutable Attempt evidence it is deterministic, so a repeated claim cannot change the assignment journal digest.

func (*RuntimeInvocationSigner) VerifyInvocationToken added in v0.1.56

func (s *RuntimeInvocationSigner) VerifyInvocationToken(token string, databaseNow time.Time) (RuntimeInvocationCapability, error)

func (*RuntimeInvocationSigner) VerifyNodeEnvelope added in v0.1.56

func (s *RuntimeInvocationSigner) VerifyNodeEnvelope(envelope string, databaseNow time.Time) (RuntimeInvocationCapability, error)

type RuntimeInvocationVerifier added in v0.1.56

type RuntimeInvocationVerifier interface {
	VerifyNodeEnvelope(string, time.Time) (RuntimeInvocationCapability, error)
	VerifyInvocationToken(string, time.Time) (RuntimeInvocationCapability, error)
}

RuntimeInvocationVerifier verifies the two domain-separated capabilities emitted with one assignment. Both must decode to the same immutable Attempt authority before an Agent may create a child Run.

type RuntimeLeaseConfig added in v0.1.56

type RuntimeLeaseConfig struct {
	OfferTTL     time.Duration
	LeaseTTL     time.Duration
	AttemptTTL   time.Duration
	HeartbeatTTL time.Duration
}

RuntimeLeaseConfig contains the database-enforced durations used by the reliable Runtime offer and execution state machine. Zero values select the protocol defaults; negative or sub-millisecond values fail closed.

func DefaultRuntimeLeaseConfig added in v0.1.56

func DefaultRuntimeLeaseConfig() RuntimeLeaseConfig

type RuntimeLeaseError added in v0.1.56

type RuntimeLeaseError struct {
	Code RuntimeLeaseErrorCode `json:"code"`
	// contains filtered or unexported fields
}

RuntimeLeaseError is transport-neutral. HTTP and WebSocket adapters map the stable code without parsing its human-readable text.

func (*RuntimeLeaseError) Error added in v0.1.56

func (e *RuntimeLeaseError) Error() string

func (*RuntimeLeaseError) Unwrap added in v0.1.56

func (e *RuntimeLeaseError) Unwrap() error

type RuntimeLeaseErrorCode added in v0.1.56

type RuntimeLeaseErrorCode string
const (
	RuntimeLeaseErrorValidationFailed RuntimeLeaseErrorCode = "VALIDATION_FAILED"
	RuntimeLeaseErrorIdentityMismatch RuntimeLeaseErrorCode = "LEASE_IDENTITY_MISMATCH"
	RuntimeLeaseErrorStaleLease       RuntimeLeaseErrorCode = "STALE_LEASE"
	RuntimeLeaseErrorLeaseExpired     RuntimeLeaseErrorCode = "LEASE_EXPIRED"
	RuntimeLeaseErrorNodeAtCapacity   RuntimeLeaseErrorCode = "NODE_AT_CAPACITY"
	RuntimeLeaseErrorRunTerminal      RuntimeLeaseErrorCode = "RUN_ALREADY_TERMINAL"
	RuntimeLeaseErrorCancelRequested  RuntimeLeaseErrorCode = "RUN_CANCEL_REQUESTED"
)

type RuntimeLeaseService added in v0.1.56

type RuntimeLeaseService struct {
	// contains filtered or unexported fields
}

RuntimeLeaseService owns offer reservation, assignment confirmation, rejection, renewal, and disconnect cleanup. PostgreSQL is the only clock and linearization point; no state decision in this service uses process time.

func NewRuntimeLeaseService added in v0.1.56

func NewRuntimeLeaseService(
	pool *pgxpool.Pool,
	coreInstanceID uuid.UUID,
	issuer RuntimeInvocationCapabilityIssuer,
	config RuntimeLeaseConfig,
) *RuntimeLeaseService

func (*RuntimeLeaseService) AckAssignment added in v0.1.56

func (*RuntimeLeaseService) ClaimOffer added in v0.1.56

ClaimOffer returns nil when no Run is currently claimable. An outstanding, unacknowledged offer for this Session is returned first and never reserves a second capacity slot. An expired offer is released in the same transaction before the next candidate is considered.

func (*RuntimeLeaseService) RejectAssignment added in v0.1.56

func (*RuntimeLeaseService) ReleaseUnackedOffer added in v0.1.56

func (s *RuntimeLeaseService) ReleaseUnackedOffer(
	ctx context.Context,
	principal RuntimeSessionPrincipal,
	reason ...string,
) error

ReleaseUnackedOffer is used by a failed send, ACK timeout, or Session close. It deliberately searches only for the Session's current unaccepted offer; an accepted executing Attempt is left untouched for resume or expiry.

func (*RuntimeLeaseService) RenewLease added in v0.1.56

type RuntimeLivenessPolicy added in v0.1.56

type RuntimeLivenessPolicy struct {
	HeartbeatInterval time.Duration
	SessionStaleAfter time.Duration
	PresenceTTL       time.Duration
}

func CurrentRuntimeLivenessPolicy added in v0.1.56

func CurrentRuntimeLivenessPolicy() RuntimeLivenessPolicy

CurrentRuntimeLivenessPolicy is the single Core-owned definition used by Session transports, PostgreSQL readiness views and Redis advisory presence.

type RuntimeMaintenanceResult added in v0.1.56

type RuntimeMaintenanceResult struct {
	ReconcileBatches    int
	CancellationBatches int
	SessionBatches      int
	Reconciled          int
	Requeued            int
	TimedOut            int
	DeadLettered        int
	CancellationsReaped int
	SessionsReaped      int
}

func RunRuntimeMaintenanceOnce added in v0.1.56

func RunRuntimeMaintenanceOnce(
	ctx context.Context,
	reconciler runtimeDeadlineReconcileWorker,
	cancellations runtimeCancellationReapWorker,
	sessions runtimeSessionReapWorker,
	cfg RuntimeMaintenanceWorkerConfig,
) (RuntimeMaintenanceResult, error)

RunRuntimeMaintenanceOnce executes lease/deadline reconciliation and cancellation deadline recovery independently. An error in one path does not suppress the other path; errors are joined after both bounded passes finish.

type RuntimeMaintenanceWorkerConfig added in v0.1.56

type RuntimeMaintenanceWorkerConfig struct {
	Interval              time.Duration
	ReconcileBatchSize    int
	CancellationBatchSize int
	SessionBatchSize      int
	MaxCatchUpBatches     int
	Observer              WorkerObserver
}

RuntimeMaintenanceWorkerConfig bounds every tick so a large stale queue cannot monopolize a Core process. A full batch triggers a small, bounded catch-up loop; the next tick continues any remaining work.

type RuntimeMessageType added in v0.1.56

type RuntimeMessageType string

RuntimeMessageType is the closed set of Runtime WebSocket message types.

const (
	RuntimeMessageHello                 RuntimeMessageType = "runtime.hello"
	RuntimeMessageReady                 RuntimeMessageType = "runtime.ready"
	RuntimeMessageRunAssigned           RuntimeMessageType = "run.assigned"
	RuntimeMessageAssignmentAck         RuntimeMessageType = "run.assignment.ack"
	RuntimeMessageAssignmentConfirmed   RuntimeMessageType = "run.assignment.confirmed"
	RuntimeMessageAssignmentReject      RuntimeMessageType = "run.assignment.reject"
	RuntimeMessageAssignmentRejected    RuntimeMessageType = "run.assignment.rejected"
	RuntimeMessageLeaseRenew            RuntimeMessageType = "run.lease.renew"
	RuntimeMessageLeaseRenewed          RuntimeMessageType = "run.lease.renewed"
	RuntimeMessageRunEvent              RuntimeMessageType = "run.event"
	RuntimeMessageRunEventAck           RuntimeMessageType = "run.event.ack"
	RuntimeMessageRunResult             RuntimeMessageType = "run.result"
	RuntimeMessageRunResultAck          RuntimeMessageType = "run.result.ack"
	RuntimeMessageRunCancel             RuntimeMessageType = "run.cancel"
	RuntimeMessageRunCancelAck          RuntimeMessageType = "run.cancel.ack"
	RuntimeMessageResume                RuntimeMessageType = "runtime.resume"
	RuntimeMessageResumeAccepted        RuntimeMessageType = "run.resume.accepted"
	RuntimeMessageLeaseRevoked          RuntimeMessageType = "run.lease.revoked"
	RuntimeMessageDrain                 RuntimeMessageType = "runtime.drain"
	RuntimeMessageBrowserViewerCommand  RuntimeMessageType = "browser.viewer.command"
	RuntimeMessageBrowserViewerFrame    RuntimeMessageType = "browser.viewer.frame"
	RuntimeMessageBrowserViewerFrameAck RuntimeMessageType = "browser.viewer.frame.ack"

	// Authenticated read-only observation. Separate from browser.viewer.* on
	// purpose: observing while the Agent keeps working is a different capability
	// from taking control away from it, and they must not share a message space.
	RuntimeMessageBrowserObserverCommand  RuntimeMessageType = "browser.observer.command"
	RuntimeMessageBrowserObserverEvent    RuntimeMessageType = "browser.observer.event"
	RuntimeMessageBrowserObserverEventAck RuntimeMessageType = "browser.observer.event.ack"
	RuntimeMessageError                   RuntimeMessageType = "runtime.error"
)

type RuntimeNodeCredentialVerifier added in v0.1.56

type RuntimeNodeCredentialVerifier interface {
	VerifyRuntimeNodeCredential(context.Context, RuntimePresentedCertificate) (RuntimeDeviceIdentity, error)
}

RuntimeNodeCredentialVerifier resolves an enrolled, active Runtime Node by exact certificate identity. The current inventory pins serial plus SPKI; the authenticator also keeps the full leaf fingerprint bound across the certificate-verifier boundary so callers cannot substitute a peer.

type RuntimeNodeListItem added in v0.1.56

type RuntimeNodeListItem struct {
	NodeID                string     `json:"node_id"`
	DisplayName           string     `json:"display_name"`
	NodeVersion           string     `json:"node_version"`
	ProtocolVersion       int32      `json:"protocol_version"`
	RuntimeContractID     string     `json:"runtime_contract_id"`
	RuntimeContractDigest string     `json:"runtime_contract_digest"`
	ContractMatch         bool       `json:"contract_match"`
	Features              []string   `json:"features"`
	Capacity              int32      `json:"capacity"`
	Inflight              int32      `json:"inflight"`
	Status                string     `json:"status"`
	LastSeenAt            *time.Time `json:"last_seen_at,omitempty"`
	DrainingAt            *time.Time `json:"draining_at,omitempty"`
	RevokedAt             *time.Time `json:"revoked_at,omitempty"`
	RevokeReason          *string    `json:"revoke_reason,omitempty"`
	CreatedAt             time.Time  `json:"created_at"`
	UpdatedAt             time.Time  `json:"updated_at"`
	ActiveSessionCount    int32      `json:"active_session_count"`
	ActiveAgentCount      int32      `json:"active_agent_count"`
}

type RuntimeNodeListResponse added in v0.1.56

type RuntimeNodeListResponse struct {
	Items                 []RuntimeNodeListItem `json:"items"`
	Total                 int32                 `json:"total"`
	Limit                 int32                 `json:"limit"`
	Offset                int32                 `json:"offset"`
	CurrentContractID     string                `json:"current_contract_id"`
	CurrentContractDigest string                `json:"current_contract_digest"`
	DatabaseTime          time.Time             `json:"database_time"`
}

RuntimeNodeListResponse is the admin-only Runtime Node inventory. The database timestamp and current contract make freshness/compatibility decisions reproducible without trusting the API host clock.

type RuntimePresence added in v0.1.56

type RuntimePresence struct {
	CoreInstanceID     uuid.UUID              `json:"core_instance_id"`
	NodeID             uuid.UUID              `json:"node_id"`
	AgentID            uuid.UUID              `json:"agent_id"`
	RuntimeSessionID   uuid.UUID              `json:"runtime_session_id"`
	ConnectionID       string                 `json:"connection_id"`
	WorkerID           string                 `json:"worker_id"`
	Capacity           int32                  `json:"capacity"`
	Inflight           int32                  `json:"inflight"`
	NodeVersion        string                 `json:"version"`
	Transport          RuntimeTransport       `json:"transport"`
	TransportReason    RuntimeTransportReason `json:"transport_reason"`
	TransportChangedAt time.Time              `json:"transport_changed_at"`
}

RuntimePresence is an expiring routing/display hint. PostgreSQL Node, Session, attachment, certificate and credential state always wins.

type RuntimePresenceStore added in v0.1.56

type RuntimePresenceStore interface {
	Refresh(context.Context, RuntimePresence, time.Duration) error
	ListByAgent(context.Context, uuid.UUID) ([]RuntimePresence, error)
	Remove(context.Context, RuntimePresence) error
}

type RuntimePresenceStoreProvider added in v0.1.56

type RuntimePresenceStoreProvider interface {
	RuntimePresenceStore() (RuntimePresenceStore, error)
}

RuntimePresenceStoreProvider lets Core reuse the Redis client owned by its signal bus without widening RuntimeSignalBus or making PostgreSQL logic depend on Redis.

type RuntimePresentedCertificate added in v0.1.56

type RuntimePresentedCertificate struct {
	Serial                    string
	FingerprintSHA256         string
	PublicKeyThumbprintSHA256 string
	NotBefore                 time.Time
	NotAfter                  time.Time
}

RuntimePresentedCertificate is the canonical certificate identity passed to the durable credential verifier. The verifier must use database time for credential expiry and must fail closed for revoked credentials.

type RuntimePrincipalBinder added in v0.1.56

type RuntimePrincipalBinder interface {
	VerifyRuntimePrincipalBinding(context.Context, uuid.UUID, RuntimeDeviceIdentity) error
	ResolveRuntimeDeviceIdentity(context.Context, uuid.UUID) (RuntimeDeviceIdentity, error)
	ResolveTokenOnlyRuntimeDeviceIdentity(context.Context, uuid.UUID, uuid.UUID) (RuntimeDeviceIdentity, error)
}

RuntimePrincipalBinder enforces the durable one-to-one relationship between an Agent Token Credential and a Runtime Node public key. Token-only transport resolves the same Node through this binding instead of trusting request data.

type RuntimeReadyMessage added in v0.1.56

type RuntimeReadyMessage = RuntimeTypedEnvelope[RuntimeReadyPayload]

type RuntimeReadyPayload added in v0.1.56

type RuntimeReadyPayload struct {
	CoreInstanceID  string    `json:"core_instance_id" runtime:"required"`
	AttachmentID    uuid.UUID `json:"attachment_id" runtime:"required"`
	Features        []string  `json:"features" runtime:"required"`
	OfferTTLSeconds int64     `json:"offer_ttl_seconds" runtime:"required"`
	LeaseTTLSeconds int64     `json:"lease_ttl_seconds" runtime:"required"`
	DatabaseTime    time.Time `json:"database_time" runtime:"required"`
}

type RuntimeReconcileBatchResult added in v0.1.56

type RuntimeReconcileBatchResult struct {
	Scanned       int `json:"scanned"`
	Reconciled    int `json:"reconciled"`
	Requeued      int `json:"requeued"`
	TimedOut      int `json:"timed_out"`
	DeadLettered  int `json:"dead_lettered"`
	ResultUnknown int `json:"result_unknown"`
}

RuntimeReconcileBatchResult reports committed state changes. Scanned is bounded by the requested limit; candidates skipped because another worker held a SKIP LOCKED row are not counted as reconciled.

type RuntimeResultAck added in v0.1.56

type RuntimeResultAck struct {
	ResultID       uuid.UUID                   `json:"result_id"`
	Classification RuntimeResultClassification `json:"classification"`
	RunStatus      string                      `json:"run_status"`
	DispatchState  string                      `json:"dispatch_state"`
	Replayed       bool                        `json:"replayed"`
	NextAttemptAt  *time.Time                  `json:"next_attempt_at,omitempty"`
}

RuntimeResultAck is durable business acknowledgement. A caller may delete its Result spool only after receiving this ACK.

type RuntimeResultClassification added in v0.1.56

type RuntimeResultClassification string

RuntimeResultClassification is Core's persisted interpretation of a Result. Timeout and dead_letter are ACK classifications derived from the final Run state; only the first three values are written to run_attempts.

const (
	RuntimeResultClassificationSuccess      RuntimeResultClassification = "success"
	RuntimeResultClassificationRetryable    RuntimeResultClassification = "retryable_failure"
	RuntimeResultClassificationNonRetryable RuntimeResultClassification = "non_retryable_failure"
	RuntimeResultClassificationTimeout      RuntimeResultClassification = "timeout"
	RuntimeResultClassificationCanceled     RuntimeResultClassification = "canceled"
	RuntimeResultClassificationDeadLetter   RuntimeResultClassification = "dead_letter"
)

type RuntimeResultError added in v0.1.56

type RuntimeResultError struct {
	Code          RuntimeResultErrorCode `json:"code"`
	MissingRanges []EventRange           `json:"missing_ranges,omitempty"`
	// contains filtered or unexported fields
}

RuntimeResultError carries a stable code and exact inclusive Event gaps.

func (*RuntimeResultError) Error added in v0.1.56

func (e *RuntimeResultError) Error() string

func (*RuntimeResultError) Unwrap added in v0.1.56

func (e *RuntimeResultError) Unwrap() error

type RuntimeResultErrorCode added in v0.1.56

type RuntimeResultErrorCode string

RuntimeResultErrorCode is stable across HTTP, WebSocket, gRPC, and MCP adapters. Human-readable transport text must never be parsed for behavior.

const (
	RuntimeResultErrorValidationFailed      RuntimeResultErrorCode = "VALIDATION_FAILED"
	RuntimeResultErrorResultIDConflict      RuntimeResultErrorCode = "RESULT_ID_CONFLICT"
	RuntimeResultErrorRunAlreadyTerminal    RuntimeResultErrorCode = "RUN_ALREADY_TERMINAL"
	RuntimeResultErrorStaleLease            RuntimeResultErrorCode = "STALE_LEASE"
	RuntimeResultErrorLeaseExpired          RuntimeResultErrorCode = "LEASE_EXPIRED"
	RuntimeResultErrorLeaseIdentityMismatch RuntimeResultErrorCode = "LEASE_IDENTITY_MISMATCH"
	RuntimeResultErrorEventsMissing         RuntimeResultErrorCode = "EVENTS_MISSING"
	RuntimeResultErrorRunCancelRequested    RuntimeResultErrorCode = "RUN_CANCEL_REQUESTED"
)

type RuntimeResultFailure added in v0.1.56

type RuntimeResultFailure struct {
	ErrorCode     string `json:"error_code"`
	Message       string `json:"message"`
	RetryableHint bool   `json:"retryable_hint,omitempty"`
}

RuntimeResultFailure is the normalized failed Result body. RetryableHint is advisory input; the server-side ResultClassifier remains authoritative.

type RuntimeResultFinalizer added in v0.1.56

type RuntimeResultFinalizer interface {
	Finalize(context.Context, RuntimeResultPrincipal, RuntimeResultRequest) (RuntimeResultAck, error)
}

type RuntimeResultPrincipal added in v0.1.56

type RuntimeResultPrincipal = RuntimeEventPrincipal

RuntimeResultPrincipal is an already-authenticated execution principal. Authentication and revocation precede envelope decoding in the Task 6 transport. It aliases the event principal so Event and Result use identical Node/Agent/worker/session ownership semantics.

type RuntimeResultRequest added in v0.1.56

type RuntimeResultRequest struct {
	AttemptIdentity     RuntimeAttemptIdentity `json:"attempt_identity"`
	ResultID            uuid.UUID              `json:"result_id"`
	Status              string                 `json:"status"`
	Output              map[string]any         `json:"output,omitempty"`
	Error               *RuntimeResultFailure  `json:"error,omitempty"`
	DurationMS          int32                  `json:"duration_ms"`
	FinalClientEventSeq int64                  `json:"final_client_event_seq"`
}

RuntimeResultRequest is the transport-neutral Runtime Result payload. The envelope message ID and transport correlation fields intentionally live outside this type and therefore outside the semantic fingerprint.

type RuntimeResumeAPI added in v0.1.56

type RuntimeResumeAPI interface {
	Resume(context.Context, RuntimeSessionPrincipal, RuntimeResumePayload) (RuntimeResumeResponse, error)
}

RuntimeResumeAPI is optional while the HTTP/WS recovery surface is being wired. A missing implementation produces a correlated runtime.error; it never acknowledges recovery that Core did not durably authorize.

type RuntimeResumeAction added in v0.1.56

type RuntimeResumeAction string
const (
	RuntimeActionContinueExecution RuntimeResumeAction = "continue_execution"
	RuntimeActionUploadEvents      RuntimeResumeAction = "upload_events"
	RuntimeActionUploadResult      RuntimeResumeAction = "upload_result"
	RuntimeActionStopExecution     RuntimeResumeAction = "stop_execution"
	RuntimeActionClearSpool        RuntimeResumeAction = "clear_spool"
)

type RuntimeResumeDecision added in v0.1.56

type RuntimeResumeDecision string
const (
	RuntimeResumeContinueExecution RuntimeResumeDecision = "continue_execution"
	RuntimeResumeUploadSpoolOnly   RuntimeResumeDecision = "upload_spool_only"
	RuntimeResumeResultAcked       RuntimeResumeDecision = "result_already_acked"
	RuntimeResumeLeaseRevoked      RuntimeResumeDecision = "lease_revoked"
)

type RuntimeResumeMessage added in v0.1.56

type RuntimeResumeMessage = RuntimeTypedEnvelope[RuntimeResumePayload]

type RuntimeResumePayload added in v0.1.56

type RuntimeResumePayload struct {
	NodeID           uuid.UUID       `json:"node_id" runtime:"required"`
	AgentID          uuid.UUID       `json:"agent_id" runtime:"required"`
	WorkerID         string          `json:"worker_id" runtime:"required"`
	RuntimeSessionID uuid.UUID       `json:"runtime_session_id" runtime:"required"`
	Attempts         []ResumeAttempt `json:"attempts" runtime:"required"`
}

type RuntimeResumeResponse added in v0.1.56

type RuntimeResumeResponse struct {
	Decisions []RunResumeAcceptedPayload `json:"decisions" runtime:"required"`
}

type RuntimeResumeService added in v0.1.56

type RuntimeResumeService struct {
	// contains filtered or unexported fields
}

RuntimeResumeService decides what a reconnecting Runtime Worker may do with each immutable Attempt identity. It never moves execution or capacity ownership to a replacement Session: cross-Session recovery is restricted to uploading an already durable local spool.

func NewRuntimeResumeService added in v0.1.56

func NewRuntimeResumeService(
	pool *pgxpool.Pool,
	coreInstanceID uuid.UUID,
	grantTTL time.Duration,
) *RuntimeResumeService

NewRuntimeResumeService creates a PostgreSQL-backed Resume service. A zero grant TTL selects the protocol default; all authorization deadlines are written and compared by PostgreSQL.

func (*RuntimeResumeService) Resume added in v0.1.56

Resume returns one decision per requested Attempt, preserving the caller's order. Locks are nevertheless acquired by Run UUID and Attempt UUID so two reconnecting processes cannot deadlock by presenting a different order.

type RuntimeRunStatus added in v0.1.56

type RuntimeRunStatus string
const (
	RuntimeRunRunning  RuntimeRunStatus = "running"
	RuntimeRunSuccess  RuntimeRunStatus = "success"
	RuntimeRunFailed   RuntimeRunStatus = "failed"
	RuntimeRunTimeout  RuntimeRunStatus = "timeout"
	RuntimeRunCanceled RuntimeRunStatus = "canceled"
)

type RuntimeSchemaContractSnapshot added in v0.1.56

type RuntimeSchemaContractSnapshot struct {
	SchemaVersion         int32  `json:"schema_version"`
	MigrationName         string `json:"migration_name"`
	RuntimeContractID     string `json:"runtime_contract_id"`
	RuntimeContractDigest string `json:"runtime_contract_digest"`
}

type RuntimeSessionAPI added in v0.1.56

type RuntimeSessionAPI interface {
	CreateOrAttachSession(context.Context, AuthenticatedRuntimePrincipal, RuntimeSessionRequest) (RuntimeSessionState, error)
	HeartbeatSession(context.Context, AuthenticatedRuntimePrincipal, RuntimeSessionHeartbeatRequest) (RuntimeSessionState, error)
	DrainSession(context.Context, AuthenticatedRuntimePrincipal, RuntimeSessionDrainRequest) (RuntimeDrainPayload, error)
	CloseSession(context.Context, AuthenticatedRuntimePrincipal, RuntimeSessionCloseRequest) (RuntimeSessionState, error)
	ResolveSessionPrincipal(context.Context, AuthenticatedRuntimePrincipal, uuid.UUID) (RuntimeSessionPrincipal, error)
	// ResolveWorkerSessionPrincipal returns the currently acting Session. The
	// source Session ID in an Attempt remains immutable across resume and must
	// not be mistaken for the authenticated uploader.
	ResolveWorkerSessionPrincipal(context.Context, AuthenticatedRuntimePrincipal, string) (RuntimeSessionPrincipal, error)
}

type RuntimeSessionCloseRequest added in v0.1.56

type RuntimeSessionCloseRequest struct {
	RuntimeSessionIdentity
	Status       string    `json:"status"`
	Reason       string    `json:"reason"`
	AttachmentID uuid.UUID `json:"-"`
}

RuntimeSessionCloseRequest closes an attachment and Session atomically. Status is either offline (reconnectable) or closed (permanent).

type RuntimeSessionDrainRequest added in v0.1.56

type RuntimeSessionDrainRequest struct {
	RuntimeSessionID uuid.UUID           `json:"-"`
	AttachmentID     uuid.UUID           `json:"-"`
	Payload          RuntimeDrainPayload `json:"-"`
}

RuntimeSessionDrainRequest carries trusted transport identity outside the JSON payload. RuntimeDrainPayload is shared by WebSocket and Pull so both adapters commit exactly the same durable state transition. Capacity must be zero; Inflight is only an untrusted client snapshot and is never persisted.

type RuntimeSessionError added in v0.1.56

type RuntimeSessionError struct {
	Code RuntimeSessionErrorCode `json:"code"`
	// contains filtered or unexported fields
}

func (*RuntimeSessionError) Error added in v0.1.56

func (e *RuntimeSessionError) Error() string

func (*RuntimeSessionError) Unwrap added in v0.1.56

func (e *RuntimeSessionError) Unwrap() error

type RuntimeSessionErrorCode added in v0.1.56

type RuntimeSessionErrorCode string
const (
	RuntimeSessionErrorAuthenticationFailed   RuntimeSessionErrorCode = "AUTHENTICATION_FAILED"
	RuntimeSessionErrorValidationFailed       RuntimeSessionErrorCode = "VALIDATION_FAILED"
	RuntimeSessionErrorAgentMismatch          RuntimeSessionErrorCode = "AGENT_MISMATCH"
	RuntimeSessionErrorDeviceMismatch         RuntimeSessionErrorCode = "DEVICE_MISMATCH"
	RuntimeSessionErrorProtocolUnsupported    RuntimeSessionErrorCode = "PROTOCOL_UNSUPPORTED"
	RuntimeSessionErrorContractMismatch       RuntimeSessionErrorCode = "CONTRACT_MISMATCH"
	RuntimeSessionErrorRequiredFeatureMissing RuntimeSessionErrorCode = "REQUIRED_FEATURE_MISSING"
	RuntimeSessionErrorSessionConflict        RuntimeSessionErrorCode = "SESSION_CONFLICT"
	RuntimeSessionErrorPrincipalInactive      RuntimeSessionErrorCode = "PRINCIPAL_INACTIVE"
	RuntimeSessionErrorNotAttached            RuntimeSessionErrorCode = "SESSION_NOT_ATTACHED"
)

type RuntimeSessionHeartbeatRequest added in v0.1.56

type RuntimeSessionHeartbeatRequest = RuntimeSessionRequest

RuntimeSessionHeartbeatRequest repeats immutable contract identity so a reconnect cannot silently downgrade features or switch a worker/session.

type RuntimeSessionIdentity added in v0.1.56

type RuntimeSessionIdentity struct {
	RuntimeSessionID uuid.UUID `json:"runtime_session_id"`
	NodeID           uuid.UUID `json:"node_id"`
	AgentID          uuid.UUID `json:"agent_id"`
	WorkerID         string    `json:"worker_id"`
	SessionEpoch     int64     `json:"session_epoch"`
}

RuntimeSessionIdentity is the complete durable identity of one Node process session for one Agent. SessionEpoch changes on every process start, while WorkerID remains installation-stable.

type RuntimeSessionLease added in v0.1.56

type RuntimeSessionLease struct {
	Version          int       `json:"version"`
	CoreInstanceID   uuid.UUID `json:"core_instance_id"`
	NodeID           uuid.UUID `json:"node_id"`
	AgentID          uuid.UUID `json:"agent_id"`
	RuntimeSessionID uuid.UUID `json:"runtime_session_id"`
	AttachmentID     uuid.UUID `json:"attachment_id"`
	ConnectionID     string    `json:"connection_id"`
	WorkerID         string    `json:"worker_id"`
	SessionEpoch     int64     `json:"session_epoch"`
	RefreshedAt      time.Time `json:"refreshed_at"`
}

RuntimeSessionLease is an expiring, advisory liveness hint. It intentionally contains no token, certificate, Run, input, output, or authorization data. PostgreSQL Session/attachment state and fencing remain authoritative.

type RuntimeSessionLeaseManager added in v0.1.56

type RuntimeSessionLeaseManager struct {
	// contains filtered or unexported fields
}

RuntimeSessionLeaseManager owns one process-level refresh loop. It replaces per-WebSocket Redis and database heartbeat loops without becoming a source of truth. Missing Redis keys are ignored until a full warmup has elapsed.

func NewRuntimeSessionLeaseManager added in v0.1.56

func NewRuntimeSessionLeaseManager(
	store RuntimeSessionLeaseStore,
	config RuntimeSessionLeaseManagerConfig,
) (*RuntimeSessionLeaseManager, error)

func (*RuntimeSessionLeaseManager) AbsenceReady added in v0.1.56

func (m *RuntimeSessionLeaseManager) AbsenceReady() bool

func (*RuntimeSessionLeaseManager) Forget added in v0.1.56

func (m *RuntimeSessionLeaseManager) Forget(
	ctx context.Context,
	runtimeSessionID uuid.UUID,
) error

func (*RuntimeSessionLeaseManager) Health added in v0.1.56

func (*RuntimeSessionLeaseManager) HealthyFor added in v0.1.56

func (m *RuntimeSessionLeaseManager) HealthyFor(connectionID string) bool

func (*RuntimeSessionLeaseManager) ListExpired added in v0.1.56

func (m *RuntimeSessionLeaseManager) ListExpired(
	ctx context.Context,
	limit int,
) ([]uuid.UUID, error)

func (*RuntimeSessionLeaseManager) Lookup added in v0.1.56

func (m *RuntimeSessionLeaseManager) Lookup(
	ctx context.Context,
	runtimeSessionID uuid.UUID,
) (RuntimeSessionLease, bool, error)

func (*RuntimeSessionLeaseManager) RefreshConnection added in v0.1.56

func (m *RuntimeSessionLeaseManager) RefreshConnection(ctx context.Context, connectionID string) error

RefreshConnection establishes the first lease before Ready is emitted. The process loop owns all later refreshes in batches; this one synchronous write prevents a newly attached Session from existing only in process memory.

func (*RuntimeSessionLeaseManager) Register added in v0.1.56

func (*RuntimeSessionLeaseManager) Run added in v0.1.56

func (*RuntimeSessionLeaseManager) ScheduleCheck added in v0.1.56

func (m *RuntimeSessionLeaseManager) ScheduleCheck(
	ctx context.Context,
	runtimeSessionID uuid.UUID,
	after time.Duration,
) error

func (*RuntimeSessionLeaseManager) Unregister added in v0.1.56

func (m *RuntimeSessionLeaseManager) Unregister(
	ctx context.Context,
	connectionID string,
	attachmentID uuid.UUID,
) (bool, error)

func (*RuntimeSessionLeaseManager) UpdatePresence added in v0.1.56

func (m *RuntimeSessionLeaseManager) UpdatePresence(connectionID string, presence RuntimePresence) bool

type RuntimeSessionLeaseManagerConfig added in v0.1.56

type RuntimeSessionLeaseManagerConfig struct {
	RefreshInterval time.Duration
	LeaseTTL        time.Duration
	PresenceTTL     time.Duration
	Warmup          time.Duration
	BatchSize       int
	DisableJitter   bool
}

type RuntimeSessionLeaseManagerHealth added in v0.1.56

type RuntimeSessionLeaseManagerHealth struct {
	Connected       bool
	AbsenceReady    bool
	Reason          string
	Generation      uint64
	Registered      int
	LastSuccessAt   time.Time
	HealthySince    time.Time
	RefreshFailures uint64
}

type RuntimeSessionLeaseRecord added in v0.1.56

type RuntimeSessionLeaseRecord struct {
	Lease    RuntimeSessionLease
	Presence RuntimePresence
}

type RuntimeSessionLeaseStore added in v0.1.56

type RuntimeSessionLeaseStoreProvider added in v0.1.56

type RuntimeSessionLeaseStoreProvider interface {
	RuntimeSessionLeaseStore() (RuntimeSessionLeaseStore, error)
}

type RuntimeSessionPrincipal added in v0.1.56

type RuntimeSessionPrincipal struct {
	RuntimeSessionID                uuid.UUID
	NodeID                          uuid.UUID
	AgentID                         uuid.UUID
	CredentialID                    uuid.UUID
	WorkerID                        string
	SessionEpoch                    int64
	RuntimeContractDigest           string
	Features                        []string
	CoreInstanceID                  uuid.UUID
	AttachmentID                    uuid.UUID
	DeviceCertificateSerial         string
	DevicePublicKeyThumbprintSHA256 string
	Status                          string
	DatabaseTime                    time.Time
}

RuntimeSessionPrincipal is the active, database-validated Session identity consumed by claim, command, Event, and Result adapters. It can be converted directly to the transport-neutral Event/Result principal.

func (RuntimeSessionPrincipal) EventPrincipal added in v0.1.56

type RuntimeSessionReaper added in v0.1.56

type RuntimeSessionReaper struct {
	// contains filtered or unexported fields
}

RuntimeSessionReaper closes transport Sessions whose database heartbeat has expired. It changes only Session/attachment state; offer and Attempt expiry remain owned by the lease/deadline reconciler.

func NewRuntimeSessionReaper added in v0.1.56

func NewRuntimeSessionReaper(pool *pgxpool.Pool, heartbeatTTL time.Duration) *RuntimeSessionReaper

func NewRuntimeSessionReaperWithLeases added in v0.1.56

func NewRuntimeSessionReaperWithLeases(
	pool *pgxpool.Pool,
	heartbeatTTL time.Duration,
	leases *RuntimeSessionLeaseManager,
) *RuntimeSessionReaper

NewRuntimeSessionReaperWithLeases keeps the database fencing transaction authoritative while allowing a healthy Redis lease to prove that a database-stale WebSocket is still connected. Redis errors and startup uncertainty can only delay a reap; they can never close a Session.

func (*RuntimeSessionReaper) ReapStaleSessions added in v0.1.56

func (r *RuntimeSessionReaper) ReapStaleSessions(ctx context.Context, limit int) (int, error)

type RuntimeSessionRequest added in v0.1.56

type RuntimeSessionRequest struct {
	RuntimeSessionIdentity
	NodeVersion           string           `json:"node_version"`
	ProtocolVersion       int32            `json:"protocol_version"`
	RuntimeContractID     string           `json:"runtime_contract_id"`
	RuntimeContractDigest string           `json:"runtime_contract_digest"`
	Features              []string         `json:"features"`
	Capacity              int32            `json:"capacity"`
	AttachmentID          uuid.UUID        `json:"-"`
	Transport             RuntimeTransport `json:"-"`
	// ReportedTransportReason is an untrusted, bounded SDK hint received only
	// from the CreateSession request or WebSocket upgrade. Core validates it
	// against Transport, the prior durable attachment and TransportPolicy.
	ReportedTransportReason RuntimeTransportReason `json:"-"`
	TransportPolicy         RuntimeTransportPolicy `json:"-"`
}

RuntimeSessionRequest is the transport-neutral hello/session request.

type RuntimeSessionService added in v0.1.56

type RuntimeSessionService struct {
	// contains filtered or unexported fields
}

RuntimeSessionService owns Session lifecycle and the global lock order: Session -> Node -> Token -> Attachment.

func NewRuntimeSessionService added in v0.1.56

func NewRuntimeSessionService(pool *pgxpool.Pool, coreInstanceID uuid.UUID) *RuntimeSessionService

func (*RuntimeSessionService) CloseSession added in v0.1.56

func (*RuntimeSessionService) CreateOrAttachSession added in v0.1.56

func (*RuntimeSessionService) DetachCutoverSessions added in v0.1.56

func (s *RuntimeSessionService) DetachCutoverSessions(ctx context.Context) (int64, error)

DetachCutoverSessions closes every zero-inflight Session owned by this Core instance. It is used only by the runtime-attach-only process after both listeners have stopped, so the normal Core can claim the same durable Session identities without waiting for stale-member reconciliation.

func (*RuntimeSessionService) DrainSession added in v0.1.56

DrainSession commits the admission fence before acknowledging either transport. The client capacity/inflight snapshot is never authoritative: capacity is forced to zero and inflight is read back from PostgreSQL.

func (*RuntimeSessionService) HeartbeatSession added in v0.1.56

func (*RuntimeSessionService) ResolveSessionPrincipal added in v0.1.56

func (s *RuntimeSessionService) ResolveSessionPrincipal(
	ctx context.Context,
	principal AuthenticatedRuntimePrincipal,
	runtimeSessionID uuid.UUID,
) (RuntimeSessionPrincipal, error)

ResolveSessionPrincipal turns an authenticated Token/device pair plus an opaque Session ID into the exact active executor identity. The repository rechecks Session, Node, Token expiry/revocation, attachment, and Core owner in one database statement; callers must still perform the Attempt lock/fence check inside their own business transaction.

func (*RuntimeSessionService) ResolveWorkerSessionPrincipal added in v0.1.56

func (s *RuntimeSessionService) ResolveWorkerSessionPrincipal(
	ctx context.Context,
	principal AuthenticatedRuntimePrincipal,
	workerID string,
) (RuntimeSessionPrincipal, error)

ResolveWorkerSessionPrincipal resolves the currently attached acting Session without trusting the immutable source Session ID carried by an Attempt. HTTP Event/Result uploads use this path so a replacement process can present a durable resume grant while the wire Attempt identity remains unchanged.

type RuntimeSessionState added in v0.1.56

type RuntimeSessionState struct {
	Session      db.RuntimeSession            `json:"session"`
	Attachment   *db.RuntimeSessionAttachment `json:"attachment,omitempty"`
	DatabaseTime time.Time                    `json:"database_time"`
	Replayed     bool                         `json:"replayed"`
	Resumed      bool                         `json:"resumed"`
}

RuntimeSessionState is durable state returned after transaction commit. DatabaseTime is always taken from a PostgreSQL-written timestamp.

type RuntimeSignal added in v0.1.56

type RuntimeSignal struct {
	SignalID         uuid.UUID                   `json:"signal_id"`
	Type             string                      `json:"type"`
	AgentID          uuid.UUID                   `json:"agent_id"`
	RunID            *uuid.UUID                  `json:"run_id,omitempty"`
	NodeID           *uuid.UUID                  `json:"node_id,omitempty"`
	TargetInstanceID *uuid.UUID                  `json:"target_instance_id,omitempty"`
	CredentialID     *uuid.UUID                  `json:"credential_id,omitempty"`
	Connections      []RuntimeConnectionIdentity `json:"connections,omitempty"`
}

RuntimeSignal is deliberately a wake-up hint, never a data transport. Its wire shape is the complete allowlist: Run input/output, token material, invocation capabilities, payloads and secrets have no field through which they can reach Redis or an in-process subscriber.

func ParseRuntimeSignal added in v0.1.56

func ParseRuntimeSignal(encoded []byte) (RuntimeSignal, error)

ParseRuntimeSignal rejects unknown fields and multiple JSON values. This is important even for trusted publishers: accepting an accidental payload field would silently widen the Redis data-classification boundary.

type RuntimeSignalBus added in v0.1.56

type RuntimeSignalBus interface {
	Publish(context.Context, RuntimeSignal) error
	Subscribe(context.Context, RuntimeSignalHandler) error
	Health(context.Context) error
	Close() error
}

RuntimeSignalBus accelerates database-backed work discovery. Subscribe is blocking and returns when ctx is canceled, the bus closes, or the transport fails. Callers must supervise and retry subscriptions; correctness must not depend on receiving any particular signal.

type RuntimeSignalHandler added in v0.1.56

type RuntimeSignalHandler func(context.Context, RuntimeSignal) error

type RuntimeSignalOutboxBatchResult added in v0.1.56

type RuntimeSignalOutboxBatchResult struct {
	Claimed   int
	Published int
	Retried   int
}

type RuntimeSignalOutboxWorker added in v0.1.56

type RuntimeSignalOutboxWorker struct {
	// contains filtered or unexported fields
}

func NewRuntimeSignalOutboxWorker added in v0.1.56

func NewRuntimeSignalOutboxWorker(queries runtimeSignalOutboxStore, bus RuntimeSignalBus) *RuntimeSignalOutboxWorker

func (*RuntimeSignalOutboxWorker) Backlog added in v0.1.56

func (*RuntimeSignalOutboxWorker) ProcessOnce added in v0.1.56

ProcessOnce claims one bounded batch. The lease owner changes for every batch so a late mark cannot complete a signal reacquired by another Core. A crash after Publish and before Mark intentionally causes a duplicate signal after lease expiry; consumers must return to PostgreSQL to dedupe.

type RuntimeSignalOutboxWorkerConfig added in v0.1.56

type RuntimeSignalOutboxWorkerConfig struct {
	Interval          time.Duration
	LeaseDuration     time.Duration
	BatchSize         int32
	MaxCatchUpBatches int
	RetryBase         time.Duration
	RetryMaximum      time.Duration
	Observer          WorkerObserver
}

type RuntimeTokenValidator added in v0.1.56

type RuntimeTokenValidator interface {
	ValidateRuntimeToken(context.Context, string, ...string) (db.AgentRuntimeToken, error)
}

RuntimeTokenValidator authenticates the Agent half of a runtime principal. Implementations must check revocation, expiry, and the requested scope.

type RuntimeTransport added in v0.1.56

type RuntimeTransport string

RuntimeTransport is the actual Server-observed transport of one attachment. unknown is reserved for detached history created before transport metadata became durable; a live attachment must use websocket or long_poll.

func (RuntimeTransport) IsLive added in v0.1.56

func (transport RuntimeTransport) IsLive() bool

type RuntimeTransportError added in v0.1.56

type RuntimeTransportError struct {
	Body RuntimeErrorBody
	// contains filtered or unexported fields
}

RuntimeTransportError preserves a stable wire error without exposing an internal cause in its JSON representation.

func MapRuntimeTransportError added in v0.1.56

func MapRuntimeTransportError(err error) *RuntimeTransportError

MapRuntimeTransportError converts EventStore and ResultFinalizer errors into the shared transport vocabulary. Unknown causes are intentionally hidden.

func NewRuntimeTransportError added in v0.1.56

func NewRuntimeTransportError(code RuntimeErrorCode, message string) *RuntimeTransportError

func (*RuntimeTransportError) Envelope added in v0.1.56

func (e *RuntimeTransportError) Envelope() RuntimeError

func (*RuntimeTransportError) Error added in v0.1.56

func (e *RuntimeTransportError) Error() string

func (*RuntimeTransportError) Unwrap added in v0.1.56

func (e *RuntimeTransportError) Unwrap() error

type RuntimeTransportPolicy added in v0.1.56

type RuntimeTransportPolicy struct {
	Version                int
	OrderedTransports      []RuntimeTransport
	DefaultTransport       string
	RetryMinimum           time.Duration
	RetryMaximum           time.Duration
	WebSocketProbeInterval time.Duration
	WebSocketProbeTimeout  time.Duration
}

func CurrentRuntimeTransportPolicy added in v0.1.56

func CurrentRuntimeTransportPolicy() RuntimeTransportPolicy

CurrentRuntimeTransportPolicy is copied into discovery responses so callers cannot mutate Core's allowlist. The order is authoritative for auto mode.

func RuntimeAttachOnlyTransportPolicy added in v0.1.56

func RuntimeAttachOnlyTransportPolicy() RuntimeTransportPolicy

RuntimeAttachOnlyTransportPolicy is the honest cutover policy for the narrow Core that accepts Session attachments before the producer boundary. Pull execution routes are deliberately absent in that mode, so advertising long-poll fallback would let an auto-mode SDK attach to an unusable transport. Timing and retry semantics remain identical to full Core.

type RuntimeTransportPolicyProvider added in v0.1.56

type RuntimeTransportPolicyProvider func() RuntimeTransportPolicy

RuntimeTransportPolicyProvider supplies the Server-owned policy at request time. Production leaves this nil and uses CurrentRuntimeTransportPolicy; tests and future dynamic policy storage may inject a concurrency-safe provider without moving admission decisions into an SDK or AgentNode.

type RuntimeTransportReason added in v0.1.56

type RuntimeTransportReason string

RuntimeTransportReason is bounded, safe-to-display evidence for why a new attachment uses its observed transport. Arbitrary client/network text never enters durable attachment or advisory presence state.

func (RuntimeTransportReason) IsValid added in v0.1.56

func (reason RuntimeTransportReason) IsValid() bool

type RuntimeTypedEnvelope added in v0.1.56

type RuntimeTypedEnvelope[P any] struct {
	RuntimeEnvelopeFields
	Payload P `json:"payload" runtime:"required"`
}

RuntimeTypedEnvelope is used by concrete messages after routing.

func DecodeRuntimeTypedMessage added in v0.1.56

func DecodeRuntimeTypedMessage[P any](reader io.Reader, expected RuntimeMessageType) (RuntimeTypedEnvelope[P], error)

DecodeRuntimeTypedMessage performs strict one-pass decoding when the caller already knows the only permitted message type.

func NewRuntimeTypedMessage added in v0.1.56

func NewRuntimeTypedMessage[P any](messageType RuntimeMessageType, replyTo *uuid.UUID, payload P) (RuntimeTypedEnvelope[P], error)

NewRuntimeTypedMessage creates an internally valid outbound envelope.

type RuntimeWakeHub added in v0.1.56

type RuntimeWakeHub struct {
	// contains filtered or unexported fields
}

RuntimeWakeHub broadcasts typed, edge-triggered hints to all local HTTP Pull and WebSocket waiters for an Agent. PostgreSQL remains authoritative; the separate channels prevent a cancellation or lifecycle event from probing the dispatch queue and prevent a new Run from polling cancellation commands.

func NewRuntimeWakeHub added in v0.1.56

func NewRuntimeWakeHub() *RuntimeWakeHub

func (*RuntimeWakeHub) ConnectionSnapshot added in v0.1.56

func (h *RuntimeWakeHub) ConnectionSnapshot() []RuntimeConnectionRegistration

func (*RuntimeWakeHub) CredentialRevalidationWake added in v0.1.56

func (h *RuntimeWakeHub) CredentialRevalidationWake() <-chan struct{}

func (*RuntimeWakeHub) RegisterConnection added in v0.1.56

func (h *RuntimeWakeHub) RegisterConnection(
	identity RuntimeConnectionIdentity,
	credentialID uuid.UUID,
) <-chan struct{}

RegisterConnection returns a close-on-revocation channel for one exact attachment generation. A recently delivered signal is retained long enough to close a connection that races registration after the database commit.

func (*RuntimeWakeHub) RegisterWebSocketDispatch added in v0.1.56

func (h *RuntimeWakeHub) RegisterWebSocketDispatch(agentID uuid.UUID) (<-chan struct{}, bool)

RegisterWebSocketDispatch registers one persistent WebSocket waiter and reports whether it is the first live Session for this Agent. Dispatch hints are coalesced tokens, so one waiter claims durable work without waking every sibling Session on the same Agent.

func (*RuntimeWakeHub) RequireCredentialRevalidation added in v0.1.56

func (h *RuntimeWakeHub) RequireCredentialRevalidation()

func (*RuntimeWakeHub) RevokeConnections added in v0.1.56

func (h *RuntimeWakeHub) RevokeConnections(identities []RuntimeConnectionIdentity) int

RevokeConnections closes only exact registered attachment generations and retains a bounded race tombstone for connections that have not registered yet. It returns the number of live registrations that were signaled.

func (*RuntimeWakeHub) RevokeCredentialConnections added in v0.1.56

func (h *RuntimeWakeHub) RevokeCredentialConnections(
	credentialID uuid.UUID,
	identities []RuntimeConnectionIdentity,
) int

func (*RuntimeWakeHub) UnregisterConnection added in v0.1.56

func (h *RuntimeWakeHub) UnregisterConnection(identity RuntimeConnectionIdentity)

func (*RuntimeWakeHub) UnregisterWebSocketDispatch added in v0.1.56

func (h *RuntimeWakeHub) UnregisterWebSocketDispatch(agentID uuid.UUID)

func (*RuntimeWakeHub) Wait added in v0.1.56

func (h *RuntimeWakeHub) Wait(agentID uuid.UUID) <-chan struct{}

Wait is the compatibility alias for dispatch waiters.

func (*RuntimeWakeHub) WaitControl added in v0.1.56

func (h *RuntimeWakeHub) WaitControl(agentID uuid.UUID) <-chan struct{}

func (*RuntimeWakeHub) WaitDispatch added in v0.1.56

func (h *RuntimeWakeHub) WaitDispatch(agentID uuid.UUID) <-chan struct{}

func (*RuntimeWakeHub) WaitNodeDispatch added in v0.1.56

func (h *RuntimeWakeHub) WaitNodeDispatch(nodeID uuid.UUID) <-chan struct{}

func (*RuntimeWakeHub) Wake added in v0.1.56

func (h *RuntimeWakeHub) Wake(agentID uuid.UUID)

Wake preserves the previous broad-wake behavior for compatibility. New signal routing must use WakeDispatch or WakeControl explicitly.

func (*RuntimeWakeHub) WakeAll added in v0.1.56

func (h *RuntimeWakeHub) WakeAll()

WakeAll broadcasts a one-shot recovery hint to every Agent that currently has a local waiter. It is called only when the signal subscription reconnects so missed Pub/Sub notifications converge without per-connection polling.

func (*RuntimeWakeHub) WakeControl added in v0.1.56

func (h *RuntimeWakeHub) WakeControl(agentID uuid.UUID)

func (*RuntimeWakeHub) WakeDispatch added in v0.1.56

func (h *RuntimeWakeHub) WakeDispatch(agentID uuid.UUID)

func (*RuntimeWakeHub) WakeDispatchIfRegistered added in v0.1.56

func (h *RuntimeWakeHub) WakeDispatchIfRegistered(agentID uuid.UUID) bool

WakeDispatchIfRegistered wakes only an Agent that has registered a local waiter. Recovery scans use this form so a durable backlog owned by another Core cannot grow this process's in-memory wake map.

func (*RuntimeWakeHub) WakeNodeDispatch added in v0.1.56

func (h *RuntimeWakeHub) WakeNodeDispatch(nodeID uuid.UUID)

WakeNodeDispatch announces that durable capacity was released on a Node. Connections without pending dispatch demand ignore it without touching the database; blocked demand on another Agent can immediately retry its Claim.

type RuntimeWebSocketConcurrencyConfig added in v0.1.56

type RuntimeWebSocketConcurrencyConfig struct {
	ConnectionMaxInflight int
	ProcessMaxInflight    int
	LaneQueueDepth        int
}

RuntimeWebSocketConcurrencyConfig bounds independent Attempt work without changing message ordering or protocol behavior.

type RuntimeWireCompatibilitySnapshot added in v0.1.56

type RuntimeWireCompatibilitySnapshot struct {
	CurrentContractDigest     string
	SupportedContractDigests  []string
	PreviousSupportedUntilRFC string
}

RuntimeWireCompatibilitySnapshot is a read-only copy of Core's bounded wire compatibility ring for public discovery. Adapter details remain internal.

func CurrentRuntimeWireCompatibility added in v0.1.56

func CurrentRuntimeWireCompatibility() RuntimeWireCompatibilitySnapshot

type Service

type Service struct {
	// contains filtered or unexported fields
}

Service invokes Agents and maintains the Runtime Run/Attempt state. Core-owned HTTP/MCP calls create and confirm a fenced Attempt before any network I/O. Progress is written through EventStore and terminal facts only through ResultFinalizer (or the database deadline reconciler after a crash).

func NewService

func NewService(pool *pgxpool.Pool, cfg *config.Config) *Service

NewService 构造 Service。HTTP client timeout 取自 cfg.RunTimeoutSeconds(默认 60s)。

func (*Service) ActivateRuntimeNode added in v0.1.56

func (s *Service) ActivateRuntimeNode(ctx context.Context, nodeID uuid.UUID) (*RuntimeNodeListItem, error)

ActivateRuntimeNode is a guarded rollback of an administrative drain. It is deliberately separate from heartbeat/session attach: only a quiescent, current, non-revoked Node can move backwards from draining to active.

func (*Service) AppendRuntimeEvent added in v0.1.56

func (s *Service) AppendRuntimeEvent(
	ctx context.Context,
	principal RuntimeEventPrincipal,
	identity RuntimeAttemptIdentity,
	req RuntimeEventRequest,
) (RuntimeEventAck, error)

AppendRuntimeEvent persists one Runtime execution Event and emits projections only for the transaction winner. A replay receives its original ACK without duplicating callbacks, messages, or artifacts.

func (*Service) BrowserHumanControl added in v0.1.56

func (s *Service) BrowserHumanControl() *BrowserHumanControl

func (*Service) BrowserObservation added in v0.1.56

func (s *Service) BrowserObservation() *BrowserObservation

func (*Service) CancelRun

func (s *Service) CancelRun(ctx context.Context, userID, runID uuid.UUID) (*RunResponse, error)

CancelRun cancels an owned Runtime Run through the durable coordinator.

func (*Service) ConfigureCoreRuntime added in v0.1.56

func (s *Service) ConfigureCoreRuntime(coreInstanceID uuid.UUID)

ConfigureCoreRuntime binds Core-owned Attempts to the process identity used by cluster membership. It must be called during startup before serving Runs.

func (*Service) DrainRuntimeNode added in v0.1.56

func (s *Service) DrainRuntimeNode(ctx context.Context, nodeID uuid.UUID) (*RuntimeNodeListItem, error)

func (*Service) DryRun

func (s *Service) DryRun(
	ctx context.Context,
	agent *db.Agent,
	input map[string]interface{},
) (map[string]interface{}, string)

DryRun 让创作者侧 endpoint 跑一次「不计费、不写 runs」的探活调用。

用于 Agent 接入和测评,返回给定输入的输出或错误信息。 HTTP/MCP 直接命中 endpoint,随机 runID 仅为响应头标识,不写 DB。 Runtime Worker 则需要创建持久化 Run,通过正常 调度通道执行并等待终态;调用者取消或诊断超时时取消该 Run。

返回 (output, errMsg):errMsg 为空字符串时表示成功。

func (*Service) FinalizeRuntimeResult added in v0.1.56

func (s *Service) FinalizeRuntimeResult(
	ctx context.Context,
	principal RuntimeResultPrincipal,
	req RuntimeResultRequest,
) (RuntimeResultAck, error)

FinalizeRuntimeResult is the transport-neutral Runtime Result entrypoint. Task 6 adapters authenticate and decode envelopes before calling it.

func (*Service) GetConversationRuns added in v0.1.56

func (s *Service) GetConversationRuns(
	ctx context.Context,
	userID, anchorRunID uuid.UUID,
) (*ConversationRunListResponse, error)

func (*Service) GetRun

func (s *Service) GetRun(ctx context.Context, userID, runID uuid.UUID) (*RunResponse, error)

GetRun 查单条调用详情;调用者本人和被调用 Agent 创作者可看。

func (*Service) GetRunCancellationEvidence added in v0.1.56

func (s *Service) GetRunCancellationEvidence(
	ctx context.Context,
	actorUserID, runID uuid.UUID,
) (RunCancellationEvidence, error)

GetRunCancellationEvidence is a read-only projection of physical stop evidence. A public Run status of canceled is deliberately not sufficient.

func (*Service) GetRunWaitStatus added in v0.1.56

func (s *Service) GetRunWaitStatus(ctx context.Context, userID, runID uuid.UUID) (string, error)

GetRunWaitStatus performs the narrow authorization and lifecycle read used by Prefer: wait. Full Run evidence is assembled only for the HTTP response.

func (*Service) ListRunArtifacts

func (s *Service) ListRunArtifacts(ctx context.Context, userID, runID uuid.UUID) ([]RunArtifactResponse, error)

ListRunArtifacts returns persisted artifacts for a run. Only the run owner can read them.

func (*Service) ListRunEvents

func (s *Service) ListRunEvents(ctx context.Context, userID, runID uuid.UUID, afterSequence, limit int32) ([]RunEventResponse, error)

ListRunEvents 查询单个 run 的事件流;仅 owner 可看。

func (*Service) ListRunEventsPage added in v0.1.56

func (s *Service) ListRunEventsPage(
	ctx context.Context,
	userID, runID uuid.UUID,
	afterSequence, limit int32,
) (*RunEventPageResponse, error)

ListRunEventsPage returns the readable event window and enough retention state for clients to distinguish an empty page from a complete history.

func (*Service) ListRunMessages

func (s *Service) ListRunMessages(ctx context.Context, userID, runID uuid.UUID) ([]RunMessageResponse, error)

ListRunMessages returns stable message replay records for a run.

func (*Service) ListRuntimeDeadLetters added in v0.1.56

func (s *Service) ListRuntimeDeadLetters(
	ctx context.Context,
	limit, offset int32,
) (*RuntimeDeadLetterListResponse, error)

ListRuntimeDeadLetters returns the admin operational inventory without Run input/output or credential-linked execution identity.

func (*Service) ListRuntimeNodes added in v0.1.56

func (s *Service) ListRuntimeNodes(
	ctx context.Context,
	limit, offset int32,
) (*RuntimeNodeListResponse, error)

ListRuntimeNodes returns a database-clock snapshot of enrolled Runtime Nodes. Session counts intentionally use the canonical Runtime liveness window as runtime availability; an attached but stale Session is not presented as online to an operator.

func (*Service) LookupRunByCreationIdentity added in v0.1.56

func (s *Service) LookupRunByCreationIdentity(
	ctx context.Context,
	actorUserID uuid.UUID,
	keyHash, fingerprint []byte,
) (*RunResponse, bool, error)

LookupRunByCreationIdentity is strictly read-only and requires the exact actor/key/fingerprint tuple persisted before launch. It never needs input.

func (*Service) LookupRunByCreationRequest added in v0.1.56

func (s *Service) LookupRunByCreationRequest(
	ctx context.Context,
	userID uuid.UUID,
	req *RunRequest,
	source string,
) (*RunResponse, bool, error)

LookupRunByCreationRequest performs the read-only half of Run creation idempotency. It deliberately resolves the immutable user/key/fingerprint identity before any mutable Agent eligibility lookup, so retiring an Agent cannot invalidate a committed Run. A missing identity is reported without creating a Run or consulting the Agent registry.

func (*Service) PrepareRunCreationIdentity added in v0.1.56

func (s *Service) PrepareRunCreationIdentity(req *RunRequest, source string) (RunCreationIdentity, error)

func (*Service) ReplayRun added in v0.1.56

func (s *Service) ReplayRun(
	ctx context.Context,
	userID, sourceRunID uuid.UUID,
	idempotencyKey, source string,
) (*RunResponse, error)

ReplayRun creates a new asynchronous Run from one owned dead-letter Run. The source Run, Attempt history, events, ledger and DLQ are immutable; only the new Run carries replay_of_run_id.

func (*Service) RevokeRuntimeNode added in v0.1.56

func (s *Service) RevokeRuntimeNode(
	ctx context.Context,
	nodeID uuid.UUID,
	reason string,
) (*RuntimeNodeListItem, error)

func (*Service) Run

func (s *Service) Run(ctx context.Context, userID uuid.UUID, req *RunRequest, source string) (*RunResponse, error)

Run 调用 Agent。

流程见 Service 注释。Core 不执行商业结算;财务字段仅保留为外部兼容记录。

source 标记调用来源:'web' / 'mcp' / 'api',写入 runs.source 以便 /usage 分类显示。 传空字符串时默认 'web',便于旧调用方零修改。

func (*Service) RunWorkflowChild added in v0.1.56

func (s *Service) RunWorkflowChild(
	ctx context.Context,
	actorUserID uuid.UUID,
	req *RunRequest,
	source string,
	fence WorkflowChildLaunchFence,
) (*RunResponse, error)

RunWorkflowChild makes the workflow parent/step launch fence part of the Run creation transaction. The created state and child Run ID commit atomically.

func (*Service) SetRunEffectHandlers added in v0.1.56

func (s *Service) SetRunEffectHandlers(
	webhook WebhookRunEffectHandler,
	delivery DeliveryRunEffectHandler,
)

SetRunEffectHandlers wires the durable terminal-effect dispatchers. It does not start the worker; coreapi starts one worker under the process root context after every handler has been constructed.

func (*Service) SetTaskCallbackEnqueuer

func (s *Service) SetTaskCallbackEnqueuer(w TaskCallbackEnqueuer)

SetTaskCallbackEnqueuer 注入 task callback 触发器。

func (*Service) StartCoreAttemptCancellationCoordinator added in v0.1.56

func (s *Service) StartCoreAttemptCancellationCoordinator(ctx context.Context)

StartCoreAttemptCancellationCoordinator starts the single Core-scoped database fallback for local HTTP/MCP Attempts. Redis run.cancel signals remain the immediate path; this coordinator only closes missed-signal gaps.

func (*Service) StartExternalRun added in v0.1.56

func (s *Service) StartExternalRun(
	ctx context.Context,
	actorUserID uuid.UUID,
	req *RunRequest,
	source string,
	fence ExternalExecutionLaunchFence,
) (*RunResponse, error)

StartExternalRun creates a Run only while the launch token is current. The external key row is locked before its execution row, which is the global key -> execution lock order shared with cancellation.

func (*Service) StartRun

func (s *Service) StartRun(ctx context.Context, userID uuid.UUID, req *RunRequest, source string) (*RunResponse, error)

StartRun 创建 running run 并在后台执行;调用方可用 GetRun/ListRunEvents/SSE 查询进度。

func (*Service) ValidateRuntimeToken

func (s *Service) ValidateRuntimeToken(ctx context.Context, plaintext string, acceptedScopes ...string) (db.AgentRuntimeToken, error)

type TaskCallbackAuthentication

type TaskCallbackAuthentication struct {
	Scheme      string `json:"scheme,omitempty"`
	Credentials string `json:"credentials,omitempty"`
}

type TaskCallbackConfig

type TaskCallbackConfig struct {
	URL             string                      `json:"url,omitempty"`
	Token           string                      `json:"token,omitempty"`
	Secret          string                      `json:"secret,omitempty"`
	Authentication  *TaskCallbackAuthentication `json:"authentication,omitempty"`
	Metadata        map[string]interface{}      `json:"metadata,omitempty"`
	EventTypes      []string                    `json:"eventTypes,omitempty"`
	EventTypesAlias []string                    `json:"event_types,omitempty"`
}

type TaskCallbackEnqueuer

type TaskCallbackEnqueuer interface {
	EnqueueRunEvent(ctx context.Context, event db.RunEvent) error
}

TaskCallbackEnqueuer 触发 task callback,payload 来自 run_events。

type WebhookRunEffectHandler added in v0.1.56

type WebhookRunEffectHandler interface {
	AttemptAgentWebhookEffect(context.Context, db.RunEffectOutbox) RunEffectAttemptResult
	AttemptTaskCallbackEffect(context.Context, db.RunEffectOutbox) RunEffectAttemptResult
	ResetWebhookEffectDelivery(context.Context, db.RunEffectOutbox) error
	EnqueueRunEventDurable(context.Context, db.RunEvent) error
}

WebhookRunEffectHandler owns Agent webhook and terminal task callback materialization. Downstream delivery IDs must equal effect.ID.

type WorkerObservation added in v0.1.56

type WorkerObservation struct {
	Category  string
	Reason    string
	BatchSize int
}

WorkerObservation is a bounded, payload-free marker for tests and load baselines. Category and reason are stable code values; BatchSize is a count, never a resource identifier. Production leaves the observer nil.

type WorkerObserver added in v0.1.56

type WorkerObserver interface {
	ObserveWorker(WorkerObservation)
}

WorkerObserver must not perform database or network I/O. Implementations used by tests should be concurrency-safe and return immediately.

type WorkerObserverFunc added in v0.1.56

type WorkerObserverFunc func(WorkerObservation)

func (WorkerObserverFunc) ObserveWorker added in v0.1.56

func (f WorkerObserverFunc) ObserveWorker(observation WorkerObservation)

type WorkflowChildLaunchFence added in v0.1.56

type WorkflowChildLaunchFence struct {
	WorkflowRunID  uuid.UUID
	WorkflowNodeID uuid.UUID
	LaunchToken    uuid.UUID
}

Jump to

Keyboard shortcuts

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