diff --git a/apps/server/src/main/kotlin/com/correx/apps/server/protocol/ClientMessage.kt b/apps/server/src/main/kotlin/com/correx/apps/server/protocol/ClientMessage.kt index 1c796f69..550322e7 100644 --- a/apps/server/src/main/kotlin/com/correx/apps/server/protocol/ClientMessage.kt +++ b/apps/server/src/main/kotlin/com/correx/apps/server/protocol/ClientMessage.kt @@ -56,6 +56,14 @@ sealed class ClientMessage { @Serializable data class DiscardIdea(val ideaId: String) : ClientMessage() + /** + * Operator request to promote an idea into the bound session's curated project profile (added as + * a `conventions` entry in `/.correx/project.toml`). Tombstones the idea like a + * discard; reply is a fresh IdeaList. + */ + @Serializable + data class PromoteIdea(val ideaId: String) : ClientMessage() + /** Operator request for the current editable config (replied to with a ConfigSnapshot). */ @Serializable data object GetConfig : ClientMessage() diff --git a/apps/server/src/main/kotlin/com/correx/apps/server/ws/GlobalStreamHandler.kt b/apps/server/src/main/kotlin/com/correx/apps/server/ws/GlobalStreamHandler.kt index 41d3c9e5..a850106b 100644 --- a/apps/server/src/main/kotlin/com/correx/apps/server/ws/GlobalStreamHandler.kt +++ b/apps/server/src/main/kotlin/com/correx/apps/server/ws/GlobalStreamHandler.kt @@ -19,6 +19,7 @@ import com.correx.core.events.events.ApprovalGrantCreatedEvent import com.correx.core.events.events.ChatSessionStartedEvent import com.correx.core.events.events.EventMetadata import com.correx.core.events.events.IdeaDiscardedEvent +import com.correx.core.events.events.IdeaPromotedEvent import com.correx.core.events.events.InitialIntentEvent import com.correx.core.events.events.NewEvent import com.correx.core.events.events.SessionWorkspaceBoundEvent @@ -30,6 +31,9 @@ import com.correx.core.events.types.SessionId import com.correx.core.events.types.StageId import com.correx.core.inference.ProviderHealth import com.correx.core.kernel.orchestration.WorkspaceContext +import com.correx.core.config.ProjectProfileLoader +import com.correx.core.config.ProjectProfileWriter +import com.correx.core.config.withConvention import com.correx.core.router.IdeaReader import com.correx.core.utils.TypeId import com.correx.infrastructure.inference.commons.ResourceProbe @@ -233,6 +237,7 @@ class GlobalStreamHandler(private val module: ServerModule) { is ClientMessage.ListArtifacts -> sendFrame(queries.listArtifacts(msg.sessionId)) is ClientMessage.ListIdeas -> sendFrame(queries.listIdeas()) is ClientMessage.DiscardIdea -> handleDiscardIdea(msg, sendFrame) + is ClientMessage.PromoteIdea -> handlePromoteIdea(msg, sendFrame) is ClientMessage.GetSessionStats -> sendFrame(queries.sessionStats(msg.sessionId)) is ClientMessage.GetConfig -> sendFrame(queries.configSnapshot(error = null, restartRequired = emptyList())) is ClientMessage.UpdateConfig -> sendFrame(queries.updateConfig(msg.patch)) @@ -336,6 +341,57 @@ class GlobalStreamHandler(private val module: ServerModule) { sendFrame(queries.listIdeas()) } + // Promote an idea into the bound session's curated project profile (a `conventions` entry in + // /.correx/project.toml), then tombstone it like a discard so it leaves the board. + // No-ops (just refresh the board) when the idea is missing or its capturing session has no bound + // workspace — promotion needs a workspace to write the profile into. + private suspend fun handlePromoteIdea( + msg: ClientMessage.PromoteIdea, + sendFrame: suspend (ServerMessage) -> Unit, + ) { + runCatching { + val reader = IdeaReader(module.eventStore) + val sessionId = reader.sessionOf(msg.ideaId) ?: return@runCatching + val text = reader.textOf(msg.ideaId) ?: return@runCatching + val workspaceRoot = workspaceRootOf(sessionId) ?: return@runCatching + + withContext(Dispatchers.IO) { + val updated = ProjectProfileLoader.load(workspaceRoot).withConvention(text) + ProjectProfileWriter.write(workspaceRoot, updated) + } + + module.eventStore.append( + NewEvent( + metadata = EventMetadata( + eventId = EventId(UUID.randomUUID().toString()), + sessionId = sessionId, + timestamp = Clock.System.now(), + schemaVersion = 1, + causationId = null, + correlationId = null, + ), + payload = IdeaPromotedEvent( + ideaId = msg.ideaId, + sessionId = sessionId, + text = text, + timestampMs = System.currentTimeMillis(), + ), + ), + ) + }.onFailure { + log.error("promoteIdea failed for ideaId={}: {}", msg.ideaId, it.message, it) + } + sendFrame(queries.listIdeas()) + } + + // The workspace a session was bound to, from its latest SessionWorkspaceBoundEvent (invariant #9), + // or null if the session never bound one. + private fun workspaceRootOf(sessionId: SessionId): String? = + module.eventStore.read(sessionId) + .mapNotNull { it.payload as? SessionWorkspaceBoundEvent } + .lastOrNull() + ?.workspaceRoot + private fun errorResponse(message: String) = ServerMessage.ProtocolError(message) private fun encodeError(message: String): Frame.Text = diff --git a/apps/server/src/test/kotlin/com/correx/apps/server/protocol/ServerMessageSerializationTest.kt b/apps/server/src/test/kotlin/com/correx/apps/server/protocol/ServerMessageSerializationTest.kt index 8be8fa5c..85ec1665 100644 --- a/apps/server/src/test/kotlin/com/correx/apps/server/protocol/ServerMessageSerializationTest.kt +++ b/apps/server/src/test/kotlin/com/correx/apps/server/protocol/ServerMessageSerializationTest.kt @@ -267,6 +267,14 @@ class ServerMessageSerializationTest { assertEquals("i1", (decoded as ClientMessage.DiscardIdea).ideaId) } + @Test + fun `PromoteIdea decodes from the client wire format`() { + val wire = """{"type":"com.correx.apps.server.protocol.ClientMessage.PromoteIdea","ideaId":"i1"}""" + val decoded = ProtocolSerializer.decodeClientMessage(wire) + assert(decoded is ClientMessage.PromoteIdea) { "expected PromoteIdea, got $decoded" } + assertEquals("i1", (decoded as ClientMessage.PromoteIdea).ideaId) + } + @Test fun `ClarificationResponse decodes from the client wire format`() { val wire = """ diff --git a/core/config/src/main/kotlin/com/correx/core/config/ConfigLoader.kt b/core/config/src/main/kotlin/com/correx/core/config/ConfigLoader.kt index 00352588..1dd5451c 100644 --- a/core/config/src/main/kotlin/com/correx/core/config/ConfigLoader.kt +++ b/core/config/src/main/kotlin/com/correx/core/config/ConfigLoader.kt @@ -57,8 +57,11 @@ internal object SimpleToml { var depth = 0 var inQuotes = false var quoteChar = ' ' + var escaped = false for (ch in value) { when { + escaped -> escaped = false + ch == '\\' && inQuotes && quoteChar == '"' -> escaped = true (ch == '"' || ch == '\'') && !inQuotes -> { inQuotes = true; quoteChar = ch } ch == quoteChar && inQuotes -> inQuotes = false ch == '[' && !inQuotes -> depth++ @@ -92,8 +95,13 @@ internal object SimpleToml { var current = StringBuilder() var inQuotes = false var quoteChar = ' ' + var escaped = false for (ch in content) { when { + // Keep escape sequences (e.g. \") intact so they reach stripQuotes verbatim; an + // escaped quote must not be treated as the string's closing delimiter. + escaped -> { current.append(ch); escaped = false } + ch == '\\' && inQuotes && quoteChar == '"' -> { current.append(ch); escaped = true } (ch == '"' || ch == '\'') && !inQuotes -> { inQuotes = true; quoteChar = ch; current.append(ch) } ch == quoteChar && inQuotes -> { inQuotes = false; current.append(ch) } ch == ',' && !inQuotes -> { result.add(stripQuotes(current.toString().trim())); current = StringBuilder() } @@ -104,10 +112,39 @@ internal object SimpleToml { return result } + /** + * Strips the surrounding quotes from [value] and, for double-quoted strings, unescapes the + * `\\` and `\"` sequences emitted by the writers (mirrors [ProjectProfileWriter]/[CorrexConfigWriter]). + */ private fun stripQuotes(value: String): String { - val isDoubleQuoted = value.startsWith("\"") && value.endsWith("\"") - val isSingleQuoted = value.startsWith("'") && value.endsWith("'") - return if (isDoubleQuoted || isSingleQuoted) value.substring(1, value.length - 1) else value + val isDoubleQuoted = value.startsWith("\"") && value.endsWith("\"") && value.length >= 2 + val isSingleQuoted = value.startsWith("'") && value.endsWith("'") && value.length >= 2 + return when { + isDoubleQuoted -> unescapeDoubleQuoted(value.substring(1, value.length - 1)) + isSingleQuoted -> value.substring(1, value.length - 1) + else -> value + } + } + + /** Reverses the writer's escaping: `\\` -> `\`, `\"` -> `"`; a trailing lone `\` is kept as-is. */ + private fun unescapeDoubleQuoted(value: String): String { + if ('\\' !in value) return value + val sb = StringBuilder(value.length) + var i = 0 + while (i < value.length) { + val ch = value[i] + if (ch == '\\' && i + 1 < value.length) { + val next = value[i + 1] + if (next == '\\' || next == '"') { + sb.append(next) + i += 2 + continue + } + } + sb.append(ch) + i++ + } + return sb.toString() } } diff --git a/core/config/src/main/kotlin/com/correx/core/config/ProjectProfileWriter.kt b/core/config/src/main/kotlin/com/correx/core/config/ProjectProfileWriter.kt new file mode 100644 index 00000000..40808484 --- /dev/null +++ b/core/config/src/main/kotlin/com/correx/core/config/ProjectProfileWriter.kt @@ -0,0 +1,64 @@ +package com.correx.core.config + +import java.nio.file.Files + +/** + * Serializes a [ProjectProfile] back to TOML text that [ProjectProfileLoader] reads identically: + * the invariant is `load(write(profile)) == profile` for every field the loader understands + * (about, conventions, commands). + * + * Like [CorrexConfigWriter] this **regenerates** the file from the profile object — it does not + * preserve hand-written comments. Empty sections are skipped: a blank `about` and an empty + * `conventions`/`commands` produce no output for that part (the loader reconstructs the defaults + * from their absence), so a [ProjectProfile.isEmpty] profile serializes to an empty string. + */ +object ProjectProfileWriter { + + /** Renders [profile] as TOML. String values escape `\` and `"`; empty sections are skipped. */ + fun serialize(profile: ProjectProfile): String { + val b = StringBuilder() + + if (profile.about.isNotBlank()) { + b.append("about = ").append(str(profile.about)).append('\n') + } + + if (profile.conventions.isNotEmpty()) { + b.append("conventions = ").append(list(profile.conventions)).append('\n') + } + + if (profile.commands.isNotEmpty()) { + if (b.isNotEmpty()) b.append('\n') + b.append("[commands]\n") + profile.commands.forEach { (key, value) -> + b.append(key).append(" = ").append(str(value)).append('\n') + } + } + + return b.toString() + } + + /** + * Writes [profile] to `/.correx/project.toml`, creating the `.correx/` directory + * if it is missing. Overwrites any existing file (the writer owns the file once a save happens). + */ + fun write(workspaceRoot: String, profile: ProjectProfile) { + val path = ProjectProfileLoader.profilePath(workspaceRoot) + Files.createDirectories(path.parent) + Files.writeString(path, serialize(profile)) + } + + private fun str(v: String): String = "\"" + v.replace("\\", "\\\\").replace("\"", "\\\"") + "\"" + + private fun list(v: List): String = "[" + v.joinToString(", ") { str(it) } + "]" +} + +/** + * Appends [text] to this profile's conventions, de-duplicated by exact string match. A blank + * [text] is ignored (returns the profile unchanged). Used when promoting an idea-board entry into + * the curated profile so repeated promotions of the same text don't pile up duplicate conventions. + */ +fun ProjectProfile.withConvention(text: String): ProjectProfile { + val trimmed = text.trim() + if (trimmed.isBlank() || trimmed in conventions) return this + return copy(conventions = conventions + trimmed) +} diff --git a/core/config/src/test/kotlin/com/correx/core/config/ProjectProfileWriterTest.kt b/core/config/src/test/kotlin/com/correx/core/config/ProjectProfileWriterTest.kt new file mode 100644 index 00000000..c8ab2ead --- /dev/null +++ b/core/config/src/test/kotlin/com/correx/core/config/ProjectProfileWriterTest.kt @@ -0,0 +1,79 @@ +package com.correx.core.config + +import org.junit.jupiter.api.Assertions.assertEquals +import org.junit.jupiter.api.Assertions.assertFalse +import org.junit.jupiter.api.Assertions.assertSame +import org.junit.jupiter.api.Assertions.assertTrue +import org.junit.jupiter.api.Test +import org.junit.jupiter.api.io.TempDir +import java.nio.file.Path + +class ProjectProfileWriterTest { + + @TempDir + lateinit var tempDir: Path + + @Test + fun `write then load round-trips about conventions and commands including embedded quotes`() { + val profile = ProjectProfile( + about = "Event-sourced kernel \"the\" orchestrator, Kotlin/JVM 21", + conventions = listOf( + "no bare try-catch", + "prefer the \"narrow\" interface", // a convention with a double-quote + ), + commands = mapOf( + "build" to "./gradlew build", + "test" to "./gradlew check", + ), + ) + + ProjectProfileWriter.write(tempDir.toString(), profile) + val loaded = ProjectProfileLoader.load(tempDir.toString()) + + assertEquals(profile.about, loaded.about) + assertEquals(profile.conventions, loaded.conventions) + assertEquals(profile.commands, loaded.commands) + assertEquals(profile, loaded) + } + + @Test + fun `an empty profile serializes to an empty string and skips empty sections`() { + val serialized = ProjectProfileWriter.serialize(ProjectProfile()) + assertEquals("", serialized) + assertFalse(serialized.contains("[commands]")) + } + + @Test + fun `blank about is omitted but conventions and commands are still written`() { + val serialized = ProjectProfileWriter.serialize( + ProjectProfile( + about = " ", + conventions = listOf("c1"), + commands = mapOf("run" to "x"), + ), + ) + assertFalse(serialized.contains("about ="), "blank about must be skipped") + assertTrue(serialized.contains("conventions = [\"c1\"]")) + assertTrue(serialized.contains("[commands]")) + } + + @Test + fun `withConvention appends a new convention`() { + val updated = ProjectProfile(conventions = listOf("a")).withConvention("b") + assertEquals(listOf("a", "b"), updated.conventions) + } + + @Test + fun `withConvention de-dupes an exact existing convention`() { + val profile = ProjectProfile(conventions = listOf("no bare try-catch")) + val updated = profile.withConvention("no bare try-catch") + assertSame(profile, updated, "an exact duplicate must return the same profile unchanged") + assertEquals(listOf("no bare try-catch"), updated.conventions) + } + + @Test + fun `withConvention ignores blank text`() { + val profile = ProjectProfile(conventions = listOf("a")) + assertSame(profile, profile.withConvention(" ")) + } +} diff --git a/core/events/src/main/kotlin/com/correx/core/events/events/IdeaEvents.kt b/core/events/src/main/kotlin/com/correx/core/events/events/IdeaEvents.kt index 6f7629c9..f7a69346 100644 --- a/core/events/src/main/kotlin/com/correx/core/events/events/IdeaEvents.kt +++ b/core/events/src/main/kotlin/com/correx/core/events/events/IdeaEvents.kt @@ -31,3 +31,19 @@ data class IdeaDiscardedEvent( val ideaId: String, val sessionId: SessionId, ) : EventPayload + +/** + * The operator promoted an idea from the cross-session board into the curated per-repo profile + * (`/.correx/project.toml`, as a `conventions` entry). Like [IdeaDiscardedEvent] this + * is a tombstone — the original [IdeaCapturedEvent] stays in the log (and replays), but the idea-board + * projection filters out any promoted [ideaId]. [sessionId] is the capturing session (so all of one + * idea's events stay co-located); [text] records the convention text written to the profile. + */ +@Serializable +@SerialName("IdeaPromoted") +data class IdeaPromotedEvent( + val ideaId: String, + val sessionId: SessionId, + val text: String, + val timestampMs: Long, +) : 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 2c394076..64a985ae 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 @@ -29,6 +29,7 @@ import com.correx.core.events.events.HealthDegradedEvent import com.correx.core.events.events.HealthRestoredEvent import com.correx.core.events.events.IdeaCapturedEvent import com.correx.core.events.events.IdeaDiscardedEvent +import com.correx.core.events.events.IdeaPromotedEvent import com.correx.core.events.events.JournalCompactedEvent import com.correx.core.events.events.FileWrittenEvent import com.correx.core.events.events.L3MemoryRetrievedEvent @@ -137,6 +138,7 @@ val eventModule = SerializersModule { subclass(WorkflowProposedEvent::class) subclass(IdeaCapturedEvent::class) subclass(IdeaDiscardedEvent::class) + subclass(IdeaPromotedEvent::class) subclass(CritiqueOutcomeCorrelatedEvent::class) subclass(StageCheckpointPassedEvent::class) subclass(StageCheckpointFailedEvent::class) diff --git a/core/events/src/test/kotlin/com/correx/core/events/EventSerializationHardeningTest.kt b/core/events/src/test/kotlin/com/correx/core/events/EventSerializationHardeningTest.kt index cc9481b6..e3a75b92 100644 --- a/core/events/src/test/kotlin/com/correx/core/events/EventSerializationHardeningTest.kt +++ b/core/events/src/test/kotlin/com/correx/core/events/EventSerializationHardeningTest.kt @@ -6,6 +6,7 @@ import com.correx.core.approvals.Tier import com.correx.core.events.events.ApprovalDecisionResolvedEvent import com.correx.core.events.events.ApprovalRequestedEvent import com.correx.core.events.events.EventPayload +import com.correx.core.events.events.IdeaPromotedEvent import com.correx.core.events.events.OrchestrationPausedEvent import com.correx.core.events.events.RiskAssessedEvent import com.correx.core.events.events.SteeringNoteAddedEvent @@ -81,6 +82,12 @@ class EventSerializationHardeningTest { terminalStageId = stageId, totalStages = 1, ), + "IdeaPromoted" to IdeaPromotedEvent( + ideaId = "idea-1", + sessionId = sessionId, + text = "cache the repo map across sessions", + timestampMs = 1_700_000_000_000, + ), ) @Test diff --git a/core/router/src/main/kotlin/com/correx/core/router/IdeaReader.kt b/core/router/src/main/kotlin/com/correx/core/router/IdeaReader.kt index f9212085..d724c88a 100644 --- a/core/router/src/main/kotlin/com/correx/core/router/IdeaReader.kt +++ b/core/router/src/main/kotlin/com/correx/core/router/IdeaReader.kt @@ -2,6 +2,7 @@ package com.correx.core.router import com.correx.core.events.events.IdeaCapturedEvent import com.correx.core.events.events.IdeaDiscardedEvent +import com.correx.core.events.events.IdeaPromotedEvent import com.correx.core.events.stores.EventStore import com.correx.core.events.types.SessionId import com.correx.core.router.model.Idea @@ -10,28 +11,36 @@ import com.correx.core.router.model.Idea * Rebuilds the operator's idea board from the event log. Cross-session by design: it folds over * [EventStore.allEvents] rather than a single session's replay, so an idea captured in one session * shows on the board (and feeds the router) in every session. A captured idea is dropped once a - * matching [IdeaDiscardedEvent] tombstones it (invariant #1 — the capture stays in the log). + * matching [IdeaDiscardedEvent] (or [IdeaPromotedEvent]) tombstones it (invariant #1 — the capture + * stays in the log). */ class IdeaReader(private val eventStore: EventStore) { - /** Active ideas (captured, not later discarded) across all sessions, newest first. */ + /** Active ideas (captured, not later discarded or promoted) across all sessions, newest first. */ fun activeIdeas(): List { - val discarded = mutableSetOf() + val tombstoned = mutableSetOf() val captured = mutableListOf() eventStore.allEvents().forEach { stored -> when (val payload = stored.payload) { - is IdeaDiscardedEvent -> discarded += payload.ideaId + is IdeaDiscardedEvent -> tombstoned += payload.ideaId + is IdeaPromotedEvent -> tombstoned += payload.ideaId is IdeaCapturedEvent -> captured += Idea(payload.ideaId, payload.text, payload.timestampMs) else -> Unit } } - return captured.filterNot { it.id in discarded }.sortedByDescending { it.capturedAtMs } + return captured.filterNot { it.id in tombstoned }.sortedByDescending { it.capturedAtMs } } - /** The session that captured [ideaId], so a discard tombstone lands alongside its capture. */ + /** The session that captured [ideaId], so a tombstone lands alongside its capture. */ fun sessionOf(ideaId: String): SessionId? = + capturedOf(ideaId)?.sessionId + + /** The captured text of [ideaId] (so a promotion can write it into the project profile), or null. */ + fun textOf(ideaId: String): String? = + capturedOf(ideaId)?.text + + private fun capturedOf(ideaId: String): IdeaCapturedEvent? = eventStore.allEvents() .mapNotNull { it.payload as? IdeaCapturedEvent } .firstOrNull { it.ideaId == ideaId } - ?.sessionId } diff --git a/core/router/src/test/kotlin/com/correx/core/router/IdeaReaderTest.kt b/core/router/src/test/kotlin/com/correx/core/router/IdeaReaderTest.kt new file mode 100644 index 00000000..0a48ae10 --- /dev/null +++ b/core/router/src/test/kotlin/com/correx/core/router/IdeaReaderTest.kt @@ -0,0 +1,100 @@ +package com.correx.core.router + +import com.correx.core.events.events.EventMetadata +import com.correx.core.events.events.EventPayload +import com.correx.core.events.events.IdeaCapturedEvent +import com.correx.core.events.events.IdeaDiscardedEvent +import com.correx.core.events.events.IdeaPromotedEvent +import com.correx.core.events.events.NewEvent +import com.correx.core.events.events.StoredEvent +import com.correx.core.events.stores.EventStore +import com.correx.core.events.types.EventId +import com.correx.core.events.types.SessionId +import kotlinx.coroutines.flow.Flow +import kotlinx.coroutines.flow.emptyFlow +import kotlinx.datetime.Instant +import org.junit.jupiter.api.Assertions.assertEquals +import org.junit.jupiter.api.Assertions.assertFalse +import org.junit.jupiter.api.Assertions.assertNull +import org.junit.jupiter.api.Assertions.assertTrue +import org.junit.jupiter.api.Test + +class IdeaReaderTest { + + private val session = SessionId("s-1") + + @Test + fun `a promoted idea is filtered out of the board like a discard`() { + val store = fakeStore( + IdeaCapturedEvent("keep", session, "cache the repo map", timestampMs = 1), + IdeaCapturedEvent("promoted", session, "no bare try-catch", timestampMs = 2), + IdeaCapturedEvent("discarded", session, "old note", timestampMs = 3), + IdeaPromotedEvent("promoted", session, "no bare try-catch", timestampMs = 99), + IdeaDiscardedEvent("discarded", session), + ) + + val active = IdeaReader(store).activeIdeas() + + assertEquals(listOf("keep"), active.map { it.id }) + assertFalse(active.any { it.id == "promoted" }, "promoted idea must leave the board") + assertFalse(active.any { it.id == "discarded" }, "discarded idea must leave the board") + } + + @Test + fun `textOf returns the captured text and null for an unknown idea`() { + val store = fakeStore( + IdeaCapturedEvent("i1", session, "events are the only source of truth", timestampMs = 1), + ) + val reader = IdeaReader(store) + + assertEquals("events are the only source of truth", reader.textOf("i1")) + assertNull(reader.textOf("missing")) + } + + @Test + fun `textOf still resolves the text after the idea is promoted`() { + val store = fakeStore( + IdeaCapturedEvent("i1", session, "prefer data classes", timestampMs = 1), + IdeaPromotedEvent("i1", session, "prefer data classes", timestampMs = 5), + ) + val reader = IdeaReader(store) + + // The capture stays in the log (invariant #1), so the text remains readable. + assertEquals("prefer data classes", reader.textOf("i1")) + assertTrue(reader.activeIdeas().isEmpty()) + } + + private fun fakeStore(vararg payloads: EventPayload): EventStore = FakeEventStore(payloads.toList()) + + /** Minimal read-only [EventStore]: only [allEvents] is exercised by [IdeaReader]. */ + private class FakeEventStore(payloads: List) : EventStore { + private val stored: List = payloads.mapIndexed { index, payload -> + StoredEvent( + metadata = EventMetadata( + eventId = EventId("e-$index"), + sessionId = SessionId("s-1"), + timestamp = Instant.fromEpochMilliseconds(index.toLong()), + schemaVersion = 1, + causationId = null, + correlationId = null, + ), + sequence = index.toLong(), + sessionSequence = index.toLong(), + payload = payload, + ) + } + + override fun allEvents(): Sequence = stored.asSequence() + + override suspend fun append(event: NewEvent): StoredEvent = throw NotImplementedError() + override suspend fun appendAll(events: List): List = throw NotImplementedError() + override fun read(sessionId: SessionId): List = throw NotImplementedError() + override fun readFrom(sessionId: SessionId, fromSequence: Long): List = + throw NotImplementedError() + override fun lastSequence(sessionId: SessionId): Long? = throw NotImplementedError() + override fun subscribe(sessionId: SessionId): Flow = emptyFlow() + override fun subscribeAll(): Flow = emptyFlow() + override suspend fun lastGlobalSequence(): Long = throw NotImplementedError() + override fun allSessionIds(): Set = throw NotImplementedError() + } +}