reasoningpreservation

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

Documentation

Overview

Package reasoningpreservation owns the official reasoning-output-preservation feature domain.

Index

Constants

View Source
const (
	HardCompressionTimeout                = 30 * time.Second
	HardCompressionMaxInputTokens         = 200000
	HardCompressionMaxInputBytes          = 4 * 1024 * 1024
	HardCompressionMaxOutputTokens        = 32000
	HardCompressionMaxOutputBytes         = HardRawOutputCeiling
	HardCompressionMaxSurrogateBytes      = 256 * 1024
	HardCompressionMaxSurrogateOverhead   = 1024
	HardCompressionMaxPendingPerSession   = 64
	HardCompressionMaxSurrogatePerSession = 4 * 1024 * 1024
	HardCompressionMaxPendingTotal        = 4096
	HardCompressionMaxSurrogateTotal      = 64 * 1024 * 1024
	HardCompressionMinSourceBytesCeiling  = 4 * 1024 * 1024
	HardCompressionMinSavedBytesCeiling   = 4 * 1024 * 1024
	CompressionSurrogateOverhead          = 1024
)
View Source
const (
	SanitizationNone     = "none"
	SanitizationRedacted = "redacted"
)
View Source
const (
	ActionObserve = "observe"
	ActionRestore = "restore"

	PolicyLogSkip = "log_skip"
	PolicyReject  = "reject"
)
View Source
const (
	EgressPurposeReasoningSemanticCompression = "reasoning_semantic_compression"
	EgressSourceClassSemanticText             = "semantic_reasoning_text"
)
View Source
const BuiltinCatalogVersion = reasoningreplay.CatalogVersion
View Source
const CompressorOutputSchema = `{"schema_version":1,"segments":[{"index":0,"text":"..."}]}`

CompressorOutputSchema is the strict versioned output contract shown to the compressor model. It contains only local segment indexes and textual surrogates.

View Source
const CompressorSystemPrompt = "" /* 371-byte string literal not displayed */

CompressorSystemPrompt is the fixed model-visible instruction for the reasoning semantic compressor. It is versioned and never includes control-plane lineage or principal data.

View Source
const HardRawOutputCeiling = 512 * 1024 // 512 KiB

HardRawOutputCeiling is the hard implementation byte limit distinct from configured max_output_bytes. It is a defense-in-depth cap applied before JSON decode regardless of operator config.

View Source
const ID = "reasoning-output-preservation"

Variables

View Source
var (
	ErrCompressionNotFound       = errors.New(ID + ": compression not found")
	ErrCompressionConflict       = errors.New(ID + ": compression conflict")
	ErrCompressionBudgetExceeded = errors.New(ID + ": compression budget exceeded")
)
View Source
var (
	ErrPreparationIneligible          = errors.New(ID + ": ineligible")
	ErrPreparationDenied              = errors.New(ID + ": egress denied")
	ErrPreparationMissingPolicy       = errors.New(ID + ": egress denied missing-policy")
	ErrPreparationSanitizerFailed     = errors.New(ID + ": sanitizer failed")
	ErrPreparationInputBytesExceeded  = errors.New(ID + ": input exceeds max_input_bytes")
	ErrPreparationInputTokensExceeded = errors.New(ID + ": input exceeds max_input_tokens")
)
View Source
var (
	ErrRawOversize       = errors.New(ID + ": raw_oversize")
	ErrRawInvalidChannel = errors.New(ID + ": raw_invalid_channel")
	ErrRawInvalidLimit   = errors.New(ID + ": raw_invalid_limit")
)
View Source
var (
	ErrSurrogateDecodeInvalid       = errors.New(ID + ": decode_invalid")
	ErrSurrogateSchemaInvalid       = errors.New(ID + ": schema_invalid")
	ErrSurrogateControlInvalid      = errors.New(ID + ": control_invalid")
	ErrSurrogateOversize            = errors.New(ID + ": surrogate_oversize")
	ErrSurrogateInsufficientSavings = errors.New(ID + ": insufficient_savings")
)
View Source
var ErrNotImplemented = errors.New("reasoningpreservation: not implemented")

Functions

func BuildCompressorAuxRequest

func BuildCompressorAuxRequest(p CompressorAuxRequestParams) (auxiliary.Request, error)

BuildCompressorAuxRequest builds one detached no-tools auxiliary request per artifact. Control-plane metadata stays in the envelope; model-visible Call.Messages contain only the fixed instruction, versioned schema, and sanitized segment JSON wrapped as untrusted quoted data. One call may carry multiple segments. It validates route, segment presence/indexes/text, output-token bound, and canonical Call invariants, and propagates json.Marshal errors instead of swallowing them.

func BuildEphemeralCandidates

func BuildEphemeralCandidates(cands []restoreCandidate, decisions map[string]ReasoningViewResult, getSurrogate func(string) (*ReasoningSurrogate, bool)) []restoreCandidate

BuildEphemeralCandidates builds defensive ephemeral candidates for restore. For each candidate where decisions is ViewSurrogate, the surrogate is fetched and the artifact's semantic Text is replaced defensively; otherwise original is cloned. Fallback (duplicate/missing/index mismatch) returns original whole artifact.

func BuildFeatureBundle

func BuildFeatureBundle(n yaml.Node) (lipfeature.FeatureBundle, error)

BuildFeatureBundle decodes YAML and returns a FeatureBundle for registry factories.

func ComputeAnchor

func ComputeAnchor(msg lipapi.Message) ([32]byte, error)

func ComputeEgressPolicyHash

func ComputeEgressPolicyHash(decision CompressionEgressDecision, route string) [32]byte

ComputeEgressPolicyHash derives a stable versioned authoritative hash from the trusted egress decision PolicyVersion + explicit route + purpose + source class. It is deterministic and provider-neutral. The separator 0x00 avoids collisions across field boundaries. Version prefix "v1" ensures stable versioning if scheme changes.

func DecodeSurrogate

DecodeSurrogate implements the strict versioned JSON decoder for raw bounded bytes. It enforces DisallowUnknownFields, exactly schema_version=1, segments expected indexes exactly once (reject duplicate/missing/unexpected), text non-empty/non-whitespace, valid UTF-8, reject disallowed controls, raw bytes already bounded but decoded aggregate and each retained surrogate enforce MaxSurrogateBytes/hard ceiling, reject trailing JSON/tokens, and builds a ReasoningSurrogate correlated with original/policy/sanitization using caller params. Savings are validated (source bytes, decoded bytes, strictly smaller, MinSavedBytes, MinSavingsRatio) with overflow-safe math. Typed content-free outcome/error taxonomy is returned. This function never claims semantic equivalence; active replay is an operator-approved lossy mode.

func EnsureCompanionRules

func EnsureCompanionRules(n yaml.Node, backendIDs []string, rulePrefix string) (yaml.Node, error)

EnsureCompanionRules appends missing backend-only rules while retaining every existing config node and its ordering. DecodeConfig performs validation before mutation, so malformed or unknown feature config is never silently repaired by composition.

func EstimateInputTokens

func EstimateInputTokens(text string) int

EstimateInputTokens is the deterministic bounded token estimator for sanitized semantic text. It reuses the repository's conservative byte-as-token pattern (extractor's "each UTF-8 byte as one token-equivalent unit") and does not add a tokenizer dependency. One byte = one token-equivalent: an upper bound on real token count, deterministic, bounded, and provider-neutral.

func EstimatedBytesForSegments

func EstimatedBytesForSegments(segments []CompressorInputSegment) int

EstimatedBytesForSegments sums byte lengths without mutation.

func EstimatedTokensForSegments

func EstimatedTokensForSegments(segments []CompressorInputSegment) int

EstimatedTokensForSegments sums EstimateInputTokens over segments without mutation.

func ExtractBoundedRaw

func ExtractBoundedRaw(collected lipapi.Collected, maxOutputBytes int) ([]byte, error)

ExtractBoundedRaw extracts raw bytes from a collected canonical response with strict bounding. It rejects terminal errors, missing finish, tool calls / non-text channels first, then enforces max_output_bytes (clamped to hard ceiling) without invoking JSON decode on oversize and without constructing a string beyond the bound. Scheduler MaxResultBytes is outer defense-in-depth only; this function enforces the feature's stricter max_output_bytes before decode per requirements 3,10 and design §9.

func FeatureBundle

func FeatureBundle(cfg Config) (lipfeature.FeatureBundle, error)

FeatureBundle builds the schema-V1 contribution for an enabled feature instance.

func FeatureBundleWithCompanionPolicy

func FeatureBundleWithCompanionPolicy(cfg Config, policy CompanionPolicy) (lipfeature.FeatureBundle, error)

func FormatSafeDiagnostic

func FormatSafeDiagnostic(outcome SafeOutcome, ruleID string, counts map[string]int) (string, error)

func IsBudgetError

func IsBudgetError(err error) bool

IsBudgetError reports whether err is a budget rejection.

func IsConflictError

func IsConflictError(err error) bool

IsConflictError reports whether err is a CAS/stale conflict.

func IsNotFoundError

func IsNotFoundError(err error) bool

IsNotFoundError reports whether err is a not-found.

func MatchEligible

func MatchEligible(kind MatchKind) bool

func NewCompanionConfig

func NewCompanionConfig(backendIDs []string, rulePrefix string) (yaml.Node, error)

NewCompanionConfig returns generic feature config for backend instance IDs. Provider-specific eligibility and trust issuance remain composition policy.

func PrepareCompressorInputWithLimits

func PrepareCompressorInputWithLimits(ctx context.Context, segments []CompressorInputSegment, decision CompressionEgressDecision, maxInputBytes, maxInputTokens int) ([]CompressorInputSegment, PreparationOutcome, error)

PrepareCompressorInputWithLimits is the bounded sanitization+budget step that enforces redaction BEFORE byte/token accounting. It is used by PrepareSemanticSegments and is also exposed for direct segment callers. Typed outcomes distinguish bytes vs tokens.

func PrepareSemanticSegments

func PrepareSemanticSegments(ctx context.Context, placements []PlacedReasoning, decision CompressionEgressDecision, maxInputBytes, maxInputTokens int) ([]CompressorInputSegment, PreparationOutcome, error)

PrepareSemanticSegments is the bounded semantic-segment preparation entry point. It extracts only ReplaySemanticText placements, applies egress-required sanitization BEFORE byte/token accounting, retains local placement index only, supports multiple eligible placements, uses the deterministic bounded token estimator, and returns typed outcomes for ineligible, denied, sanitizer failure, bytes/tokens exceeded. Inputs are preserved (defensive copies); no session/account/lineage/anchor/digest is emitted.

func PrepareSemanticSegmentsFromArtifact

func PrepareSemanticSegmentsFromArtifact(ctx context.Context, artifact TurnArtifact, decision CompressionEgressDecision, maxInputBytes, maxInputTokens int) ([]CompressorInputSegment, PreparationOutcome, error)

PrepareSemanticSegmentsFromArtifact is a convenience wrapper over PrepareSemanticSegments.

func ProjectSafeError

func ProjectSafeError(err error) (string, error)

Types

type AdoptionOutcome

type AdoptionOutcome string

AdoptionOutcome is a typed content-free outcome for the raw guard stage (5.2).

const (
	AdoptionOutcomeNone              AdoptionOutcome = "none"
	AdoptionOutcomeBoundedRaw        AdoptionOutcome = "bounded_raw"
	AdoptionOutcomeRawOversize       AdoptionOutcome = "raw_oversize"
	AdoptionOutcomeRawInvalidChannel AdoptionOutcome = "raw_invalid_channel"
	AdoptionOutcomeRawInvalidLimit   AdoptionOutcome = "raw_invalid_limit"
)

type AdoptionResult

type AdoptionResult struct {
	Outcome      AdoptionOutcome
	Candidate    *CompletedPollCandidate
	BoundedRaw   []byte
	RawByteCount int
	Err          error
}

AdoptionResult carries bounded raw bytes for the next stage (5.3) or a rejection outcome. It is typed, content-free, and holds only byte counts, not raw reasoning text. BoundedRaw is non-nil only for AdoptionOutcomeBoundedRaw.

type AttemptTransform

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

func NewAttemptTransform

func NewAttemptTransform(cfg Config, store TurnStore, tel ...*Telemetry) *AttemptTransform

func NewAttemptTransformWithCompanionPolicyServicesAndStage

func NewAttemptTransformWithCompanionPolicyServicesAndStage(cfg Config, store TurnStore, svc CompressionServices, policy CompanionPolicy, stage CompletedAdoptionStage, tel ...*Telemetry) *AttemptTransform

NewAttemptTransformWithCompanionPolicyServicesAndStage is the full composition for bundle wiring.

func (*AttemptTransform) FailureMode

func (t *AttemptTransform) FailureMode() sdkhooks.FailureMode

func (*AttemptTransform) HandleAttempt

func (*AttemptTransform) ID

func (t *AttemptTransform) ID() string

func (*AttemptTransform) Order

func (t *AttemptTransform) Order() int

type BudgetError

type BudgetError struct {
	Kind    BudgetKind
	Limit   int
	Current int
}

BudgetError is returned when a per-session or aggregate budget rejects admission.

func (*BudgetError) Error

func (e *BudgetError) Error() string

func (*BudgetError) Unwrap

func (e *BudgetError) Unwrap() error

type BudgetKind

type BudgetKind string

BudgetKind distinguishes which limit was exceeded.

const (
	BudgetPendingPerSession   BudgetKind = "pending_per_session"
	BudgetPendingTotal        BudgetKind = "pending_total"
	BudgetSurrogatePerTurn    BudgetKind = "surrogate_per_turn"
	BudgetSurrogatePerSession BudgetKind = "surrogate_per_session"
	BudgetSurrogateTotal      BudgetKind = "surrogate_total"
)

type BuiltinCatalogEntry

type BuiltinCatalogEntry struct {
	ID              string
	BackendPrefixes []string
}

func BuiltinCatalogEntries

func BuiltinCatalogEntries() ([]BuiltinCatalogEntry, error)

type CandidateIdentity

type CandidateIdentity struct {
	BackendID       string
	BackendPrefixes []string
	Model           string
}

type Classification

type Classification string
const (
	ClassMissing     Classification = "missing"
	ClassPreserved   Classification = "preserved"
	ClassConflicting Classification = "conflicting"
	ClassAmbiguous   Classification = "ambiguous"
	ClassUnmatched   Classification = "unmatched"
)

type ClassifiedTurn

type ClassifiedTurn struct {
	Classification Classification
	ArtifactID     string
}

func ClassifyAssistantTurns

func ClassifyAssistantTurns(messages []lipapi.Message, artifacts []TurnArtifact) ([]ClassifiedTurn, error)

type CompanionPolicy

type CompanionPolicy struct {
	BeforeMatch  func(*lipapi.Call, request.AttemptMeta)
	AfterRestore func(context.Context, *lipapi.Call, request.AttemptMeta, MatchResult, RestoreResult)
}

CompanionPolicy is an internal composition seam for a provider companion. The generic feature owns state matching/restoration; composition owns any provider-specific trust protocol.

type CompletedAdoptionStage

type CompletedAdoptionStage func(context.Context, AdoptionResult) AdoptionResult

CompletedAdoptionStage is a composable hook for the adoption chain (5.2 -> 5.3 -> ...). It receives the raw-guard AdoptionResult and returns the (possibly) transformed result. The transform holds this stage immutably; default is identity. Task 5.3 will inject the decoder/attach stage via constructor. No cross-attempt state is stored.

func NewDecoderAdoptionStage

func NewDecoderAdoptionStage(cfg Config, store CompressionStore, svc CompressionServices, tel *Telemetry) CompletedAdoptionStage

NewDecoderAdoptionStage returns a CompletedAdoptionStage that validates/correlation-checks, decodes, and CAS-attaches a surrogate under aggregate byte budgets. It derives ExpectedIndexes and SourceBytes from the authoritative artifact's semantic segments, enforces MaxSurrogateBytes/min savings, verifies reservation/job/artifact/ original/semantic/egress/policy, decodes typed outcomes, always Forget exactly once after terminal handling, atomically moves pending->surrogate on success, and clears expected pending (CAS) on stale/budget/decode/insufficient while preserving original. Shadow mode always keeps original; no active selection.

type CompletedPollCandidate

type CompletedPollCandidate struct {
	Partition     SessionPartition
	ArtifactID    string
	ReservationID string
	JobID         auxiliary.JobID
	Collected     lipapi.Collected
	PollState     auxiliary.PollState
}

CompletedPollCandidate is a typed completed adoption candidate for the next stage (5.2). In 5.1 the candidate is produced but original is still replayed (shadow). It carries the defensive Collected copy plus correlation IDs.

type CompressionClaim

type CompressionClaim struct {
	Partition      SessionPartition
	ArtifactID     string
	ReservationID  string
	OriginalDigest [32]byte
	PolicyRevision string
}

CompressionClaim identifies one in-flight compression reservation with value semantics. Passed as an untrusted input by value to CompressionStore mutation operations; the store fully validates all fields against its internal authoritative records.

type CompressionConfig

type CompressionConfig struct {
	Enabled                     bool            `yaml:"enabled"`
	Mode                        CompressionMode `yaml:"mode"`
	Route                       string          `yaml:"route"`
	Timeout                     time.Duration   `yaml:"timeout"`
	MaxInputTokens              int             `yaml:"max_input_tokens"`
	MaxInputBytes               int             `yaml:"max_input_bytes"`
	MaxOutputTokens             int             `yaml:"max_output_tokens"`
	MaxOutputBytes              int             `yaml:"max_output_bytes"`
	MaxSurrogateBytes           int             `yaml:"max_surrogate_bytes"`
	MinSourceBytes              int             `yaml:"min_source_bytes"`
	MinSavedBytes               int             `yaml:"min_saved_bytes"`
	MinSavingsRatio             float64         `yaml:"min_savings_ratio"`
	MaxPendingPerSession        int             `yaml:"max_pending_per_session"`
	MaxSurrogateBytesPerSession int             `yaml:"max_surrogate_bytes_per_session"`
	MaxPendingTotal             int             `yaml:"max_pending_total"`
	MaxSurrogateBytesTotal      int             `yaml:"max_surrogate_bytes_total"`
	EgressPolicyRef             string          `yaml:"egress_policy_ref"`
}

func (CompressionConfig) ToLimits

func (c CompressionConfig) ToLimits() CompressionLimits

func (CompressionConfig) Validate

func (c CompressionConfig) Validate() error

type CompressionEgressDecision

type CompressionEgressDecision struct {
	Action        EgressAction
	PolicyVersion string
	Sanitizer     TrustedTextSanitizer
}

CompressionEgressDecision is the trusted policy decision.

func EvaluateEgress

EvaluateEgress enforces the hard rules: nil/missing policy => missing-policy deny-equivalent; redact without sanitizer => deny; route string alone never approves (policy always consulted).

type CompressionEgressInput

type CompressionEgressInput struct {
	Route       string
	Purpose     string
	SourceClass string
	Principal   EgressPrincipalView
}

CompressionEgressInput is the narrow egress policy input.

type CompressionLimits

type CompressionLimits struct {
	MaxPendingPerSession        int
	MaxPendingTotal             int
	MaxSurrogateBytesPerTurn    int
	MaxSurrogateBytesPerSession int
	MaxSurrogateBytesTotal      int
}

CompressionLimits configures bounded optional-state budgets. Zero means unlimited for backward compatibility with disabled compression.

type CompressionMeasurements

type CompressionMeasurements struct {
	Counts       map[SafeOutcome]int64
	RawBytes     map[SafeOutcome]int64
	DecodedBytes map[SafeOutcome]int64
	SavedBytes   map[SafeOutcome]int64
}

CompressionMeasurements is a content-free snapshot of per-outcome byte measurements.

type CompressionMode

type CompressionMode string
const (
	CompressionShadow CompressionMode = "shadow"
	CompressionActive CompressionMode = "active"
)

type CompressionPollAttemptResult

type CompressionPollAttemptResult struct {
	Kind      PollForAttemptKind
	Candidate *CompletedPollCandidate
	State     auxiliary.PollState
	Err       error
}

CompressionPollAttemptResult is the typed result of the one-shot poll.

func PollOnceForMatchingArtifact

func PollOnceForMatchingArtifact(ctx context.Context, call *lipapi.Call, cs CompressionStore, partition SessionPartition, artifacts []TurnArtifact, support lipapi.ReasoningReplaySupport, svc CompressionServices) CompressionPollAttemptResult

PollOnceForMatchingArtifact performs a single non-blocking Poll for the first restorable artifact (ClassMissing and destination-supported) that has a pending bound JobID. It shares ClassifyAssistantTurns / dialectSet semantics with RestoreMissingReasoning via collectRestoreCandidates. Never calls Await or busy-waits; operational errors are compression-local and must NOT trigger authoritative on_state_error reject.

type CompressionServices

type CompressionServices struct {
	// Client is the generation-bound BackgroundClient (Scheduler.BindRunner result).
	Client auxiliary.BackgroundClient
	// Poller is the optional non-blocking poll capability. When compression is
	// enabled Poller must be non-nil; when disabled it must be nil (no extra goroutines).
	Poller auxiliary.BackgroundPoller
	// EgressPolicy is the trusted egress decision authority for purpose
	// reasoning_semantic_compression. When compression is enabled it must be non-nil;
	// route string alone never approves (EvaluateEgress enforces).
	EgressPolicy EgressPolicy
	// Sanitizer is the trusted redaction authority reused for redact-then-allow.
	// When compression is enabled it must be non-nil so redact decisions can succeed;
	// deny remains possible via policy, but missing sanitizer would make redaction fail closed.
	Sanitizer TrustedTextSanitizer
}

CompressionServices bundles generation-local capabilities required when compression is enabled. Capabilities are carried explicitly per construction; no global service locator, no DI container, no init() globals (Decision 6).

Single-value pattern: a process-owned scheduler implements both auxiliary.BackgroundClient and auxiliary.BackgroundPoller. Caller passes the same value for both Client and Poller fields. Two-field representation is chosen for explicitness and to keep disabled mode zero-value clean without type assertions inside the feature.

type CompressionState

type CompressionState struct {
	PolicyRevision string
	ReservationID  string
	Pending        *PendingCompression
	Surrogate      *ReasoningSurrogate
}

CompressionState is additive optional state per artifact revision.

type CompressionStats

type CompressionStats struct {
	TotalPending             int
	TotalSurrogateBytes      int
	PendingPerSession        map[string]int
	SurrogateBytesPerSession map[string]int
}

CompressionStats exposes aggregate counters for tests.

type CompressionStore

type CompressionStore interface {
	TurnStore
	ReserveCompression(ctx context.Context, partition SessionPartition, artifactID string, originalDigest [32]byte, policyRevision string, semanticDigest [32]byte, egressPolicyHash [32]byte) (CompressionClaim, error)
	UpdateReservationPolicyHash(ctx context.Context, claim CompressionClaim, expectedOldHash [32]byte, semanticDigest [32]byte, newHash [32]byte, sanitization string, routeHash [32]byte) error
	BindCompressionJob(ctx context.Context, claim CompressionClaim, jobID auxiliary.JobID) error
	AttachSurrogate(ctx context.Context, claim CompressionClaim, jobID auxiliary.JobID, surrogate ReasoningSurrogate) error
	ClearCompression(ctx context.Context, partition SessionPartition, artifactID string, expectedReservationID string) error
	GetCompressionState(ctx context.Context, partition SessionPartition, artifactID string) (CompressionState, bool, error)
	CompressionStats() CompressionStats
}

CompressionStore extends TurnStore with optional-state operations.

type CompressorAuxRequestParams

type CompressorAuxRequestParams struct {
	Route               string
	ParentTraceID       string
	ParentALegID        string
	ParentBLegID        string
	ParentBranchBinding string
	Segments            []CompressorInputSegment
	MaxOutputTokens     int
}

CompressorAuxRequestParams groups control-plane lineage plus sanitized segments.

type CompressorInputSegment

type CompressorInputSegment struct {
	Index int
	Text  string
}

CompressorInputSegment is a local index + text for one semantic placement. It carries only the local placement index and sanitized text; no session/account/lineage/anchor/digest is ever stored here.

func ExtractSemanticSegments

func ExtractSemanticSegments(placements []PlacedReasoning) []CompressorInputSegment

ExtractSemanticSegments returns ONLY placements classified as ReplaySemanticText. It excludes ordinary answer/transcript/tools/media/signatures/opaque/native. Retains the local placement index (input order) only; never copies session/account/lineage/anchor/digest. Result is a new slice; inputs are not mutated.

func ExtractSemanticSegmentsFromArtifact

func ExtractSemanticSegmentsFromArtifact(artifact TurnArtifact) []CompressorInputSegment

ExtractSemanticSegmentsFromArtifact extracts eligible segments from an artifact's reasoning. Never includes artifact ID/anchor/backend/model/lineage.

func PrepareCompressorInput

func PrepareCompressorInput(ctx context.Context, segments []CompressorInputSegment, decision CompressionEgressDecision, maxInputBytes int) ([]CompressorInputSegment, string, error)

PrepareCompressorInput applies egress decision => sanitization if required => then budgets. It enforces that redaction occurs before input-size accounting. Returns sanitized segments, outcome, error. On deny, returns nil segments, outcome "denied" or "missing-policy", and an error. Compatibility wrapper: delegates to PrepareCompressorInputWithLimits with no token bound.

type Config

type Config struct {
	Action            string            `yaml:"action"`
	UseBuiltinCatalog bool              `yaml:"use_builtin_catalog"`
	Rules             []RuleConfig      `yaml:"rules"`
	OnAmbiguous       string            `yaml:"on_ambiguous"`
	OnUnrepresentable string            `yaml:"on_unrepresentable"`
	OnStateError      string            `yaml:"on_state_error"`
	State             StateConfig       `yaml:"state"`
	Compression       CompressionConfig `yaml:"compression"`
}

func DecodeConfig

func DecodeConfig(n yaml.Node) (Config, error)

type EgressAction

type EgressAction uint8

EgressAction is the bounded egress decision.

const (
	EgressDeny EgressAction = iota
	EgressAllow
	EgressRedactThenAllow
)

func (EgressAction) String

func (a EgressAction) String() string

type EgressPolicy

type EgressPolicy interface {
	Decide(ctx context.Context, in CompressionEgressInput) (CompressionEgressDecision, error)
}

EgressPolicy is the trusted data-egress decision seam.

func NewRouteBoundEgressPolicy

func NewRouteBoundEgressPolicy(allowed map[string]struct{}, delegate EgressPolicy) EgressPolicy

NewRouteBoundEgressPolicy returns an EgressPolicy that denies when in.Route is not in allowed, otherwise delegates. An empty allowed map means no route restriction (delegate decides). A nil delegate on allowed match is treated as missing-policy deny.

type EgressPrincipalView

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

EgressPrincipalView is an opaque minimal principal-scope view not interpreted by this package. It holds the full scope.PrincipalScopeView for policy enforcement (tenant/org/workspace/project/cost-center/policy-labels) but remains opaque to preservation internals. Call Scope() for a defensive clone.

func NewEgressPrincipalScopeView

func NewEgressPrincipalScopeView(v scope.PrincipalScopeView) EgressPrincipalView

NewEgressPrincipalScopeView constructs a view from the full trusted scope. The scope is defensively cloned on construction.

func NewEgressPrincipalView

func NewEgressPrincipalView(opaque string) EgressPrincipalView

NewEgressPrincipalView constructs an opaque view from a legacy principal ID string. Compatibility is preserved by constructing a scope with Known PrincipalID.

func (EgressPrincipalView) PrincipalID

func (v EgressPrincipalView) PrincipalID() string

PrincipalID returns the string form of the principal ID (bounded accessor).

func (EgressPrincipalView) Scope

Scope returns a defensive clone of the underlying principal scope.

type EvictionSummary

type EvictionSummary struct {
	EvictedTurns int
	EvictedBytes int
	ExpiredTurns int
	ExpiredBytes int
}

type InstanceParts

type InstanceParts struct {
	Config              Config
	Store               TurnStore
	Telemetry           *Telemetry
	Transform           *AttemptTransform
	Observer            *StreamObserverFactory
	CompressionServices CompressionServices
}

InstanceParts exposes test/diagnostics handles for one enabled feature instance.

func FeatureBundleWithCompression

func FeatureBundleWithCompression(cfg Config, svc CompressionServices) (*InstanceParts, lipfeature.FeatureBundle, error)

FeatureBundleWithCompression is a convenience wrapper without companion policy.

func FeatureBundleWithParts

func FeatureBundleWithParts(cfg Config) (*InstanceParts, lipfeature.FeatureBundle, error)

FeatureBundleWithParts returns the shared store/telemetry participants plus the schema-V1 bundle. Disabled configurations must not call this constructor (D12).

func FeatureBundleWithPartsAndCompression

func FeatureBundleWithPartsAndCompression(cfg Config, svc CompressionServices, policy CompanionPolicy) (*InstanceParts, lipfeature.FeatureBundle, error)

FeatureBundleWithPartsAndCompression extends the feature composition with explicit generation-local compression capabilities. It wires validated CompressionConfig.ToLimits() into StoreOptions.CompressionLimits and validates that enabled compression has all required capabilities, while disabled mode requires none (zero delta).

func FeatureBundleWithPartsAndPolicy

func FeatureBundleWithPartsAndPolicy(cfg Config, policy CompanionPolicy) (*InstanceParts, lipfeature.FeatureBundle, error)

func (*InstanceParts) Inventory

func (p *InstanceParts) Inventory() SafeInventory

type MatchKind

type MatchKind string
const (
	MatchNone                    MatchKind = "none"
	MatchExplicitDisabledModel   MatchKind = "explicit_disabled_model"
	MatchExplicitEnabledModel    MatchKind = "explicit_enabled_model"
	MatchExplicitDisabledBackend MatchKind = "explicit_disabled_backend"
	MatchExplicitEnabledBackend  MatchKind = "explicit_enabled_backend"
	MatchBuiltin                 MatchKind = "builtin"
)

type MatchResult

type MatchResult struct {
	Kind   MatchKind
	RuleID string
}

func ResolveMatch

func ResolveMatch(cfg Config, cand CandidateIdentity) (MatchResult, error)

type PendingCompression

type PendingCompression struct {
	JobID                   auxiliary.JobID
	OriginalDigest          [32]byte
	SemanticDigest          [32]byte
	EgressPolicyHash        [32]byte
	AuthorizedRouteHash     [32]byte
	Sanitization            string
	PolicyHashAuthoritative bool
	CreatedAt               time.Time
	PolicyRevision          string
	ReservationID           string
}

PendingCompression tracks a reservation awaiting background result. SemanticDigest is the SHA-256 of the source semantic text; EgressPolicyHash identifies the egress policy revision that authorized the pending work. Both are recorded at Reserve and immutable via Bind; Attach verifies they match. Zero SemanticDigest is invalid and Attach rejects it as a conflict (real artifacts have content). Zero EgressPolicyHash is allowed but CAS-checked — a mismatch still yields a conflict. Sanitization is the content-free class (none/redacted) derived from the egress decision. AuthorizedRouteHash is sha256(route) at promotion, verified by adoption. PolicyHashAuthoritative is false at Reserve (provisional ref hash) and becomes true after successful UpdateReservationPolicyHash CAS promotion. Bind and Attach reject unless authoritative.

type PlacedReasoning

type PlacedReasoning struct {
	BeforeNonReasoningPart int
	Part                   lipapi.Part
}

func DerivePlacements

func DerivePlacements(nonReasoningCount int, reasoning []lipapi.Part) ([]PlacedReasoning, error)

DerivePlacements places every reasoning block at index 0 when only a flat reasoning slice is available (equal indexes preserve block order). Prefer DerivePlacementsFromParts when the interleaved assistant parts are known.

func DerivePlacementsFromParts

func DerivePlacementsFromParts(parts []lipapi.Part) ([]PlacedReasoning, int, error)

type PollForAttemptKind

type PollForAttemptKind string

PollForAttemptKind is a typed outcome for the one-shot poll attempt.

const (
	PollKindNoPending   PollForAttemptKind = "no_pending"
	PollKindUnavailable PollForAttemptKind = "unavailable"
	PollKindPending     PollForAttemptKind = "pending"
	PollKindFailed      PollForAttemptKind = "failed"
	PollKindNotFound    PollForAttemptKind = "not_found"
	PollKindCompleted   PollForAttemptKind = "completed"
	PollKindPollError   PollForAttemptKind = "poll_error"
)

type PostAppendCorrelation

type PostAppendCorrelation struct {
	Partition           SessionPartition
	ArtifactID          string
	Anchor              [32]byte
	OriginalDigest      [32]byte
	SemanticDigest      [32]byte
	EgressPolicyRefHash [32]byte
	SourceBytes         int
	TraceID             string
	ALegID              string
	BLegID              string
	BranchBinding       string
	Scope               scope.PrincipalScopeView
	PolicyRevision      string
}

PostAppendCorrelation carries trusted post-append attribution for compression. It contains only envelope/context data—no model payload or content telemetry. EgressPolicyRefHash is a provisional hash of the configured egress_policy_ref; task 4.3 will derive the authoritative PendingCompression.EgressPolicyHash from the trusted egress decision PolicyVersion via UpdateReservationPolicyHash before Bind. Do NOT treat EgressPolicyRefHash as authoritative — it is a stable provisional placeholder derived from config, not a policy decision hash.

type PostAppendHook

type PostAppendHook func(context.Context, PostAppendCorrelation) error

PostAppendHook is a local post-commit optimization seam invoked unlocked after authoritative Append success. It may reserve optional state but must not synchronously await provider; task 4.4 will submit after reservation. Failure must not invalidate original.

func BuildPostAppendHook

func BuildPostAppendHook(cfg Config, store TurnStore, svc CompressionServices) PostAppendHook

BuildPostAppendHook constructs the composable post-append hook chain rooted in compression. Chain is reserve -> egress -> submit (4.4). Hook is nil when compression disabled to preserve disabled-mode byte equivalency.

func BuildPostAppendHookWithEgressNext

func BuildPostAppendHookWithEgressNext(cfg Config, store TurnStore, svc CompressionServices, next PostEgressStage) PostAppendHook

BuildPostAppendHookWithEgressNext allows tests and 4.4 to inject the post-egress next stage that receives PreparedReservation (sanitized segments + decision).

func BuildPostAppendHookWithNext

func BuildPostAppendHookWithNext(cfg Config, store TurnStore, svc CompressionServices, next PostReservationStage) PostAppendHook

BuildPostAppendHookWithNext allows tests and future stages to inject the next stage for chaining. This is the reservation-level injection (pre-egress) retained for 4.2-era tests.

func BuildPostAppendHookWithTelemetry

func BuildPostAppendHookWithTelemetry(cfg Config, store TurnStore, svc CompressionServices, tel *Telemetry) PostAppendHook

BuildPostAppendHookWithTelemetry is the telemetry-aware chain builder.

func NewCompressionReservationHook

func NewCompressionReservationHook(cfg Config, store CompressionStore, next PostReservationStage) PostAppendHook

NewCompressionReservationHook returns a concrete LOCAL post-append hook that reserves optional capacity before any provider submission. It is invoked unlocked after original append success. It returns nil always (fail-open) so original is never invalidated. If reservation succeeds (OutcomeReserved) and next is non-nil, it invokes next with the ReservationResult; non-reserved outcomes never call next; next errors are fail-open.

func NewCompressionReservationHookWithTelemetry

func NewCompressionReservationHookWithTelemetry(cfg Config, store CompressionStore, next PostReservationStage, tel *Telemetry) PostAppendHook

NewCompressionReservationHookWithTelemetry is the telemetry-aware reservation hook. It records content-free outcomes: below_threshold, ineligible, budget_exceeded (per-session/total) and reservation reserved, without emitting reasoning text or IDs.

type PostEgressStage

type PostEgressStage func(context.Context, PreparedReservation) error

PostEgressStage is the next composable stage after egress+redaction. It receives the sanitized segments and decision metadata. No request build/provider is performed in this stage (task 4.3 boundary).

func NewPostEgressSubmitStage

func NewPostEgressSubmitStage(cfg Config, store CompressionStore, svc CompressionServices) PostEgressStage

NewPostEgressSubmitStage returns a PostEgressStage that builds the detached no-tools auxiliary compressor request from the prepared sanitized segments, submits it through the generation-bound BackgroundClient with timeout and versioned coalesce key, and binds the returned JobID via CAS. It never calls Await or Poll. On request-build or submit failure it clears only the expected reservation (CAS) and leaves the original intact. On accepted JobID it attempts BindCompressionJob; bind failure triggers Forget(jobID) and clears only the expected reservation. Empty JobID is treated as invalid and clears. Incurred accepted work remains accounting-owned (Forget does not cancel billing).

func NewPostEgressSubmitStageWithTelemetry

func NewPostEgressSubmitStageWithTelemetry(cfg Config, store CompressionStore, svc CompressionServices, tel *Telemetry) PostEgressStage

NewPostEgressSubmitStageWithTelemetry is the telemetry-aware submit stage. It records content-free queue/submit outcomes: submitted, coalesced, queue_saturated, submit_failed. Admission denial is asynchronous via PollFailed, not a synchronous SubmitCollect error today; OutcomeAdmissionDenied is reserved for future synchronous credit screening but currently never emitted here. No reasoning text or IDs are emitted.

type PostReservationStage

type PostReservationStage func(ctx context.Context, res ReservationResult) error

PostReservationStage is the next composable stage after reservation (egress, submit). It receives the ReservationResult (including Claim) for chaining.

func NewPostReservationEgressStage

func NewPostReservationEgressStage(cfg Config, store CompressionStore, svc CompressionServices, next PostEgressStage) PostReservationStage

NewPostReservationEgressStage returns a PostReservationStage that evaluates the trusted egress policy with explicit route/purpose/source/ principal from the trusted correlation, redacts locally before byte/token accounting, prepares sanitized segments from the authoritative artifact (via store Snapshot lookup, never from correlation/model metadata), enforces MaxInputBytes/Tokens after redaction, derives the authoritative EgressPolicyHash, performs full CAS provisional->authoritative promotion via UpdateReservationPolicyHash(returns nil, reasoningpreservation.SanitizationNone, sha256.Sum256([]byte("test-route"))) so original remains untouched. Trusted sanitizer is taken from svc.Sanitizer, not from untrusted policy decision; policy may request redact but cannot inject arbitrary sanitizer authority.

func NewPostReservationEgressStageWithTelemetry

func NewPostReservationEgressStageWithTelemetry(cfg Config, store CompressionStore, svc CompressionServices, next PostEgressStage, tel *Telemetry) PostReservationStage

NewPostReservationEgressStageWithTelemetry is the telemetry-aware egress stage. It records content-free privacy outcomes (allow/redact/deny/missing-policy) without content.

type PreparationOutcome

type PreparationOutcome string

PreparationOutcome is the typed content-free outcome of bounded preparation.

const (
	OutcomeIneligible          PreparationOutcome = "ineligible"
	OutcomeDenied              PreparationOutcome = "denied"
	OutcomeMissingPolicy       PreparationOutcome = "missing-policy"
	OutcomeSanitizerFailed     PreparationOutcome = "sanitizer_failed"
	OutcomeInputBytesExceeded  PreparationOutcome = "input_bytes_exceeded"
	OutcomeInputTokensExceeded PreparationOutcome = "input_tokens_exceeded"
	OutcomeInputOversize       PreparationOutcome = "input_oversize" // alias for bytes exceeded (compat)
	OutcomePrepared            PreparationOutcome = "prepared"
)

type PreparedReservation

type PreparedReservation struct {
	Reservation      ReservationResult
	Segments         []CompressorInputSegment
	Decision         CompressionEgressDecision
	EgressPolicyHash [32]byte
	Route            string
}

PreparedReservation is passed to the next stage after successful egress decision, sanitization, input budgeting, and authoritative CAS promotion. It is content-bearing only for the next stage's local use; it is not persisted. ReservationResult is content-free. Segments contain only local placement index + sanitized text; no session/account/lineage/anchor/digest is ever placed here.

type ReasoningSurrogate

type ReasoningSurrogate struct {
	OriginalDigest      [32]byte
	PolicyRevision      string
	Sanitization        string
	Segments            []SurrogateSegment
	Bytes               int
	SemanticDigest      [32]byte
	EgressPolicyHash    [32]byte
	AuthorizedRouteHash [32]byte
}

ReasoningSurrogate is the validated optional replacement for semantic-text placements. SemanticDigest is the SHA-256 of the canonical semantic-text payload that was compressed; EgressPolicyHash is the hash of the egress policy version that authorized submission. Both are immutable correlation digests verified by AttachSurrogate CAS. A zero SemanticDigest is invalid (real artifacts always have content) and Attach rejects it as a conflict. Sanitization is the content-free class that authorized the surrogate (none/redacted). AuthorizedRouteHash is sha256(route) at promotion, verified by Attach CAS.

type ReasoningViewResult

type ReasoningViewResult struct {
	Kind   ViewKind
	Reason string
}

ReasoningViewResult is the immutable stage result for 6.1. It is content-free and carries the surrogate eligibility before any substitution (6.2).

type ReplaySemantics

type ReplaySemantics uint8

ReplaySemantics is the bounded typed classification for reasoning replay. It is derived from canonical dialect plus structure/presence, not provider strings.

const (
	ReplayUnknown ReplaySemantics = iota
	ReplayExactRequired
	ReplaySemanticText
)

func ClassifyReasoningPart

func ClassifyReasoningPart(part lipapi.Part) ReplaySemantics

ClassifyReasoningPart returns the replay semantics for a single reasoning part. Pure function: provider-name free, no I/O, deterministic.

type ReservationOutcome

type ReservationOutcome string

ReservationOutcome is a typed content-free outcome of the reservation stage.

const (
	ReservationReserved              ReservationOutcome = "reserved"
	ReservationSkippedIneligible     ReservationOutcome = "skipped_ineligible"
	ReservationSkippedBelowThreshold ReservationOutcome = "skipped_below_threshold"
	ReservationNotFound              ReservationOutcome = "not_found"
	ReservationConflict              ReservationOutcome = "conflict"
	ReservationBudgetExceeded        ReservationOutcome = "budget_exceeded"
	ReservationError                 ReservationOutcome = "error"
)

type ReservationResult

type ReservationResult struct {
	Outcome     ReservationOutcome
	Claim       CompressionClaim
	Correlation PostAppendCorrelation
	Err         error
}

ReservationResult is passed to next hooks in the post-append chain. It is content-free: no reasoning text, no credentials, no raw hashes.

func TryReserveCompression

func TryReserveCompression(ctx context.Context, cfg Config, store CompressionStore, corr PostAppendCorrelation) ReservationResult

TryReserveCompression performs the local reservation BEFORE any egress/provider call. It checks ineligibility and MinSourceBytes, then calls ReserveCompression with the correct SemanticDigest from correlation and provisional EgressPolicyRefHash. On budget/conflict/notfound it returns typed outcome and leaves original untouched. It never evicts original reasoning and never calls provider.

func (ReservationResult) IsReserved

func (r ReservationResult) IsReserved() bool

IsReserved reports whether reservation succeeded.

type RestoreInput

type RestoreInput struct {
	Action            string
	OnUnrepresentable string
	OnStateError      string
	Call              *lipapi.Call
	Artifacts         []TurnArtifact
	ReplaySupport     lipapi.ReasoningReplaySupport
	Eligible          bool
}

type RestoreResult

type RestoreResult struct {
	Mutated       bool
	RestoredCount int
	RestoredBytes int
	Exclude       bool
	ReasonCode    string
	Outcomes      []SafeOutcome
}

func RestoreMissingReasoning

func RestoreMissingReasoning(in RestoreInput) (RestoreResult, error)

type RuleConfig

type RuleConfig struct {
	ID            string   `yaml:"id"`
	Backend       string   `yaml:"backend"`
	ModelKeywords []string `yaml:"model_keywords"`
	Enabled       *bool    `yaml:"enabled"`
}

type SafeInventory

type SafeInventory struct {
	Enabled            bool             `json:"enabled"`
	Action             string           `json:"action"`
	CatalogVersion     string           `json:"catalog_version,omitempty"`
	RuleIDs            []string         `json:"rule_ids,omitempty"`
	RuleCount          int              `json:"rule_count"`
	TTL                string           `json:"ttl"`
	MaxTurnsPerSession int              `json:"max_turns_per_session"`
	MaxBytesPerTurn    int              `json:"max_reasoning_bytes_per_turn"`
	MaxSessionBytes    int              `json:"max_session_bytes"`
	ProcessLocal       bool             `json:"process_local"`
	AggregateCounters  map[string]int64 `json:"aggregate_counters"`
}

SafeInventory is the process-local, content-safe diagnostics projection for one enabled instance.

func BuildSafeInventory

func BuildSafeInventory(cfg Config, tel *Telemetry) SafeInventory

BuildSafeInventory projects config + aggregate counters without payloads, anchors, or session partitions.

type SafeOutcome

type SafeOutcome string
const (
	OutcomeObserved            SafeOutcome = "observed"
	OutcomePreserved           SafeOutcome = "preserved"
	OutcomeMissing             SafeOutcome = "missing"
	OutcomeRestored            SafeOutcome = "restored"
	OutcomeAmbiguous           SafeOutcome = "ambiguous"
	OutcomeConflicting         SafeOutcome = "conflicting"
	OutcomeUnmatched           SafeOutcome = "unmatched"
	OutcomeUnrepresentable     SafeOutcome = "unrepresentable"
	OutcomeStateError          SafeOutcome = "state_error"
	OutcomeEvicted             SafeOutcome = "evicted"
	OutcomeOversize            SafeOutcome = "oversize"
	OutcomeBoundedRaw          SafeOutcome = "bounded_raw"
	OutcomeRawOversize         SafeOutcome = "raw_oversize"
	OutcomeRawInvalidChannel   SafeOutcome = "raw_invalid_channel"
	OutcomeRawInvalidLimit     SafeOutcome = "raw_invalid_limit"
	OutcomeDecodeInvalid       SafeOutcome = "decode_invalid"
	OutcomeSchemaInvalid       SafeOutcome = "schema_invalid"
	OutcomeControlInvalid      SafeOutcome = "control_invalid"
	OutcomeSurrogateOversize   SafeOutcome = "surrogate_oversize"
	OutcomeInsufficientSavings SafeOutcome = "insufficient_savings"
	OutcomeStale               SafeOutcome = "stale"
	OutcomeSurrogateAttached   SafeOutcome = "surrogate_attached"
	OutcomeShadowReady         SafeOutcome = "shadow_ready"
	// 5.4 taxonomy additions — content-free, no raw IDs, no reasoning text.
	OutcomeEligible                  SafeOutcome = "eligible"
	OutcomeCompIneligible            SafeOutcome = "ineligible"
	OutcomeExact                     SafeOutcome = "exact"
	OutcomeBelowThreshold            SafeOutcome = "below_threshold"
	OutcomeReservationBudgetExceeded SafeOutcome = "reservation_budget_exceeded"
	OutcomeBudgetPendingPerSession   SafeOutcome = "budget_pending_per_session"
	OutcomeBudgetPendingTotal        SafeOutcome = "budget_pending_total"
	OutcomeBudgetSurrogatePerTurn    SafeOutcome = "budget_surrogate_per_turn"
	OutcomeBudgetSurrogatePerSession SafeOutcome = "budget_surrogate_per_session"
	OutcomeBudgetSurrogateTotal      SafeOutcome = "budget_surrogate_total"
	OutcomeEgressAllow               SafeOutcome = "egress_allow"
	OutcomeEgressRedact              SafeOutcome = "egress_redact"
	OutcomeEgressDeny                SafeOutcome = "egress_deny"
	OutcomeEgressMissingPolicy       SafeOutcome = "egress_missing_policy"
	OutcomeSubmitted                 SafeOutcome = "submitted"
	OutcomeCoalesced                 SafeOutcome = "coalesced"
	OutcomeQueueSaturated            SafeOutcome = "queue_saturated"
	OutcomeAdmissionDenied           SafeOutcome = "admission_denied"
	OutcomeSubmitFailed              SafeOutcome = "submit_failed"
	OutcomePollPending               SafeOutcome = "poll_pending"
	OutcomePollCompleted             SafeOutcome = "poll_completed"
	OutcomePollFailed                SafeOutcome = "poll_failed"
	OutcomePollNotFound              SafeOutcome = "poll_not_found"
	OutcomePollUnavailable           SafeOutcome = "poll_unavailable"
	OutcomePollError                 SafeOutcome = "poll_error"
	OutcomeOriginalFallback          SafeOutcome = "original_fallback"
	OutcomeActiveUsed                SafeOutcome = "active_used"
)

type SecretguardTrustedSanitizer

type SecretguardTrustedSanitizer struct {
	Matcher secretguard.Matcher
}

SecretguardTrustedSanitizer adapts the existing pkg/lipsdk/secretguard.Matcher contract to the feature's narrow TrustedTextSanitizer seam. This reuses the repository's authoritative exact-match secret redaction authority and avoids a second heuristic detector. No new detection logic is introduced here.

Import direction: feature depends on pkg/lipsdk contract only, never on concrete secretguard feature/engine implementation. Composition (runtimebundle, later task) will inject an instance produced from MatcherResolver.

func (SecretguardTrustedSanitizer) SanitizeText

func (s SecretguardTrustedSanitizer) SanitizeText(ctx context.Context, text string) (string, error)

SanitizeText redacts exact catalog matches via the underlying Matcher. Findings are discarded at this seam; only the sanitized text is returned. A nil Matcher fails closed rather than returning unredacted input.

type SegmentSemantics

type SegmentSemantics struct {
	PlacementIndex int
	Dialect        lipapi.ReasoningDialect
	Semantics      ReplaySemantics
	SourceBytes    int
}

SegmentSemantics describes the classification of a single placed reasoning segment.

func ClassifyPlacement

func ClassifyPlacement(idx int, pr PlacedReasoning) SegmentSemantics

ClassifyPlacement returns per-placement semantics for a single placed reasoning.

func ClassifyPlacements

func ClassifyPlacements(placements []PlacedReasoning) []SegmentSemantics

ClassifyPlacements returns per-placement semantics for all placements. The PlacementIndex reflects input order, not BeforeNonReasoningPart.

type SessionPartition

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

func NewSessionPartition

func NewSessionPartition(opaque string) SessionPartition

func (SessionPartition) String

func (p SessionPartition) String() string

type ShadowEvaluation

type ShadowEvaluation struct {
	Counts       map[SafeOutcome]int64
	SourceBytes  map[SafeOutcome]int64
	RawBytes     map[SafeOutcome]int64
	DecodedBytes map[SafeOutcome]int64
	SavedBytes   map[SafeOutcome]int64
	Latency      map[SafeOutcome]time.Duration
	TotalCount   int64
	TotalSource  int64
	TotalRaw     int64
	TotalDecoded int64
	TotalSaved   int64
	// Ratio is hypothetical saved/source when TotalSource>0, else 0. No claim of semantic equivalence.
	SavingsRatio     float64
	CompressionRatio float64
	// AvgLatency is TotalLatency / TotalCount for outcomes with latency if TotalCount>0.
	AvgLatency time.Duration
}

ShadowEvaluation is the bounded content-free evaluation snapshot for shadow mode. It aggregates source/raw/decoded/saved bytes and ratios, plus per-outcome counts and latency. No economic/money calculation is performed; billing remains authoritative via existing auxiliary BillingCallID / usage surfaces. All numeric fields are bounded.

type StateConfig

type StateConfig struct {
	TTL                      time.Duration `yaml:"ttl"`
	MaxTurnsPerSession       int           `yaml:"max_turns_per_session"`
	MaxReasoningBytesPerTurn int           `yaml:"max_reasoning_bytes_per_turn"`
	MaxSessionBytes          int           `yaml:"max_session_bytes"`
}

type StoreOptions

type StoreOptions struct {
	TTL                      time.Duration
	MaxTurnsPerSession       int
	MaxReasoningBytesPerTurn int
	MaxSessionBytes          int
	Now                      func() time.Time
	CompressionLimits        CompressionLimits
}

type StreamObserverFactory

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

func NewStreamObserverFactory

func NewStreamObserverFactory(cfg Config, store TurnStore, tel ...*Telemetry) *StreamObserverFactory

func NewStreamObserverFactoryWithPostAppendHook

func NewStreamObserverFactoryWithPostAppendHook(cfg Config, store TurnStore, hook PostAppendHook, tel ...*Telemetry) *StreamObserverFactory

NewStreamObserverFactoryWithPostAppendHook creates a factory with an immutable post-append hook. The hook is assigned once at construction and never mutated thereafter, avoiding data races. When compression is disabled the hook must be nil.

func (*StreamObserverFactory) FailureMode

func (f *StreamObserverFactory) FailureMode() sdkhooks.FailureMode

func (*StreamObserverFactory) ID

func (f *StreamObserverFactory) ID() string

func (*StreamObserverFactory) LastSafeDiagnostic

func (f *StreamObserverFactory) LastSafeDiagnostic() string

func (*StreamObserverFactory) Open

func (*StreamObserverFactory) Order

func (f *StreamObserverFactory) Order() int

type SurrogateDecodeOutcome

type SurrogateDecodeOutcome = SafeOutcome

SurrogateDecodeOutcome is the typed content-free outcome of strict decoding. It is an alias to SafeOutcome so telemetry and decoder share the same string taxonomy.

const (
	OutcomeSurrogateDecoded SurrogateDecodeOutcome = "decoded"
)

type SurrogateDecodeParams

type SurrogateDecodeParams struct {
	ExpectedIndexes     []int
	SourceBytes         int
	MaxSurrogateBytes   int
	MinSavedBytes       int
	MinSavingsRatio     float64
	OriginalDigest      [32]byte
	PolicyRevision      string
	Sanitization        string
	SemanticDigest      [32]byte
	EgressPolicyHash    [32]byte
	AuthorizedRouteHash [32]byte
}

SurrogateDecodeParams correlates the decoder result with the authoritative original, policy, and sanitization. SourceBytes is the authoritative source size used for savings comparison. Caller must provide raw bytes already bounded by ExtractBoundedRaw (max_output_bytes / hard ceiling).

type SurrogateSegment

type SurrogateSegment struct {
	PlacementIndex int
	Text           string
	Bytes          int
}

SurrogateSegment is a minimal placement-indexed surrogate text.

type Telemetry

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

Telemetry accumulates content-safe aggregate outcome counters for one feature instance. All byte counters are content-free, outcome-whitelisted via isKnownOutcome, and individually clamped to hard ceilings to prevent unbounded aggregation. Concurrency: single registry keyed by SafeOutcome to *outcomeBuckets with lock-free per-outcome atomic updates; accumulation is saturation-safe.

func NewTelemetry

func NewTelemetry() *Telemetry

func (*Telemetry) BytesSnapshot

func (t *Telemetry) BytesSnapshot() map[SafeOutcome]int64

BytesSnapshot returns aggregate raw bytes per outcome (alias to RawBytesSnapshot, single count, legacy compat).

func (*Telemetry) CompressionMeasurementsSnapshot

func (t *Telemetry) CompressionMeasurementsSnapshot() CompressionMeasurements

CompressionMeasurementsSnapshot returns per-outcome counts and byte measurements.

func (*Telemetry) DecodedBytesSnapshot

func (t *Telemetry) DecodedBytesSnapshot() map[SafeOutcome]int64

DecodedBytesSnapshot returns decoded bytes per outcome.

func (*Telemetry) LatencySnapshot

func (t *Telemetry) LatencySnapshot() map[SafeOutcome]time.Duration

LatencySnapshot returns total latency per outcome (bounded).

func (*Telemetry) RawBytesSnapshot

func (t *Telemetry) RawBytesSnapshot() map[SafeOutcome]int64

RawBytesSnapshot returns raw bytes per outcome (dedicated, bounded).

func (*Telemetry) Record

func (t *Telemetry) Record(outcome SafeOutcome, counts map[string]int)

func (*Telemetry) RecordCompression

func (t *Telemetry) RecordCompression(outcome SafeOutcome, rawBytes int)

RecordCompression records a compression-safe outcome with content-free byte count. It records raw bytes once in the dedicated rawBytes bucket (single count, bounded to HardRawOutputCeiling). BytesSnapshot is an alias to RawBytesSnapshot for backward compat. New code should prefer RecordCompressionMeasurement or RecordShadowMeasurement.

func (*Telemetry) RecordCompressionMeasurement

func (t *Telemetry) RecordCompressionMeasurement(outcome SafeOutcome, rawBytes, decodedBytes, savedBytes int)

RecordCompressionMeasurement records an outcome with explicit raw/decoded/saved bytes, each clamped. Raw is stored once in rawBytes; Decoded bounded to HardCompressionMaxSurrogateBytes, saved to HardRawOutputCeiling. It is content-free and outcome-whitelisted via isKnownOutcome.

func (*Telemetry) RecordShadowMeasurement

func (t *Telemetry) RecordShadowMeasurement(outcome SafeOutcome, sourceBytes, rawBytes, decodedBytes, savedBytes int, latency time.Duration)

RecordShadowMeasurement records a full shadow evaluation sample with source/raw/decoded/saved and latency. All byte values are bounded; latency is bounded to HardCompressionTimeout per sample. No money calculation. Content-free via isKnownOutcome.

func (*Telemetry) SavedBytesSnapshot

func (t *Telemetry) SavedBytesSnapshot() map[SafeOutcome]int64

SavedBytesSnapshot returns saved bytes per outcome.

func (*Telemetry) ShadowEvaluationSnapshot

func (t *Telemetry) ShadowEvaluationSnapshot() ShadowEvaluation

ShadowEvaluationSnapshot returns a bounded evaluation computing source/raw/decoded/saved totals and ratios. Ratios are hypothetical savings only; no money calculation; latency is averaged if available.

func (*Telemetry) Snapshot

func (t *Telemetry) Snapshot() map[SafeOutcome]int64

func (*Telemetry) SourceBytesSnapshot

func (t *Telemetry) SourceBytesSnapshot() map[SafeOutcome]int64

SourceBytesSnapshot returns source bytes per outcome (bounded).

type TrustedTextSanitizer

type TrustedTextSanitizer interface {
	SanitizeText(ctx context.Context, text string) (string, error)
}

TrustedTextSanitizer is a small interface for an existing secret/redaction authority.

func NewResolverSanitizer

func NewResolverSanitizer(r sdk.MatcherResolver) TrustedTextSanitizer

NewResolverSanitizer returns a TrustedTextSanitizer that resolves the matcher from resolver on each SanitizeText call. Nil resolver returns nil so CompressionServices validation can fail closed.

func NewTrustedSanitizerFromMatcher

func NewTrustedSanitizerFromMatcher(m secretguard.Matcher) TrustedTextSanitizer

NewTrustedSanitizerFromMatcher returns a TrustedTextSanitizer backed by m. A nil Matcher yields nil so CompressionServices validation can fail closed when a redact-then-allow policy is configured but no sanitizer is wired.

type TurnArtifact

type TurnArtifact struct {
	ID             string
	Anchor         [32]byte
	SourceBackend  string
	SourceModel    string
	Reasoning      []PlacedReasoning
	CreatedAt      time.Time
	ReasoningBytes int
}

func BuildEphemeralArtifact

func BuildEphemeralArtifact(original TurnArtifact, surrogate *ReasoningSurrogate) TurnArtifact

BuildEphemeralArtifact returns a defensive copy of original with only matching semantic Reasoning.Text fields replaced from surrogate. BeforeNonReasoningPart, dialect, signature/opaque/summary/content/encrypted fields remain byte-equivalent. If surrogate is nil, or any correlation is ambiguous/missing/duplicate/index out of range or semantic subset mismatch, the whole original defensive copy is returned (fallback).

func BuildEphemeralArtifacts

func BuildEphemeralArtifacts(arts []TurnArtifact, decisions map[string]ReasoningViewResult, getSurrogate func(string) (*ReasoningSurrogate, bool)) []TurnArtifact

BuildEphemeralArtifacts builds defensive ephemeral copies for all artifacts. For each ViewSurrogate decision, the corresponding surrogate is fetched via getSurrogate; if fetch fails or fallback triggers, the original defensive copy is used. The store is never mutated; returned slice is ephemeral.

type TurnStore

func NewMemoryTurnStore

func NewMemoryTurnStore(opts StoreOptions) (TurnStore, error)

type ViewKind

type ViewKind string

ViewKind is the pure selection decision before any mutation. 6.1 may return Original/Surrogate + reason but must NOT substitute yet (6.2).

const (
	ViewOriginal  ViewKind = "original"
	ViewSurrogate ViewKind = "surrogate"
)

func SelectReasoningView

func SelectReasoningView(cfg CompressionConfig, artifact TurnArtifact, surrogate *ReasoningSurrogate, support lipapi.ReasoningReplaySupport, classification Classification) (ViewKind, string)

SelectReasoningView is a pure eligibility function over the original artifact + surrogate + candidate ReplaySupport + client classification. It is a wrapper without current-policy check for backward compatibility.

func SelectReasoningViewWithCurrentPolicy

func SelectReasoningViewWithCurrentPolicy(cfg CompressionConfig, artifact TurnArtifact, surrogate *ReasoningSurrogate, support lipapi.ReasoningReplaySupport, classification Classification, currentPolicyHash *[32]byte, currentSanitization *string) (ViewKind, string)

SelectReasoningViewWithCurrentPolicy is the policy-aware pure variant. If currentPolicyHash/currentSanitization are non-nil, they are the trusted current EgressPolicy result prepared by the stage (hash via ComputeEgressPolicyHash, sanitization via EgressAllow/Redact). When provided, they must match the stored surrogate's EgressPolicyHash and Sanitization, otherwise policy_no_longer_permits. No I/O inside pure; stage does the Decide call.

Directories

Path Synopsis

Jump to

Keyboard shortcuts

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