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
- Variables
- func BuildRuntimeInvocationProof(token string, request RuntimeInvocationProofRequest) (string, error)
- func CanonicalizeRFC8785(value any) ([]byte, error)
- func DecodeRuntimeBody[P any](reader io.Reader) (P, error)
- func DecodeRuntimeMessagePayload[P any](envelope RuntimeEnvelope, expected RuntimeMessageType) (P, error)
- func FingerprintRunCreation(input RunFingerprintInput) ([sha256.Size]byte, error)
- func HashIdempotencyKey(key string) ([sha256.Size]byte, error)
- func IsRuntimeEventError(err error, code RuntimeEventErrorCode) bool
- func IsRuntimeLeaseError(err error, code RuntimeLeaseErrorCode) bool
- func IsRuntimeResultError(err error, code RuntimeResultErrorCode) bool
- func IsRuntimeSessionError(err error, code RuntimeSessionErrorCode) bool
- func MarshalRuntimeSignal(signal RuntimeSignal) ([]byte, error)
- func RequireRuntimeClusterOperation(ctx context.Context, querier interface{ ... }, ...) error
- func RevokeAgentCredential(ctx context.Context, pool *pgxpool.Pool, creatorID uuid.UUID, ...) (bool, error)
- func RuntimeEventFingerprint(request RuntimeEventRequest) ([]byte, error)
- func RuntimeHTTPStatus(code RuntimeErrorCode) int
- func RuntimeRequiredFeatures() []string
- func RuntimeResultFingerprint(request RuntimeResultRequest) ([]byte, error)
- func RuntimeWebSocketCloseCode(code RuntimeErrorCode) (int, bool)
- func StartRunEffectWorker(ctx context.Context, svc *Service, cfg RunEffectWorkerConfig)
- func StartRunEffectWorkerWithWake(ctx context.Context, svc *Service, cfg RunEffectWorkerConfig, ...)
- func StartRuntimeCredentialRevocationWake(ctx context.Context, source eventwake.TopicSource, hub *RuntimeWakeHub)
- func StartRuntimeDispatchWakeReconciler(ctx context.Context, reconciler *RuntimeDispatchWakeReconciler, ...)
- func StartRuntimeMaintenanceWorker(ctx context.Context, reconciler runtimeDeadlineReconcileWorker, ...)
- func StartRuntimeMaintenanceWorkerWithWake(ctx context.Context, reconciler runtimeDeadlineReconcileWorker, ...)
- func StartRuntimeSignalOutboxWorker(ctx context.Context, worker *RuntimeSignalOutboxWorker, ...)
- func StartRuntimeSignalOutboxWorkerWithWake(ctx context.Context, worker *RuntimeSignalOutboxWorker, ...)
- func StartRuntimeSignalSubscriber(ctx context.Context, bus RuntimeSignalBus, instanceID uuid.UUID, ...)
- func ValidateRuntimeEnvelope(envelope RuntimeEnvelope) error
- func ValidateRuntimePayload(payload any) error
- func ValidateRuntimeReplyCorrelation(request, reply RuntimeEnvelope) error
- func ValidateRuntimeSignal(signal RuntimeSignal) error
- func VerifyRuntimeInvocationProof(token, proof string, request RuntimeInvocationProofRequest) error
- type AgentA2AContext
- type AgentError
- type AgentEvent
- type AgentRequest
- type AgentResponse
- type AttemptIdentity
- type AuthenticatedRuntimePrincipal
- type BrowserHumanControl
- func (control *BrowserHumanControl) BindCommandSender(sender browserViewerCommandSender)
- func (control *BrowserHumanControl) Claim(ctx context.Context, userID uuid.UUID, runID uuid.UUID) (BrowserHumanControlState, error)
- func (control *BrowserHumanControl) Input(ctx context.Context, userID uuid.UUID, runID uuid.UUID, ...) error
- func (control *BrowserHumanControl) PauseFromEvent(ctx context.Context, identity RuntimeAttemptIdentity, payload map[string]any) error
- func (control *BrowserHumanControl) PublishFrame(frame BrowserViewerFramePayload) error
- func (control *BrowserHumanControl) Release(ctx context.Context, userID uuid.UUID, runID uuid.UUID) (BrowserHumanControlState, error)
- func (control *BrowserHumanControl) Resume(ctx context.Context, userID uuid.UUID, runID uuid.UUID) (BrowserHumanControlState, error)
- func (control *BrowserHumanControl) RunGC(ctx context.Context)
- func (control *BrowserHumanControl) State(ctx context.Context, userID uuid.UUID, runID uuid.UUID) (BrowserHumanControlState, error)
- func (control *BrowserHumanControl) WaitFrame(ctx context.Context, userID uuid.UUID, runID uuid.UUID, after uint64) (*BrowserViewerFramePayload, error)
- type BrowserHumanControlState
- type BrowserObservation
- func (observation *BrowserObservation) AuthorizeOwner(ctx context.Context, runID uuid.UUID, callerUserID uuid.UUID) error
- func (observation *BrowserObservation) BindCommandSender(sender browserObserverCommandSender)
- func (observation *BrowserObservation) BindInstance(instance uuid.UUID)
- func (observation *BrowserObservation) CloseSessionObservations(ctx context.Context, runtimeSessionID uuid.UUID, endReason string) error
- func (observation *BrowserObservation) HandleEvent(ctx context.Context, event BrowserObserverEventPayload) (BrowserObserverEventAckPayload, error)
- func (observation *BrowserObservation) ProjectFromEvent(ctx context.Context, identity RuntimeAttemptIdentity, payload map[string]any, ...) error
- func (observation *BrowserObservation) ReconcileAbandoned(ctx context.Context)
- func (observation *BrowserObservation) ReconcileExpired(ctx context.Context) error
- func (observation *BrowserObservation) ReconcileForeignExpired(ctx context.Context) error
- func (observation *BrowserObservation) RecordFrames(ctx context.Context, leaseID uuid.UUID, count int64) error
- func (observation *BrowserObservation) ResolveIdentity(ctx context.Context, runID uuid.UUID, callerUserID uuid.UUID, isAdmin bool) (BrowserObserverIdentity, error)
- func (observation *BrowserObservation) RunGC(ctx context.Context)
- func (observation *BrowserObservation) Start(ctx context.Context, runID uuid.UUID, observerUserID uuid.UUID, isAdmin bool, ...) (BrowserObservationState, error)
- func (observation *BrowserObservation) State(ctx context.Context, runID uuid.UUID) (BrowserObservationState, error)
- func (observation *BrowserObservation) Stop(ctx context.Context, runID uuid.UUID, endReason string) error
- func (observation *BrowserObservation) WaitFrame(ctx context.Context, runID uuid.UUID, after int64) (*BrowserObservationFrame, error)
- type BrowserObservationFrame
- type BrowserObservationState
- type BrowserObserverAction
- type BrowserObserverCommandPayload
- type BrowserObserverEventAckPayload
- type BrowserObserverEventKind
- type BrowserObserverEventPayload
- type BrowserObserverFramePayload
- type BrowserObserverIdentity
- type BrowserViewerAction
- type BrowserViewerCommand
- type BrowserViewerCommandMessage
- type BrowserViewerCommandPayload
- type BrowserViewerFrameAckMessage
- type BrowserViewerFrameAckPayload
- type BrowserViewerFrameMessage
- type BrowserViewerFramePayload
- type BrowserViewerInputPayload
- type CallAgentRequest
- type CancelCommand
- type ConversationContext
- type ConversationMessage
- type ConversationRunItem
- type ConversationRunListResponse
- type DBRuntimeNodeCredentialVerifier
- type DecodedPendingCommand
- type DelegatedRunReadRequest
- type DelegatedRunView
- type Delegation
- type DeliveryRunEffectHandler
- type DrainCommand
- type EventRange
- type EventStore
- func (s *EventStore) Append(ctx context.Context, principal RuntimeEventPrincipal, ...) (RuntimeEventAck, error)
- func (s *EventStore) AppendEvent(ctx context.Context, principal RuntimeEventPrincipal, ...) (RuntimeEventAck, error)
- func (s *EventStore) MissingClientEventRanges(ctx context.Context, runID uuid.UUID, attemptID uuid.UUID, ...) ([]EventRange, error)
- func (s *EventStore) RequireCompleteClientEvents(ctx context.Context, runID uuid.UUID, attemptID uuid.UUID, ...) error
- type ExternalExecutionLaunchFence
- type Handler
- func (h *Handler) ActivateRuntimeNode(c echo.Context) error
- func (h *Handler) CancelRun(c echo.Context) error
- func (h *Handler) ClaimBrowserControl(c echo.Context) error
- func (h *Handler) DrainRuntimeNode(c echo.Context) error
- func (h *Handler) GetAdminBrowserObservation(c echo.Context) error
- func (h *Handler) GetAdminBrowserObservationFrame(c echo.Context) error
- func (h *Handler) GetBrowserControl(c echo.Context) error
- func (h *Handler) GetBrowserControlFrame(c echo.Context) error
- func (h *Handler) GetBrowserObservation(c echo.Context) error
- func (h *Handler) GetBrowserObservationFrame(c echo.Context) error
- func (h *Handler) GetConversationRuns(c echo.Context) error
- func (h *Handler) GetRun(c echo.Context) error
- func (h *Handler) GetRunArtifacts(c echo.Context) error
- func (h *Handler) GetRunEvents(c echo.Context) error
- func (h *Handler) GetRunMessages(c echo.Context) error
- func (h *Handler) ListRuntimeDeadLetters(c echo.Context) error
- func (h *Handler) ListRuntimeNodes(c echo.Context) error
- func (h *Handler) PostRun(c echo.Context) error
- func (h *Handler) PostRunAsync(c echo.Context) error
- func (h *Handler) RegisterAdmin(api *echo.Group, jwtMw, adminMw echo.MiddlewareFunc)
- func (h *Handler) RegisterAgentRuntime(api *echo.Group)
- func (h *Handler) RegisterAgentRuntimeAttachOnly(api *echo.Group)
- func (h *Handler) RegisterObservation(api *echo.Group, jwtMw echo.MiddlewareFunc)
- func (h *Handler) RegisterProtected(api *echo.Group, runMw, queryMw echo.MiddlewareFunc)
- func (h *Handler) ReleaseBrowserControl(c echo.Context) error
- func (h *Handler) ReplayRun(c echo.Context) error
- func (h *Handler) ResumeBrowserControl(c echo.Context) error
- func (h *Handler) RevokeRuntimeNode(c echo.Context) error
- func (h *Handler) RuntimeController() *RuntimeHTTPController
- func (h *Handler) SendBrowserControlInput(c echo.Context) error
- func (h *Handler) SetRunUpdateSource(source RunUpdateSource)
- func (h *Handler) SetRuntimeDependencies(dependencies RuntimeHTTPDependencies)
- func (h *Handler) SetWorkerObserver(observer WorkerObserver)
- func (h *Handler) StartAdminBrowserObservation(c echo.Context) error
- func (h *Handler) StartBrowserObservation(c echo.Context) error
- func (h *Handler) StopAdminBrowserObservation(c echo.Context) error
- func (h *Handler) StopBrowserObservation(c echo.Context) error
- func (h *Handler) StreamRunEvents(c echo.Context) error
- type IdempotencyError
- type IdempotencyErrorClass
- type LocalSignalBus
- type MTLSRuntimeDeviceAuthenticator
- type PendingCommand
- type RedisRuntimePresenceStore
- func (s *RedisRuntimePresenceStore) ListByAgent(ctx context.Context, agentID uuid.UUID) ([]RuntimePresence, error)
- func (s *RedisRuntimePresenceStore) Refresh(ctx context.Context, presence RuntimePresence, ttl time.Duration) error
- func (s *RedisRuntimePresenceStore) Remove(ctx context.Context, presence RuntimePresence) error
- type RedisRuntimeSessionLeaseStore
- func (s *RedisRuntimeSessionLeaseStore) Forget(ctx context.Context, runtimeSessionID uuid.UUID) error
- func (s *RedisRuntimeSessionLeaseStore) ListExpired(ctx context.Context, limit int) ([]uuid.UUID, error)
- func (s *RedisRuntimeSessionLeaseStore) Lookup(ctx context.Context, runtimeSessionID uuid.UUID) (RuntimeSessionLease, bool, error)
- func (s *RedisRuntimeSessionLeaseStore) RefreshBatch(ctx context.Context, records []RuntimeSessionLeaseRecord, ...) error
- func (s *RedisRuntimeSessionLeaseStore) Remove(ctx context.Context, record RuntimeSessionLeaseRecord) error
- func (s *RedisRuntimeSessionLeaseStore) ScheduleCheck(ctx context.Context, runtimeSessionID uuid.UUID, after time.Duration) error
- type RedisSignalBus
- func (b *RedisSignalBus) Check(ctx context.Context, registrations []RuntimeConnectionRegistration) ([]RuntimeCredentialProjectionResult, error)
- func (b *RedisSignalBus) Close() error
- func (b *RedisSignalBus) Health(ctx context.Context) error
- func (b *RedisSignalBus) MarkActive(ctx context.Context, registrations []RuntimeConnectionRegistration) error
- func (b *RedisSignalBus) Publish(ctx context.Context, signal RuntimeSignal) error
- func (b *RedisSignalBus) RuntimeCredentialProjectionStore() (RuntimeCredentialProjectionStore, error)
- func (b *RedisSignalBus) RuntimePresenceStore() (RuntimePresenceStore, error)
- func (b *RedisSignalBus) RuntimeSessionLeaseStore() (RuntimeSessionLeaseStore, error)
- func (b *RedisSignalBus) Subscribe(ctx context.Context, handler RuntimeSignalHandler) error
- type RedisSignalBusConfig
- type ResultClassificationInput
- type ResultClassifier
- type ResultClassifierFunc
- type ResultFinalizer
- type ResultRetryPlanner
- type ResultRetryPlannerFunc
- type ResumeAttempt
- type RevokeCommand
- type RevokeRuntimeNodeRequest
- type RunA2AContextRequest
- type RunA2AContextResponse
- type RunArtifactResponse
- type RunAssignedMessage
- type RunAssignedPayload
- type RunAssignmentAckMessage
- type RunAssignmentAckPayload
- type RunAssignmentConfirmedMessage
- type RunAssignmentConfirmedPayload
- type RunAssignmentRejectMessage
- type RunAssignmentRejectPayload
- type RunAssignmentRejectedMessage
- type RunAssignmentRejectedPayload
- type RunCancelAckMessage
- type RunCancelAckPayload
- type RunCancelMessage
- type RunCancelPayload
- type RunCancellationEvidence
- type RunCancellationState
- type RunCreationIdentity
- type RunEffectAttemptResult
- type RunEffectWorker
- func (w *RunEffectWorker) ProcessOnce(ctx context.Context, cfg RunEffectWorkerConfig) (int, error)
- func (w *RunEffectWorker) Replay(ctx context.Context, effectID uuid.UUID, actorType string, actorID *uuid.UUID, ...) (*db.RunEffectOutbox, error)
- func (w *RunEffectWorker) SetHandlers(webhook WebhookRunEffectHandler, delivery DeliveryRunEffectHandler)
- type RunEffectWorkerConfig
- type RunErrorPayload
- type RunEventAckMessage
- type RunEventAckPayload
- type RunEventMessage
- type RunEventPageMeta
- type RunEventPageResponse
- type RunEventPayload
- type RunEventResponse
- type RunEvidenceSummary
- type RunFingerprintA2A
- type RunFingerprintDelegation
- type RunFingerprintInput
- type RunFingerprintSource
- type RunFingerprintTarget
- type RunFingerprintTransport
- type RunLeaseRenewMessage
- type RunLeaseRenewPayload
- type RunLeaseRenewedMessage
- type RunLeaseRenewedPayload
- type RunLeaseRevokedMessage
- type RunLeaseRevokedPayload
- type RunMessageResponse
- type RunNextAction
- type RunRequest
- type RunRequirementEvidenceResponse
- type RunResponse
- type RunResultAckMessage
- type RunResultAckPayload
- type RunResultMessage
- type RunResultPayload
- type RunResumeAcceptedMessage
- type RunResumeAcceptedPayload
- type RunSummary
- type RunTaskCallbackResponse
- type RunUpdateHub
- type RunUpdateSource
- type RunUpdateSubscription
- type RuntimeAdmissionIdentity
- type RuntimeAdmissionLimitConfig
- type RuntimeAdmissionLimiter
- type RuntimeAssignmentRejectOutcome
- type RuntimeAssignmentRejectReason
- type RuntimeAttemptIdentity
- type RuntimeAuthenticationMode
- type RuntimeCancelState
- type RuntimeCancellationAPI
- type RuntimeCancellationCoordinator
- func (c *RuntimeCancellationCoordinator) AckCancel(ctx context.Context, principal RuntimeSessionPrincipal, ...) (RunCancellationState, error)
- func (c *RuntimeCancellationCoordinator) AcknowledgeCoreStopped(ctx context.Context, coreInstanceID uuid.UUID, identity RuntimeAttemptIdentity) (RunCancellationState, error)
- func (c *RuntimeCancellationCoordinator) CancelOwnedRun(ctx context.Context, requesterID, runID uuid.UUID, reason string) (RuntimeCancellationResult, error)
- func (c *RuntimeCancellationCoordinator) NextCommand(ctx context.Context, principal RuntimeSessionPrincipal) (*PendingCommand, time.Time, error)
- func (c *RuntimeCancellationCoordinator) PollCommands(ctx context.Context, principal RuntimeSessionPrincipal) (RuntimeCommandsResponse, error)
- func (c *RuntimeCancellationCoordinator) ReapExpiredCancellation(ctx context.Context) (*RunCancellationState, error)
- func (c *RuntimeCancellationCoordinator) ReapExpiredCancellations(ctx context.Context, limit int) (int, error)
- type RuntimeCancellationResult
- type RuntimeClaimRequest
- type RuntimeClusterControlSnapshot
- type RuntimeClusterCoordinator
- type RuntimeClusterIdentity
- type RuntimeClusterMemberSnapshot
- type RuntimeClusterMode
- type RuntimeClusterOperation
- type RuntimeClusterReadiness
- type RuntimeClusterRepository
- type RuntimeClusterSnapshot
- type RuntimeCommand
- type RuntimeCommandsResponse
- type RuntimeConnectionIdentity
- type RuntimeConnectionRegistration
- type RuntimeCredentialConnectionValidator
- type RuntimeCredentialProjectionResult
- type RuntimeCredentialProjectionState
- type RuntimeCredentialProjectionStore
- type RuntimeCredentialProjectionStoreProvider
- type RuntimeCredentialReconciler
- type RuntimeCredentialReconcilerConfig
- type RuntimeCredentialValidationResult
- type RuntimeDeadLetterListItem
- type RuntimeDeadLetterListResponse
- type RuntimeDeadlineReconciler
- type RuntimeDelegationAPI
- type RuntimeDelegationAuthorization
- type RuntimeDelegationService
- func (s *RuntimeDelegationService) CallAgent(ctx context.Context, authorization RuntimeDelegationAuthorization) (RunSummary, error)
- func (s *RuntimeDelegationService) ReadDelegatedRun(ctx context.Context, authorization RuntimeDelegationAuthorization) (DelegatedRunView, error)
- func (s *RuntimeDelegationService) ResolveInvocationDevice(ctx context.Context, invocationToken string) (RuntimeDeviceIdentity, error)
- type RuntimeDeviceAuthenticator
- type RuntimeDeviceIdentity
- type RuntimeDispatchState
- type RuntimeDispatchWakeReconcileResult
- type RuntimeDispatchWakeReconciler
- type RuntimeDispatchWakeReconcilerConfig
- type RuntimeDrainMessage
- type RuntimeDrainPayload
- type RuntimeEnvelope
- type RuntimeEnvelopeFields
- type RuntimeError
- type RuntimeErrorBody
- type RuntimeErrorCode
- type RuntimeErrorMessage
- type RuntimeEventAck
- type RuntimeEventError
- type RuntimeEventErrorCode
- type RuntimeEventPrincipal
- type RuntimeEventProjector
- type RuntimeEventRequest
- type RuntimeHTTPController
- func (h *RuntimeHTTPController) AckAssignment(c echo.Context) error
- func (h *RuntimeHTTPController) AckCancel(c echo.Context) error
- func (h *RuntimeHTTPController) AppendEvent(c echo.Context) error
- func (h *RuntimeHTTPController) AuthenticateAgentRequest(c echo.Context) (AuthenticatedRuntimePrincipal, *RuntimeTransportError)
- func (h *RuntimeHTTPController) CallAgent(c echo.Context) error
- func (h *RuntimeHTTPController) ClaimRun(c echo.Context) error
- func (h *RuntimeHTTPController) CloseSession(c echo.Context) error
- func (h *RuntimeHTTPController) CreateSession(c echo.Context) error
- func (h *RuntimeHTTPController) DrainSession(c echo.Context) error
- func (h *RuntimeHTTPController) FinalizeResult(c echo.Context) error
- func (h *RuntimeHTTPController) HeartbeatSession(c echo.Context) error
- func (h *RuntimeHTTPController) PollCommands(c echo.Context) error
- func (h *RuntimeHTTPController) ReadDelegatedRun(c echo.Context) error
- func (h *RuntimeHTTPController) Register(api *echo.Group)
- func (h *RuntimeHTTPController) RegisterAttachOnly(api *echo.Group)
- func (h *RuntimeHTTPController) RejectAssignment(c echo.Context) error
- func (h *RuntimeHTTPController) RenewLease(c echo.Context) error
- func (h *RuntimeHTTPController) ResumeRuns(c echo.Context) error
- func (h *RuntimeHTTPController) SendBrowserObserverCommand(runtimeSessionID uuid.UUID, payload BrowserObserverCommandPayload) error
- func (h *RuntimeHTTPController) SendBrowserViewerCommand(runtimeSessionID uuid.UUID, payload BrowserViewerCommandPayload) error
- func (h *RuntimeHTTPController) Shutdown(ctx context.Context) error
- func (h *RuntimeHTTPController) WebSocket(c echo.Context) error
- type RuntimeHTTPDependencies
- type RuntimeHelloMessage
- type RuntimeHelloPayload
- type RuntimeInvocationCapability
- type RuntimeInvocationCapabilityIssuer
- type RuntimeInvocationProofRequest
- type RuntimeInvocationSigner
- func NewRuntimeInvocationSigner(secret string) (*RuntimeInvocationSigner, error)
- func NewRuntimeInvocationSignerKeyring(activeKeyID string, secrets map[string]string) (*RuntimeInvocationSigner, error)
- func NewRuntimeInvocationSignerWithPrevious(activeKeyID, activeSecret, previousKeyID, previousSecret string) (*RuntimeInvocationSigner, error)
- func (s *RuntimeInvocationSigner) Issue(capability RuntimeInvocationCapability) (nodeEnvelope, invocationToken string, err error)
- func (s *RuntimeInvocationSigner) VerifyInvocationToken(token string, databaseNow time.Time) (RuntimeInvocationCapability, error)
- func (s *RuntimeInvocationSigner) VerifyNodeEnvelope(envelope string, databaseNow time.Time) (RuntimeInvocationCapability, error)
- type RuntimeInvocationVerifier
- type RuntimeLeaseAPI
- type RuntimeLeaseConfig
- type RuntimeLeaseError
- type RuntimeLeaseErrorCode
- type RuntimeLeaseService
- func (s *RuntimeLeaseService) AckAssignment(ctx context.Context, principal RuntimeSessionPrincipal, ...) (RunAssignmentConfirmedPayload, error)
- func (s *RuntimeLeaseService) ClaimOffer(ctx context.Context, principal RuntimeSessionPrincipal) (*RunAssignedPayload, error)
- func (s *RuntimeLeaseService) RejectAssignment(ctx context.Context, principal RuntimeSessionPrincipal, ...) (RunAssignmentRejectedPayload, error)
- func (s *RuntimeLeaseService) ReleaseUnackedOffer(ctx context.Context, principal RuntimeSessionPrincipal, reason ...string) error
- func (s *RuntimeLeaseService) RenewLease(ctx context.Context, principal RuntimeSessionPrincipal, ...) (RunLeaseRenewedPayload, error)
- type RuntimeLivenessPolicy
- type RuntimeMaintenanceResult
- type RuntimeMaintenanceWorkerConfig
- type RuntimeMessageType
- type RuntimeNodeCredentialVerifier
- type RuntimeNodeListItem
- type RuntimeNodeListResponse
- type RuntimePresence
- type RuntimePresenceStore
- type RuntimePresenceStoreProvider
- type RuntimePresentedCertificate
- type RuntimePrincipalBinder
- type RuntimeReadyMessage
- type RuntimeReadyPayload
- type RuntimeReconcileBatchResult
- type RuntimeResultAck
- type RuntimeResultClassification
- type RuntimeResultError
- type RuntimeResultErrorCode
- type RuntimeResultFailure
- type RuntimeResultFinalizer
- type RuntimeResultPrincipal
- type RuntimeResultRequest
- type RuntimeResumeAPI
- type RuntimeResumeAction
- type RuntimeResumeDecision
- type RuntimeResumeMessage
- type RuntimeResumePayload
- type RuntimeResumeResponse
- type RuntimeResumeService
- type RuntimeRunStatus
- type RuntimeSchemaContractSnapshot
- type RuntimeSessionAPI
- type RuntimeSessionCloseRequest
- type RuntimeSessionDrainRequest
- type RuntimeSessionError
- type RuntimeSessionErrorCode
- type RuntimeSessionHeartbeatRequest
- type RuntimeSessionIdentity
- type RuntimeSessionLease
- type RuntimeSessionLeaseManager
- func (m *RuntimeSessionLeaseManager) AbsenceReady() bool
- func (m *RuntimeSessionLeaseManager) Forget(ctx context.Context, runtimeSessionID uuid.UUID) error
- func (m *RuntimeSessionLeaseManager) Health() RuntimeSessionLeaseManagerHealth
- func (m *RuntimeSessionLeaseManager) HealthyFor(connectionID string) bool
- func (m *RuntimeSessionLeaseManager) ListExpired(ctx context.Context, limit int) ([]uuid.UUID, error)
- func (m *RuntimeSessionLeaseManager) Lookup(ctx context.Context, runtimeSessionID uuid.UUID) (RuntimeSessionLease, bool, error)
- func (m *RuntimeSessionLeaseManager) RefreshConnection(ctx context.Context, connectionID string) error
- func (m *RuntimeSessionLeaseManager) Register(record RuntimeSessionLeaseRecord) error
- func (m *RuntimeSessionLeaseManager) Run(ctx context.Context) error
- func (m *RuntimeSessionLeaseManager) ScheduleCheck(ctx context.Context, runtimeSessionID uuid.UUID, after time.Duration) error
- func (m *RuntimeSessionLeaseManager) Unregister(ctx context.Context, connectionID string, attachmentID uuid.UUID) (bool, error)
- func (m *RuntimeSessionLeaseManager) UpdatePresence(connectionID string, presence RuntimePresence) bool
- type RuntimeSessionLeaseManagerConfig
- type RuntimeSessionLeaseManagerHealth
- type RuntimeSessionLeaseRecord
- type RuntimeSessionLeaseStore
- type RuntimeSessionLeaseStoreProvider
- type RuntimeSessionPrincipal
- type RuntimeSessionReaper
- type RuntimeSessionRequest
- type RuntimeSessionService
- func (s *RuntimeSessionService) CloseSession(ctx context.Context, principal AuthenticatedRuntimePrincipal, ...) (RuntimeSessionState, error)
- func (s *RuntimeSessionService) CreateOrAttachSession(ctx context.Context, principal AuthenticatedRuntimePrincipal, ...) (RuntimeSessionState, error)
- func (s *RuntimeSessionService) DetachCutoverSessions(ctx context.Context) (int64, error)
- func (s *RuntimeSessionService) DrainSession(ctx context.Context, principal AuthenticatedRuntimePrincipal, ...) (RuntimeDrainPayload, error)
- func (s *RuntimeSessionService) HeartbeatSession(ctx context.Context, principal AuthenticatedRuntimePrincipal, ...) (RuntimeSessionState, error)
- func (s *RuntimeSessionService) ResolveSessionPrincipal(ctx context.Context, principal AuthenticatedRuntimePrincipal, ...) (RuntimeSessionPrincipal, error)
- func (s *RuntimeSessionService) ResolveWorkerSessionPrincipal(ctx context.Context, principal AuthenticatedRuntimePrincipal, workerID string) (RuntimeSessionPrincipal, error)
- type RuntimeSessionState
- type RuntimeSignal
- type RuntimeSignalBus
- type RuntimeSignalHandler
- type RuntimeSignalOutboxBatchResult
- type RuntimeSignalOutboxWorker
- type RuntimeSignalOutboxWorkerConfig
- type RuntimeTokenValidator
- type RuntimeTransport
- type RuntimeTransportError
- type RuntimeTransportPolicy
- type RuntimeTransportPolicyProvider
- type RuntimeTransportReason
- type RuntimeTypedEnvelope
- type RuntimeWakeHub
- func (h *RuntimeWakeHub) ConnectionSnapshot() []RuntimeConnectionRegistration
- func (h *RuntimeWakeHub) CredentialRevalidationWake() <-chan struct{}
- func (h *RuntimeWakeHub) RegisterConnection(identity RuntimeConnectionIdentity, credentialID uuid.UUID) <-chan struct{}
- func (h *RuntimeWakeHub) RegisterWebSocketDispatch(agentID uuid.UUID) (<-chan struct{}, bool)
- func (h *RuntimeWakeHub) RequireCredentialRevalidation()
- func (h *RuntimeWakeHub) RevokeConnections(identities []RuntimeConnectionIdentity) int
- func (h *RuntimeWakeHub) RevokeCredentialConnections(credentialID uuid.UUID, identities []RuntimeConnectionIdentity) int
- func (h *RuntimeWakeHub) UnregisterConnection(identity RuntimeConnectionIdentity)
- func (h *RuntimeWakeHub) UnregisterWebSocketDispatch(agentID uuid.UUID)
- func (h *RuntimeWakeHub) Wait(agentID uuid.UUID) <-chan struct{}
- func (h *RuntimeWakeHub) WaitControl(agentID uuid.UUID) <-chan struct{}
- func (h *RuntimeWakeHub) WaitDispatch(agentID uuid.UUID) <-chan struct{}
- func (h *RuntimeWakeHub) WaitNodeDispatch(nodeID uuid.UUID) <-chan struct{}
- func (h *RuntimeWakeHub) Wake(agentID uuid.UUID)
- func (h *RuntimeWakeHub) WakeAll()
- func (h *RuntimeWakeHub) WakeControl(agentID uuid.UUID)
- func (h *RuntimeWakeHub) WakeDispatch(agentID uuid.UUID)
- func (h *RuntimeWakeHub) WakeDispatchIfRegistered(agentID uuid.UUID) bool
- func (h *RuntimeWakeHub) WakeNodeDispatch(nodeID uuid.UUID)
- type RuntimeWebSocketConcurrencyConfig
- type RuntimeWireCompatibilitySnapshot
- type Service
- func (s *Service) ActivateRuntimeNode(ctx context.Context, nodeID uuid.UUID) (*RuntimeNodeListItem, error)
- func (s *Service) AppendRuntimeEvent(ctx context.Context, principal RuntimeEventPrincipal, ...) (RuntimeEventAck, error)
- func (s *Service) BrowserHumanControl() *BrowserHumanControl
- func (s *Service) BrowserObservation() *BrowserObservation
- func (s *Service) CancelRun(ctx context.Context, userID, runID uuid.UUID) (*RunResponse, error)
- func (s *Service) ConfigureCoreRuntime(coreInstanceID uuid.UUID)
- func (s *Service) DrainRuntimeNode(ctx context.Context, nodeID uuid.UUID) (*RuntimeNodeListItem, error)
- func (s *Service) DryRun(ctx context.Context, agent *db.Agent, input map[string]interface{}) (map[string]interface{}, string)
- func (s *Service) FinalizeRuntimeResult(ctx context.Context, principal RuntimeResultPrincipal, ...) (RuntimeResultAck, error)
- func (s *Service) GetConversationRuns(ctx context.Context, userID, anchorRunID uuid.UUID) (*ConversationRunListResponse, error)
- func (s *Service) GetRun(ctx context.Context, userID, runID uuid.UUID) (*RunResponse, error)
- func (s *Service) GetRunCancellationEvidence(ctx context.Context, actorUserID, runID uuid.UUID) (RunCancellationEvidence, error)
- func (s *Service) GetRunWaitStatus(ctx context.Context, userID, runID uuid.UUID) (string, error)
- func (s *Service) ListRunArtifacts(ctx context.Context, userID, runID uuid.UUID) ([]RunArtifactResponse, error)
- func (s *Service) ListRunEvents(ctx context.Context, userID, runID uuid.UUID, afterSequence, limit int32) ([]RunEventResponse, error)
- func (s *Service) ListRunEventsPage(ctx context.Context, userID, runID uuid.UUID, afterSequence, limit int32) (*RunEventPageResponse, error)
- func (s *Service) ListRunMessages(ctx context.Context, userID, runID uuid.UUID) ([]RunMessageResponse, error)
- func (s *Service) ListRuntimeDeadLetters(ctx context.Context, limit, offset int32) (*RuntimeDeadLetterListResponse, error)
- func (s *Service) ListRuntimeNodes(ctx context.Context, limit, offset int32) (*RuntimeNodeListResponse, error)
- func (s *Service) LookupRunByCreationIdentity(ctx context.Context, actorUserID uuid.UUID, keyHash, fingerprint []byte) (*RunResponse, bool, error)
- func (s *Service) LookupRunByCreationRequest(ctx context.Context, userID uuid.UUID, req *RunRequest, source string) (*RunResponse, bool, error)
- func (s *Service) PrepareRunCreationIdentity(req *RunRequest, source string) (RunCreationIdentity, error)
- func (s *Service) ReplayRun(ctx context.Context, userID, sourceRunID uuid.UUID, ...) (*RunResponse, error)
- func (s *Service) RevokeRuntimeNode(ctx context.Context, nodeID uuid.UUID, reason string) (*RuntimeNodeListItem, error)
- func (s *Service) Run(ctx context.Context, userID uuid.UUID, req *RunRequest, source string) (*RunResponse, error)
- func (s *Service) RunWorkflowChild(ctx context.Context, actorUserID uuid.UUID, req *RunRequest, source string, ...) (*RunResponse, error)
- func (s *Service) SetRunEffectHandlers(webhook WebhookRunEffectHandler, delivery DeliveryRunEffectHandler)
- func (s *Service) SetTaskCallbackEnqueuer(w TaskCallbackEnqueuer)
- func (s *Service) StartCoreAttemptCancellationCoordinator(ctx context.Context)
- func (s *Service) StartExternalRun(ctx context.Context, actorUserID uuid.UUID, req *RunRequest, source string, ...) (*RunResponse, error)
- func (s *Service) StartRun(ctx context.Context, userID uuid.UUID, req *RunRequest, source string) (*RunResponse, error)
- func (s *Service) ValidateRuntimeToken(ctx context.Context, plaintext string, acceptedScopes ...string) (db.AgentRuntimeToken, error)
- type TaskCallbackAuthentication
- type TaskCallbackConfig
- type TaskCallbackEnqueuer
- type WebhookRunEffectHandler
- type WorkerObservation
- type WorkerObserver
- type WorkerObserverFunc
- type WorkflowChildLaunchFence
Constants ¶
const ( DefaultRuntimeHTTPRequestsPerSecond = 100 DefaultRuntimeHTTPBurst = 200 DefaultRuntimeWebSocketMessagesPerSecond = 200 DefaultRuntimeWebSocketMessageBurst = 400 DefaultRuntimeWebSocketsPerIdentity = 16 )
const ( BrowserObserverMinFrameIntervalMS = 100 BrowserObserverMaxFrameIntervalMS = 5000 BrowserObserverMaxFrameBytes = 1 << 20 )
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 )
const ( RuntimeProtocolVersion = 2 RuntimeContractID = "openlinker.runtime.v2" RuntimeContractDigest = "4be9b2fe09eeedf0e37119075134064be88f93b301c502cdfa21a6cb978c6481" )
const ( RunEffectTypeAgentWebhook = "run.agent_webhook" RunEffectTypeTaskCallback = "run.task_callback" RunEffectTypeDefaultDelivery = "run.default_delivery" RunEffectTypeParentCompletion = "run.parent_completion" )
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 )
const ( RuntimeWSCloseAuthenticationFailed = 4401 RuntimeWSCloseClientUpgradeRequired = 4406 RuntimeWSCloseSessionConflict = 4409 RuntimeWSCloseRequiredFeatureMissing = 4412 RuntimeWSCloseProtocolError = 1002 RuntimeWSCloseInternalError = 1011 )
const ( RuntimeAttachmentIDHeader = "OpenLinker-Runtime-Attachment" RuntimeNodeIDHeader = "OpenLinker-Runtime-Node" RuntimeFallbackReasonHeader = "OpenLinker-Runtime-Fallback-Reason" )
const ( RuntimeTransportWebSocket RuntimeTransport = "websocket" RuntimeTransportLongPoll RuntimeTransport = "long_poll" RuntimeTransportUnknown RuntimeTransport = "unknown" RuntimeTransportReasonExplicit RuntimeTransportReason = "explicit" 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" )
const BrowserObservationFeature = "browser_authenticated_observation.v1"
BrowserObservationFeature is the optional Runtime capability that gates this whole surface.
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.
const MaxRuntimeSignalConnections = 64
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" )
const RuntimeBrowserExecutionProfileFeature = "browser_execution_profile.v1"
const RuntimeBrowserFullInteractionFeature = "browser_full_interaction.v1"
const RuntimeCredentialProjectionTTL = 15 * time.Minute
const (
RuntimeCredentialRevocationWakeTopic = "runtime.credential.revoke"
)
const (
RuntimeDelegatedRunReadFeature = "delegated_run_read.v1"
)
const ( // RuntimeResultFingerprintVersion makes the immutable Result fingerprint // domain explicit. Changing it is a breaking runtime protocol change. RuntimeResultFingerprintVersion = "openlinker.runtime-result.v1" )
Variables ¶
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") )
var ( ErrExternalExecutionLaunchFenceRejected = errors.New("external execution launch fence rejected") ErrWorkflowChildLaunchFenceRejected = errors.New("workflow child launch fence rejected") )
var ( ErrInvalidRuntimeInvocation = errors.New("invalid runtime invocation capability") ErrExpiredRuntimeInvocation = errors.New("runtime invocation capability expired") )
var ( ErrRuntimeReconcilerNotConfigured = errors.New("Runtime deadline reconciler is not configured") ErrRuntimeReconcileBatchInvalid = errors.New("Runtime deadline reconcile batch must be between 1 and 1000") )
var ( ErrRuntimeSignalBusClosed = errors.New("runtime signal bus is closed") ErrRuntimeSignalInvalid = errors.New("runtime signal is invalid") )
var ( // ErrInvalidRuntimeEvent is returned before touching PostgreSQL when the // request cannot satisfy the Runtime event contract. ErrInvalidRuntimeEvent = errors.New("invalid runtime event") )
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.
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.
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.
var ErrObservationForbidden = errors.New("browser observation is not permitted for this caller")
ErrObservationForbidden is returned when the caller may not observe this Run.
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.
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.
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.
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.
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.
var ErrRuntimeDispatchWakeReconcilerNotConfigured = errors.New("Runtime dispatch wake reconciler is not configured")
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
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
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
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
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 ¶
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)
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 (*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
func (observation *BrowserObservation) HandleEvent( ctx context.Context, event BrowserObserverEventPayload, ) (BrowserObserverEventAckPayload, error)
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 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 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 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 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 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 ¶
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
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
func (s *EventStore) Append( ctx context.Context, principal RuntimeEventPrincipal, identity RuntimeAttemptIdentity, request RuntimeEventRequest, ) (RuntimeEventAck, error)
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 Handler ¶
type Handler struct {
// contains filtered or unexported fields
}
Handler 调用执行 HTTP 入口。
func NewHandler ¶
NewHandler 构造 Handler。cfg 可选(测试可省略)。
func (*Handler) ActivateRuntimeNode ¶ added in v0.1.56
func (*Handler) CancelRun ¶ added in v0.1.41
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 (*Handler) DrainRuntimeNode ¶ added in v0.1.56
func (*Handler) GetAdminBrowserObservation ¶ added in v0.1.56
func (*Handler) GetAdminBrowserObservationFrame ¶ added in v0.1.56
func (*Handler) GetBrowserControl ¶ added in v0.1.56
func (*Handler) GetBrowserControlFrame ¶ added in v0.1.56
func (*Handler) GetBrowserObservation ¶ added in v0.1.56
func (*Handler) GetBrowserObservationFrame ¶ added in v0.1.56
func (*Handler) GetConversationRuns ¶ added in v0.1.56
func (*Handler) GetRunArtifacts ¶
GetRunArtifacts 查询 run 持久化产物。只返回给 run owner。
func (*Handler) GetRunEvents ¶
GetRunEvents 查询 run 事件流。SSE 接口后续会复用同一 service 方法。
func (*Handler) GetRunMessages ¶
GetRunMessages 查询 run 的稳定消息回放。只返回给 run owner。
func (*Handler) ListRuntimeDeadLetters ¶ added in v0.1.56
func (*Handler) ListRuntimeNodes ¶ added in v0.1.56
func (*Handler) PostRun ¶
PostRun 调用 Agent。
Endpoint 连接模式会同步等待 Agent 返回;其他运行模式由各自的调度路径处理。 失败 / 超时 / 取消 → status='failed' or 'timeout' or 'canceled',已退款。
func (*Handler) PostRunAsync ¶
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 ¶
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
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 (*Handler) ReplayRun ¶ added in v0.1.56
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 (*Handler) RevokeRuntimeNode ¶ added in v0.1.56
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 (*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 (*Handler) StartBrowserObservation ¶ added in v0.1.56
func (*Handler) StopAdminBrowserObservation ¶ added in v0.1.56
func (*Handler) StopBrowserObservation ¶ added in v0.1.56
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
func (a *MTLSRuntimeDeviceAuthenticator) AuthenticateHTTP( ctx context.Context, req *http.Request, ) (RuntimeDeviceIdentity, error)
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 (s *RedisRuntimePresenceStore) Refresh(ctx context.Context, presence RuntimePresence, ttl time.Duration) error
func (*RedisRuntimePresenceStore) Remove ¶ added in v0.1.56
func (s *RedisRuntimePresenceStore) Remove(ctx context.Context, presence RuntimePresence) error
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) ListExpired ¶ added in v0.1.56
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 (s *RedisRuntimeSessionLeaseStore) Remove( ctx context.Context, record RuntimeSessionLeaseRecord, ) error
func (*RedisRuntimeSessionLeaseStore) ScheduleCheck ¶ added in v0.1.56
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 (b *RedisSignalBus) Check( ctx context.Context, registrations []RuntimeConnectionRegistration, ) ([]RuntimeCredentialProjectionResult, error)
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 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
func (fn ResultClassifierFunc) ClassifyResult(input ResultClassificationInput) RuntimeResultClassification
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
func (f *ResultFinalizer) Finalize( ctx context.Context, principal RuntimeResultPrincipal, request RuntimeResultRequest, ) (RuntimeResultAck, error)
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
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
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 RunCancellationEvidence ¶ added in v0.1.56
type RunCancellationState ¶ added in v0.1.56
type RunCreationIdentity ¶ added in v0.1.56
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) SetHandlers ¶ added in v0.1.56
func (w *RunEffectWorker) SetHandlers( webhook WebhookRunEffectHandler, delivery DeliveryRunEffectHandler, )
type RunEffectWorkerConfig ¶ added in v0.1.56
type RunErrorPayload ¶ added in v0.1.56
type RunEventAckMessage ¶ added in v0.1.56
type RunEventAckMessage = RuntimeTypedEnvelope[RunEventAckPayload]
type RunEventAckPayload ¶ added in v0.1.56
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
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
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
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 RuntimeAdmissionIdentity ¶ added in v0.1.56
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 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 RuntimeCancellationAPI ¶ added in v0.1.56
type RuntimeCancellationAPI interface {
NextCommand(context.Context, RuntimeSessionPrincipal) (*PendingCommand, time.Time, error)
PollCommands(context.Context, RuntimeSessionPrincipal) (RuntimeCommandsResponse, error)
AckCancel(context.Context, RuntimeSessionPrincipal, RunCancelAckPayload) (RunCancellationState, error)
}
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
func (c *RuntimeCancellationCoordinator) AckCancel( ctx context.Context, principal RuntimeSessionPrincipal, request RunCancelAckPayload, ) (RunCancellationState, error)
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
func (c *RuntimeCancellationCoordinator) NextCommand( ctx context.Context, principal RuntimeSessionPrincipal, ) (*PendingCommand, time.Time, error)
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 (c *RuntimeCancellationCoordinator) PollCommands( ctx context.Context, principal RuntimeSessionPrincipal, ) (RuntimeCommandsResponse, error)
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 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
func (c *RuntimeClusterCoordinator) Close(ctx context.Context) error
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 (c *RuntimeClusterCoordinator) Readiness(ctx context.Context) RuntimeClusterReadiness
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 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 NewRuntimeCredentialReconciler ¶ added in v0.1.56
func NewRuntimeCredentialReconciler( hub *RuntimeWakeHub, projection RuntimeCredentialProjectionStore, validator RuntimeCredentialConnectionValidator, config RuntimeCredentialReconcilerConfig, ) *RuntimeCredentialReconciler
type RuntimeCredentialReconcilerConfig ¶ added in v0.1.56
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 (s *RuntimeDelegationService) CallAgent( ctx context.Context, authorization RuntimeDelegationAuthorization, ) (RunSummary, error)
func (*RuntimeDelegationService) ReadDelegatedRun ¶ added in v0.1.59
func (s *RuntimeDelegationService) ReadDelegatedRun(ctx context.Context, authorization RuntimeDelegationAuthorization) (DelegatedRunView, error)
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
func (r *RuntimeDispatchWakeReconciler) ReconcileOnce( ctx context.Context, limit int, ) (RuntimeDispatchWakeReconcileResult, error)
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 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" 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" 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" 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
func (h *RuntimeHTTPController) AuthenticateAgentRequest(c echo.Context) (AuthenticatedRuntimePrincipal, *RuntimeTransportError)
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 RuntimeLeaseAPI ¶ added in v0.1.56
type RuntimeLeaseAPI interface {
ClaimOffer(context.Context, RuntimeSessionPrincipal) (*RunAssignedPayload, error)
AckAssignment(context.Context, RuntimeSessionPrincipal, RunAssignmentAckPayload) (RunAssignmentConfirmedPayload, error)
RejectAssignment(context.Context, RuntimeSessionPrincipal, RunAssignmentRejectPayload) (RunAssignmentRejectedPayload, error)
RenewLease(context.Context, RuntimeSessionPrincipal, RunLeaseRenewPayload) (RunLeaseRenewedPayload, error)
ReleaseUnackedOffer(context.Context, RuntimeSessionPrincipal, ...string) error
}
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 (s *RuntimeLeaseService) AckAssignment( ctx context.Context, principal RuntimeSessionPrincipal, request RunAssignmentAckPayload, ) (RunAssignmentConfirmedPayload, error)
func (*RuntimeLeaseService) ClaimOffer ¶ added in v0.1.56
func (s *RuntimeLeaseService) ClaimOffer( ctx context.Context, principal RuntimeSessionPrincipal, ) (*RunAssignedPayload, error)
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 (s *RuntimeLeaseService) RejectAssignment( ctx context.Context, principal RuntimeSessionPrincipal, request RunAssignmentRejectPayload, ) (RunAssignmentRejectedPayload, error)
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
func (s *RuntimeLeaseService) RenewLease( ctx context.Context, principal RuntimeSessionPrincipal, request RunLeaseRenewPayload, ) (RunLeaseRenewedPayload, error)
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 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
func (s *RuntimeResumeService) Resume( ctx context.Context, target RuntimeSessionPrincipal, payload RuntimeResumePayload, ) (RuntimeResumeResponse, error)
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 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) Health ¶ added in v0.1.56
func (m *RuntimeSessionLeaseManager) Health() RuntimeSessionLeaseManagerHealth
func (*RuntimeSessionLeaseManager) HealthyFor ¶ added in v0.1.56
func (m *RuntimeSessionLeaseManager) HealthyFor(connectionID string) bool
func (*RuntimeSessionLeaseManager) ListExpired ¶ added in v0.1.56
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 (m *RuntimeSessionLeaseManager) Register(record RuntimeSessionLeaseRecord) error
func (*RuntimeSessionLeaseManager) Run ¶ added in v0.1.56
func (m *RuntimeSessionLeaseManager) Run(ctx context.Context) error
func (*RuntimeSessionLeaseManager) ScheduleCheck ¶ added in v0.1.56
func (*RuntimeSessionLeaseManager) Unregister ¶ added in v0.1.56
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 RuntimeSessionLeaseManagerHealth ¶ added in v0.1.56
type RuntimeSessionLeaseRecord ¶ added in v0.1.56
type RuntimeSessionLeaseRecord struct {
Lease RuntimeSessionLease
Presence RuntimePresence
}
type RuntimeSessionLeaseStore ¶ added in v0.1.56
type RuntimeSessionLeaseStore interface {
RefreshBatch(context.Context, []RuntimeSessionLeaseRecord, time.Duration, time.Duration) error
Lookup(context.Context, uuid.UUID) (RuntimeSessionLease, bool, error)
Remove(context.Context, RuntimeSessionLeaseRecord) error
ListExpired(context.Context, int) ([]uuid.UUID, error)
ScheduleCheck(context.Context, uuid.UUID, time.Duration) error
Forget(context.Context, uuid.UUID) error
}
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
func (p RuntimeSessionPrincipal) EventPrincipal() RuntimeEventPrincipal
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
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 (s *RuntimeSessionService) CloseSession( ctx context.Context, principal AuthenticatedRuntimePrincipal, request RuntimeSessionCloseRequest, ) (RuntimeSessionState, error)
func (*RuntimeSessionService) CreateOrAttachSession ¶ added in v0.1.56
func (s *RuntimeSessionService) CreateOrAttachSession( ctx context.Context, principal AuthenticatedRuntimePrincipal, request RuntimeSessionRequest, ) (RuntimeSessionState, error)
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
func (s *RuntimeSessionService) DrainSession( ctx context.Context, principal AuthenticatedRuntimePrincipal, request RuntimeSessionDrainRequest, ) (RuntimeDrainPayload, error)
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 (s *RuntimeSessionService) HeartbeatSession( ctx context.Context, principal AuthenticatedRuntimePrincipal, request RuntimeSessionHeartbeatRequest, ) (RuntimeSessionState, error)
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 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 (w *RuntimeSignalOutboxWorker) Backlog(ctx context.Context) (int32, error)
func (*RuntimeSignalOutboxWorker) ProcessOnce ¶ added in v0.1.56
func (w *RuntimeSignalOutboxWorker) ProcessOnce( ctx context.Context, cfg RuntimeSignalOutboxWorkerConfig, ) (RuntimeSignalOutboxBatchResult, error)
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 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 ¶
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) ConfigureCoreRuntime ¶ added in v0.1.56
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 (*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 (*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
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 (*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
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 ¶
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
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)
Source Files
¶
- admin_nodes.go
- admission_limiter.go
- agent_request_auth.go
- browser_execution_profile.go
- browser_human_control.go
- browser_observation.go
- browser_observation_frames.go
- browser_observation_payloads.go
- cancellation.go
- cancellation_evidence.go
- cancellation_retry.go
- capacity_signal.go
- cluster.go
- contract.go
- conversation_runs.go
- core_executor.go
- creation_idempotency.go
- credential_projection.go
- credential_reconciler.go
- credential_revocation.go
- dead_letter.go
- delegated_run.go
- delegation.go
- dispatch_wake_reconciler.go
- dry_run_runtime.go
- dto.go
- effect_worker.go
- event_page.go
- event_store.go
- external_launch.go
- finalizer.go
- handler.go
- idempotency.go
- invocation.go
- lease.go
- maintenance_worker.go
- presence.go
- presence_runtime.go
- principal.go
- protocol.go
- protocol_error.go
- protocol_validation.go
- reconciler.go
- replay.go
- replay_capability.go
- requirements.go
- resume.go
- resume_authorization.go
- run_input_schema.go
- run_update_hub.go
- runtime.go
- runtime_http.go
- safe_int.go
- service.go
- session.go
- session_lease.go
- session_reaper.go
- signal_bus.go
- signal_bus_local.go
- signal_bus_redis.go
- signal_outbox_worker.go
- transport_policy.go
- wake_hub.go
- websocket.go
- websocket_scheduler.go
- wire_compatibility.go
- worker_event_schedule.go
- worker_observer.go