import type { CompressionBlock, CompressionMember, CompressionProtectedFragment, DcpState, MessageIdMeta, } from "./state.js" import type { DcpConfig } from "./config.js" import { estimateTokens } from "./pruner-metadata.js" import { isToolRecordProtected } from "./pruner-tools.js" import { toolRecordContinuity } from "./protected-continuity.js" import { compareConversationStableIds, buildExactRangeMembership } from "./conversation-index.js" import { createHash } from "node:crypto" import { open, realpath } from "node:fs/promises" import { isAbsolute, resolve, sep } from "node:path" const MAX_PROTECTED_ARTIFACT_BYTES = 50_000 const MAX_PROTECTED_ARTIFACT_TOTAL_BYTES = 200_000 export interface CompressionBlockCreationResult { block: CompressionBlock removedTokenEstimate: number summaryTokenEstimate: number } export interface CreateRangeCompressionBlockOptions { topic: string summary: string startTimestamp: number endTimestamp: number startMessageId?: string endMessageId?: string state: DcpState config: DcpConfig anchorTimestamp?: number anchorMessageId?: string createdByToolCallId?: string mode?: "range" | "message" version?: 2 replacementMode?: "range" | "message-body" validatePlaceholders?: boolean expandPlaceholders?: boolean /** E07 protected fragments prepared with bounded async artifact reads. */ preparedProtectedFragments?: CompressionProtectedFragment[] /** Exact current-projection source membership prepared before summary generation. */ sourceMembers?: CompressionMember[] /** Exact canonical raw mutation membership prepared with sourceMembers. */ mutationMembers?: CompressionMember[] } export interface ResolvedCompressionBoundary { timestamp: number stableId?: string } export interface ResolvedCompressionAnchor extends ResolvedCompressionBoundary {} /** * Resolve a model-visible range against the exact latest projection. New * callers must persist this plan with their block instead of reconstructing a * timestamp interval during a later context pass. */ export function resolveExactRangeMembership( startId: string, endId: string, state: DcpState, messageBody = false, ): { sourceMembers: CompressionMember[]; mutationMembers: CompressionMember[] } | undefined { const plan = buildExactRangeMembership( state.conversationIndexSnapshot, startId, endId, state, messageBody, ) if (!plan) return undefined return { sourceMembers: plan.sourceMembers, mutationMembers: plan.mutationMembers, } } function hasExactMembers( sourceMembers: CompressionMember[] | undefined, mutationMembers: CompressionMember[] | undefined, ): boolean { const valid = (members: CompressionMember[] | undefined): members is CompressionMember[] => Array.isArray(members) && members.length > 0 && members.every((member) => typeof member?.stableId === "string" && member.stableId.length > 0 && typeof member.hash === "string" && /^[a-f0-9]{64}$/i.test(member.hash), ) && new Set(members.map((member) => member.stableId)).size === members.length return valid(sourceMembers) && valid(mutationMembers) } function copyMembers(members: CompressionMember[]): CompressionMember[] { return members.map((member) => ({ stableId: member.stableId, hash: member.hash.toLowerCase() })) } function messageOrdinal(stableId: string | undefined, state: DcpState): number | undefined { if (!stableId) return undefined const visibleId = state.messageIdsByStableId.get(stableId) ?? [...state.messageMetaSnapshot.entries()].find(([, meta]) => meta.stableId === stableId)?.[0] const match = visibleId?.match(/^m(\d+)$/i) if (!match) return undefined const ordinal = Number.parseInt(match[1]!, 10) return Number.isFinite(ordinal) ? ordinal : undefined } /** * Compare two conversation boundaries. Current canonical branch order wins * whenever both stable identities resolve; timestamps are used only for * transient in-memory inputs that have not received stable identities yet. */ export function compareCompressionBoundaries( left: ResolvedCompressionBoundary, right: ResolvedCompressionBoundary, state: DcpState, ): number { const canonicalOrder = compareConversationStableIds( state.conversationIndexSnapshot, left.stableId, right.stableId, state, ) if (canonicalOrder !== undefined) return canonicalOrder if (left.timestamp !== right.timestamp) return left.timestamp - right.timestamp const leftOrdinal = messageOrdinal(left.stableId, state) const rightOrdinal = messageOrdinal(right.stableId, state) if (leftOrdinal !== undefined && rightOrdinal !== undefined && leftOrdinal !== rightOrdinal) { return leftOrdinal - rightOrdinal } return 0 } export function isCompressionBoundaryWithinRange( boundary: ResolvedCompressionBoundary, start: ResolvedCompressionBoundary, end: ResolvedCompressionBoundary, state: DcpState, ): boolean { return compareCompressionBoundaries(boundary, start, state) >= 0 && compareCompressionBoundaries(boundary, end, state) <= 0 } function rangeBoundaries( startTimestamp: number, endTimestamp: number, ids: { startMessageId?: string; endMessageId?: string } = {}, ): { start: ResolvedCompressionBoundary; end: ResolvedCompressionBoundary } { return { start: { timestamp: startTimestamp, stableId: ids.startMessageId }, end: { timestamp: endTimestamp, stableId: ids.endMessageId }, } } function inferExactMembershipFromBoundaries( startMessageId: string | undefined, endMessageId: string | undefined, state: DcpState, messageBody = false, ): { sourceMembers: CompressionMember[]; mutationMembers: CompressionMember[] } | undefined { if (!startMessageId || !endMessageId) return undefined const entries = state.conversationIndexSnapshot const aliasFor = (stableId: string): string | undefined => { const direct = entries.find((entry) => entry.stableId === stableId) if (direct?.visibleId) return direct.visibleId if (direct?.blockId !== undefined) return `b${direct.blockId}` const block = state.compressionBlocks.find((candidate) => candidate.active && (candidate.startMessageId === stableId || candidate.endMessageId === stableId), ) if (!block) return undefined return entries.some((entry) => entry.blockId === block.id) ? `b${block.id}` : undefined } const startId = aliasFor(startMessageId) const endId = aliasFor(endMessageId) return startId && endId ? resolveExactRangeMembership(startId, endId, state, messageBody) : undefined } const BLOCK_PLACEHOLDER_RE = /\(b(\d+)\)/gi const MAX_DIAGNOSTIC_IDS = 24 function idSortKey(id: string): [number, string] { const match = id.match(/^(\D+)(\d+)$/) if (!match) return [Number.MAX_SAFE_INTEGER, id] return [Number.parseInt(match[2]!, 10), match[1]!.toLowerCase()] } function sortIds(ids: string[]): string[] { return [...ids].sort((a, b) => { const [aNumber, aPrefix] = idSortKey(a) const [bNumber, bPrefix] = idSortKey(b) return aNumber - bNumber || aPrefix.localeCompare(bPrefix) || a.localeCompare(b) }) } function compactList(items: string[], maxItems = MAX_DIAGNOSTIC_IDS): string { if (items.length === 0) return "none" if (items.length <= maxItems) return items.join(", ") const headCount = Math.max(1, Math.floor(maxItems * 0.65)) const tailCount = Math.max(1, maxItems - headCount) return `${items.slice(0, headCount).join(", ")}, …, ${items.slice(-tailCount).join(", ")} (${items.length} total)` } function quoteTopic(topic: string): string { const singleLine = topic.replace(/\s+/g, " ").trim() if (!singleLine) return "" const truncated = singleLine.length > 48 ? `${singleLine.slice(0, 47)}…` : singleLine return ` "${truncated.replace(/"/g, "'")}"` } export function formatCompressionIdDiagnostics(state: DcpState): string { const rawIds = sortIds([ ...new Set([ ...state.messageIdSnapshot.keys(), ...state.messageMetaSnapshot.keys(), ]), ]) const blockIds = state.compressionBlocks .filter((block) => block.active) .sort((a, b) => a.id - b.id) .map((block) => `b${block.id}${quoteTopic(block.topic)}`) return [ `Current raw message IDs: ${compactList(rawIds)}.`, `Current active block IDs: ${compactList(blockIds)}.`, "Use only IDs from the latest visible context. If a raw range was already compressed, use the corresponding bN block ID instead of stale mNNN IDs.", "Retry at most once with a safe closed range from those IDs, or skip compression if none is safe.", ].join("\n") } export function unknownCompressionIdError(rawId: string, state: DcpState): Error { const id = rawId.trim() return new Error( `Unknown message ID: ${id}.\n` + "The ID is not present in the current DCP snapshot; it may be stale after compression, pruning, reload, or session switching.\n" + formatCompressionIdDiagnostics(state), ) } function formatRestoredBlock(block: CompressionBlock): string { return `[Previously compressed: ${block.topic}]\n${block.summary}` } /** * Replace `(bN)` placeholders in a summary with the stored content of the * referenced compression block. Unrecognised placeholders are left as-is. */ export function expandBlockPlaceholders(summary: string, state: DcpState): string { const consumed = new Set() return summary.replace(BLOCK_PLACEHOLDER_RE, (match, idStr) => { const id = parseInt(idStr, 10) const block = state.compressionBlocks.find((b) => b.id === id) if (block && consumed.has(id)) return "" if (block) consumed.add(id) return block ? formatRestoredBlock(block) : match }) } function expandBlockPlaceholdersWithRecovery( summary: string, coveredBlocks: CompressionBlock[], state: DcpState, ): string { const coveredIds = new Set(coveredBlocks.map((block) => block.id)) const consumed = new Set() const expanded = summary.replace(BLOCK_PLACEHOLDER_RE, (_match, idStr) => { const id = parseInt(idStr, 10) if (!coveredIds.has(id) || consumed.has(id)) return "" const block = state.compressionBlocks.find((b) => b.id === id) if (!block) return "" consumed.add(id) return formatRestoredBlock(block) }) const missing = coveredBlocks.filter((block) => !consumed.has(block.id)) if (missing.length === 0) return expanded const recoveryHeading = "\n\nThe following previously compressed summaries were also part of this conversation section and were preserved automatically:" const recovery = missing .map((block) => `\n\n### b${block.id}\n${formatRestoredBlock(block)}`) .join("") return expanded + recoveryHeading + recovery } function assertVerifiableBlockPlaceholders(summary: string, state: DcpState): void { const unknown = new Set() for (const match of summary.matchAll(BLOCK_PLACEHOLDER_RE)) { const id = Number.parseInt(match[1] ?? "", 10) if (!Number.isInteger(id) || !state.compressionBlocks.some((block) => block.id === id)) unknown.add(id) } if (unknown.size > 0) { throw new Error( `Compression summary contains unverifiable block placeholder(s): ${[...unknown].map((id) => `b${id}`).join(", ")}. ` + "Expand the referenced content explicitly or retry with current block IDs.", ) } } function preparePlaceholderSummary( summary: string, coveredBlocks: CompressionBlock[], state: DcpState, options: { validatePlaceholders: boolean; expandPlaceholders: boolean }, ): string { if (!options.expandPlaceholders) return summary if (options.validatePlaceholders) { return expandBlockPlaceholdersWithRecovery(summary, coveredBlocks, state) } return expandBlockPlaceholders(summary, state) } export function findCoveredAndPartialBlocks( startTimestamp: number, endTimestamp: number, state: DcpState, boundaries: { startMessageId?: string; endMessageId?: string } = {}, ): { coveredBlocks: CompressionBlock[]; partialBlocks: CompressionBlock[] } { const coveredBlocks: CompressionBlock[] = [] const partialBlocks: CompressionBlock[] = [] const requestedStart = { timestamp: startTimestamp, stableId: boundaries.startMessageId } const requestedEnd = { timestamp: endTimestamp, stableId: boundaries.endMessageId } for (const existing of state.compressionBlocks) { if (!existing.active) continue if (!Number.isFinite(existing.startTimestamp) || !Number.isFinite(existing.endTimestamp)) continue const existingStart = { timestamp: existing.startTimestamp, stableId: existing.startMessageId } const existingEnd = { timestamp: existing.endTimestamp, stableId: existing.endMessageId } const overlaps = compareCompressionBoundaries(requestedStart, existingEnd, state) <= 0 && compareCompressionBoundaries(existingStart, requestedEnd, state) <= 0 if (!overlaps) continue const fullyCovered = compareCompressionBoundaries(requestedStart, existingStart, state) <= 0 && compareCompressionBoundaries(existingEnd, requestedEnd, state) <= 0 if (fullyCovered) coveredBlocks.push(existing) else partialBlocks.push(existing) } return { coveredBlocks, partialBlocks } } function firstAddressableRawIdAfterBlock(block: CompressionBlock, state: DcpState): string | undefined { const endOrdinal = messageOrdinal(block.endMessageId, state) if (endOrdinal === undefined) return undefined let best: { id: string; ordinal: number } | undefined for (const [id, meta] of state.messageMetaSnapshot) { if (meta.blockId !== undefined) continue const match = id.match(/^m(\d+)$/i) if (!match) continue const ordinal = Number.parseInt(match[1]!, 10) if (!Number.isFinite(ordinal) || ordinal <= endOrdinal) continue if (!best || ordinal < best.ordinal) best = { id, ordinal } } return best?.id } export function getMessageMeta(id: string, state: DcpState): MessageIdMeta | undefined { return state.messageMetaSnapshot.get(id.trim()) } function extractProtectTagTexts(text: string): string[] { const protectedTexts: string[] = [] for (const match of text.matchAll(/([\s\S]*?)<\/protect>/gi)) { const protectedText = match[1]?.trim() if (protectedText) protectedTexts.push(protectedText) } return protectedTexts } export function appendProtectedPromptInfo( summary: string, startTimestamp: number, endTimestamp: number, state: DcpState, config: DcpConfig, ids: { startMessageId?: string; endMessageId?: string } = {}, ): string { if (!config.compress.protectTags) return summary const protectedTexts: string[] = [] const seen = new Set() const boundaries = rangeBoundaries(startTimestamp, endTimestamp, ids) for (const meta of state.messageMetaSnapshot.values()) { if (meta.role !== "user") continue if (meta.blockId !== undefined) continue if (!Number.isFinite(meta.timestamp)) continue if (!isCompressionBoundaryWithinRange( { timestamp: meta.timestamp, stableId: meta.stableId }, boundaries.start, boundaries.end, state, )) continue for (const text of extractProtectTagTexts(meta.text ?? "")) { if (seen.has(text)) continue seen.add(text) protectedTexts.push(text) } } if (protectedTexts.length === 0) return summary const heading = "\n\nThe following protected prompt information appeared in the selected user message(s) and must be preserved verbatim:" return summary + heading + protectedTexts.map((text) => `\n${text}`).join("") } export function appendProtectedUserMessages( summary: string, startTimestamp: number, endTimestamp: number, state: DcpState, enabled: boolean, ids: { startMessageId?: string; endMessageId?: string } = {}, ): string { if (!enabled) return summary const userTexts: string[] = [] const seen = new Set() const boundaries = rangeBoundaries(startTimestamp, endTimestamp, ids) for (const meta of state.messageMetaSnapshot.values()) { if (meta.role !== "user") continue if (meta.blockId !== undefined) continue if (!Number.isFinite(meta.timestamp)) continue if (!isCompressionBoundaryWithinRange( { timestamp: meta.timestamp, stableId: meta.stableId }, boundaries.start, boundaries.end, state, )) continue const text = meta.text?.trim() if (!text || seen.has(text)) continue seen.add(text) userTexts.push(text) } if (userTexts.length === 0) return summary const heading = "\n\nThe following user messages from this compressed range were preserved verbatim:" return summary + heading + userTexts.map((text) => `\n${text}`).join("") } async function readBoundedProtectedSubagentArtifact( details: unknown, fallbackText: string, cwd: string | undefined, ): Promise<{ artifactPath: string; text: string } | undefined> { if (!cwd || !details || typeof details !== "object") return undefined const record = details as any const rawArtifactPath = record.artifacts?.resultMd if (typeof rawArtifactPath !== "string" || rawArtifactPath.trim().length === 0) return undefined const artifactPath = rawArtifactPath.trim() let root: string let candidate: string try { root = await realpath(cwd) const requested = isAbsolute(artifactPath) ? artifactPath : resolve(root, artifactPath) candidate = await realpath(requested) } catch { // Missing artifacts are optional recovery material when the tool result text // itself is still available. Do not invent a replacement source. return undefined } if (candidate !== root && !candidate.startsWith(root + sep)) { throw new Error(`Protected artifact resolves outside the session cwd: ${artifactPath}`) } const handle = await open(candidate, "r") try { const info = await handle.stat() if (!info.isFile()) throw new Error(`Protected artifact is not a regular file: ${artifactPath}`) if (info.size > MAX_PROTECTED_ARTIFACT_BYTES) { throw new Error( `Protected artifact exceeds the ${MAX_PROTECTED_ARTIFACT_BYTES}-byte recovery budget: ${artifactPath}`, ) } const buffer = Buffer.alloc(MAX_PROTECTED_ARTIFACT_BYTES + 1) const { bytesRead } = await handle.read(buffer, 0, buffer.length, 0) if (bytesRead > MAX_PROTECTED_ARTIFACT_BYTES) { throw new Error( `Protected artifact grew beyond the ${MAX_PROTECTED_ARTIFACT_BYTES}-byte recovery budget: ${artifactPath}`, ) } const text = buffer.subarray(0, bytesRead).toString("utf8").trim() if (!text || fallbackText.includes(text)) return undefined return { artifactPath, text } } finally { await handle.close() } } function compressionProtectedFragment( kind: CompressionProtectedFragment["kind"], origin: string, text: string, ): CompressionProtectedFragment { return { kind, origin, hash: createHash("sha256").update(text).digest("hex"), text, } } function mergeProtectedFragments(...groups: CompressionProtectedFragment[][]): CompressionProtectedFragment[] { const merged: CompressionProtectedFragment[] = [] const seen = new Set() for (const fragment of groups.flat()) { const key = `${fragment.origin}:${fragment.hash}` if (seen.has(key)) continue seen.add(key) merged.push(fragment) } return merged } function collectCurrentProtectedFragments( startTimestamp: number, endTimestamp: number, state: DcpState, config: DcpConfig, mode: "range" | "message", ids: { startMessageId?: string; endMessageId?: string }, ): CompressionProtectedFragment[] { const fragments: CompressionProtectedFragment[] = [] const boundaries = rangeBoundaries(startTimestamp, endTimestamp, ids) for (const [visibleId, meta] of state.messageMetaSnapshot) { if (meta.blockId !== undefined || !Number.isFinite(meta.timestamp)) continue if (!isCompressionBoundaryWithinRange( { timestamp: meta.timestamp, stableId: meta.stableId }, boundaries.start, boundaries.end, state, )) continue const originBase = meta.stableId ?? visibleId ?? String(meta.timestamp) if (mode === "range" && config.compress.protectUserMessages && meta.role === "user") { const text = meta.text?.trim() if (text) fragments.push(compressionProtectedFragment("user", `user:${originBase}`, text)) } if (config.compress.protectTags && meta.role === "user") { for (const [index, text] of extractProtectTagTexts(meta.text ?? "").entries()) { fragments.push(compressionProtectedFragment("prompt", `prompt:${originBase}:${index}`, text)) } } if (meta.role !== "toolResult" && meta.role !== "bashExecution") continue const toolCallId = meta.toolCallId if (!toolCallId) continue const record = state.toolCalls.get(toolCallId) if (!record || record.toolName === "compress" || !isToolRecordProtected(record, config)) continue const continuity = toolRecordContinuity(record, config) const text = continuity.text?.trim() if (!text) continue fragments.push(compressionProtectedFragment("tool", `tool:${toolCallId}`, text)) } return mergeProtectedFragments(fragments) } export interface PrepareCompressionProtectedFragmentsOptions { startTimestamp: number endTimestamp: number startMessageId?: string endMessageId?: string state: DcpState config: DcpConfig mode: "range" | "message" cwd?: string } /** * Prepare protected continuity fragments before the synchronous block commit. * Raw transcript/state remains the primary source. Optional subagent artifacts * are read only through bounded async I/O rooted at the session cwd. */ export async function prepareCompressionProtectedFragments( options: PrepareCompressionProtectedFragmentsOptions, ): Promise { const { startTimestamp, endTimestamp, startMessageId, endMessageId, state, config, mode, cwd } = options const ids = { startMessageId, endMessageId } const { coveredBlocks } = findCoveredAndPartialBlocks(startTimestamp, endTimestamp, state, ids) const fragments = mergeProtectedFragments( collectInheritedProtectedFragments(coveredBlocks), collectCurrentProtectedFragments(startTimestamp, endTimestamp, state, config, mode, ids), ) if (!cwd) return fragments const boundaries = rangeBoundaries(startTimestamp, endTimestamp, ids) let totalArtifactBytes = 0 for (const meta of state.messageMetaSnapshot.values()) { if (meta.blockId !== undefined || !Number.isFinite(meta.timestamp)) continue if (!isCompressionBoundaryWithinRange( { timestamp: meta.timestamp, stableId: meta.stableId }, boundaries.start, boundaries.end, state, )) continue if (meta.role !== "toolResult" && meta.role !== "bashExecution") continue const toolCallId = meta.toolCallId if (!toolCallId) continue const record = state.toolCalls.get(toolCallId) if (!record || !isToolRecordProtected(record, config)) continue if (record.toolName !== "subagents" && record.toolName !== "async_subagents_result") continue const output = (record.outputText ?? meta.text ?? "").trim() if (!output) continue const artifact = await readBoundedProtectedSubagentArtifact(record.outputDetails, output, cwd) if (!artifact) continue const artifactBytes = Buffer.byteLength(artifact.text, "utf8") totalArtifactBytes += artifactBytes if (totalArtifactBytes > MAX_PROTECTED_ARTIFACT_TOTAL_BYTES) { throw new Error( `Protected artifacts exceed the ${MAX_PROTECTED_ARTIFACT_TOTAL_BYTES}-byte total recovery budget`, ) } fragments.push(compressionProtectedFragment( "tool", `artifact:${toolCallId}:${artifact.artifactPath}`, `### Expanded subagent result: ${artifact.artifactPath}\n${artifact.text}`, )) } return mergeProtectedFragments(fragments) } function collectInheritedProtectedFragments(coveredBlocks: CompressionBlock[]): CompressionProtectedFragment[] { return mergeProtectedFragments( ...coveredBlocks.map((block) => block.protectedFragments ?? []), ) } function appendProtectedFragmentLedger(summary: string, fragments: CompressionProtectedFragment[]): string { const missing = fragments.filter((fragment) => !summary.includes(fragment.text)) if (missing.length === 0) return summary return summary + "\n\nThe following protected continuity fragments were preserved verbatim:" + missing.map((fragment) => `\n\n${fragment.text}`).join("") } export function estimateVisibleRangeTokens( startTimestamp: number, endTimestamp: number, coveredBlocks: CompressionBlock[], state: DcpState, ids: { startMessageId?: string; endMessageId?: string } = {}, ): number { let total = 0 const boundaries = rangeBoundaries(startTimestamp, endTimestamp, ids) for (const meta of state.messageMetaSnapshot.values()) { if (meta.blockId !== undefined) continue if (!Number.isFinite(meta.timestamp)) continue if (!isCompressionBoundaryWithinRange( { timestamp: meta.timestamp, stableId: meta.stableId }, boundaries.start, boundaries.end, state, )) continue total += Math.max(0, Math.round(meta.tokenEstimate ?? 0)) } for (const block of coveredBlocks) { total += Math.max(0, Math.round(block.summaryTokenEstimate ?? 0)) } return total } /** * Resolve a user-supplied ID string (e.g. "m001" or "b3") to an actual * message timestamp. */ export function resolveIdToTimestamp( rawId: string, field: "startTimestamp" | "endTimestamp", state: DcpState, ): number { return resolveIdToBoundary(rawId, field, state).timestamp } export function resolveIdToBoundary( rawId: string, field: "startTimestamp" | "endTimestamp", state: DcpState, ): ResolvedCompressionBoundary { const id = rawId.trim() const blockMatch = id.match(/^b(\d+)$/i) if (blockMatch) { const blockId = parseInt(blockMatch[1]!, 10) const block = state.compressionBlocks.find((b) => b.id === blockId && b.active) if (!block) throw unknownCompressionIdError(id, state) return { timestamp: block[field], stableId: field === "startTimestamp" ? block.startMessageId : block.endMessageId, } } const meta = state.messageMetaSnapshot.get(id) if (meta) return resolveMetaBoundary(meta, field, state) const ts = state.messageIdSnapshot.get(id) if (ts !== undefined) return { timestamp: ts } throw unknownCompressionIdError(id, state) } /** * Resolve a `MessageIdMeta` to a compression boundary. When the meta is a * synthetic placeholder for an active compression block (`meta.blockId` is * set), resolve to that block's stored boundary so the caller rolls the * block up instead of nesting a new block on top of the placeholder. This * A model-visible mNNN may itself represent a previously compressed section. */ function resolveMetaBoundary( meta: MessageIdMeta, field: "startTimestamp" | "endTimestamp", state: DcpState, ): ResolvedCompressionBoundary { if (meta.blockId !== undefined) { const block = state.compressionBlocks.find((b) => b.id === meta.blockId && b.active) if (block) { return { timestamp: block[field], stableId: field === "startTimestamp" ? block.startMessageId : block.endMessageId, } } } return { timestamp: meta.timestamp, stableId: meta.stableId } } /** * Determine the anchor timestamp for a compression block — the timestamp of * the first raw message that appears strictly after `endTimestamp`. */ export function resolveAnchorTimestamp(endTimestamp: number, state: DcpState): number { return resolveAnchorBoundary(endTimestamp, state).timestamp } export function resolveAnchorBoundary( endTimestamp: number, state: DcpState, endMessageId?: string, ): ResolvedCompressionAnchor { const endBoundary = { timestamp: endTimestamp, stableId: endMessageId } let anchor: ResolvedCompressionAnchor | null = null for (const meta of state.messageMetaSnapshot.values()) { if (meta.blockId !== undefined || !Number.isFinite(meta.timestamp)) continue const candidate = { timestamp: meta.timestamp, stableId: meta.stableId } if (compareCompressionBoundaries(candidate, endBoundary, state) <= 0) continue if (anchor === null || compareCompressionBoundaries(candidate, anchor, state) < 0) anchor = candidate } if (anchor === null) { for (const ts of state.messageIdSnapshot.values()) { if (ts > endTimestamp && (anchor === null || ts < anchor.timestamp)) anchor = { timestamp: ts } } } return anchor ?? { timestamp: endTimestamp + 1 } } export function createRangeCompressionBlock( options: CreateRangeCompressionBlockOptions, ): CompressionBlockCreationResult { const { topic, summary, startTimestamp, endTimestamp, startMessageId, endMessageId, state, config, anchorTimestamp: requestedAnchorTimestamp, anchorMessageId: requestedAnchorMessageId, createdByToolCallId, mode = "range", version = 2, replacementMode, validatePlaceholders = true, expandPlaceholders = true, preparedProtectedFragments = [], sourceMembers: requestedSourceMembers, mutationMembers: requestedMutationMembers, } = options const defaultAnchor = resolveAnchorBoundary(endTimestamp, state, endMessageId) const anchorTimestamp = requestedAnchorTimestamp ?? defaultAnchor.timestamp const anchorMessageId = requestedAnchorMessageId ?? defaultAnchor.stableId const ids = { startMessageId, endMessageId } const inferredMembership = !requestedSourceMembers && !requestedMutationMembers ? inferExactMembershipFromBoundaries(startMessageId, endMessageId, state, replacementMode === "message-body") : undefined const sourceMembers = requestedSourceMembers ?? inferredMembership?.sourceMembers const mutationMembers = requestedMutationMembers ?? inferredMembership?.mutationMembers if (compareCompressionBoundaries( { timestamp: startTimestamp, stableId: startMessageId }, { timestamp: endTimestamp, stableId: endMessageId }, state, ) > 0) { throw new Error("Compression range start must appear before end in the conversation") } if (!Number.isFinite(startTimestamp)) { throw new Error(`Compression range start resolved to a non-finite timestamp (${startTimestamp})`) } if (!Number.isFinite(endTimestamp)) { throw new Error(`Compression range end resolved to a non-finite timestamp (${endTimestamp})`) } const { coveredBlocks, partialBlocks } = findCoveredAndPartialBlocks( startTimestamp, endTimestamp, state, { startMessageId, endMessageId }, ) if (partialBlocks.length > 0) { const blockList = partialBlocks.map((block) => `b${block.id} "${block.topic}"`).join(", ") const freeHints = partialBlocks .map((block) => { const firstFreeId = firstAddressableRawIdAfterBlock(block, state) return firstFreeId ? `first raw ID after b${block.id} is ${firstFreeId}` : undefined }) .filter((hint): hint is string => typeof hint === "string") const freeHint = freeHints.length > 0 ? ` ${freeHints.join("; ")}.` : "" throw new Error( `Compression range partially overlaps existing block(s): ${blockList}. ` + `Select the whole block or choose non-overlapping boundaries.${freeHint}`, ) } if (version !== 2 || !hasExactMembers(sourceMembers, mutationMembers)) { throw new Error( "Compression requires exact canonical source and mutation membership. " + "Refresh the current DCP snapshot and retry with visible IDs.", ) } // Every block in the new format carries a protected-fragment ledger. Do not // silently fall back to the old recursive-summary behavior: a missing ledger // means the in-memory state violates the journal contract and the roll-up // must fail closed. if (coveredBlocks.some((block) => block.version !== 2 || block.protectedFragments === undefined)) { throw new Error("Compression roll-up requires exact v2 blocks with protected-fragment ledgers") } const ledgerRollup = coveredBlocks.length > 0 const hasExplicitBlockReferences = [...summary.matchAll(BLOCK_PLACEHOLDER_RE)].length > 0 if (ledgerRollup) assertVerifiableBlockPlaceholders(summary, state) const placeholderSummary = preparePlaceholderSummary( summary, coveredBlocks, state, { validatePlaceholders: ledgerRollup ? false : validatePlaceholders, // Explicit references are deliberately expanded once, never left to // dangle after their source blocks are deactivated. Otherwise modern // ledger roll-ups replace old summary prose instead of recursively // appending it. expandPlaceholders: ledgerRollup ? hasExplicitBlockReferences : expandPlaceholders, }, ) const userPreservedSummary = mode === "range" ? appendProtectedUserMessages( placeholderSummary, startTimestamp, endTimestamp, state, config.compress.protectUserMessages, ids, ) : placeholderSummary const promptPreservedSummary = appendProtectedPromptInfo( userPreservedSummary, startTimestamp, endTimestamp, state, config, ids, ) const protectedFragments = mergeProtectedFragments( collectInheritedProtectedFragments(coveredBlocks), collectCurrentProtectedFragments(startTimestamp, endTimestamp, state, config, mode, ids), preparedProtectedFragments, ) const ledgerPreservedSummary = appendProtectedFragmentLedger(promptPreservedSummary, protectedFragments) const block: CompressionBlock = { id: state.nextBlockId++, topic, summary: ledgerPreservedSummary, startTimestamp, endTimestamp, startMessageId, endMessageId, anchorTimestamp, anchorMessageId, createdByToolCallId, active: true, summaryTokenEstimate: estimateTokens(ledgerPreservedSummary), createdAt: Date.now(), coveredBlockIds: coveredBlocks.map((covered) => covered.id), mode, version: 2, replacementMode, sourceMembers: copyMembers(sourceMembers!), mutationMembers: copyMembers(mutationMembers!), protectedFragments, } state.compressionBlocks.push(block) for (const covered of coveredBlocks) { covered.active = false } return { block, removedTokenEstimate: estimateVisibleRangeTokens(startTimestamp, endTimestamp, coveredBlocks, state, ids), summaryTokenEstimate: Math.max(0, Math.round(block.summaryTokenEstimate ?? 0)), } }