From e0c222392c21e564c3efb9ed9b77264898399edf Mon Sep 17 00:00:00 2001 From: kami Date: Wed, 1 Jul 2026 14:03:47 +0400 Subject: [PATCH] =?UTF-8?q?feat(context,kernel):=20compression=20pipeline?= =?UTF-8?q?=20stage=203=20=E2=80=94=20replay-safe=20tier=20summarizer?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ContextSummarizedEvent (registered) keyed on stable session sequence, not the ephemeral pack ordinal, so replay reuses the recorded summary artifact and never re-calls the LLM (invariants #6/#8). TierContextSummarizer mirrors JournalCompactionService for freeform context (tool outputs/docs/retrieved text) the journal path doesn't cover. ContextState + DefaultContextReducer fold the event; reducer replay-safety unit-tested. --- .../core/context/DefaultContextReducer.kt | 11 +++- .../correx/core/context/state/ContextState.kt | 7 ++ .../test/kotlin/ContextSummaryReducerTest.kt | 45 +++++++++++++ .../events/events/ContextSummarizedEvent.kt | 23 +++++++ .../events/serialization/Serialization.kt | 2 + .../orchestration/TierContextSummarizer.kt | 64 +++++++++++++++++++ 6 files changed, 151 insertions(+), 1 deletion(-) create mode 100644 core/context/src/test/kotlin/ContextSummaryReducerTest.kt create mode 100644 core/events/src/main/kotlin/com/correx/core/events/events/ContextSummarizedEvent.kt create mode 100644 core/kernel/src/main/kotlin/com/correx/core/kernel/orchestration/TierContextSummarizer.kt diff --git a/core/context/src/main/kotlin/com/correx/core/context/DefaultContextReducer.kt b/core/context/src/main/kotlin/com/correx/core/context/DefaultContextReducer.kt index 1a26d162..1b7431ce 100644 --- a/core/context/src/main/kotlin/com/correx/core/context/DefaultContextReducer.kt +++ b/core/context/src/main/kotlin/com/correx/core/context/DefaultContextReducer.kt @@ -1,8 +1,17 @@ package com.correx.core.context import com.correx.core.context.state.ContextState +import com.correx.core.events.events.ContextSummarizedEvent import com.correx.core.events.events.StoredEvent class DefaultContextReducer : ContextReducer { - override fun reduce(state: ContextState, event: StoredEvent): ContextState = state + override fun reduce(state: ContextState, event: StoredEvent): ContextState = + when (val p = event.payload) { + is ContextSummarizedEvent -> state.copy( + summarizedThroughSequence = p.coversThroughSequence, + summaryArtifactId = p.summaryArtifactId, + summaryTier = p.tier, + ) + else -> state + } } diff --git a/core/context/src/main/kotlin/com/correx/core/context/state/ContextState.kt b/core/context/src/main/kotlin/com/correx/core/context/state/ContextState.kt index 356e23b6..c860621f 100644 --- a/core/context/src/main/kotlin/com/correx/core/context/state/ContextState.kt +++ b/core/context/src/main/kotlin/com/correx/core/context/state/ContextState.kt @@ -1,5 +1,6 @@ package com.correx.core.context.state +import com.correx.core.events.types.ArtifactId import com.correx.core.events.types.ContextPackId import kotlinx.serialization.Serializable @@ -8,4 +9,10 @@ data class ContextState( val builtPackIds: List = emptyList(), val buildingInProgress: Boolean = false, val interrupted: Boolean = false, + // Recursive-summarization coverage (pipeline stage 5). Freeform context up through this + // stable event sequence is represented by [summaryArtifactId]; assembly drops the covered + // raw entries and injects the summary instead. -1 = nothing summarized yet. + val summarizedThroughSequence: Long = -1L, + val summaryArtifactId: ArtifactId? = null, + val summaryTier: Int = 0, ) diff --git a/core/context/src/test/kotlin/ContextSummaryReducerTest.kt b/core/context/src/test/kotlin/ContextSummaryReducerTest.kt new file mode 100644 index 00000000..5e818aa9 --- /dev/null +++ b/core/context/src/test/kotlin/ContextSummaryReducerTest.kt @@ -0,0 +1,45 @@ +import com.correx.core.context.DefaultContextReducer +import com.correx.core.context.state.ContextState +import com.correx.core.events.events.ContextSummarizedEvent +import com.correx.core.events.events.EventMetadata +import com.correx.core.events.events.StoredEvent +import com.correx.core.events.types.ArtifactId +import com.correx.core.events.types.EventId +import com.correx.core.events.types.SessionId +import kotlinx.datetime.Instant +import org.junit.jupiter.api.Assertions.assertEquals +import org.junit.jupiter.api.Test + +class ContextSummaryReducerTest { + + private val reducer = DefaultContextReducer() + + private fun stored(payload: ContextSummarizedEvent, seq: Long) = StoredEvent( + metadata = EventMetadata( + eventId = EventId("e$seq"), + sessionId = SessionId("s"), + timestamp = Instant.fromEpochSeconds(0), + schemaVersion = 1, + causationId = null, + correlationId = null, + ), + sequence = seq, + sessionSequence = seq, + payload = payload, + ) + + @Test + fun `latest summary event wins and records stable coverage`() { + // Replaying the same events must yield the same state — coverage is keyed on the stable + // event sequence, so no live LLM call is needed to rebuild (invariant #8). + var state = ContextState() + state = reducer.reduce(state, stored( + ContextSummarizedEvent(SessionId("s"), 10L, ArtifactId("a1"), tier = 2, entriesOmittedCount = 5), 1L)) + state = reducer.reduce(state, stored( + ContextSummarizedEvent(SessionId("s"), 25L, ArtifactId("a2"), tier = 3, entriesOmittedCount = 8), 2L)) + + assertEquals(25L, state.summarizedThroughSequence) + assertEquals(ArtifactId("a2"), state.summaryArtifactId) + assertEquals(3, state.summaryTier) + } +} diff --git a/core/events/src/main/kotlin/com/correx/core/events/events/ContextSummarizedEvent.kt b/core/events/src/main/kotlin/com/correx/core/events/events/ContextSummarizedEvent.kt new file mode 100644 index 00000000..e4a1d798 --- /dev/null +++ b/core/events/src/main/kotlin/com/correx/core/events/events/ContextSummarizedEvent.kt @@ -0,0 +1,23 @@ +package com.correx.core.events.events + +import com.correx.core.events.types.ArtifactId +import com.correx.core.events.types.SessionId +import kotlinx.serialization.SerialName +import kotlinx.serialization.Serializable + +/** + * Records that freeform context up through session sequence [coversThroughSequence] was + * summarized into [summaryArtifactId] at tier [tier]. Keyed on the stable event sequence (not + * the ephemeral per-assembly pack ordinal, which is re-stamped each build) so replay reuses + * the recorded summary deterministically and never re-calls the LLM (invariants #6/#8). + * [tier] 2 = light single-pass summary, 3 = summary-of-summaries. + */ +@Serializable +@SerialName("ContextSummarized") +data class ContextSummarizedEvent( + val sessionId: SessionId, + val coversThroughSequence: Long, + val summaryArtifactId: ArtifactId, + val tier: Int, + val entriesOmittedCount: Int, +) : EventPayload diff --git a/core/events/src/main/kotlin/com/correx/core/events/serialization/Serialization.kt b/core/events/src/main/kotlin/com/correx/core/events/serialization/Serialization.kt index 07f28ea0..9f07ba22 100644 --- a/core/events/src/main/kotlin/com/correx/core/events/serialization/Serialization.kt +++ b/core/events/src/main/kotlin/com/correx/core/events/serialization/Serialization.kt @@ -61,6 +61,7 @@ import com.correx.core.events.events.StageCompletedEvent import com.correx.core.events.events.PreemptRedirectBlockedEvent import com.correx.core.events.events.PreemptRedirectEvent import com.correx.core.events.events.StageFailedEvent +import com.correx.core.events.events.ContextSummarizedEvent import com.correx.core.events.events.SteeringNoteAddedEvent import com.correx.core.events.events.ToolCallAssessedEvent import com.correx.core.events.events.ToolExecutionCompletedEvent @@ -107,6 +108,7 @@ val eventModule = SerializersModule { subclass(ToolCallAssessedEvent::class) subclass(FileWrittenEvent::class) subclass(SteeringNoteAddedEvent::class) + subclass(ContextSummarizedEvent::class) subclass(PreemptRedirectEvent::class) subclass(PreemptRedirectBlockedEvent::class) subclass(InitialIntentEvent::class) diff --git a/core/kernel/src/main/kotlin/com/correx/core/kernel/orchestration/TierContextSummarizer.kt b/core/kernel/src/main/kotlin/com/correx/core/kernel/orchestration/TierContextSummarizer.kt new file mode 100644 index 00000000..c984bf19 --- /dev/null +++ b/core/kernel/src/main/kotlin/com/correx/core/kernel/orchestration/TierContextSummarizer.kt @@ -0,0 +1,64 @@ +package com.correx.core.kernel.orchestration + +import com.correx.core.artifactstore.ArtifactStore +import com.correx.core.context.model.ContextEntry +import com.correx.core.events.events.ContextSummarizedEvent +import com.correx.core.events.events.EventPayload +import com.correx.core.events.types.SessionId + +/** + * Stage 5 (RECURSIVE_SUMMARIZE) for freeform context — the tier the decision-journal + * compaction path does NOT cover: long tool outputs, docs and retrieved text. Mirrors + * [JournalCompactionService]: summarize via injected inference, store the summary in the + * artifact store, emit [ContextSummarizedEvent] referencing it. Replay reads the recorded + * artifact and never re-calls the LLM (invariants #6/#8). + */ +class TierContextSummarizer( + private val artifactStore: ArtifactStore, + private val summarize: suspend (String) -> String, + private val tokenThreshold: () -> Int = { DEFAULT_TOKEN_THRESHOLD }, +) { + companion object { + private const val DEFAULT_TOKEN_THRESHOLD = 2000 + private const val CHARS_PER_TOKEN = 4 + } + + /** + * Summarize [freeformEntries] (tier 2/3 candidates) when their combined size crosses the + * threshold. Returns true when a summary event was emitted. [tier] selects light summary + * (2) vs summary-of-summaries (3); protected spans are already preserved in the fact sheet, + * so the prompt only asks to keep intent and load-bearing facts. + */ + suspend fun summarizeIfNeeded( + sessionId: SessionId, + freeformEntries: List, + coversThroughSequence: Long, + tier: Int, + emit: suspend (EventPayload) -> Unit, + ): Boolean { + if (freeformEntries.isEmpty()) return false + val tokenEstimate = freeformEntries.sumOf { it.content.length } / CHARS_PER_TOKEN + if (tokenEstimate < tokenThreshold()) return false + + val prompt = buildString { + appendLine("Summarize the following context concisely (≤200 words).") + appendLine("Preserve user intent, decisions, constraints and every load-bearing fact.") + appendLine() + freeformEntries.forEach { appendLine("- ${it.content}") } + } + val summaryText = summarize(prompt) + val summaryArtifactId = artifactStore.put(summaryText.toByteArray(Charsets.UTF_8)) + artifactStore.flushBefore { + emit( + ContextSummarizedEvent( + sessionId = sessionId, + coversThroughSequence = coversThroughSequence, + summaryArtifactId = summaryArtifactId, + tier = tier, + entriesOmittedCount = freeformEntries.size, + ), + ) + } + return true + } +}