sessiontree

package
v2.2.0 Latest Latest
Warning

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

Go to latest
Published: Jul 29, 2026 License: MIT Imports: 28 Imported by: 0

Documentation

Overview

Package sessiontree stores durable conversation journals.

Repo implementations own thread metadata, append-only entries, forks, leaf movement, and provider-visible context reconstruction. FileRepo and MemoryRepo are safe for concurrent use. Repos used by agentharness should implement TurnLeaseRepo so active-turn serialization is durable per ThreadID and shared across harness instances that use the same backend. Different ThreadIDs, including forked threads, may run concurrently.

Index

Constants

View Source
const (
	ApprovalReasonUserRejected             = "user_rejected"
	ApprovalReasonPolicyDenied             = "policy_denied"
	ApprovalReasonAuthorizationUnavailable = "authorization_unavailable"
	ApprovalReasonAuthorizationContract    = "authorization_contract"
	ApprovalReasonCancelled                = "cancelled"
	ApprovalReasonTimedOut                 = "timed_out"
)
View Source
const (
	SubAgentInputIDMetadataKey                     = "subagent_input_id"
	SubAgentUserMessageOriginMetadataKey           = "subagent_user_message_origin"
	SubAgentUserMessageOriginDelegatedMission      = "delegated_mission"
	SubAgentUserMessageOriginInput                 = "subagent_input"
	SubAgentUserMessageOriginPendingToolCompletion = "pending_tool_completion"
)
View Source
const (
	InterruptedTurnRecoveryKindKey                 = "authority_kind"
	InterruptedTurnRecoveryKind                    = "interrupted_turn_recovery"
	InterruptedTurnRecoveryFingerprintKey          = "authority_fingerprint"
	InterruptedTurnRecoveryParentKey               = "authority_parent_thread_id"
	InterruptedTurnRecoveryApprovalRootKey         = "authority_approval_root_thread_id"
	InterruptedTurnRecoveryApprovalGenerationKey   = "authority_approval_queue_generation"
	InterruptedTurnRecoveryApprovalRevisionKey     = "authority_approval_queue_revision"
	InterruptedTurnRecoverySourceFailureEntryKey   = "authority_source_failure_entry_id"
	InterruptedTurnRecoverySourceFailureRawHashKey = "authority_source_failure_raw_hash"
	InterruptedTurnFailureMessage                  = "turn interrupted during previous process"
	InterruptedTurnEffectOutcomeUnknownMessage     = "effect outcome is unknown because the turn was interrupted after dispatch"
	BranchBoundaryTurnFailureMessage               = "turn interrupted by branch boundary"
)
View Source
const (
	PendingToolSettlementKindKey        = "authority_kind"
	PendingToolSettlementKind           = "pending_tool_settlement"
	PendingToolSettlementFingerprintKey = "authority_fingerprint"
	PendingToolEffectAttemptIDKey       = "effect_attempt_id"
)
View Source
const (
	RetrySourceTurnIDMetadataKey  = "retry_source_turn_id"
	RetrySourceEntryIDMetadataKey = "retry_source_entry_id"
)
View Source
const (
	TurnFailureCodeMetadataKey = "failure_code"

	TurnFailureCancelled                = "cancelled"
	TurnFailureInterrupted              = "interrupted"
	TurnFailureProvider                 = "provider"
	TurnFailureToolDispatch             = "tool_dispatch"
	TurnFailureEffectOutcomeUnknown     = "effect_outcome_unknown"
	TurnFailureAuthorizationUnavailable = "authorization_unavailable"
	TurnFailureAuthorizationContract    = "authorization_contract"
	TurnFailureStorage                  = "storage"
	TurnFailureEngineContract           = "engine_contract"
	TurnFailureLegacyUnclassified       = "legacy_unclassified"
)
View Source
const CompactionMutationKind = "compaction"
View Source
const MaxThreadTitleRunes = 200

MaxThreadTitleRunes is the canonical title admission and projection limit.

Variables

View Source
var (
	ErrArtifactNotFound           = errors.New("session tree artifact not found")
	ErrSubAgentNotFound           = errors.New("session tree subagent not found")
	ErrSubAgentParentRequired     = errors.New("session tree subagent requires parent authority")
	ErrUnsupportedStoreCapability = errors.New("session tree store capability is unsupported")
)
View Source
var (
	ErrRequestConflict         = errors.New("session tree authority request conflicts with persisted request")
	ErrSubAgentRequestConflict = fmt.Errorf("subagent request identity conflicts with persisted request: %w", ErrRequestConflict)
	ErrSubAgentInputNotFound   = errors.New("session tree subagent input not found")
	ErrEffectAttemptNotFound   = errors.New("session tree effect attempt not found")
	ErrEffectOutcomeUnknown    = errors.New("session tree effect outcome is unknown")
)
View Source
var (
	ErrPendingToolTurnNotFound = errors.New("session tree pending tool turn not found")
	ErrPendingToolRunNotFound  = errors.New("session tree pending tool run not found")
	ErrPendingToolNotFound     = errors.New("session tree pending tool not found")
	ErrPendingToolNotPending   = errors.New("session tree pending tool is not pending")
)
View Source
var (
	ErrThreadNotFound           = errors.New("session tree thread not found")
	ErrEntryNotFound            = errors.New("session tree entry not found")
	ErrInvalidParent            = errors.New("session tree invalid parent")
	ErrActiveTurn               = errors.New("session tree thread already has an active turn")
	ErrThreadExists             = errors.New("session tree thread already exists")
	ErrInvalidThreadAuthority   = errors.New("session tree invalid thread authority")
	ErrThreadAuthorityBusy      = errors.New("session tree thread authority is busy")
	ErrThreadClosed             = errors.New("session tree thread is closed")
	ErrThreadDeleted            = errors.New("session tree thread is deleted")
	ErrSubAgentClosing          = errors.New("session tree subagent is closing")
	ErrStaleAuthority           = errors.New("session tree authority proof is stale")
	ErrRecoveryTargetResolved   = errors.New("session tree interrupted recovery target is resolved")
	ErrAuthorityCorrupt         = errors.New("session tree authority state is corrupt")
	ErrForkDestinationConflict  = errors.New("session tree fork destination conflicts with operation marker")
	ErrAgentTodoVersionConflict = errors.New("session tree agent todo version conflict")
	ErrStaleCanonicalTurnCursor = errors.New("session tree canonical turn cursor is stale")
)
View Source
var DefaultLeasePolicy = LeasePolicy{
	TTL:                30 * time.Second,
	RenewInterval:      10 * time.Second,
	ClockSkewAllowance: 2 * time.Second,
}
View Source
var ErrApprovalNotFound = errors.New("session tree approval not found")
View Source
var ErrCanonicalTurnNotFound = errors.New("session tree canonical turn not found")
View Source
var ErrProviderStateNotFound = errors.New("provider state not found")

Functions

func ActivePathHash

func ActivePathHash(path []Entry) string

func AddInterruptedTurnApprovalQueueProof

func AddInterruptedTurnApprovalQueueProof(metadata map[string]string, proof *InterruptedTurnApprovalQueueProof) error

func AddInterruptedTurnRecoverySourceFailureProof

func AddInterruptedTurnRecoverySourceFailureProof(metadata map[string]string, proof *InterruptedTurnRecoveryFailureProof) error

func ApprovalBatchCancellationFingerprint

func ApprovalBatchCancellationFingerprint(fingerprint string) string

func ApprovalCancellationEntryID

func ApprovalCancellationEntryID(cancellationFingerprint, approvalID string) string

func ApprovalDispatchEntryID

func ApprovalDispatchEntryID(decisionID, approvalID string) string

func ApprovalDispatchEntryRequestMatches

func ApprovalDispatchEntryRequestMatches(stored, requested Entry) bool

func ApprovalEffectAttemptID

func ApprovalEffectAttemptID(invocation EffectInvocationIdentity) string

func ApprovalFinalizationEntryID

func ApprovalFinalizationEntryID(decisionID, approvalID string) string

func ApprovalPreflightMatchesEffect

func ApprovalPreflightMatchesEffect(item ApprovalPreflightItem, lease TurnLease, attempt EffectAttempt) bool

func ApprovalPreflightMatchesRecord

func ApprovalPreflightMatchesRecord(record ApprovalRecord, rootID, parentID string, item ApprovalPreflightItem, attempt EffectAttempt) bool

func ApprovalQueueVisible

func ApprovalQueueVisible(state ApprovalState) bool

func ApprovalRejectedEntryID

func ApprovalRejectedEntryID(decisionID, approvalID string) string

func ApprovalRequestedEntryID

func ApprovalRequestedEntryID(approvalID string) string

func BuildContext

func BuildContext(path []Entry, opts ContextOptions) []session.Message

func BuildContextChecked

func BuildContextChecked(path []Entry, opts ContextOptions) ([]session.Message, error)

func ContextWithTurnLease

func ContextWithTurnLease(ctx context.Context, lease TurnLease) context.Context

ContextWithTurnLease binds the exact durable mutation owner to journal writes.

func CreateRootFingerprint

func CreateRootFingerprint(req CreateRootRequest) string

CreateRootFingerprint is the stable identity used by root-create replay.

func CreateRootReplayMatches

func CreateRootReplayMatches(meta ThreadMeta, threadID string) bool

CreateRootReplayMatches reports whether a live row is the exact canonical root shape eligible for root-create replay.

func EffectResultRequestMatches

func EffectResultRequestMatches(committed, requested Entry, effectAttemptID string) bool

func InterruptedTurnRecoveryApprovalCancellationID

func InterruptedTurnRecoveryApprovalCancellationID(recoveryFingerprint, approvalID string) string

func InterruptedTurnRecoveryCancelledEffectFingerprint

func InterruptedTurnRecoveryCancelledEffectFingerprint(recoveryFingerprint string) string

func InterruptedTurnRecoveryFingerprint

func InterruptedTurnRecoveryFingerprint(
	expectedLease TurnLease,
	parentThreadID string,
	runID string,
	status TurnMarkerStatus,
	failureCode string,
	failureMessage string,
	sourceFailure *InterruptedTurnRecoveryFailureProof,
	effects []InterruptedTurnRecoveryEffect,
) (string, error)

func InterruptedTurnRecoveryUnknownEffectFingerprint

func InterruptedTurnRecoveryUnknownEffectFingerprint(recoveryFingerprint string) string

func InterruptedTurnToolResult

func InterruptedTurnToolResult(call session.Message, effectState EffectAttemptState) session.Message

func MatchesForkDestinationMeta

func MatchesForkDestinationMeta(meta ThreadMeta, destination *ForkDestinationMeta) bool

MatchesForkDestinationMeta reports whether persisted ownership metadata exactly matches the fork plan. It is used for replay idempotency.

func NormalizeApprovalCancellationEntries

func NormalizeApprovalCancellationEntries(req CancelApprovalBatchRequest) (map[string]Entry, error)

func NormalizeSubAgentCloseIntent

func NormalizeSubAgentCloseIntent(operationID, parentThreadID, targetThreadID, reason string) (string, string, string, string, string, error)

NormalizeSubAgentCloseIntent validates the storage-kernel close identity and returns its normalized fields plus immutable intent fingerprint.

func PublishSubAgentChildMatches

func PublishSubAgentChildMatches(req PublishSubAgentRequest, child ThreadMeta) bool

PublishSubAgentChildMatches verifies the durable child authority produced by or replayed for one exact publication request.

func RawForEntry

func RawForEntry(entry Entry) string

func RetryPathHasRetryEligibleDurableInput

func RetryPathHasRetryEligibleDurableInput(path []Entry) (bool, error)

RetryPathHasRetryEligibleDurableInput evaluates the newest canonical user input at or before a retry source. Callers must provide an authority-validated ancestor path ending at the source entry.

func RetrySourceHasRetryEligibleDurableInput

func RetrySourceHasRetryEligibleDurableInput(path []Entry, sourceTurnID, sourceEntryID string) (int, bool, error)

RetrySourceHasRetryEligibleDurableInput validates the structural retry source and reports whether its canonical user input contains durable text or an attachment. References alone require ephemeral supplemental context and are therefore not replayable.

func SameApprovalIdentity

func SameApprovalIdentity(left, right ApprovalIdentity) bool

func SameThreadAuthority

func SameThreadAuthority(left, right ThreadMeta) bool

SameThreadAuthority reports whether an update preserves immutable ownership and lineage identity.

func SameThreadTitleState

func SameThreadTitleState(left, right ThreadMeta) bool

func SameTurnLease

func SameTurnLease(left, right TurnLease) bool

func SortThreadsByCreatedAtDesc

func SortThreadsByCreatedAtDesc(threads []ThreadMeta)

func StableHash

func StableHash(value string) string

func SubAgentCloseRequestFingerprint

func SubAgentCloseRequestFingerprint(intentFingerprint string, nodes []SubAgentCloseNode) (string, error)

SubAgentCloseRequestFingerprint binds one close intent to the subtree membership derived by the storage transaction.

func SubAgentUserMessageOrigin

func SubAgentUserMessageOrigin(kind SubAgentRequestKind) (string, error)

func ThreadAuthorityTreeIDs

func ThreadAuthorityTreeIDs(threads []ThreadMeta, rootThreadID string) ([]string, error)

ThreadAuthorityTreeIDs returns one root and all descendants owned through ParentThreadID after validating the complete authority graph.

func TurnAdmissionRequestFingerprint

func TurnAdmissionRequestFingerprint(req AdmitTurnRequest) (string, error)

func UpdateTurnLeaseContext

func UpdateTurnLeaseContext(ctx context.Context, previous, renewed TurnLease) error

UpdateTurnLeaseContext advances one context binding after a successful durable renewal. It cannot replace a different owner or generation.

func ValidTurnFailureCode

func ValidTurnFailureCode(code string) bool

func ValidateAdmitPendingToolCompletionEnvelope

func ValidateAdmitPendingToolCompletionEnvelope(req AdmitPendingToolCompletionRequest) error

func ValidateAdmitPendingToolCompletionReplayRequest

func ValidateAdmitPendingToolCompletionReplayRequest(req AdmitPendingToolCompletionRequest) error

func ValidateAdmitPendingToolCompletionRequest

func ValidateAdmitPendingToolCompletionRequest(req AdmitPendingToolCompletionRequest) error

func ValidateAdmitTurnReplayRequest

func ValidateAdmitTurnReplayRequest(req AdmitTurnRequest) error

ValidateAdmitTurnReplayRequest preserves the attachment shape accepted by historical admissions while retaining every other request-shape check.

func ValidateAdmitTurnRequest

func ValidateAdmitTurnRequest(req AdmitTurnRequest) error

func ValidateAdmitTurnRequestEnvelope

func ValidateAdmitTurnRequestEnvelope(req AdmitTurnRequest) error

ValidateAdmitTurnRequestEnvelope validates the authority fields required to look up an existing admission before applying new-admission input limits.

func ValidateApprovalDecisionReceiptAuthority

func ValidateApprovalDecisionReceiptAuthority(receipt ApprovalDecisionReceipt, record ApprovalRecord, queue ApprovalQueue) error

func ValidateApprovalRecord

func ValidateApprovalRecord(record ApprovalRecord) error

func ValidateBeginAutomaticThreadTitleRequest

func ValidateBeginAutomaticThreadTitleRequest(req BeginAutomaticThreadTitleRequest) error

func ValidateBeginCompactionRequest

func ValidateBeginCompactionRequest(req BeginCompactionRequest) error

func ValidateCancelApprovalBatchRequest

func ValidateCancelApprovalBatchRequest(req CancelApprovalBatchRequest) error

func ValidateCanonicalThreadTitle

func ValidateCanonicalThreadTitle(title string) error

ValidateCanonicalThreadTitle validates committed title text independently from its authority state.

func ValidateCanonicalTurnEntries

func ValidateCanonicalTurnEntries(entries []Entry, threadID, turnID, runID string) error

func ValidateCanonicalTurnReadAuthority

func ValidateCanonicalTurnReadAuthority(turn CanonicalTurn, admission CanonicalTurnAdmissionFact) error

func ValidateCanonicalTurnReadStructure

func ValidateCanonicalTurnReadStructure(turn CanonicalTurn, threadID string) error

ValidateCanonicalTurnReadStructure validates journal shape independently of optional execution-admission authority.

func ValidateCommitApprovalDispatchRequest

func ValidateCommitApprovalDispatchRequest(req CommitApprovalDispatchRequest) error

func ValidateCompleteAutomaticThreadTitleRequest

func ValidateCompleteAutomaticThreadTitleRequest(req CompleteAutomaticThreadTitleRequest) error

func ValidateCreateRootRequest

func ValidateCreateRootRequest(req CreateRootRequest) error

ValidateCreateRootRequest validates the exact root-create contract.

func ValidateEffectLeaseSuccessor

func ValidateEffectLeaseSuccessor(proof, current TurnLease) error

ValidateEffectLeaseSuccessor permits an effect authority operation to use a proof captured before one or more durable heartbeat renewals. Ownership, generation, acquisition identity, and monotonic lease time must remain exact.

func ValidateEntryIntegrity

func ValidateEntryIntegrity(entry Entry) error

func ValidateEntryMessageAttachments

func ValidateEntryMessageAttachments(entry Entry) error

func ValidateEntryMessageReferences

func ValidateEntryMessageReferences(entry Entry) error

func ValidateFailAutomaticThreadTitleRequest

func ValidateFailAutomaticThreadTitleRequest(req FailAutomaticThreadTitleRequest) error

func ValidateFinalizeApprovalRequest

func ValidateFinalizeApprovalRequest(req FinalizeApprovalRequest) error

func ValidateFinalizeApprovalRequestedEntry

func ValidateFinalizeApprovalRequestedEntry(record ApprovalRecord, entry Entry) error

func ValidateFinalizeApprovalResultAuthority

func ValidateFinalizeApprovalResultAuthority(req FinalizeApprovalRequest, result FinalizeApprovalResult) error

func ValidateFinalizeApprovalSourceAuthority

func ValidateFinalizeApprovalSourceAuthority(req FinalizeApprovalRequest, record ApprovalRecord, effect EffectAttempt, queue ApprovalQueue) error

func ValidateFinishTurnRequest

func ValidateFinishTurnRequest(req FinishTurnRequest) error

func ValidateForkPrepareState

func ValidateForkPrepareState(rootThreadID string, nodes []ForkOptions, states []ForkPrepareThreadState) error

ValidateForkPrepareState rejects a plan whose pinned source snapshot or terminal-child set no longer matches the canonical state at claim time.

func ValidateForkRetryAuthorityPath

func ValidateForkRetryAuthorityPath(path []Entry, threadID string) error

ValidateForkRetryAuthorityPath validates the complete staged destination path before a fork publishes entries or retry admission facts.

func ValidateInterruptedTurnAdmissionPath

func ValidateInterruptedTurnAdmissionPath(path []Entry, threadID, turnID, runID, turnStartedID string) error

func ValidateInterruptedTurnLeaseSuccessor

func ValidateInterruptedTurnLeaseSuccessor(target, current TurnLease) error

ValidateInterruptedTurnLeaseSuccessor verifies that current is a canonical monotonic successor of the exact recovery target proof.

func ValidateInterruptedTurnRecoveryEffectAttempts

func ValidateInterruptedTurnRecoveryEffectAttempts(attempts []EffectAttempt, threadID, turnID, runID string) error

func ValidateInterruptedTurnRecoveryPath

func ValidateInterruptedTurnRecoveryPath(path []Entry, turnID, runID string) error

func ValidateInterruptedTurnStartedEntry

func ValidateInterruptedTurnStartedEntry(entry Entry, threadID, turnID, runID, turnStartedID string) error

func ValidateListCanonicalTurnsOptions

func ValidateListCanonicalTurnsOptions(opts ListCanonicalTurnsOptions) error

func ValidateNewEntryMessageAttachments

func ValidateNewEntryMessageAttachments(entry Entry) error

func ValidatePendingToolCompletionPath

func ValidatePendingToolCompletionPath(path []Entry) error

func ValidatePublishSubAgentIdentity

func ValidatePublishSubAgentIdentity(req PublishSubAgentRequest) error

ValidatePublishSubAgentIdentity validates the durable request-ledger key and conflict identity before a backend interprets the requested child shape.

func ValidatePublishSubAgentInputEnvelope

func ValidatePublishSubAgentInputEnvelope(req PublishSubAgentInputRequest) error

func ValidatePublishSubAgentInputReplayRequest

func ValidatePublishSubAgentInputReplayRequest(req PublishSubAgentInputRequest) error

func ValidatePublishSubAgentInputRequest

func ValidatePublishSubAgentInputRequest(req PublishSubAgentInputRequest) error

func ValidatePublishSubAgentPendingToolCompletionEnvelope

func ValidatePublishSubAgentPendingToolCompletionEnvelope(req PublishSubAgentPendingToolCompletionRequest) error

func ValidatePublishSubAgentPendingToolCompletionReplayRequest

func ValidatePublishSubAgentPendingToolCompletionReplayRequest(req PublishSubAgentPendingToolCompletionRequest) error

func ValidatePublishSubAgentPendingToolCompletionRequest

func ValidatePublishSubAgentPendingToolCompletionRequest(req PublishSubAgentPendingToolCompletionRequest) error

func ValidatePublishSubAgentReplayRequest

func ValidatePublishSubAgentReplayRequest(req PublishSubAgentRequest) error

func ValidatePublishSubAgentRequest

func ValidatePublishSubAgentRequest(req PublishSubAgentRequest) error

ValidatePublishSubAgentRequest enforces the exact parent, child, fork, and first-input identity before a backend starts the atomic publication.

func ValidateRecoverInterruptedTurnRequest

func ValidateRecoverInterruptedTurnRequest(req RecoverInterruptedTurnRequest) error

func ValidateResolveApprovalReplayAuthority

func ValidateResolveApprovalReplayAuthority(
	expectedDecisionID string,
	decision ApprovalDecision,
	expectedRootThreadID string,
	expectedGeneration int64,
	expectedRevision int64,
	expectedCurrent ApprovalIdentity,
	expectedApprovalRevision int64,
	receipt ApprovalDecisionReceipt,
	record ApprovalRecord,
	effect EffectAttempt,
	queue ApprovalQueue,
) error

func ValidateResolveApprovalRequest

func ValidateResolveApprovalRequest(req ResolveApprovalRequest) error

func ValidateRetrySourcePath

func ValidateRetrySourcePath(path []Entry, sourceTurnID, sourceEntryID string) (int, error)

func ValidateRetryStartedEntry

func ValidateRetryStartedEntry(entry Entry, threadID, turnID, runID, startedEntryID string, source CanonicalTurnRetrySource) error

func ValidateSetThreadTitleRequest

func ValidateSetThreadTitleRequest(req SetThreadTitleRequest) error

func ValidateThreadAuthorityGraph

func ValidateThreadAuthorityGraph(threads []ThreadMeta) error

ValidateThreadAuthorityGraph requires every SubAgent parent chain to be acyclic and terminate at an existing root thread.

func ValidateThreadAuthoritySnapshot

func ValidateThreadAuthoritySnapshot(meta ThreadMeta, path []Entry, lease *TurnLease, claimOperationID string, leaseGeneration int64) error

func ValidateThreadAuthorityState

func ValidateThreadAuthorityState(path []Entry, lease *TurnLease, claimOperationID string) error

func ValidateThreadMetaAuthority

func ValidateThreadMetaAuthority(meta ThreadMeta) error

ValidateThreadMetaAuthority enforces the durable distinction between independent root threads and parent-owned SubAgent threads.

func ValidateThreadTitleProjection

func ValidateThreadTitleProjection(projection ThreadTitleProjection) error

ValidateThreadTitleProjection validates the observable canonical title state shared by durable authority and public runtime projections.

func ValidateThreadTitleState

func ValidateThreadTitleState(meta ThreadMeta) error

Types

type AdmitPendingToolCompletionRequest

type AdmitPendingToolCompletionRequest struct {
	CompletionRequestID   string
	RequestFingerprint    string
	SettlementFingerprint string
	Target                PendingToolSettlementTarget
	Settlement            Entry
	ContinuationTurnID    string
	ContinuationRunID     string
	OwnerID               string
	Input                 session.Message
	Now                   time.Time
}

type AdmitPendingToolCompletionResult

type AdmitPendingToolCompletionResult struct {
	Settlement         Entry
	SettlementReplayed bool
	Admission          AdmitTurnResult
	Replayed           bool
}

type AdmitSubAgentInputRequest

type AdmitSubAgentInputRequest struct {
	ParentThreadID string
	ChildThreadID  string
	TurnID         string
	RunID          string
	OwnerID        string
	Now            time.Time
}

type AdmitSubAgentInputResult

type AdmitSubAgentInputResult struct {
	Input       SubAgentInputRecord
	Lease       TurnLease
	TurnStarted Entry
	UserMessage Entry
	Replayed    bool
}

type AdmitTurnRequest

type AdmitTurnRequest struct {
	ThreadID           string
	TurnID             string
	RunID              string
	OwnerID            string
	Input              session.Message
	RetrySourceTurnID  string
	RetrySourceEntryID string
	RequestFingerprint string
	Now                time.Time
}

type AdmitTurnResult

type AdmitTurnResult struct {
	Lease            TurnLease
	BoundaryTerminal Entry
	TurnStarted      Entry
	UserMessage      Entry
	BaseLeafID       string
	Terminal         *TurnTerminalOutcome
	Replayed         bool
}

type AgentTodoItem

type AgentTodoItem struct {
	ID      string          `json:"id"`
	Content string          `json:"content"`
	Status  AgentTodoStatus `json:"status"`
}

type AgentTodoState

type AgentTodoState struct {
	ThreadID          string          `json:"thread_id"`
	Version           int64           `json:"version"`
	Items             []AgentTodoItem `json:"items"`
	UpdatedAt         time.Time       `json:"updated_at,omitempty"`
	UpdatedByTurnID   string          `json:"updated_by_turn_id,omitempty"`
	UpdatedByRunID    string          `json:"updated_by_run_id,omitempty"`
	UpdatedByToolCall string          `json:"updated_by_tool_call_id,omitempty"`
}

type AgentTodoStateRepo

type AgentTodoStateRepo interface {
	ReadAgentTodoState(context.Context, string) (AgentTodoState, error)
	CompareAndSwapAgentTodoState(context.Context, AgentTodoState, int64) (AgentTodoState, error)
}

type AgentTodoStatus

type AgentTodoStatus string
const (
	AgentTodoPending    AgentTodoStatus = "pending"
	AgentTodoInProgress AgentTodoStatus = "in_progress"
	AgentTodoCompleted  AgentTodoStatus = "completed"
)

type AppendCommittedError

type AppendCommittedError struct {
	Err error
}

func (AppendCommittedError) Error

func (e AppendCommittedError) Error() string

func (AppendCommittedError) Unwrap

func (e AppendCommittedError) Unwrap() error

type AppendOptions

type AppendOptions struct {
	ID       string
	ParentID string
	Now      time.Time
}

type ApprovalCancellationEntry

type ApprovalCancellationEntry struct {
	ApprovalID string
	Entry      Entry
}

type ApprovalDecision

type ApprovalDecision string
const (
	ApprovalDecisionApprove ApprovalDecision = "approve"
	ApprovalDecisionReject  ApprovalDecision = "reject"
)

type ApprovalDecisionReceipt

type ApprovalDecisionReceipt struct {
	DecisionID             string
	ApprovalID             string
	RootThreadID           string
	Decision               ApprovalDecision
	State                  ApprovalState
	Reason                 string
	AuthorizationProofHash string
	QueueGeneration        int64
	QueueRevision          int64
	ApprovalRevision       int64
	SubmittedAt            time.Time
	ResolvedAt             time.Time
}

func CanonicalApprovalDecisionReceipt

func CanonicalApprovalDecisionReceipt(record ApprovalRecord, queue ApprovalQueue) (ApprovalDecisionReceipt, error)

type ApprovalIdentity

type ApprovalIdentity struct {
	ApprovalID      string
	ThreadID        string
	TurnID          string
	RunID           string
	ToolCallID      string
	EffectAttemptID string
}

type ApprovalPreflightItem

type ApprovalPreflightItem struct {
	EffectAttemptID            string
	EffectRequestFingerprint   string
	ApprovalRequestFingerprint string
	Invocation                 EffectInvocationIdentity
	RequestedEntry             Entry
	ToolKind                   string
	Step                       int
	BatchIndex                 int
	BatchSize                  int
	Resources                  []ApprovalResource
	Effects                    []string
	Labels                     map[string]string
	HostContext                map[string]string
	ReadOnly                   bool
	Destructive                bool
	OpenWorld                  bool
}

func NormalizeApprovalPreflightBatch

func NormalizeApprovalPreflightBatch(req PrepareApprovalBatchRequest) ([]ApprovalPreflightItem, error)

type ApprovalQueue

type ApprovalQueue struct {
	RootThreadID      string
	Generation        int64
	Revision          int64
	CurrentApprovalID string
	Items             []ApprovalRecord
	GeneratedAt       time.Time
}

type ApprovalRecord

type ApprovalRecord struct {
	ApprovalID             string
	RootThreadID           string
	ParentThreadID         string
	ThreadID               string
	TurnID                 string
	RunID                  string
	ToolCallID             string
	EffectAttemptID        string
	ToolName               string
	ToolKind               string
	Step                   int
	BatchIndex             int
	BatchSize              int
	ArgsHash               string
	Resources              []ApprovalResource
	Effects                []string
	Labels                 map[string]string
	HostContext            map[string]string
	ReadOnly               bool
	Destructive            bool
	OpenWorld              bool
	RequestFingerprint     string
	State                  ApprovalState
	Revision               int64
	QueueSequence          int64
	DecisionID             string
	Reason                 string
	AuthorizationProofHash string
	RequestedAt            time.Time
	UpdatedAt              time.Time
	ResolvedAt             time.Time
}

func ApprovalRecordFromPreflight

func ApprovalRecordFromPreflight(rootID, parentID string, item ApprovalPreflightItem, attempt EffectAttempt, sequence int64, now time.Time) ApprovalRecord

func (ApprovalRecord) Identity

func (a ApprovalRecord) Identity() ApprovalIdentity

type ApprovalResource

type ApprovalResource struct {
	Kind  string
	Value string
}

type ApprovalState

type ApprovalState string
const (
	ApprovalRequested         ApprovalState = "requested"
	ApprovalDecisionSubmitted ApprovalState = "decision_submitted"
	ApprovalApproved          ApprovalState = "approved"
	ApprovalRejected          ApprovalState = "rejected"
	ApprovalFailed            ApprovalState = "failed"
	ApprovalTimedOut          ApprovalState = "timed_out"
	ApprovalCancelled         ApprovalState = "cancelled"
)

type ArtifactAuthorityRepo

type ArtifactAuthorityRepo interface {
	ReadArtifact(context.Context, ArtifactReadRequest) (ArtifactContent, error)
	ArtifactClosure(context.Context, ArtifactClosureRequest) (artifact.Closure, error)
}

type ArtifactClosureRequest

type ArtifactClosureRequest struct {
	SourceThreadID      string
	DestinationThreadID string
	EntryIDs            []string
}

type ArtifactContent

type ArtifactContent struct {
	Ref  artifact.Ref
	Text string
}

type ArtifactReadRequest

type ArtifactReadRequest struct {
	ParentThreadID string
	ThreadID       string
	ArtifactID     string
}

type BackendRepo

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

BackendRepo executes the canonical session-tree semantics inside Backend snapshot and serializable transactions.

func NewBackendRepo

func NewBackendRepo(ctx context.Context, backend backendspi.Backend, policy LeasePolicy, now func() time.Time) (*BackendRepo, error)

NewBackendRepo initializes or validates the canonical session-tree state.

func (*BackendRepo) AcquireThreadAuthorityClaim

func (repo *BackendRepo) AcquireThreadAuthorityClaim(ctx context.Context, operationID string, requiredSourceThreadIDs, authorityThreadIDs []string) error

func (*BackendRepo) AcquireTurnLease

func (repo *BackendRepo) AcquireTurnLease(ctx context.Context, request TurnLease) (TurnLease, error)

func (*BackendRepo) ActiveTurnLease

func (repo *BackendRepo) ActiveTurnLease(ctx context.Context, threadID string) (TurnLease, bool, error)

func (*BackendRepo) AdmitPendingToolCompletion

func (*BackendRepo) AdmitSubAgentInput

func (*BackendRepo) AdmitTurn

func (repo *BackendRepo) AdmitTurn(ctx context.Context, req AdmitTurnRequest) (AdmitTurnResult, error)

func (*BackendRepo) Append

func (repo *BackendRepo) Append(ctx context.Context, entry Entry, opts AppendOptions) (Entry, error)

func (*BackendRepo) Approval

func (repo *BackendRepo) Approval(ctx context.Context, approvalID string) (ApprovalRecord, error)

func (*BackendRepo) ArtifactClosure

func (repo *BackendRepo) ArtifactClosure(ctx context.Context, req ArtifactClosureRequest) (artifact.Closure, error)

func (*BackendRepo) AuthorityLeasePolicy

func (repo *BackendRepo) AuthorityLeasePolicy() LeasePolicy

AuthorityLeasePolicy returns the immutable persisted lease policy.

func (*BackendRepo) BeginAutomaticThreadTitle

func (repo *BackendRepo) BeginAutomaticThreadTitle(ctx context.Context, req BeginAutomaticThreadTitleRequest) (ThreadTitleMutationResult, error)

func (*BackendRepo) BeginCompaction

func (repo *BackendRepo) BeginCompaction(ctx context.Context, req BeginCompactionRequest) (BeginCompactionResult, error)

func (*BackendRepo) BeginEffectDispatch

func (repo *BackendRepo) BeginEffectDispatch(ctx context.Context, req BeginEffectDispatchRequest) (EffectAttempt, error)

func (*BackendRepo) CancelApprovalBatch

func (*BackendRepo) CanonicalTurnEntries

func (repo *BackendRepo) CanonicalTurnEntries(ctx context.Context, threadID, turnID, runID string) ([]Entry, bool, error)

func (*BackendRepo) CommitApprovalDispatch

func (*BackendRepo) CompareAndSwapAgentTodoState

func (repo *BackendRepo) CompareAndSwapAgentTodoState(ctx context.Context, state AgentTodoState, expectedVersion int64) (AgentTodoState, error)

func (*BackendRepo) CompleteAutomaticThreadTitle

func (repo *BackendRepo) CompleteAutomaticThreadTitle(ctx context.Context, req CompleteAutomaticThreadTitleRequest) (ThreadTitleMutationResult, error)

func (*BackendRepo) CreateRoot

func (repo *BackendRepo) CreateRoot(ctx context.Context, req CreateRootRequest) (CreateRootResult, error)

func (*BackendRepo) CreateThread

func (repo *BackendRepo) CreateThread(ctx context.Context, meta ThreadMeta) (ThreadMeta, error)

func (*BackendRepo) CreateThreadWithInitialEntry

func (repo *BackendRepo) CreateThreadWithInitialEntry(ctx context.Context, meta ThreadMeta, initial Entry) (ThreadMeta, Entry, error)

func (*BackendRepo) DeleteProviderState

func (repo *BackendRepo) DeleteProviderState(ctx context.Context, threadID string) error

func (*BackendRepo) DeleteRootTree

func (repo *BackendRepo) DeleteRootTree(ctx context.Context, rootThreadID string) (DeleteRootTreeResult, error)

func (*BackendRepo) Entries

func (repo *BackendRepo) Entries(ctx context.Context, threadID string) ([]Entry, error)

func (*BackendRepo) Entry

func (repo *BackendRepo) Entry(ctx context.Context, threadID, entryID string) (Entry, error)

func (*BackendRepo) FailAutomaticThreadTitle

func (repo *BackendRepo) FailAutomaticThreadTitle(ctx context.Context, req FailAutomaticThreadTitleRequest) (ThreadTitleMutationResult, error)

func (*BackendRepo) FinalizeApproval

func (repo *BackendRepo) FinalizeApproval(ctx context.Context, req FinalizeApprovalRequest) (FinalizeApprovalResult, error)

func (*BackendRepo) FinishCompaction

func (repo *BackendRepo) FinishCompaction(ctx context.Context, req FinishCompactionRequest) (FinishCompactionResult, error)

func (*BackendRepo) FinishEffectDispatch

func (*BackendRepo) FinishSubAgentClose

func (*BackendRepo) FinishTurn

func (repo *BackendRepo) FinishTurn(ctx context.Context, req FinishTurnRequest) (FinishTurnResult, error)

func (*BackendRepo) Fork

func (repo *BackendRepo) Fork(ctx context.Context, opts ForkOptions) (ThreadMeta, error)

func (*BackendRepo) ForkWithInitialEntry

func (repo *BackendRepo) ForkWithInitialEntry(ctx context.Context, opts ForkOptions, initial Entry) (ThreadMeta, Entry, error)

func (*BackendRepo) InspectSubAgentThreadAuthority

func (repo *BackendRepo) InspectSubAgentThreadAuthority(ctx context.Context, parentThreadID, childThreadID string) (SubAgentThreadAuthoritySnapshot, error)

func (*BackendRepo) InspectThreadAuthority

func (repo *BackendRepo) InspectThreadAuthority(ctx context.Context, threadID string) (ThreadAuthoritySnapshot, error)

func (*BackendRepo) ListCanonicalTurns

func (repo *BackendRepo) ListCanonicalTurns(ctx context.Context, opts ListCanonicalTurnsOptions) (CanonicalTurnsPage, error)

func (*BackendRepo) ListSubAgentInputs

func (repo *BackendRepo) ListSubAgentInputs(ctx context.Context, childThreadID string, state SubAgentInputState) ([]SubAgentInputRecord, error)

func (*BackendRepo) ListThreads

func (repo *BackendRepo) ListThreads(ctx context.Context, opts ListThreadsOptions) ([]ThreadMeta, error)

func (*BackendRepo) MarkEffectUnknown

func (repo *BackendRepo) MarkEffectUnknown(ctx context.Context, req MarkEffectUnknownRequest) (EffectAttempt, error)

func (*BackendRepo) MoveLeaf

func (repo *BackendRepo) MoveLeaf(ctx context.Context, threadID, entryID string) error

func (*BackendRepo) Path

func (repo *BackendRepo) Path(ctx context.Context, threadID, leafID string) ([]Entry, error)

func (*BackendRepo) PathPage

func (repo *BackendRepo) PathPage(ctx context.Context, threadID, leafID, beforeEntryID string, limit int) (PathPage, error)

func (*BackendRepo) PendingAutomaticThreadTitles

func (repo *BackendRepo) PendingAutomaticThreadTitles(ctx context.Context) ([]ThreadMeta, error)

func (*BackendRepo) PrepareApprovalBatch

func (*BackendRepo) PrepareEffectAttempt

func (*BackendRepo) PrepareForkClaim

func (repo *BackendRepo) PrepareForkClaim(ctx context.Context, operationID, rootThreadID string, nodes []ForkOptions) error

func (*BackendRepo) PrepareSubAgentClose

func (*BackendRepo) ProviderState

func (repo *BackendRepo) ProviderState(ctx context.Context, threadID string) (ProviderStateRecord, error)

func (*BackendRepo) PublishSubAgent

func (repo *BackendRepo) PublishSubAgent(ctx context.Context, req PublishSubAgentRequest) (PublishSubAgentResult, error)

func (*BackendRepo) PublishSubAgentInput

func (repo *BackendRepo) PublishSubAgentInput(ctx context.Context, req PublishSubAgentInputRequest) (SubAgentInputRecord, bool, error)

func (*BackendRepo) PutProviderState

func (repo *BackendRepo) PutProviderState(ctx context.Context, record ProviderStateRecord) error

func (*BackendRepo) ReadAgentTodoState

func (repo *BackendRepo) ReadAgentTodoState(ctx context.Context, threadID string) (AgentTodoState, error)

func (*BackendRepo) ReadApprovalQueue

func (repo *BackendRepo) ReadApprovalQueue(ctx context.Context, threadID string) (ApprovalQueue, error)

func (*BackendRepo) ReadArtifact

func (repo *BackendRepo) ReadArtifact(ctx context.Context, req ArtifactReadRequest) (ArtifactContent, error)

func (*BackendRepo) ReadCanonicalTurn

func (repo *BackendRepo) ReadCanonicalTurn(ctx context.Context, threadID, turnID string) (CanonicalTurnRead, error)

func (*BackendRepo) ReadCompaction

func (repo *BackendRepo) ReadCompaction(ctx context.Context, threadID, requestID string) (CompactionOperation, bool, error)

func (*BackendRepo) ReadSubAgentInput

func (repo *BackendRepo) ReadSubAgentInput(ctx context.Context, inputID string) (SubAgentInputRecord, bool, error)

func (*BackendRepo) ReadTurnAdmission

func (repo *BackendRepo) ReadTurnAdmission(ctx context.Context, threadID, turnID, runID string) (AdmitTurnResult, bool, error)

func (*BackendRepo) RecoverInterruptedTurn

func (*BackendRepo) RejectEffectAttempt

func (repo *BackendRepo) RejectEffectAttempt(ctx context.Context, req RejectEffectAttemptRequest) (EffectAttempt, error)

func (*BackendRepo) ReleaseTurnLease

func (repo *BackendRepo) ReleaseTurnLease(ctx context.Context, proof TurnLease) error

func (*BackendRepo) RenewTurnLease

func (repo *BackendRepo) RenewTurnLease(ctx context.Context, proof TurnLease) (TurnLease, error)

func (*BackendRepo) ResolveApproval

func (repo *BackendRepo) ResolveApproval(ctx context.Context, req ResolveApprovalRequest) (ResolveApprovalResult, error)

func (*BackendRepo) SetThreadTitle

func (*BackendRepo) SettlePendingToolRecovery

func (*BackendRepo) TakeOverCompaction

func (repo *BackendRepo) TakeOverCompaction(ctx context.Context, req TakeOverCompactionRequest) (BeginCompactionResult, error)

func (*BackendRepo) Thread

func (repo *BackendRepo) Thread(ctx context.Context, threadID string) (ThreadMeta, error)

func (*BackendRepo) ThreadTombstone

func (repo *BackendRepo) ThreadTombstone(ctx context.Context, threadID string) (ThreadTombstone, error)

func (*BackendRepo) UpdateDomain

func (repo *BackendRepo) UpdateDomain(ctx context.Context, mutate func(*MemoryRepo, backendspi.WriteTx) error) error

UpdateDomain executes one session-tree mutation and related domain writes in the same serializable backend transaction.

func (*BackendRepo) UpdateThread

func (repo *BackendRepo) UpdateThread(ctx context.Context, meta ThreadMeta) error

func (*BackendRepo) ValidateArtifactForkDestination

func (repo *BackendRepo) ValidateArtifactForkDestination(ctx context.Context, closure artifact.Closure) error

func (*BackendRepo) ValidateInterruptedTurnResolution

func (repo *BackendRepo) ValidateInterruptedTurnResolution(ctx context.Context, req RecoverInterruptedTurnRequest) error

func (*BackendRepo) ViewDomain

func (repo *BackendRepo) ViewDomain(ctx context.Context, read func(*MemoryRepo, backendspi.ReadTx) error) error

ViewDomain executes one read against an exact backend snapshot.

func (*BackendRepo) WaitApprovalDecision

func (repo *BackendRepo) WaitApprovalDecision(ctx context.Context, approvalID string) (result WaitApprovalDecisionResult, err error)

WaitApprovalDecision waits without holding a backend transaction open.

type BeginAutomaticThreadTitleRequest

type BeginAutomaticThreadTitleRequest struct {
	ThreadID string
	Token    string
	Now      time.Time
}

type BeginCompactionRequest

type BeginCompactionRequest struct {
	ThreadID             string
	RequestID            string
	RequestFingerprint   string
	Source               string
	SourceLeafID         string
	ActivePathHash       string
	SummarySchemaVersion string
	PromptIdentity       string
	RequestPayloadHash   string
	OwnerID              string
	Now                  time.Time
}

type BeginCompactionResult

type BeginCompactionResult struct {
	Operation        CompactionOperation
	Owner            bool
	TakeoverEligible bool
	Replayed         bool
}

type BeginEffectDispatchRequest

type BeginEffectDispatchRequest struct {
	Lease                  TurnLease
	EffectAttemptID        string
	RequestFingerprint     string
	ObservedHeartbeat      int64
	AuthorizationProofHash string
	Now                    time.Time
}

type CancelApprovalBatchRequest

type CancelApprovalBatchRequest struct {
	Lease                   TurnLease
	RunID                   string
	CancellationFingerprint string
	CancellationEntries     []ApprovalCancellationEntry
	Now                     time.Time
}

type CancelApprovalBatchResult

type CancelApprovalBatchResult struct {
	Queue               ApprovalQueue
	Approvals           []ApprovalRecord
	Effects             []EffectAttempt
	CancellationEntries []Entry
	Replayed            bool
}

type CanonicalTurn

type CanonicalTurn struct {
	TurnID         string
	RunID          string
	StartedEntryID string
	StartedOrdinal int64
	RetrySource    *CanonicalTurnRetrySource
	Entries        []CanonicalTurnPathEntry
}

type CanonicalTurnAdmissionFact

type CanonicalTurnAdmissionFact struct {
	ThreadID      string
	TurnID        string
	RunID         string
	TurnStartedID string
	UserMessageID string
	BaseLeafID    string
}

CanonicalTurnAdmissionFact validates persisted execution authority when it exists for a normal turn and whenever a retry turn omits a user entry.

type CanonicalTurnBeforeCursor

type CanonicalTurnBeforeCursor struct {
	EntryID string
}

type CanonicalTurnPageRepo

type CanonicalTurnPageRepo interface {
	ListCanonicalTurns(context.Context, ListCanonicalTurnsOptions) (CanonicalTurnsPage, error)
}

type CanonicalTurnPathEntry

type CanonicalTurnPathEntry struct {
	Entry   Entry
	Ordinal int64
}

type CanonicalTurnRead

type CanonicalTurnRead struct {
	Turn           CanonicalTurn
	LatestTurn     CanonicalTurn
	ThroughOrdinal int64
	HasRetryTarget bool
}

CanonicalTurnRead is one admitted canonical turn on the current active path.

type CanonicalTurnReadRepo

type CanonicalTurnReadRepo interface {
	ReadCanonicalTurn(context.Context, string, string) (CanonicalTurnRead, error)
}

CanonicalTurnReadRepo reads one admitted turn without scanning a public page.

type CanonicalTurnRepo

type CanonicalTurnRepo interface {
	CanonicalTurnEntries(context.Context, string, string, string) ([]Entry, bool, error)
}

CanonicalTurnRepo reads all journal entries associated with one exact turn. Execution replay authority remains owned by TurnAuthorityRepo.

type CanonicalTurnRetrySource

type CanonicalTurnRetrySource struct {
	TurnID  string
	EntryID string
}

func CanonicalTurnRetrySourceForStartedEntry

func CanonicalTurnRetrySourceForStartedEntry(entry Entry) (*CanonicalTurnRetrySource, error)

func ValidateCanonicalRetrySourceTurn

func ValidateCanonicalRetrySourceTurn(entries []Entry, threadID string, source CanonicalTurnRetrySource) (bool, *CanonicalTurnRetrySource, string, string, error)

type CanonicalTurnSinceCursor

type CanonicalTurnSinceCursor struct {
	EntryID string
}

type CanonicalTurnsPage

type CanonicalTurnsPage struct {
	Turns          []CanonicalTurn
	BeforeCursor   *CanonicalTurnBeforeCursor
	SinceCursor    CanonicalTurnSinceCursor
	HasMore        bool
	ThroughOrdinal int64
	LatestTurnID   string
	HasRetryTarget bool
}

type CommitApprovalDispatchRequest

type CommitApprovalDispatchRequest struct {
	DecisionID               string
	ExpectedRootThreadID     string
	ExpectedGeneration       int64
	ExpectedCurrent          ApprovalIdentity
	ExpectedApprovalRevision int64
	EffectAttemptID          string
	Lease                    TurnLease
	AuthorizationProofHash   string
	ApprovedEntry            Entry
	Now                      time.Time
}

type CommitApprovalDispatchResult

type CommitApprovalDispatchResult struct {
	Receipt       ApprovalDecisionReceipt
	Queue         ApprovalQueue
	Approval      ApprovalRecord
	Effect        EffectAttempt
	ApprovedEntry Entry
	Replayed      bool
}

type CompactionOperation

type CompactionOperation struct {
	ThreadID             string
	RequestID            string
	RequestFingerprint   string
	Source               string
	SourceLeafID         string
	ActivePathHash       string
	SummarySchemaVersion string
	PromptIdentity       string
	RequestPayloadHash   string
	State                CompactionOperationState
	Lease                TurnLease
	ResultEntryID        string
	ErrorCode            string
	ErrorMessage         string
	OutcomeFingerprint   string
	FinishedOwnerID      string
	FinishedGeneration   int64
	CreatedAt            time.Time
	UpdatedAt            time.Time
	FinishedAt           time.Time
}

type CompactionOperationState

type CompactionOperationState string
const (
	CompactionOperationPrepared  CompactionOperationState = "prepared"
	CompactionOperationCompleted CompactionOperationState = "completed"
	CompactionOperationFailed    CompactionOperationState = "failed"
)

type CompleteAutomaticThreadTitleRequest

type CompleteAutomaticThreadTitleRequest struct {
	ThreadID   string
	Generation int64
	Token      string
	Title      string
	Now        time.Time
}

type ContextOptions

type ContextOptions struct{}

type ContextProjection

type ContextProjection struct {
	Messages []session.Message  `json:"messages"`
	Segments []ProjectedSegment `json:"segments,omitempty"`
}

func BuildContextProjection

func BuildContextProjection(path []Entry, opts ContextProjectionOptions) ContextProjection

func BuildContextProjectionChecked

func BuildContextProjectionChecked(path []Entry, opts ContextProjectionOptions) (ContextProjection, error)

type ContextProjectionOptions

type ContextProjectionOptions struct {
	Purpose ProjectionPurpose
}

type CreateRootRequest

type CreateRootRequest struct {
	ThreadID        string
	CreateIntentID  string
	ContractVersion string
	Meta            ThreadMeta
}

type CreateRootResult

type CreateRootResult struct {
	Thread   ThreadMeta
	Replayed bool
}

type DeleteRootTreeResult

type DeleteRootTreeResult struct {
	ThreadIDs []string
	Replayed  bool
}

type EffectAttempt

type EffectAttempt struct {
	EffectAttemptID     string
	Invocation          EffectInvocationIdentity
	RequestFingerprint  string
	State               EffectAttemptState
	RejectionCode       string
	TerminalFingerprint string
	ResultEntryID       string
	OwnerID             string
	Generation          int64
	CreatedAt           time.Time
	UpdatedAt           time.Time
}

type EffectAttemptState

type EffectAttemptState string
const (
	EffectAttemptPrepared    EffectAttemptState = "prepared"
	EffectAttemptDispatching EffectAttemptState = "dispatching"
	EffectAttemptCompleted   EffectAttemptState = "completed"
	EffectAttemptFailed      EffectAttemptState = "failed"
	EffectAttemptRejected    EffectAttemptState = "rejected"
	EffectAttemptUnknown     EffectAttemptState = "unknown"
	EffectAttemptCancelled   EffectAttemptState = "cancelled"
)

type EffectInvocationIdentity

type EffectInvocationIdentity struct {
	ThreadID     string
	TurnID       string
	RunID        string
	ToolCallID   string
	ToolName     string
	ArgumentHash string
}

type Entry

type Entry struct {
	ID                      string              `json:"id"`
	ThreadID                string              `json:"thread_id"`
	ParentID                string              `json:"parent_id,omitempty"`
	PathDepth               int64               `json:"path_depth"`
	Type                    EntryType           `json:"type"`
	TurnID                  string              `json:"turn_id,omitempty"`
	CreatedAt               time.Time           `json:"created_at"`
	Message                 session.Message     `json:"message,omitempty"`
	Raw                     string              `json:"raw,omitempty"`
	RawHash                 string              `json:"raw_hash,omitempty"`
	TurnStatus              TurnMarkerStatus    `json:"turn_status,omitempty"`
	Provider                string              `json:"provider,omitempty"`
	Model                   string              `json:"model,omitempty"`
	CompactionID            string              `json:"compaction_id,omitempty"`
	PreviousCompactionID    string              `json:"previous_compaction_id,omitempty"`
	CompactedThroughEntryID string              `json:"compacted_through_entry_id,omitempty"`
	SummarySchemaVersion    string              `json:"summary_schema_version,omitempty"`
	CompactionGeneration    int                 `json:"compaction_generation,omitempty"`
	CompactionWindowID      string              `json:"compaction_window_id,omitempty"`
	FirstKeptEntryID        string              `json:"first_kept_entry_id,omitempty"`
	KeptUserEntryIDs        []string            `json:"kept_user_entry_ids,omitempty"`
	Summary                 string              `json:"summary,omitempty"`
	CompactionTrigger       string              `json:"compaction_trigger,omitempty"`
	CompactionReason        string              `json:"compaction_reason,omitempty"`
	CompactionPhase         string              `json:"compaction_phase,omitempty"`
	CompactionOperationID   string              `json:"compaction_operation_id,omitempty"`
	CompactionRequestID     string              `json:"compaction_request_id,omitempty"`
	CompactionSource        string              `json:"compaction_source,omitempty"`
	TokensBefore            int64               `json:"tokens_before,omitempty"`
	TokensAfterEstimate     int64               `json:"tokens_after_estimate,omitempty"`
	ContextUsageBefore      contextpolicy.Usage `json:"context_usage_before,omitempty"`
	ContextUsageAfter       contextpolicy.Usage `json:"context_usage_after,omitempty"`
	Error                   string              `json:"error,omitempty"`
	Metadata                map[string]string   `json:"metadata,omitempty"`
}

func AppendActiveTools

func AppendActiveTools(ctx context.Context, repo JournalRepo, threadID string, metadata map[string]string) (Entry, error)

func AppendCompaction

func AppendCompaction(ctx context.Context, repo JournalRepo, threadID, turnID string, result compaction.Result) (Entry, error)

func AppendFailure

func AppendFailure(ctx context.Context, repo JournalRepo, threadID, turnID string, message string) (Entry, error)

func AppendMessage

func AppendMessage(ctx context.Context, repo JournalRepo, threadID, turnID string, msg session.Message) (Entry, error)

func AppendMessageAt

func AppendMessageAt(ctx context.Context, repo JournalRepo, threadID, turnID string, msg session.Message, observedAt time.Time) (Entry, error)

func AppendTurnMarker

func AppendTurnMarker(ctx context.Context, repo JournalRepo, threadID, turnID string, status TurnMarkerStatus, metadata map[string]string) (Entry, error)

func AppendTurnMarkerWithID

func AppendTurnMarkerWithID(ctx context.Context, repo JournalRepo, threadID, turnID, entryID string, status TurnMarkerStatus, metadata map[string]string) (Entry, error)

func CanonicalTurnEntriesForRead

func CanonicalTurnEntriesForRead(entries []Entry) []Entry

CanonicalTurnEntriesForRead excludes retry/fork structural closures when the same turn already has its execution terminal. A copied unfinished turn keeps its branch boundary as the only canonical terminal.

func CompactionEntry

func CompactionEntry(threadID, turnID string, result compaction.Result) (Entry, error)

func PrepareBranchBoundaryEntry

func PrepareBranchBoundaryEntry(path []Entry, threadID, parentEntryID, entryID, reason string, now time.Time) (Entry, error)

PrepareBranchBoundaryEntry closes one copied or rewound unfinished turn so a fork or retry path is idle before another turn authority is admitted.

func PrepareEntry

func PrepareEntry(entry Entry) Entry

func PrepareSubAgentCloseLifecycleEntry

func PrepareSubAgentCloseLifecycleEntry(operation SubAgentCloseOperation, threadID, parentEntryID, entryID string, now time.Time) Entry

PrepareSubAgentCloseLifecycleEntry builds the canonical lifecycle entry that FinishSubAgentClose persists atomically with terminal child state.

func UnresolvedInterruptedTurnCalls

func UnresolvedInterruptedTurnCalls(path []Entry, turnID string) []Entry

type EntryType

type EntryType string
const (
	EntryThreadInfo       EntryType = "thread_info"
	EntryTurnMarker       EntryType = "turn_marker"
	EntryUserMessage      EntryType = "user_message"
	EntryAssistantMessage EntryType = "assistant_message"
	EntryToolCall         EntryType = "tool_call"
	EntryToolResult       EntryType = "tool_result"
	EntryModelChange      EntryType = "model_change"
	EntryActiveTools      EntryType = "active_tools_change"
	EntryCompaction       EntryType = "compaction"
	EntryBranchSummary    EntryType = "branch_summary"
	EntryRunFailure       EntryType = "run_failure"
	EntryCustom           EntryType = "custom"
)

type FailAutomaticThreadTitleRequest

type FailAutomaticThreadTitleRequest struct {
	ThreadID   string
	Generation int64
	Token      string
	Error      string
	Now        time.Time
}

type FileRepo

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

func NewFileRepo

func NewFileRepo(root string) *FileRepo

func (*FileRepo) AcquireTurnLease

func (r *FileRepo) AcquireTurnLease(ctx context.Context, lease TurnLease) error

func (*FileRepo) ActiveTurnLease

func (r *FileRepo) ActiveTurnLease(ctx context.Context, threadID string) (TurnLease, bool, error)

func (*FileRepo) Append

func (r *FileRepo) Append(ctx context.Context, entry Entry, opts AppendOptions) (Entry, error)

func (*FileRepo) ArtifactClosure

func (*FileRepo) BeginAutomaticThreadTitle

func (r *FileRepo) BeginAutomaticThreadTitle(ctx context.Context, req BeginAutomaticThreadTitleRequest) (ThreadTitleMutationResult, error)

func (*FileRepo) CanonicalTurnEntries

func (r *FileRepo) CanonicalTurnEntries(ctx context.Context, threadID, turnID, runID string) ([]Entry, bool, error)

func (*FileRepo) ClearExpiredTurnLease

func (r *FileRepo) ClearExpiredTurnLease(ctx context.Context, threadID string, cutoff time.Time) (TurnLease, bool, error)

func (*FileRepo) CompareAndSwapAgentTodoState

func (r *FileRepo) CompareAndSwapAgentTodoState(ctx context.Context, state AgentTodoState, expectedVersion int64) (AgentTodoState, error)

func (*FileRepo) CompleteAutomaticThreadTitle

func (r *FileRepo) CompleteAutomaticThreadTitle(ctx context.Context, req CompleteAutomaticThreadTitleRequest) (ThreadTitleMutationResult, error)

func (*FileRepo) CreateThread

func (r *FileRepo) CreateThread(ctx context.Context, meta ThreadMeta) (ThreadMeta, error)

func (*FileRepo) Entries

func (r *FileRepo) Entries(ctx context.Context, threadID string) ([]Entry, error)

func (*FileRepo) Entry

func (r *FileRepo) Entry(ctx context.Context, threadID, entryID string) (Entry, error)

func (*FileRepo) FailAutomaticThreadTitle

func (*FileRepo) Fork

func (r *FileRepo) Fork(ctx context.Context, opts ForkOptions) (ThreadMeta, error)

func (*FileRepo) ListCanonicalTurns

func (r *FileRepo) ListCanonicalTurns(ctx context.Context, opts ListCanonicalTurnsOptions) (CanonicalTurnsPage, error)

func (*FileRepo) ListThreads

func (r *FileRepo) ListThreads(ctx context.Context, opts ListThreadsOptions) ([]ThreadMeta, error)

func (*FileRepo) MoveLeaf

func (r *FileRepo) MoveLeaf(ctx context.Context, threadID, entryID string) error

func (*FileRepo) Path

func (r *FileRepo) Path(ctx context.Context, threadID, leafID string) ([]Entry, error)

func (*FileRepo) PathPage

func (r *FileRepo) PathPage(ctx context.Context, threadID, leafID, beforeEntryID string, limit int) (PathPage, error)

func (*FileRepo) PendingAutomaticThreadTitles

func (r *FileRepo) PendingAutomaticThreadTitles(ctx context.Context) ([]ThreadMeta, error)

func (*FileRepo) ReadAgentTodoState

func (r *FileRepo) ReadAgentTodoState(ctx context.Context, threadID string) (AgentTodoState, error)

func (*FileRepo) ReadArtifact

func (*FileRepo) ReleaseTurnLease

func (r *FileRepo) ReleaseTurnLease(ctx context.Context, lease TurnLease) error

func (*FileRepo) SetThreadTitle

func (*FileRepo) Thread

func (r *FileRepo) Thread(ctx context.Context, threadID string) (ThreadMeta, error)

func (*FileRepo) UpdateThread

func (r *FileRepo) UpdateThread(ctx context.Context, meta ThreadMeta) error

type FinalizeApprovalRequest

type FinalizeApprovalRequest struct {
	ResolutionID             string
	ExpectedRootThreadID     string
	ExpectedGeneration       int64
	ExpectedCurrent          ApprovalIdentity
	ExpectedApprovalRevision int64
	State                    ApprovalState
	Reason                   string
	FinalizedEntry           Entry
	Now                      time.Time
}

type FinalizeApprovalResult

type FinalizeApprovalResult struct {
	Receipt        ApprovalDecisionReceipt
	Queue          ApprovalQueue
	Approval       ApprovalRecord
	Effect         EffectAttempt
	FinalizedEntry Entry
	Replayed       bool
}

type FinishCompactionRequest

type FinishCompactionRequest struct {
	Lease              TurnLease
	RequestID          string
	RequestFingerprint string
	OutcomeFingerprint string
	Result             *Entry
	ErrorCode          string
	ErrorMessage       string
	Now                time.Time
}

type FinishCompactionResult

type FinishCompactionResult struct {
	Operation CompactionOperation
	Entry     *Entry
	Replayed  bool
}

type FinishEffectDispatchRequest

type FinishEffectDispatchRequest struct {
	Lease              TurnLease
	EffectAttemptID    string
	RequestFingerprint string
	OutcomeFingerprint string
	Failed             bool
	Result             Entry
	FullOutput         *artifact.FullOutput
	Now                time.Time
}

type FinishEffectDispatchResult

type FinishEffectDispatchResult struct {
	Attempt  EffectAttempt
	Result   Entry
	Artifact *artifact.Ref
	Replayed bool
}

type FinishSubAgentCloseRequest

type FinishSubAgentCloseRequest struct {
	CloseOperationID string
	ParentThreadID   string
	TargetThreadID   string
	Reason           string
	Now              time.Time
}

type FinishSubAgentCloseResult

type FinishSubAgentCloseResult struct {
	Operation         SubAgentCloseOperation
	Threads           []ThreadMeta
	Entries           []Entry
	CancelledInputIDs []string
	Replayed          bool
}

type FinishTurnRequest

type FinishTurnRequest struct {
	Lease              TurnLease
	RunID              string
	TerminalEntryID    string
	Status             TurnMarkerStatus
	Metadata           map[string]string
	FailureMessage     string
	ProviderState      *ProviderStateRecord
	ClearProviderState bool
	OutcomeFingerprint string
	Now                time.Time
}

type FinishTurnResult

type FinishTurnResult struct {
	Failure  *Entry
	Terminal Entry
	Replayed bool
}

type ForkDestinationMeta

type ForkDestinationMeta struct {
	ParentThreadID  string          `json:"parent_thread_id"`
	ParentTurnID    string          `json:"parent_turn_id,omitempty"`
	TaskName        string          `json:"task_name,omitempty"`
	TaskDescription string          `json:"task_description,omitempty"`
	AgentPath       string          `json:"agent_path,omitempty"`
	HostProfileRef  string          `json:"host_profile_ref,omitempty"`
	ForkMode        string          `json:"fork_mode,omitempty"`
	Lifecycle       ThreadLifecycle `json:"lifecycle,omitempty"`
}

ForkDestinationMeta is the child ownership metadata written atomically with a fork destination. A nil value creates an independent root fork.

type ForkEntryIdentity

type ForkEntryIdentity struct {
	SourceThreadID      string
	DestinationThreadID string
	TurnIDMap           map[string]string
	RunIDMap            map[string]string
}

type ForkOptions

type ForkOptions struct {
	SourceThreadID       string
	EntryID              string
	EntryIDPinned        bool
	ExpectedSourceLeafID string
	Position             ForkPosition
	NewThreadID          string
	OperationID          string
	OperationNodeID      string
	Now                  time.Time
	TurnIDMap            map[string]string
	RunIDMap             map[string]string
	DestinationMeta      *ForkDestinationMeta
	ArtifactClosure      artifact.Closure
	RewriteEntry         func(Entry, ForkEntryIdentity) (Entry, error)
}

type ForkPosition

type ForkPosition string
const (
	ForkAt     ForkPosition = "at"
	ForkBefore ForkPosition = "before"
)

type ForkPrepareThreadState

type ForkPrepareThreadState struct {
	Meta              ThreadMeta
	Path              []Entry
	PinnedPath        []Entry
	PendingInputCount int
}

ForkPrepareThreadState is the transaction-local source state used to validate one complete replayable fork plan before any structural claim is published.

type InterruptedTurnApprovalQueueProof

type InterruptedTurnApprovalQueueProof struct {
	RootThreadID string
	Generation   int64
	Revision     int64
}

func InterruptedTurnApprovalQueueProofFromMetadata

func InterruptedTurnApprovalQueueProofFromMetadata(metadata map[string]string) (*InterruptedTurnApprovalQueueProof, error)

type InterruptedTurnRecoveryEffect

type InterruptedTurnRecoveryEffect struct {
	EffectAttemptID string             `json:"effect_attempt_id"`
	ToolCallID      string             `json:"tool_call_id"`
	State           EffectAttemptState `json:"state"`
}

func InterruptedTurnRecoveryEffects

func InterruptedTurnRecoveryEffects(attempts []EffectAttempt, committedFingerprint string) ([]InterruptedTurnRecoveryEffect, error)

type InterruptedTurnRecoveryFailureProof

type InterruptedTurnRecoveryFailureProof struct {
	EntryID string `json:"entry_id"`
	Message string `json:"message"`
	RawHash string `json:"raw_hash"`
}

func InterruptedTurnRecoverySourceFailureProofFromMetadata

func InterruptedTurnRecoverySourceFailureProofFromMetadata(metadata map[string]string) (*InterruptedTurnRecoveryFailureProof, error)

type InterruptedTurnRecoveryPlan

type InterruptedTurnRecoveryPlan struct {
	RunID              string
	Status             TurnMarkerStatus
	FailureCode        string
	FailureMessage     string
	SourceFailure      *InterruptedTurnRecoveryFailureProof
	OutcomeFingerprint string
	TerminalEntryID    string
	Effects            []InterruptedTurnRecoveryEffect
}

func DeriveInterruptedTurnRecoveryPlan

func DeriveInterruptedTurnRecoveryPlan(path []Entry, expectedLease TurnLease, parentThreadID string, effects []InterruptedTurnRecoveryEffect) (InterruptedTurnRecoveryPlan, error)

type InterruptedTurnRecoveryRepo

type InterruptedTurnRecoveryRepo interface {
	RecoverInterruptedTurn(context.Context, RecoverInterruptedTurnRequest) (RecoverInterruptedTurnResult, error)
}

type InterruptedTurnResolutionValidationRepo

type InterruptedTurnResolutionValidationRepo interface {
	ValidateInterruptedTurnResolution(context.Context, RecoverInterruptedTurnRequest) error
}

type JournalRepo

JournalRepo is the durable journal capability used by normal Agent execution. It intentionally excludes lifecycle creation, deletion, fork, and metadata replacement capabilities.

type LeasePolicy

type LeasePolicy struct {
	TTL                time.Duration
	RenewInterval      time.Duration
	ClockSkewAllowance time.Duration
}

func (LeasePolicy) Validate

func (p LeasePolicy) Validate() error

type LeasePolicyRepo

type LeasePolicyRepo interface {
	AuthorityLeasePolicy() LeasePolicy
}

type ListCanonicalTurnsOptions

type ListCanonicalTurnsOptions struct {
	ThreadID     string
	BeforeCursor *CanonicalTurnBeforeCursor
	SinceCursor  *CanonicalTurnSinceCursor
	Tail         int
	Limit        int
}

type ListThreadsOptions

type ListThreadsOptions struct {
	IncludeArchived bool
	RootOnly        bool
	Limit           int
	AfterCreatedAt  time.Time
	AfterID         string
}

type MarkEffectUnknownRequest

type MarkEffectUnknownRequest struct {
	Lease              TurnLease
	EffectAttemptID    string
	RequestFingerprint string
	OutcomeFingerprint string
	Now                time.Time
}

type MemoryRepo

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

func DecodeMemoryState

func DecodeMemoryState(data []byte, now func() time.Time) (*MemoryRepo, error)

DecodeMemoryState constructs a repo from one exact encoded state.

func NewMemoryRepo

func NewMemoryRepo() *MemoryRepo

func NewMemoryRepoWithLeasePolicy

func NewMemoryRepoWithLeasePolicy(policy LeasePolicy, now func() time.Time) (*MemoryRepo, error)

func (*MemoryRepo) AcquireThreadAuthorityClaim

func (r *MemoryRepo) AcquireThreadAuthorityClaim(_ context.Context, operationID string, requiredSourceThreadIDs, authorityThreadIDs []string) error

AcquireThreadAuthorityClaim reserves identities for one replayable structural operation. Required source threads must exist and have no active turn lease.

func (*MemoryRepo) AcquireTurnLease

func (r *MemoryRepo) AcquireTurnLease(_ context.Context, request TurnLease) (TurnLease, error)

func (*MemoryRepo) ActiveTurnLease

func (r *MemoryRepo) ActiveTurnLease(_ context.Context, threadID string) (TurnLease, bool, error)

func (*MemoryRepo) AdmitSubAgentInput

func (*MemoryRepo) AdmitTurn

func (*MemoryRepo) Append

func (r *MemoryRepo) Append(ctx context.Context, entry Entry, opts AppendOptions) (Entry, error)

func (*MemoryRepo) Approval

func (r *MemoryRepo) Approval(_ context.Context, approvalID string) (ApprovalRecord, error)

func (*MemoryRepo) ArtifactClosure

func (r *MemoryRepo) ArtifactClosure(_ context.Context, req ArtifactClosureRequest) (artifact.Closure, error)

func (*MemoryRepo) AuthorityLeasePolicy

func (r *MemoryRepo) AuthorityLeasePolicy() LeasePolicy

func (*MemoryRepo) BeginAutomaticThreadTitle

func (*MemoryRepo) BeginCompaction

func (*MemoryRepo) BeginEffectDispatch

func (r *MemoryRepo) BeginEffectDispatch(_ context.Context, req BeginEffectDispatchRequest) (EffectAttempt, error)

func (*MemoryRepo) CancelApprovalBatch

func (*MemoryRepo) CanonicalTurnEntries

func (r *MemoryRepo) CanonicalTurnEntries(_ context.Context, threadID, turnID, runID string) ([]Entry, bool, error)

func (*MemoryRepo) CommitApprovalDispatch

func (*MemoryRepo) CommitForkBatch

func (r *MemoryRepo) CommitForkBatch(ctx context.Context, operationID string, nodes []ForkOptions, commit func() error) ([]ThreadMeta, error)

CommitForkBatch publishes every destination and releases the operation's complete claim set in one MemoryRepo critical section. The callback persists the terminal operation record before readers can observe the destinations.

func (*MemoryRepo) CompareAndSwapAgentTodoState

func (r *MemoryRepo) CompareAndSwapAgentTodoState(ctx context.Context, state AgentTodoState, expectedVersion int64) (AgentTodoState, error)

func (*MemoryRepo) CompleteAutomaticThreadTitle

func (*MemoryRepo) CreateRoot

func (*MemoryRepo) CreateThread

func (r *MemoryRepo) CreateThread(_ context.Context, meta ThreadMeta) (ThreadMeta, error)

func (*MemoryRepo) CreateThreadWithInitialEntry

func (r *MemoryRepo) CreateThreadWithInitialEntry(ctx context.Context, meta ThreadMeta, initial Entry) (ThreadMeta, Entry, error)

func (*MemoryRepo) DeleteProviderState

func (r *MemoryRepo) DeleteProviderState(ctx context.Context, threadID string) error

func (*MemoryRepo) DeleteRootTree

func (r *MemoryRepo) DeleteRootTree(_ context.Context, rootThreadID string) (DeleteRootTreeResult, error)

func (*MemoryRepo) EncodeMemoryState

func (repo *MemoryRepo) EncodeMemoryState() ([]byte, error)

EncodeMemoryState returns a detached strict representation of the complete session-tree domain state. The caller owns the returned bytes.

func (*MemoryRepo) Entries

func (r *MemoryRepo) Entries(_ context.Context, threadID string) ([]Entry, error)

func (*MemoryRepo) Entry

func (r *MemoryRepo) Entry(_ context.Context, threadID, entryID string) (Entry, error)

func (*MemoryRepo) FailAutomaticThreadTitle

func (*MemoryRepo) FailForkClaim

func (r *MemoryRepo) FailForkClaim(operationID string, sourceThreadIDs, authorityThreadIDs []string, commit func() error) error

FailForkClaim records one deterministic pre-publication failure and releases the complete claim set without exposing an unclaimed prepared operation.

func (*MemoryRepo) FinalizeApproval

func (*MemoryRepo) FinishCompaction

func (*MemoryRepo) FinishEffectDispatch

func (*MemoryRepo) FinishSubAgentClose

func (*MemoryRepo) FinishTurn

func (*MemoryRepo) Fork

func (r *MemoryRepo) Fork(ctx context.Context, opts ForkOptions) (ThreadMeta, error)

func (*MemoryRepo) ForkWithInitialEntry

func (r *MemoryRepo) ForkWithInitialEntry(ctx context.Context, opts ForkOptions, initial Entry) (ThreadMeta, Entry, error)

func (*MemoryRepo) InspectSubAgentThreadAuthority

func (r *MemoryRepo) InspectSubAgentThreadAuthority(_ context.Context, parentThreadID, childThreadID string) (SubAgentThreadAuthoritySnapshot, error)

func (*MemoryRepo) InspectThreadAuthority

func (r *MemoryRepo) InspectThreadAuthority(_ context.Context, threadID string) (ThreadAuthoritySnapshot, error)

func (*MemoryRepo) ListCanonicalTurns

func (r *MemoryRepo) ListCanonicalTurns(_ context.Context, opts ListCanonicalTurnsOptions) (CanonicalTurnsPage, error)

func (*MemoryRepo) ListSubAgentInputs

func (r *MemoryRepo) ListSubAgentInputs(_ context.Context, childThreadID string, state SubAgentInputState) ([]SubAgentInputRecord, error)

func (*MemoryRepo) ListThreads

func (r *MemoryRepo) ListThreads(_ context.Context, opts ListThreadsOptions) ([]ThreadMeta, error)

func (*MemoryRepo) MarkEffectUnknown

func (r *MemoryRepo) MarkEffectUnknown(_ context.Context, req MarkEffectUnknownRequest) (EffectAttempt, error)

func (*MemoryRepo) MoveLeaf

func (r *MemoryRepo) MoveLeaf(_ context.Context, threadID, entryID string) error

func (*MemoryRepo) Path

func (r *MemoryRepo) Path(_ context.Context, threadID, leafID string) ([]Entry, error)

func (*MemoryRepo) PathPage

func (r *MemoryRepo) PathPage(_ context.Context, threadID, leafID, beforeEntryID string, limit int) (PathPage, error)

func (*MemoryRepo) PendingAutomaticThreadTitles

func (r *MemoryRepo) PendingAutomaticThreadTitles(_ context.Context) ([]ThreadMeta, error)

func (*MemoryRepo) PrepareApprovalBatch

func (*MemoryRepo) PrepareEffectAttempt

func (*MemoryRepo) PrepareForkClaim

func (r *MemoryRepo) PrepareForkClaim(ctx context.Context, operationID, rootThreadID string, nodes []ForkOptions) error

PrepareForkClaim validates and publishes one complete Memory authority claim in the same critical section.

func (*MemoryRepo) PrepareSubAgentClose

func (*MemoryRepo) ProviderState

func (r *MemoryRepo) ProviderState(_ context.Context, threadID string) (ProviderStateRecord, error)

func (*MemoryRepo) PublishSubAgent

func (*MemoryRepo) PublishSubAgentInput

func (*MemoryRepo) PutProviderState

func (r *MemoryRepo) PutProviderState(ctx context.Context, record ProviderStateRecord) error

func (*MemoryRepo) ReadAgentTodoState

func (r *MemoryRepo) ReadAgentTodoState(_ context.Context, threadID string) (AgentTodoState, error)

func (*MemoryRepo) ReadApprovalQueue

func (r *MemoryRepo) ReadApprovalQueue(_ context.Context, threadID string) (ApprovalQueue, error)

func (*MemoryRepo) ReadArtifact

func (*MemoryRepo) ReadCanonicalTurn

func (r *MemoryRepo) ReadCanonicalTurn(_ context.Context, threadID, turnID string) (CanonicalTurnRead, error)

func (*MemoryRepo) ReadCompaction

func (r *MemoryRepo) ReadCompaction(_ context.Context, threadID, requestID string) (CompactionOperation, bool, error)

func (*MemoryRepo) ReadSubAgentInput

func (r *MemoryRepo) ReadSubAgentInput(_ context.Context, inputID string) (SubAgentInputRecord, bool, error)

func (*MemoryRepo) ReadTurnAdmission

func (r *MemoryRepo) ReadTurnAdmission(_ context.Context, threadID, turnID, runID string) (AdmitTurnResult, bool, error)

func (*MemoryRepo) RejectEffectAttempt

func (r *MemoryRepo) RejectEffectAttempt(_ context.Context, req RejectEffectAttemptRequest) (EffectAttempt, error)

func (*MemoryRepo) ReleaseThreadAuthorityClaim

func (r *MemoryRepo) ReleaseThreadAuthorityClaim(_ context.Context, operationID string)

ReleaseThreadAuthorityClaim releases every identity held by one operation.

func (*MemoryRepo) ReleaseTurnLease

func (r *MemoryRepo) ReleaseTurnLease(_ context.Context, proof TurnLease) error

func (*MemoryRepo) RenewTurnLease

func (r *MemoryRepo) RenewTurnLease(_ context.Context, proof TurnLease) (TurnLease, error)

func (*MemoryRepo) ResolveApproval

func (*MemoryRepo) SetThreadTitle

func (*MemoryRepo) TakeOverCompaction

func (*MemoryRepo) Thread

func (r *MemoryRepo) Thread(_ context.Context, threadID string) (ThreadMeta, error)

func (*MemoryRepo) ThreadTombstone

func (r *MemoryRepo) ThreadTombstone(_ context.Context, threadID string) (ThreadTombstone, error)

func (*MemoryRepo) UpdateThread

func (r *MemoryRepo) UpdateThread(ctx context.Context, meta ThreadMeta) error

func (*MemoryRepo) ValidateArtifactForkDestination

func (r *MemoryRepo) ValidateArtifactForkDestination(_ context.Context, closure artifact.Closure) error

ValidateArtifactForkDestination verifies the complete copied artifact set without relying on the source thread still being live.

func (*MemoryRepo) ValidateInterruptedTurnResolution

func (r *MemoryRepo) ValidateInterruptedTurnResolution(_ context.Context, req RecoverInterruptedTurnRequest) error

func (*MemoryRepo) WaitApprovalDecision

func (r *MemoryRepo) WaitApprovalDecision(ctx context.Context, approvalID string) (WaitApprovalDecisionResult, error)

type PathPage

type PathPage struct {
	Entries     []Entry
	NextEntryID string
	HasMore     bool
	// NewestOrdinal is the active-path ordinal of Entries[0]. Entries are
	// returned newest first, so later entries decrement this value by one.
	NewestOrdinal int64
}

type PendingToolRecoveryRepo

type PendingToolRecoveryRepo interface {
	SettlePendingToolRecovery(context.Context, SettlePendingToolRecoveryRequest) (SettlePendingToolRecoveryResult, error)
}

type PendingToolSettlementTarget

type PendingToolSettlementTarget struct {
	ThreadID        string
	TurnID          string
	RunID           string
	ToolCallID      string
	ToolName        string
	Handle          string
	EffectAttemptID string
}

type PrepareApprovalBatchRequest

type PrepareApprovalBatchRequest struct {
	Lease TurnLease
	Items []ApprovalPreflightItem
	Now   time.Time
}

type PrepareApprovalBatchResult

type PrepareApprovalBatchResult struct {
	Queue            ApprovalQueue
	Effects          []EffectAttempt
	Approvals        []ApprovalRecord
	RequestedEntries []Entry
	Replayed         bool
}

type PrepareEffectAttemptRequest

type PrepareEffectAttemptRequest struct {
	Lease              TurnLease
	Invocation         EffectInvocationIdentity
	RequestFingerprint string
	Now                time.Time
}

type PrepareEffectAttemptResult

type PrepareEffectAttemptResult struct {
	Attempt  EffectAttempt
	Replayed bool
}

type PrepareSubAgentCloseRequest

type PrepareSubAgentCloseRequest struct {
	CloseOperationID string
	ParentThreadID   string
	TargetThreadID   string
	Reason           string
	TargetLease      *TurnLease
	Now              time.Time
}

type PrepareSubAgentCloseResult

type PrepareSubAgentCloseResult struct {
	Operation SubAgentCloseOperation
	Replayed  bool
}

type ProjectedSegment

type ProjectedSegment struct {
	EntryID       string         `json:"entry_id,omitempty"`
	EntryType     EntryType      `json:"entry_type,omitempty"`
	MessageIndex  int            `json:"message_index"`
	Role          session.Role   `json:"role,omitempty"`
	ToolCallID    string         `json:"tool_call_id,omitempty"`
	ToolName      string         `json:"tool_name,omitempty"`
	TokenEstimate int64          `json:"token_estimate,omitempty"`
	ArtifactRefs  []artifact.Ref `json:"artifact_refs,omitempty"`
	UIPreview     string         `json:"ui_preview,omitempty"`
}

type ProjectionPurpose

type ProjectionPurpose string
const (
	ProjectionProviderRequest ProjectionPurpose = "provider_request"
	ProjectionCompaction      ProjectionPurpose = "compaction"
	ProjectionTestUI          ProjectionPurpose = "test_ui"
)

type ProviderStateReader

type ProviderStateReader interface {
	ProviderState(context.Context, string) (ProviderStateRecord, error)
}

type ProviderStateRecord

type ProviderStateRecord struct {
	ThreadID         string
	LeafEntryID      string
	CompatibilityKey string
	State            provider.State
	CreatedByRunID   string
	CreatedByTurnID  string
	UpdatedAt        time.Time
}

type ProviderStateStore

type ProviderStateStore interface {
	ProviderStateReader
	PutProviderState(context.Context, ProviderStateRecord) error
	DeleteProviderState(context.Context, string) error
}

type PublishSubAgentInputRequest

type PublishSubAgentInputRequest struct {
	InputRequestID     string
	RequestFingerprint string
	ParentThreadID     string
	ChildThreadID      string
	Message            session.Message
	HostLabels         map[string]string
	CorrelationLabels  map[string]string
	Interrupt          bool
	Now                time.Time
}

type PublishSubAgentPendingToolCompletionRequest

type PublishSubAgentPendingToolCompletionRequest struct {
	InputRequestID        string
	RequestFingerprint    string
	SettlementFingerprint string
	ParentThreadID        string
	ChildThreadID         string
	Target                PendingToolSettlementTarget
	Settlement            Entry
	Message               session.Message
	HostLabels            map[string]string
	CorrelationLabels     map[string]string
	Now                   time.Time
}

type PublishSubAgentPendingToolCompletionResult

type PublishSubAgentPendingToolCompletionResult struct {
	Settlement         Entry
	SettlementReplayed bool
	Input              SubAgentInputRecord
	Replayed           bool
}

type PublishSubAgentRequest

type PublishSubAgentRequest struct {
	PublicationID      string
	RequestFingerprint string
	ParentThreadID     string
	ChildMeta          ThreadMeta
	ForkOptions        *ForkOptions
	ArtifactClosure    artifact.Closure
	Message            session.Message
	HostLabels         map[string]string
	CorrelationLabels  map[string]string
	Now                time.Time
}

type PublishSubAgentResult

type PublishSubAgentResult struct {
	Thread   ThreadMeta
	Input    SubAgentInputRecord
	Replayed bool
}

type RecoverInterruptedTurnRequest

type RecoverInterruptedTurnRequest struct {
	ExpectedLease  TurnLease
	ParentThreadID string
	Now            time.Time
}

type RecoverInterruptedTurnResult

type RecoverInterruptedTurnResult struct {
	RunID              string
	Status             TurnMarkerStatus
	OutcomeFingerprint string
	Failure            *Entry
	ToolResults        []Entry
	Terminal           Entry
	Generation         int64
	Replayed           bool
}

type RejectEffectAttemptRequest

type RejectEffectAttemptRequest struct {
	Lease                TurnLease
	EffectAttemptID      string
	RequestFingerprint   string
	RejectionCode        string
	RejectionFingerprint string
	Now                  time.Time
}

type Repo

type Repo interface {
	JournalRepo
	CreateThread(context.Context, ThreadMeta) (ThreadMeta, error)
	UpdateThread(context.Context, ThreadMeta) error
	MoveLeaf(context.Context, string, string) error
	Fork(context.Context, ForkOptions) (ThreadMeta, error)
}

Repo is the internal storage implementation contract. Production runtime actors receive narrower capabilities such as JournalRepo instead.

type ResolveApprovalRequest

type ResolveApprovalRequest struct {
	DecisionID               string
	ExpectedRootThreadID     string
	ExpectedGeneration       int64
	ExpectedRevision         int64
	ExpectedCurrent          ApprovalIdentity
	ExpectedApprovalRevision int64
	Decision                 ApprovalDecision
	RejectedEntry            Entry
	Now                      time.Time
}

type ResolveApprovalResult

type ResolveApprovalResult struct {
	Receipt       ApprovalDecisionReceipt
	Queue         ApprovalQueue
	Approval      ApprovalRecord
	Effect        EffectAttempt
	RejectedEntry Entry
	Replayed      bool
}

type RootAuthorityRepo

type RootAuthorityRepo interface {
	ThreadTombstoneRepo
	CreateRoot(context.Context, CreateRootRequest) (CreateRootResult, error)
	DeleteRootTree(context.Context, string) (DeleteRootTreeResult, error)
}

type SetThreadTitleRequest

type SetThreadTitleRequest struct {
	ThreadID string
	Title    string
	Now      time.Time
}

type SettlePendingToolRecoveryRequest

type SettlePendingToolRecoveryRequest struct {
	Target             PendingToolSettlementTarget
	RequestFingerprint string
	Settlement         Entry
	Now                time.Time
}

type SettlePendingToolRecoveryResult

type SettlePendingToolRecoveryResult struct {
	Entry    Entry
	Replayed bool
}

type SubAgentCloseAuthorityRepo

type SubAgentCloseAuthorityRepo interface {
	PrepareSubAgentClose(context.Context, PrepareSubAgentCloseRequest) (PrepareSubAgentCloseResult, error)
	FinishSubAgentClose(context.Context, FinishSubAgentCloseRequest) (FinishSubAgentCloseResult, error)
}

type SubAgentCloseNode

type SubAgentCloseNode struct {
	ThreadID string `json:"thread_id"`
	WasOpen  bool   `json:"was_open"`
}

type SubAgentCloseOperation

type SubAgentCloseOperation struct {
	CloseOperationID   string
	ParentThreadID     string
	TargetThreadID     string
	Reason             string
	IntentFingerprint  string
	RequestFingerprint string
	State              SubAgentCloseState
	Nodes              []SubAgentCloseNode
	ResultEntryIDs     []string
	PreparedAt         time.Time
	FinishedAt         time.Time
}

type SubAgentCloseState

type SubAgentCloseState string
const (
	SubAgentClosePrepared  SubAgentCloseState = "prepared"
	SubAgentCloseCompleted SubAgentCloseState = "completed"
)

type SubAgentInputReadRepo

type SubAgentInputReadRepo interface {
	ReadSubAgentInput(context.Context, string) (SubAgentInputRecord, bool, error)
}

type SubAgentInputRecord

type SubAgentInputRecord struct {
	SubAgentInputID    string
	ParentThreadID     string
	ChildThreadID      string
	RequestKind        SubAgentRequestKind
	RequestID          string
	RequestFingerprint string
	Sequence           int64
	State              SubAgentInputState
	Message            session.Message
	HostLabels         map[string]string
	CorrelationLabels  map[string]string
	AdmittedTurnID     string
	AdmittedRunID      string
	CreatedAt          time.Time
	AdmittedAt         time.Time
	CancelledAt        time.Time
}

type SubAgentInputState

type SubAgentInputState string
const (
	SubAgentInputPending   SubAgentInputState = "pending"
	SubAgentInputAdmitted  SubAgentInputState = "admitted"
	SubAgentInputCancelled SubAgentInputState = "cancelled"
)

type SubAgentRequestKind

type SubAgentRequestKind string
const (
	SubAgentRequestPublication           SubAgentRequestKind = "publication"
	SubAgentRequestInput                 SubAgentRequestKind = "input"
	SubAgentRequestPendingToolCompletion SubAgentRequestKind = "pending_tool_completion"
)

type SubAgentThreadAuthorityInspectionRepo

type SubAgentThreadAuthorityInspectionRepo interface {
	InspectSubAgentThreadAuthority(context.Context, string, string) (SubAgentThreadAuthoritySnapshot, error)
}

type SubAgentThreadAuthoritySnapshot

type SubAgentThreadAuthoritySnapshot struct {
	Parent ThreadMeta
	Child  ThreadAuthoritySnapshot
}

type TakeOverCompactionRequest

type TakeOverCompactionRequest struct {
	ThreadID           string
	RequestID          string
	RequestFingerprint string
	ExpectedLease      TurnLease
	OwnerID            string
	Now                time.Time
}

type ThreadAuthorityInspectionRepo

type ThreadAuthorityInspectionRepo interface {
	InspectThreadAuthority(context.Context, string) (ThreadAuthoritySnapshot, error)
}

type ThreadAuthoritySnapshot

type ThreadAuthoritySnapshot struct {
	Thread           ThreadMeta
	Lease            *TurnLease
	ClaimOperationID string
	LeaseGeneration  int64
}

type ThreadLifecycle

type ThreadLifecycle string

ThreadLifecycle is the durable canonical lifecycle of a thread identity. A deleted identity is represented by a tombstone rather than by absence.

const (
	ThreadLifecycleOpen    ThreadLifecycle = "open"
	ThreadLifecycleClosing ThreadLifecycle = "closing"
	ThreadLifecycleClosed  ThreadLifecycle = "closed"
	ThreadLifecycleDeleted ThreadLifecycle = "deleted"
)

func (ThreadLifecycle) Valid

func (l ThreadLifecycle) Valid() bool

type ThreadListRepo

type ThreadListRepo interface {
	ListThreads(context.Context, ListThreadsOptions) ([]ThreadMeta, error)
}

type ThreadMeta

type ThreadMeta struct {
	ID                  string            `json:"id"`
	LeafID              string            `json:"leaf_id,omitempty"`
	ParentThreadID      string            `json:"parent_thread_id,omitempty"`
	ParentTurnID        string            `json:"parent_turn_id,omitempty"`
	ForkedFromThreadID  string            `json:"forked_from_thread_id,omitempty"`
	ForkedFromEntryID   string            `json:"forked_from_entry_id,omitempty"`
	ForkOperationID     string            `json:"fork_operation_id,omitempty"`
	ForkOperationNodeID string            `json:"fork_operation_node_id,omitempty"`
	TaskName            string            `json:"task_name,omitempty"`
	TaskDescription     string            `json:"task_description,omitempty"`
	AgentPath           string            `json:"agent_path,omitempty"`
	HostProfileRef      string            `json:"host_profile_ref,omitempty"`
	ForkMode            string            `json:"fork_mode,omitempty"`
	Lifecycle           ThreadLifecycle   `json:"lifecycle,omitempty"`
	CloseOperationID    string            `json:"close_operation_id,omitempty"`
	Archived            bool              `json:"archived,omitempty"`
	Title               string            `json:"title,omitempty"`
	TitleStatus         ThreadTitleStatus `json:"title_status,omitempty"`
	TitleSource         ThreadTitleSource `json:"title_source,omitempty"`
	TitleUpdatedAt      time.Time         `json:"title_updated_at,omitempty"`
	TitleError          string            `json:"title_error,omitempty"`
	TitleGeneration     int64             `json:"title_generation,omitempty"`
	TitleToken          string            `json:"title_token,omitempty"`
	CreatedAt           time.Time         `json:"created_at"`
	UpdatedAt           time.Time         `json:"updated_at"`
	LastViewedAt        time.Time         `json:"last_viewed_at,omitempty"`
}

func ApplyThreadListOptions

func ApplyThreadListOptions(threads []ThreadMeta, opts ListThreadsOptions) []ThreadMeta

func ListThreads

func ListThreads(ctx context.Context, repo JournalRepo, opts ListThreadsOptions) ([]ThreadMeta, error)

func (ThreadMeta) CanonicalLifecycle

func (m ThreadMeta) CanonicalLifecycle() (ThreadLifecycle, error)

func (ThreadMeta) IsClosed

func (m ThreadMeta) IsClosed() bool

func (ThreadMeta) IsClosing

func (m ThreadMeta) IsClosing() bool

type ThreadPublishRepo

type ThreadPublishRepo interface {
	CreateThreadWithInitialEntry(context.Context, ThreadMeta, Entry) (ThreadMeta, Entry, error)
	ForkWithInitialEntry(context.Context, ForkOptions, Entry) (ThreadMeta, Entry, error)
}

type ThreadTitleMutationResult

type ThreadTitleMutationResult struct {
	Thread  ThreadMeta
	Changed bool
}

type ThreadTitleProjection

type ThreadTitleProjection struct {
	Title      string
	Status     ThreadTitleStatus
	Source     ThreadTitleSource
	UpdatedAt  time.Time
	Error      string
	Generation int64
}

ThreadTitleProjection is the observable subset of canonical title authority.

type ThreadTitleSource

type ThreadTitleSource string
const (
	ThreadTitleSourceProvider ThreadTitleSource = "provider"
	ThreadTitleSourceHost     ThreadTitleSource = "host"
)

type ThreadTitleStatus

type ThreadTitleStatus string
const (
	ThreadTitlePending ThreadTitleStatus = "pending"
	ThreadTitleReady   ThreadTitleStatus = "ready"
	ThreadTitleFailed  ThreadTitleStatus = "failed"
)

type ThreadTombstone

type ThreadTombstone struct {
	ThreadID            string
	RootThreadID        string
	ParentThreadID      string
	CreateIntentID      string
	ForkOperationID     string
	ForkOperationNodeID string
	ForkedFromThreadID  string
	ForkedFromEntryID   string
	DeletedAt           time.Time
}

ThreadTombstone retains identity provenance after queryable Agent state is deleted. It is intentionally not a ThreadMeta and is never returned as a normal thread read.

type ThreadTombstoneRepo

type ThreadTombstoneRepo interface {
	ThreadTombstone(context.Context, string) (ThreadTombstone, error)
}

type TurnAuthorityRepo

type TurnAuthorityRepo interface {
	AdmitTurn(context.Context, AdmitTurnRequest) (AdmitTurnResult, error)
	ReadTurnAdmission(context.Context, string, string, string) (AdmitTurnResult, bool, error)
	FinishTurn(context.Context, FinishTurnRequest) (FinishTurnResult, error)
}

type TurnLease

type TurnLease struct {
	ThreadID     string           `json:"thread_id"`
	Purpose      TurnLeasePurpose `json:"purpose"`
	TurnID       string           `json:"turn_id,omitempty"`
	MutationID   string           `json:"mutation_id,omitempty"`
	MutationKind string           `json:"mutation_kind,omitempty"`
	OwnerID      string           `json:"owner_id"`
	Generation   int64            `json:"generation"`
	Heartbeat    int64            `json:"heartbeat"`
	AcquiredAt   time.Time        `json:"acquired_at"`
	RenewedAt    time.Time        `json:"renewed_at"`
	ExpiresAt    time.Time        `json:"expires_at"`
}

func TurnLeaseFromContext

func TurnLeaseFromContext(ctx context.Context) (TurnLease, bool)

TurnLeaseFromContext returns the durable mutation owner bound to ctx.

func (TurnLease) Fresh

func (l TurnLease) Fresh(now time.Time) bool

func (TurnLease) TakeoverEligible

func (l TurnLease) TakeoverEligible(now time.Time, policy LeasePolicy) bool

func (TurnLease) Validate

func (l TurnLease) Validate() error

type TurnLeasePurpose

type TurnLeasePurpose string
const (
	TurnLeasePurposeTurn     TurnLeasePurpose = "turn"
	TurnLeasePurposeMutation TurnLeasePurpose = "mutation"
)

func (TurnLeasePurpose) Normalize

func (p TurnLeasePurpose) Normalize() (TurnLeasePurpose, error)

type TurnLeaseRepo

type TurnLeaseRepo interface {
	AcquireTurnLease(context.Context, TurnLease) (TurnLease, error)
	RenewTurnLease(context.Context, TurnLease) (TurnLease, error)
	ReleaseTurnLease(context.Context, TurnLease) error
	ActiveTurnLease(context.Context, string) (TurnLease, bool, error)
}

type TurnMarkerStatus

type TurnMarkerStatus string
const (
	TurnStarted   TurnMarkerStatus = "started"
	TurnSavePoint TurnMarkerStatus = "save_point"
	TurnCompleted TurnMarkerStatus = "completed"
	TurnWaiting   TurnMarkerStatus = "waiting"
	TurnFailed    TurnMarkerStatus = "failed"
	TurnAborted   TurnMarkerStatus = "aborted"
)

type TurnTerminalOutcome

type TurnTerminalOutcome struct {
	Failure  *Entry
	Terminal Entry
}

type WaitApprovalDecisionResult

type WaitApprovalDecisionResult struct {
	Receipt  ApprovalDecisionReceipt
	Queue    ApprovalQueue
	Approval ApprovalRecord
}

Directories

Path Synopsis
cmd
genbackendrepo command

Jump to

Keyboard shortcuts

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