continuation

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

Documentation

Overview

Package continuation defines protocol-neutral proxy-owned response continuation contracts: opaque response IDs, scoped stores, persistence policy, terminal recording, and bounded materialization. Wire profile types and upstream continuation IDs stay outside this package.

Index

Constants

View Source
const (
	// MinResponseIDEntropyBytes is the minimum raw entropy required for proxy IDs.
	MinResponseIDEntropyBytes = 16
	// MaxResponseIDLength is the maximum allowed string length for proxy continuation IDs.
	MaxResponseIDLength = 512
	// ResponseIDPrefix is the stable external prefix for proxy-issued response IDs.
	ResponseIDPrefix = "resp_"
)

Variables

View Source
var (
	// ErrPreviousResponseNotFound is returned for missing, expired, unauthorized,
	// evicted, or incompatible previous_response_id lookups. Callers must treat
	// all causes identically on the wire.
	ErrPreviousResponseNotFound = errors.New("continuation: previous response not found")

	// ErrChainDepthExceeded rejects continuation chains beyond configured depth.
	ErrChainDepthExceeded = errors.New("continuation: chain depth exceeded")

	// ErrCycleDetected rejects cyclic previous_response_id chains.
	ErrCycleDetected = errors.New("continuation: cycle detected")

	// ErrLineageMismatch rejects a chain whose provider-bound links disagree.
	ErrLineageMismatch = errors.New("continuation: lineage mismatch")

	// ErrMaterializedSizeExceeded rejects reconstructed context above byte bounds.
	ErrMaterializedSizeExceeded = errors.New("continuation: materialized size exceeded")

	// ErrMaterializedItemsExceeded rejects reconstructed context above item bounds.
	ErrMaterializedItemsExceeded = errors.New("continuation: materialized items exceeded")

	// ErrRecordNotReady rejects lookup of a reserved but non-terminal record.
	ErrRecordNotReady = errors.New("continuation: record not terminal")

	// ErrInvalidPolicy rejects malformed persistence policy values.
	ErrInvalidPolicy = errors.New("continuation: invalid storage policy")

	// ErrStorageLimitExceeded rejects records or stores above configured bounds.
	ErrStorageLimitExceeded = errors.New("continuation: storage limit exceeded")

	// ErrIncompleteNotEligible rejects incomplete responses when policy disallows them.
	ErrIncompleteNotEligible = errors.New("continuation: incomplete record is not eligible")

	// ErrRecordNotEligible rejects terminal records that cannot be continued.
	ErrRecordNotEligible = errors.New("continuation: record is not eligible")

	// ErrStorageFailure classifies an unavailable or failed persistence boundary.
	ErrStorageFailure = errors.New("continuation: storage failure")

	// ErrNativeReferencesUnprotected rejects durable writes when provider-native
	// evidence has no configured encryption/protection boundary.
	ErrNativeReferencesUnprotected = errors.New("continuation: native references require protected storage")

	// ErrStoreClosed rejects operations attempted after a store has been closed.
	ErrStoreClosed = errors.New("continuation: store closed")
)

Functions

func CloneItems

func CloneItems(in []lipapi.Item) []lipapi.Item

CloneItems returns defensive copies of trajectory slices, including nested JSON.

func CloneRequirements

CloneRequirements returns an independent requirements value.

func EstimateItemsBytes

func EstimateItemsBytes(items []lipapi.Item) int64

func EstimateRecordBytes

func EstimateRecordBytes(rec ContinuationRecord) int64

func RecordSize

func RecordSize(record ContinuationRecord) int64

RecordSize returns a conservative serialized size for storage accounting.

Types

type Bounds

type Bounds struct {
	MaxChainDepth        int
	MaxMaterializedItems int
	MaxMaterializedBytes int64
}

Bounds configures deterministic traversal and materialization limits.

func DefaultBounds

func DefaultBounds() Bounds

DefaultBounds returns conservative Task 1.5 contract defaults aligned with design.

type ContinuationRecord

type ContinuationRecord struct {
	ID                 ResponseID
	Scope              Scope
	PreviousID         ResponseID
	ProfileID          string
	Lineage            Lineage
	InputItems         []lipapi.Item
	OutputItems        []lipapi.Item
	Requirements       lipapi.ProtocolRequirements
	Policy             StoragePolicy
	ExpiresAt          time.Time
	Terminal           bool
	Status             RecordStatus
	NativeRefs         []NativeReference
	NativeRequirements []NativeRequirement
	MaterializedBytes  int64
	ChainDepth         int
}

ContinuationRecord is the protocol-neutral persisted continuation payload.

func CloneRecord

func CloneRecord(in ContinuationRecord) ContinuationRecord

CloneRecord returns a deep copy suitable for crossing the store boundary.

func Lookup

func Lookup(ctx context.Context, store Store, scope Scope, id ResponseID) (ContinuationRecord, error)

Lookup performs a scoped get and maps every miss to ErrPreviousResponseNotFound.

type Lineage

type Lineage struct {
	ProfileID     string
	Model         string
	RouteSelector string
	// ProviderBound prevents a record containing provider-specific semantics from
	// silently moving to a different candidate.
	ProviderBound bool
	ProviderID    string
	CandidateKey  string
}

Lineage captures model/route context required for portable reroute or pinning.

type MaterializeInput

type MaterializeInput struct {
	Store      Store
	Scope      Scope
	StartID    ResponseID
	NewInput   []lipapi.Item
	Bounds     Bounds
	Now        func() int64
	EstimateFn func(ContinuationRecord) int64
}

MaterializeInput describes one bounded continuation materialization request.

type MaterializedTrajectory

type MaterializedTrajectory struct {
	Items              []lipapi.Item
	InputItems         []lipapi.Item
	OutputItems        []lipapi.Item
	NewInput           []lipapi.Item
	ChainDepth         int
	TotalBytes         int64
	Lineage            Lineage
	Requirements       lipapi.ProtocolRequirements
	NativeRequirements []NativeRequirement
}

MaterializedTrajectory is prior input, prior output, then new input in order.

func Materialize

Materialize walks the previous_response_id chain, enforces depth/cycle/byte bounds, and returns the concatenated semantic order without invoking backends.

func MaterializeCall

MaterializeCall resolves proxy-owned continuation state into a fresh provider-facing call.

type MemoryStore

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

MemoryStore is a bounded in-memory implementation of the continuation Store port. It is intentionally protocol-neutral and safe for concurrent use.

func NewMemoryStore

func NewMemoryStore() *MemoryStore

NewMemoryStore returns an empty bounded in-memory continuation store.

func NewMemoryStoreWithLimits

func NewMemoryStoreWithLimits(limits StorageLimits) *MemoryStore

NewMemoryStoreWithLimits creates a bounded in-memory continuation store.

func (*MemoryStore) Close

func (s *MemoryStore) Close() error

Close idempotently closes the store and clears records and reservations.

func (*MemoryStore) Delete

func (s *MemoryStore) Delete(ctx context.Context, scope Scope, id ResponseID) error

func (*MemoryStore) Get

func (*MemoryStore) PutTerminal

func (s *MemoryStore) PutTerminal(ctx context.Context, record ContinuationRecord) error

func (*MemoryStore) Reserve

func (s *MemoryStore) Reserve(ctx context.Context, scope Scope, policy StoragePolicy) (ResponseID, error)

func (*MemoryStore) SetClock

func (s *MemoryStore) SetClock(now func() time.Time)

SetClock replaces the store clock for deterministic tests.

type NativeReference

type NativeReference struct {
	Provider string
	Kind     string
	ID       string
	Opaque   []byte
}

NativeReference is private provider evidence retained for exact-lineage optimizations. It is never a proxy ID and must not be forwarded as client continuation state.

func (NativeReference) GoString

func (r NativeReference) GoString() string

GoString returns a redacted representation for fmt %#v formatting.

func (NativeReference) String

func (r NativeReference) String() string

String returns a redacted string representation to prevent sensitive native state leakage.

type NativeRequirement

type NativeRequirement struct {
	BackendID   string
	Model       string
	Kind        string
	Dialect     string
	Implementor string
}

NativeRequirement describes private provider evidence that is valid only for one exact backend/model lineage. It is never a provider request field.

type PersistenceMode

type PersistenceMode string

PersistenceMode selects whether a record survives reconnect/process boundaries.

const (
	PersistencePersistent PersistenceMode = "persistent"
	PersistenceConnection PersistenceMode = "connection_local"
)

type RecordStatus

type RecordStatus string

RecordStatus describes whether a terminal record can be used as a parent.

const (
	RecordStatusCompleted  RecordStatus = "completed"
	RecordStatusIncomplete RecordStatus = "incomplete"
	RecordStatusFailed     RecordStatus = "failed"
)

func EffectiveStatus

func EffectiveStatus(record ContinuationRecord) RecordStatus

EffectiveStatus treats records written by the Phase 1 contract as completed.

type Recorder

type Recorder interface {
	RecordTerminal(ctx context.Context, record ContinuationRecord) error
}

Recorder accepts incremental terminal output for one reserved response ID.

type Resolver

type Resolver interface {
	ResolveParent(ctx context.Context, scope Scope, parentID string, baseCall lipapi.Call) (lipapi.Call, ContinuationRecord, error)
}

Resolver resolves a client-visible previous response ID into a canonical call.

func NewResolver

func NewResolver(store Store, bounds Bounds) Resolver

NewResolver constructs a resolver backed by a protocol-neutral continuation store.

type ResponseID

type ResponseID string

ResponseID is a proxy-issued opaque continuation identifier.

func NewResponseID

func NewResponseID(ctx context.Context) (ResponseID, error)

NewResponseID returns a high-entropy proxy response ID suitable for a continuation store. The SDK owns this protocol-neutral primitive so frontend plugins do not depend on an internal core implementation.

func (ResponseID) IsZero

func (id ResponseID) IsZero() bool

IsZero reports whether the ID is unset.

func (ResponseID) String

func (id ResponseID) String() string

String returns the wire-safe ID string.

func (ResponseID) Validate

func (id ResponseID) Validate() error

Validate checks prefix and minimum encoded entropy for externally supplied IDs.

type Scope

type Scope struct {
	TenantID     string
	PrincipalID  string
	SessionID    string
	ConnectionID string // non-empty for connection-local store:false records
}

Scope isolates continuation records to an authoritative client/session identity. Client-controlled metadata must not widen or substitute a scope.

func (Scope) Equal

func (s Scope) Equal(o Scope) bool

Equal reports whether two scopes denote the same isolation boundary.

func (Scope) IsZero

func (s Scope) IsZero() bool

IsZero reports whether the scope carries no authoritative identity.

func (Scope) String

func (s Scope) String() string

String returns a stable diagnostic label without embedding secrets.

type StorageLimits

type StorageLimits struct {
	MaxRecords     int
	MaxBytes       int64
	MaxRecordBytes int64
	MaxChainDepth  int
}

StorageLimits bounds one store. Zero values select the store defaults.

func DefaultStorageLimits

func DefaultStorageLimits() StorageLimits

DefaultStorageLimits returns finite production-safe storage limits.

type StoragePolicy

type StoragePolicy struct {
	Mode            PersistenceMode
	TTL             time.Duration
	AllowIncomplete bool
	Limits          StorageLimits
}

StoragePolicy governs reservation and retention for one continuation chain link.

func (StoragePolicy) Validate

func (p StoragePolicy) Validate() error

Validate rejects policies that could weaken scope or retention guarantees.

type Store

type Store interface {
	Reserve(ctx context.Context, scope Scope, policy StoragePolicy) (ResponseID, error)
	PutTerminal(ctx context.Context, record ContinuationRecord) error
	Get(ctx context.Context, scope Scope, id ResponseID) (ContinuationRecord, error)
	Delete(ctx context.Context, scope Scope, id ResponseID) error
}

Store reserves proxy response IDs, stores terminal records, and performs scoped lookup.

type StreamObserver

type StreamObserver interface {
	Observe(ctx context.Context, event lipapi.Event)
	Close() error
}

StreamObserver is a best-effort canonical-stream observation seam. Observe must not return an error: persistence is intentionally outside downstream commitment and failover decisions.

type StreamRecorder

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

StreamRecorder observes canonical events and attempts one terminal write. Storage failures are deliberately retained as diagnostics only; they never become stream errors and therefore cannot trigger retry after commitment.

func NewStreamRecorder

func NewStreamRecorder(recorder Recorder, record ContinuationRecord, cleanup func()) *StreamRecorder

NewStreamRecorder creates an observer for one reserved response record.

func (*StreamRecorder) Close

func (r *StreamRecorder) Close() error

Close is idempotent and does not turn a partial/cancelled stream into a stored continuation record.

func (*StreamRecorder) ContinuationReservationCleanupConsumed

func (r *StreamRecorder) ContinuationReservationCleanupConsumed() bool

ContinuationReservationCleanupConsumed reports whether the recorder detached its cleanup callback.

func (*StreamRecorder) FinalizeIncomplete

func (r *StreamRecorder) FinalizeIncomplete(ctx context.Context) error

FinalizeIncomplete persists the events observed so far as an incomplete terminal record. It is idempotent.

func (*StreamRecorder) Observe

func (r *StreamRecorder) Observe(ctx context.Context, event lipapi.Event)

Observe records a defensive copy of each event and stores only after the canonical terminal event. Non-terminal failures and cancellation are never eligible for persistence.

func (*StreamRecorder) OwnsContinuationReservation

func (r *StreamRecorder) OwnsContinuationReservation() bool

OwnsContinuationReservation identifies the recorder as the owner of the reservation cleanup callback passed at construction time.

func (*StreamRecorder) ReleaseContinuationReservation

func (r *StreamRecorder) ReleaseContinuationReservation()

ReleaseContinuationReservation consumes the reservation cleanup callback without attempting to persist a terminal record.

func (*StreamRecorder) StorageError

func (r *StreamRecorder) StorageError() error

StorageError exposes a best-effort persistence failure for diagnostics/tests.

type TerminalRecorder

type TerminalRecorder struct {
	Store Store
}

TerminalRecorder adapts a Store to the Recorder port for one response lifecycle.

func (TerminalRecorder) RecordTerminal

func (r TerminalRecorder) RecordTerminal(ctx context.Context, record ContinuationRecord) error

RecordTerminal validates terminal state and persists through the store.

Jump to

Keyboard shortcuts

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