import { type CompactionPreparation, type CompactionResult, estimateContextTokens } from "@caupulican/pi-agent-core/compaction/compaction"; import { type SessionContext, SessionManager } from "@caupulican/pi-agent-core/session"; import type { Message, Usage } from "@caupulican/pi-ai"; import type { AgentResumeContext, AttemptUsageSnapshot, ResourcePointer } from "../orchestration/contracts.ts"; import { type WorkerContextForkReference } from "../orchestration/worker-context-fork-reference.ts"; import { type WorkerConversationFileRevision, type WorkerSessionFileHead } from "./worker-conversation-revision.ts"; export declare const MAX_WORKER_TRANSCRIPT_PAGE_MESSAGES = 64; export declare const MAX_WORKER_TRANSCRIPT_PAGE_BYTES: number; interface WorkerConversationTranscriptPageOptions { /** Opaque raw-entry offset returned by the previous page. */ cursor?: number; maxMessages?: number; maxBytes?: number; } interface WorkerConversationTranscriptPage { /** Opaque raw-entry offset used for this page; it is not a message count. */ cursor: number; messages: Message[]; nextCursor?: number; /** Messages consumed but not cloned because one message exceeded the complete page byte ceiling. */ omittedMessages: number; /** Exact UTF-8 JSON byte length of `messages`, including array framing. */ serializedBytes: number; } interface WorkerControlTranscriptExpectation { messageId: string; content: string; } interface WorkerControlTranscriptReconciliation { delivered: boolean; appended: boolean; } export interface CreateWorkerConversationOptions { agentDir: string; parentSessionId: string; /** Durable logical identity for the worker/lane. It never becomes a path segment directly. */ logicalAgentId: string; cwd: string; orchestrationProfileId?: string; modelRef?: string; resourceProfileNames: readonly string[]; contextPointers: readonly ResourcePointer[]; /** Immutable sanitized parent context captured before this logical agent is admitted. */ birthContextForkReference?: WorkerContextForkReference; } interface OpenWorkerConversationOptions { agentDir: string; resumeContext: AgentResumeContext; expectedLogicalAgentId?: string; } /** * Explicit, token-based retention for a durable worker transcript. * * This deliberately has no default: the worker model/context policy belongs to orchestration, not * the transcript store. Call only at a safe worker turn boundary, after all messages from that turn * have been durably committed. The execution controller provides `generateVerifiedCompaction` * through the shared model-aware pipeline; this store owns only preparation, append-only apply, * and its deterministic verified fallback. */ export interface WorkerConversationRetentionPolicy { /** Maximum provider-visible context tokens after a checkpoint is applied. */ maxContextTokens: number; /** Shared compaction's retained recent-context target. Must be lower than the maximum. */ keepRecentTokens: number; /** * Generate a verified shared compaction result. Failures, malformed results, and omitted * generators fall back to `createDeterministicCompaction`; raw transcript entries are never * replaced or removed. */ generateVerifiedCompaction?: (preparation: CompactionPreparation) => Promise; /** * Cumulative provider usage spent by the current verified generation if it failed before it * could return a CompactionResult. This usage is attached to the deterministic checkpoint so * recovery never loses rejected-summary spend. */ getFailedCompactionUsage?: () => Usage | undefined; } interface WorkerConversationRetentionOutcome { status: "within_limit" | "compacted_verified" | "compacted_deterministic" | "cannot_compact"; context: SessionContext; contextUsage: ReturnType; } interface WorkerConversationMetadataState { logicalAgentId: string; parentSessionId?: string; birthContextForkReference?: WorkerContextForkReference; usageAccountingVersion?: 1; } interface WorkerConversationMetadataBinding extends WorkerConversationMetadataState { file: string; agentDir: string; } interface WorkerConversationCore { sessionManager: SessionManager; head?: WorkerSessionFileHead; metadataRevision?: WorkerConversationFileRevision; metadataState?: WorkerConversationMetadataState; invalid: boolean; generation: number; activeTranscriptCursors: number; } export interface WorkerTranscriptCommitCursor { readonly kind: "worker-transcript-suffix-v1"; } /** * Durable SessionManager-backed transcript for exactly one logical worker lane. * * The canonical session-file lock plus a raw-byte revision head fences competing processes. One * store shares the parsed core across lightweight resume-context views; unexpected durable changes * fail closed for active owners and only strict append-only recovery may replace the core. */ export declare class WorkerConversation { private readonly core; private readonly resumeContext; private readonly agentDir?; private readonly metadataFile?; constructor(sessionManager: SessionManager, resumeContext: AgentResumeContext, metadata?: WorkerConversationMetadataBinding, sharedCore?: WorkerConversationCore); private get sessionManager(); private set sessionManager(value); /** Persisted conversations share one adopted metadata state; direct in-memory instances are current. */ private get usageAccountingVersion(); /** Immutable parent-context identity installed before this logical worker's first attempt. */ getBirthContextForkReference(): WorkerContextForkReference | undefined; /** Resolve current provider-visible messages lazily through SessionManager. */ getProviderContext(): SessionContext; /** Convert the current compacted projection at its transcript owner, never at each consumer. */ getProviderMessages(): Message[]; /** True when provider context is a compacted projection rather than the raw transcript prefix. */ hasProviderCompaction(): boolean; /** * The immutable, append-only raw worker messages, including messages compacted out of provider * context. This is recovery/audit data; it is never loaded into a provider request implicitly. */ getRawTranscript(): Message[]; /** * Read one bounded raw-transcript page without materializing the complete session entry list. * A message larger than the complete page ceiling is consumed as an omission so every valid * non-terminal cursor advances. The opaque cursor is a raw-entry offset, so each entry is visited * at most once across a complete pagination pass and a page may contain no transcript messages. * `serializedBytes` measures the returned `messages` JSON array. */ getRawTranscriptPage(options?: WorkerConversationTranscriptPageOptions): WorkerConversationTranscriptPage; /** Last durable provider message owned by one attempt, never an earlier persistent-worker turn. */ getLastAttemptMessage(attemptId: string): Message | undefined; /** Durable host-observed mutation progress. Custom entries never enter provider context. */ recordChangedFile(attemptId: string, filePath: string): void; /** Rehydrate the bounded mutation set across owner-session disposal and worker resume. */ getChangedFiles(attemptId: string): string[]; /** * Mark the first transcript entry owned by one durable attempt. Persistent logical workers share * a conversation across tasks, so recovery accounting must not replay an earlier task's spend as * the new attempt's baseline. The marker is idempotent and never enters provider context. */ beginAttemptUsage(attemptId: string): void; /** * Persist the authoritative task prompt once inside its durable attempt boundary. The boundary, * rather than provider-history length, is the idempotency owner: inherited birth messages and * queued mailbox controls may already exist, while a crash after append must not duplicate the * prompt on resume. */ ensureAttemptUserPrompt(attemptId: string, prompt: string): void; /** Whether missing attempt markers have versioned, fail-closed accounting semantics. */ usesAttemptUsageBoundaries(): boolean; /** * Upgrade a legacy idle conversation before its next durable task is prepared. Persisting this * evidence first closes the crash window where the task exists but its transcript marker does not. */ enableAttemptUsageBoundaries(): void; /** * Recover cumulative accounting from raw entry metadata without cloning or resolving message * payloads that compaction deliberately moved out of the provider-visible working set. Attempts * created before usage boundaries existed conservatively fall back to the complete transcript. */ getRawTranscriptUsage(attemptId?: string): AttemptUsageSnapshot; /** * Locate a bounded set of durable mailbox delivery commits without materializing or cloning the * complete raw transcript. Recovery uses this only to close the narrow crash window after a * control message was appended but before its mailbox acknowledgement was persisted. */ findDeliveredWorkerControlMessageIds(expectations: readonly WorkerControlTranscriptExpectation[]): Set; /** * Atomically prove one exact worker-control identity and optionally append it. Reopening inside * the canonical session-file lock makes two stale recovery owners serialize on the latest durable * transcript instead of both observing absence and appending a duplicate. */ reconcileWorkerControlMessage(expectation: WorkerControlTranscriptExpectation, message: Message, appendIfMissing: boolean): WorkerControlTranscriptReconciliation; /** * Append one shared-compaction checkpoint when the current provider-visible context exceeds the * explicit worker policy. The source transcript remains append-only; only the projection sent to * the next provider call changes. The method is idempotent while no new messages are appended. */ compactProviderContext(policy: WorkerConversationRetentionPolicy, signal?: AbortSignal): Promise; /** Append one already-authorized worker message to the canonical transcript. */ appendMessage(message: Message): string; private appendSessionEntry; private appendSessionEntryLocked; private reopenVerifiedSessionManagerLocked; private synchronizeCoreLocked; /** Store-only refresh while the canonical session-file lock is already held. */ refreshCachedCoreLocked(): void; private captureTranscriptCommitCursorLocked; /** Atomically capture the provider projection and the raw-entry cursor that immediately follows it. */ beginTranscriptCommit(): { history: Message[]; cursor: WorkerTranscriptCommitCursor; }; private withCanonicalSessionLock; /** Commit only the bounded child-loop suffix captured by an opaque raw-entry cursor. */ captureTranscriptCommitCursor(): WorkerTranscriptCommitCursor; abortTranscriptCommit(cursor: WorkerTranscriptCommitCursor): void; commitTranscript(cursor: WorkerTranscriptCommitCursor, suffix: readonly Message[], options?: { appendMissing?: boolean; }): number; /** Return an isolated copy suitable for durable orchestration/process-resume state. */ getResumeContext(): AgentResumeContext; } /** Creates and reopens canonical SessionManager transcripts for logical Pi worker lanes. */ export declare class WorkerConversationStore { private readonly cachedCores; clearCache(): void; create(options: CreateWorkerConversationOptions): WorkerConversation; private createLocked; /** Open the one canonical transcript or atomically create it with the exact same durable identity. */ ensure(options: CreateWorkerConversationOptions): WorkerConversation; open(options: OpenWorkerConversationOptions): WorkerConversation; private openExisting; } export {}; //# sourceMappingURL=worker-conversation-store.d.ts.map