Documentation
¶
Overview ¶
Package reasoningpreservation owns the official reasoning-output-preservation feature domain.
Index ¶
- Constants
- Variables
- func BuildCompressorAuxRequest(p CompressorAuxRequestParams) (auxiliary.Request, error)
- func BuildEphemeralCandidates(cands []restoreCandidate, decisions map[string]ReasoningViewResult, ...) []restoreCandidate
- func BuildFeatureBundle(n yaml.Node) (lipfeature.FeatureBundle, error)
- func ComputeAnchor(msg lipapi.Message) ([32]byte, error)
- func ComputeEgressPolicyHash(decision CompressionEgressDecision, route string) [32]byte
- func DecodeSurrogate(raw []byte, params SurrogateDecodeParams) (ReasoningSurrogate, SurrogateDecodeOutcome, error)
- func EnsureCompanionRules(n yaml.Node, backendIDs []string, rulePrefix string) (yaml.Node, error)
- func EstimateInputTokens(text string) int
- func EstimatedBytesForSegments(segments []CompressorInputSegment) int
- func EstimatedTokensForSegments(segments []CompressorInputSegment) int
- func ExtractBoundedRaw(collected lipapi.Collected, maxOutputBytes int) ([]byte, error)
- func FeatureBundle(cfg Config) (lipfeature.FeatureBundle, error)
- func FeatureBundleWithCompanionPolicy(cfg Config, policy CompanionPolicy) (lipfeature.FeatureBundle, error)
- func FormatSafeDiagnostic(outcome SafeOutcome, ruleID string, counts map[string]int) (string, error)
- func IsBudgetError(err error) bool
- func IsConflictError(err error) bool
- func IsNotFoundError(err error) bool
- func MatchEligible(kind MatchKind) bool
- func NewCompanionConfig(backendIDs []string, rulePrefix string) (yaml.Node, error)
- func PrepareCompressorInputWithLimits(ctx context.Context, segments []CompressorInputSegment, ...) ([]CompressorInputSegment, PreparationOutcome, error)
- func PrepareSemanticSegments(ctx context.Context, placements []PlacedReasoning, ...) ([]CompressorInputSegment, PreparationOutcome, error)
- func PrepareSemanticSegmentsFromArtifact(ctx context.Context, artifact TurnArtifact, decision CompressionEgressDecision, ...) ([]CompressorInputSegment, PreparationOutcome, error)
- func ProjectSafeError(err error) (string, error)
- type AdoptionOutcome
- type AdoptionResult
- type AttemptTransform
- type BudgetError
- type BudgetKind
- type BuiltinCatalogEntry
- type CandidateIdentity
- type Classification
- type ClassifiedTurn
- type CompanionPolicy
- type CompletedAdoptionStage
- type CompletedPollCandidate
- type CompressionClaim
- type CompressionConfig
- type CompressionEgressDecision
- type CompressionEgressInput
- type CompressionLimits
- type CompressionMeasurements
- type CompressionMode
- type CompressionPollAttemptResult
- type CompressionServices
- type CompressionState
- type CompressionStats
- type CompressionStore
- type CompressorAuxRequestParams
- type CompressorInputSegment
- func ExtractSemanticSegments(placements []PlacedReasoning) []CompressorInputSegment
- func ExtractSemanticSegmentsFromArtifact(artifact TurnArtifact) []CompressorInputSegment
- func PrepareCompressorInput(ctx context.Context, segments []CompressorInputSegment, ...) ([]CompressorInputSegment, string, error)
- type Config
- type EgressAction
- type EgressPolicy
- type EgressPrincipalView
- type EvictionSummary
- type InstanceParts
- func FeatureBundleWithCompression(cfg Config, svc CompressionServices) (*InstanceParts, lipfeature.FeatureBundle, error)
- func FeatureBundleWithParts(cfg Config) (*InstanceParts, lipfeature.FeatureBundle, error)
- func FeatureBundleWithPartsAndCompression(cfg Config, svc CompressionServices, policy CompanionPolicy) (*InstanceParts, lipfeature.FeatureBundle, error)
- func FeatureBundleWithPartsAndPolicy(cfg Config, policy CompanionPolicy) (*InstanceParts, lipfeature.FeatureBundle, error)
- type MatchKind
- type MatchResult
- type PendingCompression
- type PlacedReasoning
- type PollForAttemptKind
- type PostAppendCorrelation
- type PostAppendHook
- func BuildPostAppendHook(cfg Config, store TurnStore, svc CompressionServices) PostAppendHook
- func BuildPostAppendHookWithEgressNext(cfg Config, store TurnStore, svc CompressionServices, next PostEgressStage) PostAppendHook
- func BuildPostAppendHookWithNext(cfg Config, store TurnStore, svc CompressionServices, ...) PostAppendHook
- func BuildPostAppendHookWithTelemetry(cfg Config, store TurnStore, svc CompressionServices, tel *Telemetry) PostAppendHook
- func NewCompressionReservationHook(cfg Config, store CompressionStore, next PostReservationStage) PostAppendHook
- func NewCompressionReservationHookWithTelemetry(cfg Config, store CompressionStore, next PostReservationStage, tel *Telemetry) PostAppendHook
- type PostEgressStage
- type PostReservationStage
- type PreparationOutcome
- type PreparedReservation
- type ReasoningSurrogate
- type ReasoningViewResult
- type ReplaySemantics
- type ReservationOutcome
- type ReservationResult
- type RestoreInput
- type RestoreResult
- type RuleConfig
- type SafeInventory
- type SafeOutcome
- type SecretguardTrustedSanitizer
- type SegmentSemantics
- type SessionPartition
- type ShadowEvaluation
- type StateConfig
- type StoreOptions
- type StreamObserverFactory
- func (f *StreamObserverFactory) FailureMode() sdkhooks.FailureMode
- func (f *StreamObserverFactory) ID() string
- func (f *StreamObserverFactory) LastSafeDiagnostic() string
- func (f *StreamObserverFactory) Open(_ context.Context, meta response.StreamMeta, _ response.Services) (response.StreamObserver, error)
- func (f *StreamObserverFactory) Order() int
- type SurrogateDecodeOutcome
- type SurrogateDecodeParams
- type SurrogateSegment
- type Telemetry
- func (t *Telemetry) BytesSnapshot() map[SafeOutcome]int64
- func (t *Telemetry) CompressionMeasurementsSnapshot() CompressionMeasurements
- func (t *Telemetry) DecodedBytesSnapshot() map[SafeOutcome]int64
- func (t *Telemetry) LatencySnapshot() map[SafeOutcome]time.Duration
- func (t *Telemetry) RawBytesSnapshot() map[SafeOutcome]int64
- func (t *Telemetry) Record(outcome SafeOutcome, counts map[string]int)
- func (t *Telemetry) RecordCompression(outcome SafeOutcome, rawBytes int)
- func (t *Telemetry) RecordCompressionMeasurement(outcome SafeOutcome, rawBytes, decodedBytes, savedBytes int)
- func (t *Telemetry) RecordShadowMeasurement(outcome SafeOutcome, sourceBytes, rawBytes, decodedBytes, savedBytes int, ...)
- func (t *Telemetry) SavedBytesSnapshot() map[SafeOutcome]int64
- func (t *Telemetry) ShadowEvaluationSnapshot() ShadowEvaluation
- func (t *Telemetry) Snapshot() map[SafeOutcome]int64
- func (t *Telemetry) SourceBytesSnapshot() map[SafeOutcome]int64
- type TrustedTextSanitizer
- type TurnArtifact
- type TurnStore
- type ViewKind
Constants ¶
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 )
const ( SanitizationNone = "none" SanitizationRedacted = "redacted" )
const ( ActionObserve = "observe" ActionRestore = "restore" PolicyLogSkip = "log_skip" PolicyReject = "reject" )
const ( EgressPurposeReasoningSemanticCompression = "reasoning_semantic_compression" EgressSourceClassSemanticText = "semantic_reasoning_text" )
const BuiltinCatalogVersion = reasoningreplay.CatalogVersion
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.
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.
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.
const ID = "reasoning-output-preservation"
Variables ¶
var ( ErrCompressionNotFound = errors.New(ID + ": compression not found") ErrCompressionConflict = errors.New(ID + ": compression conflict") ErrCompressionBudgetExceeded = errors.New(ID + ": compression budget exceeded") )
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") )
var ( ErrRawOversize = errors.New(ID + ": raw_oversize") ErrRawInvalidChannel = errors.New(ID + ": raw_invalid_channel") ErrRawInvalidLimit = errors.New(ID + ": raw_invalid_limit") )
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") )
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 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 ¶
func DecodeSurrogate(raw []byte, params SurrogateDecodeParams) (ReasoningSurrogate, SurrogateDecodeOutcome, error)
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 ¶
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 ¶
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 ¶
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 IsBudgetError ¶
IsBudgetError reports whether err is a budget rejection.
func IsConflictError ¶
IsConflictError reports whether err is a CAS/stale conflict.
func IsNotFoundError ¶
IsNotFoundError reports whether err is a not-found.
func MatchEligible ¶
func NewCompanionConfig ¶
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 ¶
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 (t *AttemptTransform) HandleAttempt(ctx context.Context, call *lipapi.Call, meta request.AttemptMeta, _ request.Services) (request.AttemptDecision, error)
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 ¶
func BuiltinCatalogEntries ¶
func BuiltinCatalogEntries() ([]BuiltinCatalogEntry, error)
type CandidateIdentity ¶
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 ¶
func EvaluateEgress(ctx context.Context, policy EgressPolicy, in CompressionEgressInput) CompressionEgressDecision
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 ¶
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"`
}
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 ¶
func (v EgressPrincipalView) Scope() scope.PrincipalScopeView
Scope returns a defensive clone of the underlying principal scope.
type EvictionSummary ¶
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 ¶
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 ¶
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" 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 ¶
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 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" 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 ¶
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 StoreOptions ¶
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 (f *StreamObserverFactory) Open(_ context.Context, meta response.StreamMeta, _ response.Services) (response.StreamObserver, error)
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 ¶
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) 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 ¶
type TurnStore interface {
Append(context.Context, SessionPartition, TurnArtifact) (EvictionSummary, error)
Snapshot(context.Context, SessionPartition) ([]TurnArtifact, error)
Delete(context.Context, SessionPartition, ...string) error
}
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).
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.
Source Files
¶
- anchor.go
- artifact.go
- bundle.go
- catalog.go
- classify.go
- companion.go
- companion_policy.go
- compression_adoption_stage.go
- compression_attempt_poll.go
- compression_cleanup.go
- compression_config.go
- compression_egress_stage.go
- compression_ephemeral_view.go
- compression_reservation.go
- compression_selection.go
- compression_services.go
- compression_store.go
- compression_submit_stage.go
- compressor_input.go
- compressor_request.go
- config.go
- continuity_marker.go
- doc.go
- egress.go
- egress_route.go
- id.go
- inert_observer.go
- observer.go
- observer_compression.go
- outcome.go
- partition.go
- raw_extractor.go
- resolver_sanitizer.go
- restore.go
- restore_candidates.go
- secretguard_adapter.go
- semantics.go
- store.go
- surrogate_decoder.go
- telemetry.go
- transform.go