runtime

package
v0.1.0 Latest Latest
Warning

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

Go to latest
Published: Oct 5, 2026 License: Apache-2.0 Imports: 104 Imported by: 0

Documentation

Overview

Bounded, attempt-local capture of the one proxy-owned control call for task 4.1 of agent-loop-explicit-completion-protocol (spec: .kiro/specs/agent-loop-explicit-completion-protocol, design Response Interception / Capture State; requirements 3.2, 3.5, 5.1-5.6, 8.5, 10.3).

This file owns capture only. It is not wired into the response pipeline, it never calls Provider.Handle, and it names no concrete feature: ownership comes from the trusted activation the request path already froze, so the only control truth it reads is that activation's immutable projection. Wording, provider keys, and schema stay the control provider's own policy, so nothing here interprets an argument member or a value.

The capture is single-owner, exactly like the ordinary tool-call assembler it will run ahead of: one backend Recv loop drives it and no other goroutine may touch it, so it owns no lock, no timer, and no background work.

Package runtime owns the top-level application assembly and request execution lifecycle wiring.

Private, attempt-local control-call interception for agent-loop-explicit-completion-protocol (spec: .kiro/specs/agent-loop-explicit-completion-protocol, design Response Interception / Placement, Handler Semantics, and Concurrency and Lifecycle; requirements 3.5, 4.7, 5.1-5.7, 8.2-8.7, 10.3, 11.1-11.3, 12.3-12.5).

This file owns the interception seam only. It runs on the backend-event boundary after BTP and provider usage observation and before the ordinary tool-call assembler, so a proxy-owned control call never reaches ordinary tool-call finalizers, tool policy, tool reactors, response hooks, the client recorder, or client release, while operator BTP capture and provider accounting keep seeing the legitimate upstream control traffic. Ordinary events are returned unchanged, undelayed, and unbuffered.

Ownership comes exclusively from the trusted activation the request path already froze. Nothing here re-reads a runtime snapshot, a provider ID() or Spec(), capability eligibility, request Extensions, or a mutable context value to decide response ownership.

The capture and its normalized outcome are attempt-local mutable state, owned by the attempt session under one short-held control lock (attemptSession.controlMu, documented at its declaration). The single backend Recv loop observes and hands off; attempt lifecycle cleanup — cancellation, Close, attempt loss, and replacement — disposes that same state from whichever goroutine owns the transition. The lock therefore covers capture observation, handoff bookkeeping, outcome storage, and disposal, and nothing else: no Provider.Handle call, no backend Recv/Cancel/Close, and no terminal effect ever runs while it is held, so cleanup can never wait on an in-flight handler.

Disposal is one-way. A released attempt keeps its immutable activation and a visibly released capture, so a late event fails closed instead of reopening the capture or falling through to ordinary client tool execution. This file adds no timer, no goroutine, no global map, and no public extension metadata.

Pending completion-result publication for agent-loop-explicit-completion-protocol (spec: .kiro/specs/agent-loop-explicit-completion-protocol, design Completion Evidence and Pending Result / Response Interception Placement / Concurrency and Lifecycle; requirements 6.3-6.7, 11.3, 11.5, 12.5).

This file owns one private, bounded terminal-publication value for a valid proxy-owned completion. It is not a gate engine, a terminal owner, a registry, or a side channel to a frontend: the accepted request-terminal owner calls it from inside the winning request effects, and every effect it performs runs through the existing response observation seam.

Ownership of each fact:

  • The bounded lexical result is a copy taken from the winning attempt's private control state under that attempt's control lock, before attempt cleanup disposes it. It is dropped by every losing, errored, continued, cancelled, or closed path.
  • The canonical events are ordinary legal response/message/text events. They run the existing response-part hooks exactly once each and participate in the existing completion-gate chain exactly once. No synthetic client tool call, tool result, or item is ever produced.
  • Eligibility uses trimmed released assistant text plus post-hook assistant text already inside the gate buffer. The general output-committed bit is deliberately not used, because reasoning, tool, and media output also commit while proving nothing about an assistant answer existing.
  • An unresolved, incomplete, ambiguous, or over-capacity ordinary tool boundary in the held candidate suppresses the automatic result even when the terminal provider allowed the stop. Ordinary canonical output is never suppressed with it.

The prepared value is not released while it is prepared. Every effect it needs runs exactly once at the accepted publication boundary, and the private release drain below exists only to deliver already-observed events to the client in canonical order.

Index

Constants

View Source
const (
	IntentSuccess              attemptTerminalIntent = "success"
	IntentSwallowedFailure     attemptTerminalIntent = "swallowed_failure"
	IntentSurfacedFailure      attemptTerminalIntent = "surfaced_failure"
	IntentCancellation         attemptTerminalIntent = "cancellation"
	IntentTimeout              attemptTerminalIntent = "timeout"
	IntentReplacement          attemptTerminalIntent = "replacement"
	IntentParallelLoser        attemptTerminalIntent = "parallel_loser"
	IntentOpenReadinessFailure attemptTerminalIntent = "open_readiness_failure"
	IntentPublicationDenied    attemptTerminalIntent = "publication_denied"
	IntentPreReturnAbort       attemptTerminalIntent = "pre_return_abort"
)
View Source
const (
	CancellationCauseExplicit    CancellationCauseClass = "explicit"
	CancellationCauseClientGone  CancellationCauseClass = "client_gone"
	CancellationCauseContextDone CancellationCauseClass = "context_done"
	CancellationCauseRaceLoser   CancellationCauseClass = "race_loser"
	CancellationCauseNone        CancellationCauseClass = "none"
	CancellationCauseOther       CancellationCauseClass = "other"

	CancellationModeNone      CancellationModeClass = "none"
	CancellationModeProvider  CancellationModeClass = "provider"
	CancellationModeTransport CancellationModeClass = "transport"
	CancellationModeCloseOnly CancellationModeClass = "close_only"
	CancellationModeOther     CancellationModeClass = "other"

	CancellationPhaseRequested CancellationPhase = "requested"
	CancellationPhaseOutcome   CancellationPhase = "outcome"
	CancellationPhaseForced    CancellationPhase = "forced"
	CancellationPhaseTerminal  CancellationPhase = "terminal"
	CancellationPhaseNone      CancellationPhase = "none"
	CancellationPhaseOther     CancellationPhase = "other"

	CancellationFallbackNegotiated CancellationFallback = "negotiated"
	CancellationFallbackLegacy     CancellationFallback = "legacy"
	CancellationFallbackNone       CancellationFallback = "none"
	CancellationFallbackOther      CancellationFallback = "other"
)

Variables

View Source
var (
	ErrBillingAdmissionDenied         = errors.New("executor: billing admission denied")
	ErrBillingCreditScreenDenied      = errors.New("executor: cheap credit screen denied")
	ErrBillingCreditScreenUnavailable = errors.New("executor: cheap credit screen unavailable")
)
View Source
var ErrNilConfig = errors.New("runtime: config is required")
View Source
var ErrNilLogger = errors.New("runtime: logger is required")

ErrNilLogger is returned by New when Options.Logger is nil.

Functions

func BuildAuthorityCoordinators

BuildAuthorityCoordinators wires thin adapters into request/attempt coordinators when a usage-authority service and/or concurrency provider is present.

func BuildClientTurnRecordInputFromShape

func BuildClientTurnRecordInputFromShape(
	now time.Time,
	traceID string,
	br app.BeginResult,
	shape largebody.ClientTurnShape,
	maxFactBytes int64,
) (app.ClientTurnRecordInput, error)

BuildClientTurnRecordInputFromShape builds a bounded ClientTurnRecordInput from a ClientTurnShape without prompt text materialization (Requirements 14.3, 14.5). It validates the shape under maxFactBytes and returns an error wrapping largebody.ErrSemanticFactBudgetExceeded if the semantic fact budget is exceeded (Requirement 14.4).

func HashOpaqueIDForLog

func HashOpaqueIDForLog(s string) string

HashOpaqueIDForLog returns a short stable hex digest of opaque strings (session ids, user ids) for log attributes. It does not replace access control; it limits raw identifier leakage in structured logs.

func StandardLaneDomainPolicies

func StandardLaneDomainPolicies() map[string]LaneDomainPolicy

StandardLaneDomainPolicies returns the default domain policies for standard frontend lanes.

Types

type AccountingRuntime

type AccountingRuntime struct {
	Preflight                    *accountingpreflight.Checker
	StreamUsage                  *accountingstream.Reconstructor
	TokenAccountingObservability *accountingobs.Stats
	AdminCountService            *accountingapp.Service
	UsageAuthority               UsageAuthorityService
	UsageAuthorityCleanupTimeout time.Duration
	// ConcurrencyProvider is the optional Phase 8 logical-request lease authority.
	ConcurrencyProvider authority.ConcurrencyProvider
	// ConcurrencyLeaseTTL / ConcurrencyRenewBefore are defaults used by heartbeat
	// when the admit decision omits rule-level values.
	ConcurrencyLeaseTTL    time.Duration
	ConcurrencyRenewBefore time.Duration
	// ConcurrencyAuxiliaryLeasePolicy is inherit (default) or acquire_own (10.10).
	ConcurrencyAuxiliaryLeasePolicy string
	// MeteringRecorder is the optional Phase 3 metering journal port. Nil means
	// checkpoints are retained in-request only (no durable append until Phase 5).
	MeteringRecorder metering.Recorder
	// MeteringObservationSink is the optional V2 evidence journal port. It is
	// deliberately separate from MeteringRecorder because pre-terminal economic
	// checkpoints persist immutable observations only; they never rate or mutate
	// money from a receive callback. Runtime checkpoint flushing requires the
	// sink to also implement metering.AtomicObservationSink; a plain
	// ObservationSink is retained for source compatibility but is not retried or
	// invoked by that path.
	MeteringObservationSink metering.ObservationSink
	// RequestCoordinator admits customer/logical-request authority once per request (Phase 6).
	RequestCoordinator *authoritycoord.RequestCoordinator
	// AttemptCoordinator admits operator/attempt authority per B-leg (Phase 6).
	AttemptCoordinator *authoritycoord.AttemptCoordinator
	// SnapshotGeneration is the atomic policy/rating generation publisher (Phase 9.3).
	// Admit binds Current() usage/concurrency/rating refs when present.
	SnapshotGeneration *snapshotgen.Publisher
	// TerminalWork accepts durable settle/release intents on post-output failures
	// (Phase 4.5; requirements 7.7, 8.3; design D9).
	TerminalWork *terminalworkapp.IntentService
}

AccountingRuntime carries token-accounting admission, usage reconstruction, and ledger hooks.

type App

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

App is the bootstrap composition root for the standard distribution.

func New

func New(opts Options) (*App, error)

New validates bootstrap wiring without starting the HTTP server (see cmd/lipstd and stdhttp.RunWithGenerationHost). It does not validate coreconfig.Config field semantics; load YAML and run config.Validate upstream.

func (*App) HookBus

func (a *App) HookBus() *hooks.Bus

HookBus returns the configured hook bus (never nil after New).

func (*App) Registrations

func (a *App) Registrations() []lipsdk.Registration

Registrations returns a copy of bootstrap plugin registrations (may be empty).

func (*App) Shutdown

func (a *App) Shutdown(ctx context.Context)

Shutdown stops plugin lifecycles in reverse registration order.

func (*App) Start

func (a *App) Start(ctx context.Context) error

Start logs hook chain lengths and starts plugin lifecycles. The bundled HTTP server is started by stdhttp.RunWithGenerationHost from cmd/lipstd.

type BillingAdmissionInput

type BillingAdmissionInput = BillingRoutePlanInput

type BillingCreditGate

type BillingCreditGate interface {
	Check(context.Context, string) error
}

type BillingExposureAdmission

type BillingExposureAdmission interface {
	Admit(context.Context, BillingExposureAdmissionInput) (billing.CallExposure, error)
}

type BillingExposureAdmissionInput

type BillingExposureAdmissionInput struct {
	BillingAdmissionInput
	CallID string
}

type BillingIdentity

type BillingIdentity struct {
	AccountID func(context.Context, lipapi.Call) string
	// StoreID resolves the trusted durable billing/metering store identity for a
	// request. It is captured with the terminal billing facts and is intentionally
	// absent from public runtime options.
	StoreID            func(context.Context) string
	CustomerPricingRef func(context.Context, lipapi.Call) billing.VersionRef
	ChargePolicyRef    func(context.Context, lipapi.Call) billing.VersionRef
	OperatorRateRef    func(context.Context, string, string) billing.VersionRef

	// WireBounded indicates whether this identity bundle is backed by stock/bounded
	// facts (such as PrincipalSessionIdentity) rather than custom Call-inspecting callbacks.
	// Custom callbacks without an explicit bounded contract remain wire blockers (Requirements 15.6, 19).
	WireBounded bool

	// WireAccountID optionally resolves an account identity directly from bounded scope
	// without constructing or inspecting a lipapi.Call.
	WireAccountID func(context.Context, scope.PrincipalScopeView) string

	// WireCustomerPricingRef optionally resolves the customer pricing snapshot ref from bounded facts.
	WireCustomerPricingRef func(context.Context) billing.VersionRef

	// WireChargePolicyRef optionally resolves the charge policy snapshot ref from bounded facts.
	WireChargePolicyRef func(context.Context) billing.VersionRef
}

BillingIdentity is the composition-root identity bundle shared by billing exposure admission and terminal call-closure stamping so those seams cannot drift.

func (BillingIdentity) HasCustomCallCallbacks

func (b BillingIdentity) HasCustomCallCallbacks() bool

HasCustomCallCallbacks reports whether BillingIdentity has Call-shaped callbacks that are NOT certified as wire-bounded (Requirements 15.6, 19).

type BillingLegObserver

type BillingLegObserver interface {
	ObserveBillingLeg(context.Context, billing.CallLegUsageRecord)
}

type BillingLegObserverFunc

type BillingLegObserverFunc func(context.Context, billing.CallLegUsageRecord)

func (BillingLegObserverFunc) ObserveBillingLeg

func (f BillingLegObserverFunc) ObserveBillingLeg(ctx context.Context, record billing.CallLegUsageRecord)

type BillingRoutePlanInput

type BillingRoutePlanInput struct {
	Call            lipapi.Call
	TraceID         string
	ALegID          string
	BillingCallID   string
	Route           *routing.Selector
	RequestSize     routing.RequestSizeEstimate
	Scope           scope.PrincipalScopeView
	SessionID       string
	MaxOutputTokens *int
	AccountID       string
}

type BillingRuntime

type BillingRuntime struct {
	// BillingCreditGate is the cheap settled-credit screen. Authoritative
	// billing requires it; when wired, runtime invokes it after identity
	// preparation and before route expansion.
	BillingCreditGate BillingCreditGate
	// BillingExposureAdmission is the authoritative operational-exposure seam.
	// When set, it replaces hold admission for this executor generation.
	BillingExposureAdmission BillingExposureAdmission
	// BillingLegObserver receives one recorded current call-leg usage record at terminal ownership.
	// It is observational only and never participates in authorization,
	// settlement, retry, cancellation, or client-visible output.
	BillingLegObserver BillingLegObserver
	TerminalUsageSink  billing.TerminalUsageSink
	// BillingIdentity is the composition identity bundle for exposure admission
	// and terminal call-closure stamping. It contains no hold/authorization identity.
	BillingIdentity BillingIdentity
}

BillingRuntime carries the runtime seams for two-stage exposure admission and durable terminal usage. Token quota and usage-authority stay on AccountingRuntime.

type CancellationCauseClass

type CancellationCauseClass string // explicit | client_gone | context_done | race_loser | none | other

type CancellationFallback

type CancellationFallback string // negotiated | legacy | none | other

type CancellationModeClass

type CancellationModeClass string // none | provider | transport | close_only | other

type CancellationObservation

type CancellationObservation struct {
	CauseClass CancellationCauseClass
	Mode       CancellationModeClass
	Phase      CancellationPhase
	Fallback   CancellationFallback
}

type CancellationPhase

type CancellationPhase string // requested | outcome | forced | terminal | none | other

type CatalogResolver

type CatalogResolver interface {
	Resolve(
		ctx context.Context,
		candidate routing.AttemptCandidate,
		call lipapi.Call,
		backend lipapi.BackendCaps,
	) modelcatalog.EffectiveFacts
}

CatalogResolver merges administrator overrides, catalog facts, and backend capability maps for one routing candidate. Defined here (consumer package) per ports-and-adapters guidance.

type CompactionDetector

CompactionDetector is the repository-internal consumer port for session compaction detection. It exposes only the four observation operations consumed by the core runtime, using canonical lipapi and compaction SDK types.

type CompactionRuntime

type CompactionRuntime struct {
	Detector CompactionDetector
	// BackgroundAux is the generation-bound process scheduler client. The
	// scheduler itself remains process-owned; this interface lets callbacks
	// submit work against the executor's frozen generation binding.
	BackgroundAux auxiliary.BackgroundClient
}

CompactionRuntime carries the process-owned compaction detector reference. The detector is shared across generations and never owned by the executor; nil is safe and disables compaction observation entirely. Detection is observational only: it never alters routing, prompts, responses, retries, accounting, or client framing.

type CompactionWireDetector

type CompactionWireDetector interface {
	CompactionDetector
	RequestOpenedFacts(compaction.PreservationMeta, compactionfacts.RequestFacts) []compaction.Event
}

CompactionWireDetector extends CompactionDetector with exact-facts request observation for wire streaming execution without materializing a full lipapi.Call.

type ConversationViewObserver

type ConversationViewObserver interface {
	OnProjection(stage string, summary conversationprojection.ProjectionSummary)
	OnProjectionFailure(stage string)
	OnAnchorFallback(stage string, policy conversationprojection.AnchorMissingPolicy)
	OnAnchorFailure(policy conversationprojection.AnchorMissingPolicy)
}

ConversationViewObserver is optional narrow diagnostics for bounded conversation-view projection/anchor/steering metrics. Nil is no-op. Labels are bounded enums only (placement, operation, policy, stage).

type ConversationViewTagger

type ConversationViewTagger interface {
	TagNeverBackend(ctx context.Context, aLegID string, tags []TagRequest) (TagResult, error)
}

ConversationViewTagger is the narrow tagger interface consumed by runtime.

type CoreRuntime

type CoreRuntime struct {
	Store                b2bua.Store
	Backends             map[string]execbackend.Backend
	ALegLifecycle        *leglifecycle.Coordinator
	Rand                 routing.Rng
	Now                  func() time.Time
	MaxPendingWireEvents int
	StreamRecovery       streamrecovery.Config
	// PromptCacheMaintenance is the optional generation-owned provider-neutral maintenance port.
	PromptCacheMaintenance PromptCacheMaintenance
	// ConversationViewReader is an optional narrow snapshot port. When set,
	// runtime preserves the single-snapshot per-turn invariant (task 3.2).
	ConversationViewReader conversationprojection.Reader
	// ConversationViewTagger is the optional narrow tagger port for local-turn
	// tag-before-release.
	ConversationViewTagger ConversationViewTagger
	// SteeringWriterFactory constructs an authoritative steering writer for a given A-leg.
	SteeringWriterFactory SteeringWriterFactory
	// LargeBodyAssessor evaluates frontend proof for candidate fast-path execution.
	LargeBodyAssessor largebody.LargeBodyAssessor
	// LargeBodyGenerationID is the optional generation identity for largebody execution validation.
	LargeBodyGenerationID string
	// LargeBodyCandidateDomainGeneration is the optional generation identity for late-route proof validation.
	LargeBodyCandidateDomainGeneration string
}

CoreRuntime carries continuity store, backends, lifecycle coordination, and clocks.

type EligibilityResolver

type EligibilityResolver interface {
	Check(
		ctx context.Context,
		candidate routing.AttemptCandidate,
		call lipapi.Call,
		facts modelcatalog.EffectiveFacts,
	) modelcatalog.EligibilityDecision
}

EligibilityResolver runs context-limit checks after successful capability negotiation.

type Executor

Executor orchestrates hooks, capability negotiation, routing, B2BUA, and backend attempts. Fields are grouped by concern; promoted fields preserve the historical flat access API.

func NewExecutor

func NewExecutor(cfg ExecutorConfig) *Executor

NewExecutor constructs an Executor from grouped runtime configuration.

func TestExecutor

func TestExecutor() *Executor

func (*Executor) AppendWireExposureAbort

func (e *Executor) AppendWireExposureAbort(ctx context.Context, args WireExposureAbortArgs) error

AppendWireExposureAbort records an exposure abort closure for a failed wire attempt sharing exact record sealing and sink handoff logic with canonical aborts (Requirements 15.6, 19).

func (*Executor) AssertWireAttemptNotWidened

func (e *Executor) AssertWireAttemptNotWidened(
	holder *checkpoint.RequestHolder,
	attemptID string,
	current checkpoint.WireAttemptEvidence,
) error

AssertWireAttemptNotWidened verifies that current bounded attempt evidence has not widened beyond the authorized backend-ingress freeze for attemptID (Requirements 10, 15.3, 19).

func (*Executor) AssessLargeBody

func (e *Executor) AssessLargeBody(ctx context.Context, proof largebody.Proof) (largebody.Assessment, error)

AssessLargeBody evaluates frontend proof for candidate fast-path execution (Phase 4). If LargeBodyAssessor is not configured or interleaved thinking is enabled, it returns a declined assessment (fail closed to canonical; Item 4).

func (*Executor) AssessWireBilling

AssessWireBilling evaluates billing eligibility under the held decode permit (Requirements 15.6, 19). It accepts whenever no custom BillingIdentity Call callbacks block the wire path, regardless of whether credit/exposure gates are configured (unconfigured gates accept the same way on the canonical path). If custom BillingIdentity Call callbacks exist without a wire-bounded contract, it declines with (AssessmentDecisionDecline, DeclineReasonAuthorityBlocker).

func (*Executor) AssessWirePreflight

AssessWirePreflight evaluates token preflight admission for a wire request under the held decode permit without materializing prompt text or constructing a lipapi.Call. If preflight is disabled or nil, it succeeds with AssessmentDecisionAccept. If preflight requires counting and only CountCall is supported, or exact tokenizer semantics do not exist, or counting is expensive/unbounded, it returns (AssessmentDecisionDecline, DeclineReasonCountingUnsupported, decision). If preflight token limits are exceeded, it returns (AssessmentDecisionDecline, DeclineReasonAuthorityBlocker, decision).

func (*Executor) AuthorizeWireBilling

func (e *Executor) AuthorizeWireBilling(ctx context.Context, args WireBillingExposureArgs) (billing.CallExposure, error)

AuthorizeWireBilling admits operational exposure from bounded wire facts post-quote without constructing or inspecting a lipapi.Call (Requirements 15.6, 19).

func (*Executor) CancelALeg

func (e *Executor) CancelALeg(ctx context.Context, req lipapi.ALegCancelRequest) error

CancelALeg explicitly cancels an active A-leg and all registered B-legs.

func (*Executor) CaptureWireBackendIngress

func (e *Executor) CaptureWireBackendIngress(
	ctx context.Context,
	holder *checkpoint.RequestHolder,
	args WireBackendIngressArgs,
) (checkpoint.Snapshot, error)

CaptureWireBackendIngress captures an immutable backend-attempt checkpoint from bounded wire facts, inheriting correlated session/request facts from holder.FrontendIngress when available (Requirements 10, 15.1–15.3, 19).

func (*Executor) CheckWireCheapCredit

func (e *Executor) CheckWireCheapCredit(ctx context.Context, args WireBillingCreditArgs) error

CheckWireCheapCredit performs the cheap settled-credit screen on bounded wire facts without materializing or inspecting a lipapi.Call (Requirements 15.6, 19).

func (*Executor) EnrichWireBackendIngressQuantities

func (e *Executor) EnrichWireBackendIngressQuantities(holder *checkpoint.RequestHolder, attemptID string, count largebody.WireCountResult)

EnrichWireBackendIngressQuantities merges exact measured wire token counting results into the stored BackendIngress snapshot for attemptID without mutating or retaining a Call (Requirements 15.4, 15.5).

func (*Executor) EnrichWireFrontendIngressQuantities

func (e *Executor) EnrichWireFrontendIngressQuantities(holder *checkpoint.RequestHolder, count largebody.WireCountResult)

EnrichWireFrontendIngressQuantities merges exact measured wire token counting results into the stored FrontendIngress snapshot without mutating or retaining a Call (Requirements 15.4, 15.5).

func (*Executor) Execute

func (e *Executor) Execute(ctx context.Context, call *lipapi.Call) (_ lipapi.EventStream, err error)

func (*Executor) ExecuteLargeBody

func (e *Executor) ExecuteLargeBody(
	ctx context.Context,
	accepted largebody.Assessment,
	src largebody.Source,
) (largebody.ExecutionResult, error)

ExecuteLargeBody implements largebody.LargeBodyWireExecutor (Requirements 6, 7, 14, 15, 18, 19; Task 13.1). It crosses the one-way wire commit barrier, runs the wire secure-session preparation and exactly one BeginTurn/A-leg lifecycle, reads the post-BeginTurn live route override constrained to the assessed domain, applies request authority and economic admission, and returns an ExecutionResult containing the canonical EventStream, bounded ResponseFacts, and the sensitive SessionResponseCarrier. Under no circumstance does this method fall back to canonical Execute.

func (*Executor) LargeBodyStaticDisposition

func (e *Executor) LargeBodyStaticDisposition(profileID string) (largebody.StaticWireDisposition, largebody.StaticWireReason)

LargeBodyStaticDisposition returns an O(1) static wire disposition for a profile (Phase 4). If LargeBodyAssessor implements LargeBodyStaticDispositionProvider, it delegates to it. Otherwise, it returns DefinitelyCanonical with StaticBlocker.

func (*Executor) PersistBackendIngressFact

func (e *Executor) PersistBackendIngressFact(ctx context.Context, holder *checkpoint.RequestHolder, attemptID string) (string, error)

PersistBackendIngressFact appends the operator BE-ingress journal fact for attemptID when a MeteringRecorder is configured and binds its FactID for rating/admission. If holder is nil, it falls back to the holder stored in ctx (Requirements 10, 15.1–15.3, 19).

func (*Executor) PersistFrontendIngressFact

func (e *Executor) PersistFrontendIngressFact(ctx context.Context, holder *checkpoint.RequestHolder) (string, error)

PersistFrontendIngressFact appends the customer FE-ingress journal fact when a MeteringRecorder is configured and binds its FactID for rating/admission. If holder is nil, it falls back to the holder stored in ctx (Requirements 15.1–15.3, 19).

func (*Executor) PrepareSecureSession

func (e *Executor) PrepareSecureSession(ctx context.Context, in SecureSessionPrepInput) (*PreparedSecureSession, error)

PrepareSecureSession prepares fact-based inputs for secure-session execution. It executes scope resolution, session openers, and workspace resolution, but strictly DOES NOT call BeginTurn or mutate session/store state (Requirements 6.2, 14.1, 19).

func (*Executor) QuarantinePersistenceFaulted

func (e *Executor) QuarantinePersistenceFaulted() bool

func (*Executor) SetToolCallFinalizers

func (e *Executor) SetToolCallFinalizers(fs []toolcall.Finalizer, maxArgsBytes int)

SetToolCallFinalizers installs completed-call finalizers used by the per-B-leg assembler. Composition roots normally merge these via the request runtime snapshot; harnesses may call this directly when they do not run feature-bundle merge.

func (*Executor) WallClock

func (e *Executor) WallClock() func() time.Time

type ExecutorConfig

type ExecutorConfig struct {
	Core          CoreRuntime
	Billing       BillingRuntime
	Routing       RoutingRuntime
	Security      SecurityRuntime
	Accounting    AccountingRuntime
	Observability ObservabilityRuntime
	Extension     ExtensionRuntime
	Interleaved   InterleavedRuntime
	Compaction    CompactionRuntime
}

ExecutorConfig groups executor dependencies for explicit construction at the composition root and in tests. Use NewExecutor to obtain a runnable executor.

type ExtensionRuntime

type ExtensionRuntime struct {
	Bus             *hooks.Bus
	RuntimeSnapshot *extensions.RequestRuntimeSnapshot

	// TerminalPolicyReader resolves session-scoped terminal decision policy overrides
	// at request admission (Task 7.2).
	TerminalPolicyReader TerminalPolicyReader

	// ToolCallFinalizationMaxArgsBytes is the assembler buffer cap from merged
	// feature bundles (0 means default at assembler construction).
	ToolCallFinalizationMaxArgsBytes int
	// contains filtered or unexported fields
}

ExtensionRuntime carries the hook bus and frozen per-build extension snapshot.

type InterleavedMemo

type InterleavedMemo struct {
	Text       string
	Reference  string
	Version    int64
	HadContent bool
}

InterleavedMemo carries minimal evidence of captured thinker output.

func (InterleavedMemo) IsEmpty

func (m InterleavedMemo) IsEmpty() bool

IsEmpty reports whether the memo evidence is empty.

type InterleavedProcessor

type InterleavedProcessor interface {
	BeginTurn(ctx context.Context, in InterleavedTurnInput) (InterleavedTurn, error)
	IsMemoVisibleToClient(ctx context.Context, aLegID string) bool
	// MemoSteeringPutRequest builds the feature-owned steering mutation that
	// persists a captured memo.
	MemoSteeringPutRequest(memo string) steering.PutRequest
	// MemoSteeringOverlayID returns the feature-owned stable overlay identity
	// for the thinker memo.
	MemoSteeringOverlayID() steering.OverlayID
	// IsMemoSteeringOverlay reports whether overlayID carries the thinker memo.
	IsMemoSteeringOverlay(overlayID string) bool
}

InterleavedProcessor is the runtime-owned consumer interface for interleaved thinking. Memo steering policy (rendering, overlay identity, placement/fallback selection, memo filtering) is feature-owned: the processor answers every memo-policy question so core never hardcodes feature semantics.

type InterleavedRuntime

type InterleavedRuntime struct {
	Processor InterleavedProcessor
}

InterleavedRuntime carries the interleaved-thinking consumer processor port.

type InterleavedTurn

type InterleavedTurn interface {
	ShapeThinker(call lipapi.Call) (lipapi.Call, error)
	ObserveThinkerEvent(ev lipapi.Event) ([]lipapi.Event, error)
	FinalizeThinker(ctx context.Context) (InterleavedMemo, error)
	FinalizeThinkerStatus(ctx context.Context, interrupted bool, visibleCommitted bool) (InterleavedMemo, error)
	ShapeExecutor(ctx context.Context, call lipapi.Call, memo InterleavedMemo) (lipapi.Call, error)
	Visible() bool
	CanContinue() bool
	ShapeDiagnostics() (outcome string, turnsRemaining int)
	CommitExecutor(ctx context.Context) (remaining int, err error)
	FlushVisible() []lipapi.Event
}

InterleavedTurn is the runtime-owned per-turn lifecycle contract.

type InterleavedTurnInput

type InterleavedTurnInput struct {
	ALegID              string
	Selector            string
	Backend             string
	Model               string
	RequestID           string
	StreamToClient      string
	SuppressVisibleMemo bool
}

InterleavedTurnInput carries per-turn facts required by interleaved thinking.

type LaneDomainPolicy

type LaneDomainPolicy struct {
	// UniversalOnly requires backend universal model certification (AnyAcceptedModel: true)
	// and evaluates universal late-route domain proof (e.g. Lane 3 OpenResponses).
	UniversalOnly bool
	// CandidateModels lists finite candidate models certified for the lane when not UniversalOnly.
	CandidateModels []string
}

LaneDomainPolicy configures late-route domain constraints per frontend lane.

type MetricsSink

type MetricsSink interface {
	OnAttemptRecorded(outcome lipapi.AttemptOutcome, backend string)
	OnBackendOpenDuration(backend string, seconds float64)
	OnTransportNegotiation(operation lipapi.Operation, mode lipapi.TransportMode, outcome string)
	OnCancellation(obs CancellationObservation)
}

MetricsSink receives coarse executor-level observations (Prometheus, etc.) when non-nil.

type ObservabilityRuntime

type ObservabilityRuntime struct {
	Log                        *slog.Logger
	Metrics                    MetricsSink
	ExtensionMetrics           extensions.StageMetrics
	SecretGuardDecisionMetrics extensions.SecretGuardDecisionMetrics
	RouteTrace                 *diag.RouteTraceBuffer
	PolicyDiagnosticsEnabled   bool
	CompletionBufferLimits     completion.BufferLimits
	// ConversationViewObserver is optional narrow diagnostics for bounded conversation-view
	// projection/anchor/steering metrics. Nil is no-op. Labels are bounded enums only (placement, operation, policy, stage).
	ConversationViewObserver ConversationViewObserver
}

ObservabilityRuntime carries structured logging, metrics, and diagnostics toggles.

type Options

type Options struct {
	Config *coreconfig.Config
	Logger *slog.Logger

	// Registrations enumerates configured plugins at bootstrap. When non-empty or Mandatory
	// is set, duplicates and mandatory coverage are validated in New.
	Registrations []lipsdk.Registration
	Mandatory     []lipsdk.Requirement

	// Hooks configures submit, part, and tool-reactor chains (zero value means empty chains).
	Hooks hooks.Config

	// Lifecycles are started after validation and stopped on shutdown (reverse order).
	Lifecycles []lipplugin.Lifecycle
}

Options carries bootstrap-only runtime dependencies for New.

We use an explicit struct rather than functional options: the surface is small, fields are easy to read at call sites, and the project prefers explicit construction over indirection (see repository AGENTS.md). Add fields here as needed instead of a generic options closure chain.

Config must be non-nil (otherwise ErrNilConfig). Logger must be non-nil (otherwise ErrNilLogger). Nil entries in Lifecycles are ignored by App.Start and App.Shutdown.

type OrchestrationScenarioSpec

type OrchestrationScenarioSpec struct {
	ID               string
	InvariantSummary string
	TestName         string // existing *_test.go function name in package runtime_test
}

OrchestrationScenarioSpec links a stable scenario identifier to steering-level invariants and the primary regression test that exercises it (specification bundle).

func SpecBundleOrchestrationScenarios

func SpecBundleOrchestrationScenarios() []OrchestrationScenarioSpec

SpecBundleOrchestrationScenarios lists core-owned orchestration invariants. Keep aligned with .kiro/steering/routing-and-orchestration.md and the referenced tests.

type PreparedSecureSession

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

PreparedSecureSession holds the pre-turn resolution state before BeginTurn commit. Preparation performs scope/principal resolution, context diagnostic binding, session openers, and workspace resolution, but strictly DOES NOT invoke BeginTurn or mutate session/store state (Requirements 6.2, 14.1, 19).

func (*PreparedSecureSession) BeginInput

func (p *PreparedSecureSession) BeginInput() app.BeginInput

func (*PreparedSecureSession) BindSession

func (*PreparedSecureSession) BuildClientTurnRecordInput

func (p *PreparedSecureSession) BuildClientTurnRecordInput(
	br app.BeginResult,
	shape largebody.ClientTurnShape,
	maxFactBytes int64,
) (app.ClientTurnRecordInput, error)

BuildClientTurnRecordInput builds an app.ClientTurnRecordInput from this prepared session and a ClientTurnShape under maxFactBytes without materializing prompt text (Requirements 14.3, 14.5).

func (*PreparedSecureSession) CaptureFrontendIngressCheckpoint

func (p *PreparedSecureSession) CaptureFrontendIngressCheckpoint(
	ctx context.Context,
	requestID string,
	br app.BeginResult,
	aLeg b2bua.ALegRecord,
	maxOutputTokens *int,
) (context.Context, *checkpoint.RequestHolder, error)

CaptureFrontendIngressCheckpoint captures an immutable FE-ingress checkpoint from bounded wire facts and post-BeginTurn session/a-leg correlation, sharing exact canonical checkpoint helpers without cloning or retaining a lipapi.Call (Requirements 15.1–15.3, 16.1–16.6, 19).

func (*PreparedSecureSession) Context

func (p *PreparedSecureSession) Context() context.Context

func (*PreparedSecureSession) EnrichWireFrontendIngressQuantities

func (p *PreparedSecureSession) EnrichWireFrontendIngressQuantities(holder *checkpoint.RequestHolder, count largebody.WireCountResult)

EnrichWireFrontendIngressQuantities merges exact measured wire token counting results into the stored FrontendIngress snapshot for this prepared session (Requirements 15.4, 15.5).

func (*PreparedSecureSession) ExecuteBeginTurn

func (p *PreparedSecureSession) ExecuteBeginTurn(ctx context.Context) (app.BeginResult, error)

func (*PreparedSecureSession) HasPrincipal

func (p *PreparedSecureSession) HasPrincipal() bool

func (*PreparedSecureSession) PersistFrontendIngressFact

func (p *PreparedSecureSession) PersistFrontendIngressFact(ctx context.Context, holder *checkpoint.RequestHolder) (string, error)

PersistFrontendIngressFact appends the customer FE-ingress journal fact when a MeteringRecorder is configured and binds its FactID (Requirements 15.1–15.3, 19).

func (*PreparedSecureSession) PreSession

func (p *PreparedSecureSession) PreSession() session.SessionView

func (*PreparedSecureSession) Principal

func (*PreparedSecureSession) RecordClientTurnWithShape

func (p *PreparedSecureSession) RecordClientTurnWithShape(
	ctx context.Context,
	br app.BeginResult,
	shape largebody.ClientTurnShape,
	maxFactBytes int64,
) error

RecordClientTurnWithShape records accepted client turn facts directly from a bounded ClientTurnShape without prompt text materialization (Requirements 14.3, 14.5). Semantic-fact budget overflow returns an error wrapping largebody.ErrSemanticFactBudgetExceeded so the caller can trigger pre-commit canonical fallback (Requirement 14.4).

func (*PreparedSecureSession) ResolveALeg

func (p *PreparedSecureSession) ResolveALeg(ctx context.Context, alegID string) (b2bua.ALegRecord, routeAuthoritySnapshot, error)

func (*PreparedSecureSession) ResponseCarrier

func (*PreparedSecureSession) Scope

func (*PreparedSecureSession) SecureTurn

func (*PreparedSecureSession) SessionInput

func (p *PreparedSecureSession) SessionInput() largebody.SessionInput

func (*PreparedSecureSession) TraceID

func (p *PreparedSecureSession) TraceID() string

func (*PreparedSecureSession) Workspace

type ProductionLargeBodyAssessor

type ProductionLargeBodyAssessor struct {
	GenerationID              string
	CandidateDomainGeneration string
	AuthorityGate             *largebody.AuthorityAssessmentGate
	WireProofGate             *largebody.BackendWireProofGate
	LaneDomainPolicies        map[string]LaneDomainPolicy
}

ProductionLargeBodyAssessor unifies authority assessment, backend wire proof, and route override/late selector gates into a side-effect-free proof assessor for runtime.Executor (Phase 4, Requirements 5, 6, 7, 8, 9, 13, 14, 15, 19).

func NewProductionLargeBodyAssessor

func NewProductionLargeBodyAssessor(
	generationID string,
	candidateDomainGen string,
	authGate *largebody.AuthorityAssessmentGate,
	wireProofGate *largebody.BackendWireProofGate,
	lanePolicies map[string]LaneDomainPolicy,
) *ProductionLargeBodyAssessor

NewProductionLargeBodyAssessor constructs a ProductionLargeBodyAssessor.

func (*ProductionLargeBodyAssessor) AssessLargeBody

AssessLargeBody evaluates frontend proof across authority and wire proof gates. Invariants:

  • Streaming-only delivery gate (proof.Delivery == DeliveryModeStreaming; 15.4 carry).
  • Universal-vs-finite domain policy per lane.
  • Same-permit decline discipline: always returns (declined, nil) on decline, never an error, so caller continues canonical processing under held decode permit.

func (*ProductionLargeBodyAssessor) LargeBodyStaticDisposition

LargeBodyStaticDisposition returns an O(1) static wire disposition for a profile.

type PromptCacheCommittedTurn

type PromptCacheCommittedTurn struct {
	ALegID              string
	BLegID              string
	CommittedSuccessful bool
	ToolEvents          []lipapi.ToolEvent
	Observations        []promptcache.Observation
	BackendInstanceID   string
	CanonicalModelID    string
	Controller          promptcache.Controller
}

PromptCacheCommittedTurn carries attempt-local prompt-cache observations and committed tool evidence using canonical/SDK DTOs only (design §6, Requirement 6.3).

type PromptCacheMaintenance

type PromptCacheMaintenance interface {
	BeginRealTurn(aLegID string)
	EndSession(aLegID string)
	ArmCommittedTurn(turn PromptCacheCommittedTurn)
}

PromptCacheMaintenance is the narrow core runtime consumer port for lifecycle facts that only core authoritatively emits (Requirement 6.3, design §6).

type RequestTokenEstimator

type RequestTokenEstimator interface {
	EstimateRequestTokens(ctx context.Context, call lipapi.Call) modelcatalog.SizeEstimate
}

RequestTokenEstimator supplies a provider-neutral request-size estimate for routing constraints. Unavailable estimates fail open in the routing planner.

type RoutingRuntime

type RoutingRuntime struct {
	MaxAttempts             int
	DefaultBackend          string
	SelectorAliases         *routing.AliasResolver
	CapsResolver            capabilities.Resolver
	CatalogResolver         CatalogResolver
	EligibilityResolver     EligibilityResolver
	RequestTokenEstimator   RequestTokenEstimator
	CandidateHealth         policy.CandidateHealth
	RouteObserver           lipsdk.RouteObserver
	AffinityStore           affinity.Store
	AffinityMissingIdentity affinity.MissingIdentityPolicy
	TransportFallbackPolicy lipapi.TransportFallbackPolicy
	RouteOverrideReader     routeoverride.Reader

	ExecutionCompositionPolicy config.ExecutionCompositionPolicy
	BackendExecutionResolver   routing.BackendExecutionResolver
}

RoutingRuntime carries selector parsing, planning, negotiation, and affinity policy.

type SecureSessionMetrics

type SecureSessionMetrics interface {
	ObserveBeginTurnNew()
	ObserveBeginTurnResume()
	// ObserveBeginTurnDenied increments when BeginTurn fails (code is lipapi SessionDenialCode string or "unknown").
	ObserveBeginTurnDenied(code string)
	ObserveStorageUnavailable()
	ObserveActivityTouch(seconds float64)
	ObserveRecorderClientTurnFailed(mandatory bool)
	ObserveRecorderStreamEventFailed(committed bool, mandatory bool)
}

SecureSessionMetrics receives secure-session observability signals (optional; nil skips all).

type SecureSessionPrepInput

type SecureSessionPrepInput struct {
	TraceID       string
	Session       largebody.SessionInput
	ContinuityKey string
}

SecureSessionPrepInput carries fact-based inputs required to prepare a secure-session turn for either canonical or wire execution without requiring a canonical lipapi.Call.

type SecureSessionRecorder

type SecureSessionRecorder = app.GateRecording

SecureSessionRecorder is an alias for the secure-session gate recording port implemented by app.Recorder.

type SecurityRuntime

type SecurityRuntime struct {
	SecureSession                           *app.Manager
	SyntheticLocalPrincipal                 bool
	SecureSessionRecorder                   app.GateRecording
	SecureSessionRecordingMandatory         bool
	SessionDenialMapper                     func(error) error
	SecureSessionMetrics                    SecureSessionMetrics
	SecureSessionRequireWorkspaceID         bool
	SecureSessionWorkspaceResolveFailClosed bool
	AuthEvents                              *auth.EventDispatcher
	SessionAuditPolicy                      auth.SessionAuditPolicy
}

SecurityRuntime carries secure-session gates, auth events, and session audit policy.

type SteeringWriterFactory

type SteeringWriterFactory func(ctx context.Context, aLegID string, resolver SteeringWriterResolver) (steering.Writer, error)

SteeringWriterFactory constructs an authoritative steering writer for a given A-leg.

type SteeringWriterResolver

type SteeringWriterResolver func(ctx context.Context) (lipapi.Call, conversationprojection.Snapshot, error)

SteeringWriterResolver resolves the current call and snapshot for anchor calculation.

type TagRequest

type TagRequest struct {
	Identity conversationprojection.MessageIdentity `json:"identity"`
	Reason   conversationprojection.ReasonCode      `json:"reason"`
}

TagRequest is one element of a TagNeverBackend batch.

func (TagRequest) Validate

func (r TagRequest) Validate() error

type TagResult

type TagResult struct {
	StateRevision uint64                       `json:"state_revision"`
	Tags          []conversationprojection.Tag `json:"tags"`
}

TagResult is returned from a successful TagNeverBackend call.

type TerminalPolicyQuery

type TerminalPolicyQuery struct {
	SecureSessionIncarnation string
	ALegID                   string
	FeatureID                string
	GenerationDefault        bool
}

TerminalPolicyQuery provides the scoped session identity and generation default needed to evaluate the effective terminal decision policy for an incoming turn.

type TerminalPolicyReader

type TerminalPolicyReader interface {
	Effective(ctx context.Context, in TerminalPolicyQuery) (TerminalPolicySnapshot, error)
}

TerminalPolicyReader resolves session-scoped terminal decision policy overrides at request admission. The interface is consumer-owned: core execution never mutates policy and never observes actor-specific keys or store internals.

type TerminalPolicySnapshot

type TerminalPolicySnapshot struct {
	EffectiveEnabled bool
	Revision         uint64
}

TerminalPolicySnapshot carries only the immutable effective decision and revision resolved at request admission.

type UsageAuthorityService

UsageAuthorityService is the runtime-owned boundary for accounting authority admission, settlement, release, advisory usage application, and bounded query access.

type WireBackendIngressArgs

type WireBackendIngressArgs struct {
	RequestID       string
	TraceID         string
	AttemptID       string
	BLegID          string
	ALegID          string
	SessionID       string
	Scope           scope.PrincipalScopeView
	BackendID       string
	Model           string
	MaxOutputTokens *int
	Now             time.Time

	SourceDigest  [32]byte
	RewriteDigest [32]byte
	AttemptDigest [32]byte
}

WireBackendIngressArgs carries bounded facts required to freeze an immutable backend-attempt checkpoint on the wire fast-path (Requirements 10, 15.1–15.3, 19).

type WireBillingAssessmentArgs

type WireBillingAssessmentArgs struct {
	Scope scope.PrincipalScopeView
}

WireBillingAssessmentArgs carries facts needed to assess billing compatibility under held permit (Req 15.6, 19).

type WireBillingCreditArgs

type WireBillingCreditArgs struct {
	Scope     scope.PrincipalScopeView
	AccountID string // optional explicit account ID override
}

WireBillingCreditArgs carries bounded facts for the cheap pre-route credit screen (Req 15.6).

type WireBillingExposureArgs

type WireBillingExposureArgs struct {
	BillingCallID   billing.BillingCallID
	TraceID         string
	ALegID          string
	SessionID       string
	Scope           scope.PrincipalScopeView
	AccountID       string
	Route           *routing.Selector
	RequestSize     routing.RequestSizeEstimate
	MaxOutputTokens *int
}

WireBillingExposureArgs carries bounded facts for post-quote atomic exposure admission (Req 15.6, 19).

type WireExposureAbortArgs

type WireExposureAbortArgs struct {
	BillingCallID   billing.BillingCallID
	SubmissionID    string
	Exposure        billing.CallExposure
	ALegID          string
	SessionID       string
	ExpectedBLegIDs []string
	Now             time.Time
}

WireExposureAbortArgs carries bounded facts for recording an exposure abort closure (Req 15.6, 19).

type WireFrontendIngressArgs

type WireFrontendIngressArgs struct {
	RequestID       string
	TraceID         string
	Scope           scope.PrincipalScopeView
	ALegID          string
	SessionID       string
	MaxOutputTokens *int
	Now             time.Time
}

WireFrontendIngressArgs carries the bounded facts required to capture an immutable frontend-ingress checkpoint on the wire fast-path (Requirements 15.1–15.3, 16.1–16.6, 19).

type WirePreflightAssessmentArgs

type WirePreflightAssessmentArgs struct {
	Backend                  string
	Model                    string
	CallID                   string
	ProfileID                string
	Source                   largebody.Source
	RequestedMaxOutputTokens *int
	Facts                    modelcatalog.ModelFacts
	Semantics                largebody.ExactTokenizerSemantics
	MaxScanBytes             int64
	MaxPermitHoldCPU         time.Duration
}

WirePreflightAssessmentArgs carries the inputs needed to evaluate token preflight during dynamic assessment under the held decode permit (Requirements 6.1–6.3, 15.5, 21; Task 10.4).

Source Files

Jump to

Keyboard shortcuts

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