Remove :storage-inmemory module, tests, and related code.
ci / JVM build + tests (push) Successful in 6m46s

This commit is contained in:
2026-09-23 03:48:38 +03:00
parent 84f5fd84f3
commit e447525059
105 changed files with 2732 additions and 1978 deletions
+38
View File
@@ -0,0 +1,38 @@
plugins {
alias(libs.plugins.kotlin.multiplatform)
}
// KMP-реализация :reflection-api (ReflectionStore) поверх ksqlite.
//
// Минимальная — только таблица `reflection` + 2 индекса по ней.
// Остальные таблицы (`conversation`, `message`, `working_memory`) живут в
// других ksqlite-модулях; этот модуль не претендует на полную схему агента.
//
// Цели сборки — jvm() + linuxX64() + mingwX64(); Apple targets auto-disabled
// на Linux (ksqlite 0.1.2 не публикует native артефакты для Apple).
kotlin {
jvmToolchain(21)
jvm()
linuxX64()
mingwX64()
sourceSets {
commonMain.dependencies {
// ksqlite 0.1.2 опубликован в Maven Central — обычный
// `mavenCentral()` в settings.gradle.kts его подтянет.
implementation("pw.binom.db:ksqlite:0.1.2")
api(project(":reflection-api"))
// :reflection-api ссылается на :journal-api типы в сигнатурах
// (MessageRecord и т.п. для enum'ов). Фиксируем явно чтобы
// downstream-консьюмеры не должны были декларировать ещё раз.
api(project(":journal-api"))
}
commonTest.dependencies {
implementation(kotlin("test"))
implementation(libs.kotlinx.coroutines.test)
}
}
}
@@ -0,0 +1,251 @@
package pw.binom.agentik.reflection.ksqlite
import pw.binom.agentik.reflection.Reflection
import pw.binom.agentik.reflection.ReflectionEvent
import pw.binom.agentik.reflection.ReflectionStore
import pw.binom.db.ksqlite.SQLiteConnection
import pw.binom.db.ksqlite.SQLitePreparedStatement
import pw.binom.db.ksqlite.SQLiteResultSet
import kotlin.time.Clock
import kotlin.time.Instant
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.MutableSharedFlow
import kotlinx.coroutines.flow.asSharedFlow
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
import kotlinx.coroutines.withContext
/**
* ksqlite-реализация [ReflectionStore].
*
* ## Lifecycle соединения
*
* Семантика владения connection'ом идентична
* `pw.binom.agentik.journal.ksqlite.KsqliteJournalStore`:
* - `KsqliteReflectionStore(connection)` — внешнее соединение, store НЕ
* закрывает его в [close] и НЕ прогоняет миграцию (shared-connection bundle).
* - `KsqliteReflectionStore(path)` — открывает файловое соединение,
* прогоняет [Schema.migrate], закрывает соединение в [close].
* - `KsqliteReflectionStore.memory(name)` — in-memory, мигрирует,
* закрывает в [close].
*/
class KsqliteReflectionStore private constructor(
private val connection: SQLiteConnection,
private val ownsConnection: Boolean,
private val clock: Clock = Clock.System,
) : ReflectionStore {
constructor(path: String, clock: Clock = Clock.System) : this(
connection = SQLiteConnection.open(path = path),
ownsConnection = true,
clock = clock,
)
constructor(connection: SQLiteConnection, clock: Clock = Clock.System) : this(
connection = connection,
ownsConnection = false,
clock = clock,
)
init {
Schema.migrate(connection)
}
private val mutex = Mutex()
private val ev = MutableSharedFlow<ReflectionEvent>(extraBufferCapacity = 16)
// pre-prepare (см. KsqliteMessageStore KDoc — почему это критично против
// SIGSEGV в StmtHolder.finalize на закрытой connection).
private val insertStmt = connection.prepare(
"""
INSERT INTO ${Schema.TABLE_REFLECTION}
(${Schema.COL_ID}, ${Schema.COL_CONVERSATION_ID}, ${Schema.COL_CREATED_AT},
${Schema.COL_TURNS_ANALYZED}, ${Schema.COL_SCORE},
${Schema.COL_SUMMARY}, ${Schema.COL_WEAK_SPOTS_JSON})
VALUES (?, ?, ?, ?, ?, ?, ?)
""".trimIndent()
)
private val getStmt = connection.prepare(
"""
SELECT ${Schema.COL_ID}, ${Schema.COL_CONVERSATION_ID}, ${Schema.COL_CREATED_AT},
${Schema.COL_TURNS_ANALYZED}, ${Schema.COL_SCORE},
${Schema.COL_SUMMARY}, ${Schema.COL_WEAK_SPOTS_JSON}
FROM ${Schema.TABLE_REFLECTION}
WHERE ${Schema.COL_ID} = ?
""".trimIndent()
)
private val listRecentStmt = connection.prepare(
"""
SELECT ${Schema.COL_ID}, ${Schema.COL_CONVERSATION_ID}, ${Schema.COL_CREATED_AT},
${Schema.COL_TURNS_ANALYZED}, ${Schema.COL_SCORE},
${Schema.COL_SUMMARY}, ${Schema.COL_WEAK_SPOTS_JSON}
FROM ${Schema.TABLE_REFLECTION}
ORDER BY ${Schema.COL_CREATED_AT} DESC
LIMIT ?
""".trimIndent()
)
private val listForConvStmt = connection.prepare(
"""
SELECT ${Schema.COL_ID}, ${Schema.COL_CONVERSATION_ID}, ${Schema.COL_CREATED_AT},
${Schema.COL_TURNS_ANALYZED}, ${Schema.COL_SCORE},
${Schema.COL_SUMMARY}, ${Schema.COL_WEAK_SPOTS_JSON}
FROM ${Schema.TABLE_REFLECTION}
WHERE ${Schema.COL_CONVERSATION_ID} = ?
ORDER BY ${Schema.COL_CREATED_AT} DESC
LIMIT ?
""".trimIndent()
)
private val deleteOlderThanStmt = connection.prepare(
"DELETE FROM ${Schema.TABLE_REFLECTION} WHERE ${Schema.COL_CREATED_AT} < ?"
)
private val countStmt = connection.prepare(
"SELECT COUNT(*) FROM ${Schema.TABLE_REFLECTION}"
)
override suspend fun insert(reflection: Reflection): Unit = withContext(Dispatchers.Default) {
mutex.withLock {
insertStmt.reset()
insertStmt.clearBindings()
insertStmt.bindText(1, reflection.id)
val convId = reflection.conversationId
if (convId != null) insertStmt.bindText(2, convId) else insertStmt.bindNull(2)
insertStmt.bindLong(3, reflection.createdAt.toEpochMilliseconds())
insertStmt.bindLong(4, reflection.turnsAnalyzed.toLong())
insertStmt.bindLong(5, reflection.score.toLong())
insertStmt.bindText(6, reflection.summary)
insertStmt.bindText(7, encodeStringArray(reflection.weakSpots))
insertStmt.executeUpdate()
ev.tryEmit(ReflectionEvent.Created(reflection))
}
}
override suspend fun get(id: String): Reflection? = withContext(Dispatchers.Default) {
mutex.withLock {
getStmt.reset()
getStmt.clearBindings()
getStmt.bindText(1, id)
getStmt.executeQuery().use { rs ->
if (rs.next()) rs.toDomain() else null
}
}
}
override suspend fun listRecent(limit: Int): List<Reflection> = withContext(Dispatchers.Default) {
mutex.withLock {
listRecentStmt.reset()
listRecentStmt.clearBindings()
listRecentStmt.bindLong(1, limit.toLong())
val out = mutableListOf<Reflection>()
listRecentStmt.executeQuery().use { rs ->
while (rs.next()) out.add(rs.toDomain())
}
out
}
}
override suspend fun listForConversation(conversationId: String, limit: Int): List<Reflection> = withContext(Dispatchers.Default) {
mutex.withLock {
listForConvStmt.reset()
listForConvStmt.clearBindings()
listForConvStmt.bindText(1, conversationId)
listForConvStmt.bindLong(2, limit.toLong())
val out = mutableListOf<Reflection>()
listForConvStmt.executeQuery().use { rs ->
while (rs.next()) out.add(rs.toDomain())
}
out
}
}
override suspend fun deleteOlderThan(cutoff: Instant): Unit = withContext(Dispatchers.Default) {
mutex.withLock {
deleteOlderThanStmt.reset()
deleteOlderThanStmt.clearBindings()
deleteOlderThanStmt.bindLong(1, cutoff.toEpochMilliseconds())
deleteOlderThanStmt.executeUpdate()
}
}
override suspend fun count(): Int = withContext(Dispatchers.Default) {
mutex.withLock {
countStmt.reset()
countStmt.clearBindings()
countStmt.executeQuery().use { rs ->
if (rs.next()) (rs.getLong(0) ?: 0L).toInt() else 0
}
}
}
override fun events(): Flow<ReflectionEvent> = ev.asSharedFlow()
override fun close() {
insertStmt.close()
getStmt.close()
listRecentStmt.close()
listForConvStmt.close()
deleteOlderThanStmt.close()
countStmt.close()
if (ownsConnection) connection.close()
}
companion object {
fun memory(name: String? = null, clock: Clock = Clock.System): KsqliteReflectionStore =
KsqliteReflectionStore(
connection = SQLiteConnection.memory(name),
ownsConnection = true,
clock = clock,
)
}
private fun SQLiteResultSet.toDomain(): Reflection = Reflection(
id = getText(0)!!,
conversationId = getText(1),
createdAt = Instant.fromEpochMilliseconds(getLong(2)!!),
turnsAnalyzed = getLong(3)!!.toInt(),
score = getLong(4)!!.toInt(),
summary = getText(5)!!,
weakSpots = decodeStringArray(getText(6)!!),
)
}
/** Hand-rolled JSON-encode/decode для List<String>. */
internal fun encodeStringArray(items: List<String>): String = buildString {
append('[')
items.forEachIndexed { i, s ->
if (i > 0) append(',')
append('"')
for (c in s) {
when (c) {
'\\' -> append("\\\\")
'"' -> append("\\\"")
else -> append(c)
}
}
append('"')
}
append(']')
}
internal fun decodeStringArray(raw: String): List<String> {
val s = raw.trim()
if (!s.startsWith('[') || !s.endsWith(']')) return emptyList()
val inner = s.substring(1, s.length - 1)
if (inner.isBlank()) return emptyList()
val out = mutableListOf<String>()
var i = 0
val cur = StringBuilder()
var inStr = false
var escape = false
while (i < inner.length) {
val c = inner[i]
when {
escape -> { cur.append(c); escape = false }
c == '\\' -> { escape = true }
c == '"' -> { if (inStr) out.add(cur.toString()); inStr = !inStr; cur.clear() }
inStr -> cur.append(c)
}
i++
}
return out
}
@@ -0,0 +1,71 @@
package pw.binom.agentik.reflection.ksqlite
import pw.binom.db.ksqlite.SQLiteConnection
/**
* Имена таблиц/колонок/индексов для ksqlite-бэкенда `:reflection-api`.
*
* Минимум — только то, что относится к `reflection` (реализация
* [KsqliteReflectionStore]). Остальные таблицы агента (`conversation`,
* `message`, `working_memory`) живут в других ksqlite-модулях.
*
* Все DDL/DML в этом модуле должны ссылаться на эти константы — никаких
* хардкоженных литералов в `prepare("SELECT ... FROM foo ...")` в store'е.
*/
internal object Schema {
/** Версия схемы модуля. Увеличивать при ЛЮБОМ изменении DDL. */
const val CURRENT_VERSION: Int = 1
// ───── Таблица ─────
const val TABLE_REFLECTION = "reflection"
// ───── Колонки ─────
const val COL_ID = "id"
const val COL_CONVERSATION_ID = "conversation_id"
const val COL_CREATED_AT = "created_at"
const val COL_TURNS_ANALYZED = "turns_analyzed"
const val COL_SCORE = "score"
const val COL_SUMMARY = "summary"
const val COL_WEAK_SPOTS_JSON = "weak_spots_json"
// ───── Индексы ─────
const val IDX_REFLECTION_CREATED = "idx_reflection_created"
const val IDX_REFLECTION_CONV = "idx_reflection_conv"
private val v1Ddl = """
CREATE TABLE IF NOT EXISTS $TABLE_REFLECTION (
$COL_ID TEXT NOT NULL PRIMARY KEY,
$COL_CONVERSATION_ID TEXT,
$COL_CREATED_AT INTEGER NOT NULL,
$COL_TURNS_ANALYZED INTEGER NOT NULL,
$COL_SCORE INTEGER NOT NULL,
$COL_SUMMARY TEXT NOT NULL,
$COL_WEAK_SPOTS_JSON TEXT NOT NULL DEFAULT '[]'
);
""".trimIndent()
private val v1IndexesDdl = """
CREATE INDEX IF NOT EXISTS $IDX_REFLECTION_CREATED
ON $TABLE_REFLECTION($COL_CREATED_AT DESC);
CREATE INDEX IF NOT EXISTS $IDX_REFLECTION_CONV
ON $TABLE_REFLECTION($COL_CONVERSATION_ID, $COL_CREATED_AT DESC);
""".trimIndent()
/**
* Прогоняет миграцию схемы до [CURRENT_VERSION] на пустой или существующей БД.
*
* Идемпотентен: повторный вызов на уже мигрированной БД — no-op.
*/
fun migrate(conn: SQLiteConnection) {
conn.exec("BEGIN")
try {
conn.exec(v1Ddl)
conn.exec(v1IndexesDdl)
conn.exec("COMMIT")
} catch (t: Throwable) {
runCatching { conn.exec("ROLLBACK") }
throw t
}
}
}
@@ -0,0 +1,93 @@
package pw.binom.agentik.reflection.ksqlite
import kotlinx.coroutines.test.runTest
import pw.binom.agentik.reflection.Reflection
import kotlin.test.AfterTest
import kotlin.test.BeforeTest
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertNotNull
import kotlin.test.assertNull
import kotlin.time.Instant
class KsqliteReflectionStoreTest {
private lateinit var store: KsqliteReflectionStore
@BeforeTest
fun setup() {
store = KsqliteReflectionStore.memory("refl-${kotlin.random.Random.nextLong()}")
}
@AfterTest
fun tearDown() = store.close()
private fun refl(
id: String,
conv: String? = null,
ts: Instant = Instant.parse("2026-09-15T10:00:00Z"),
score: Int = 3,
summary: String = "ok",
weak: List<String> = listOf("weak1", "weak2"),
) = Reflection(
id = id,
conversationId = conv,
createdAt = ts,
turnsAnalyzed = 5,
score = score,
summary = summary,
weakSpots = weak,
)
@Test
fun testInsertAndGetRoundtrip() = runTest {
store.insert(refl("r1"))
val got = store.get("r1")
assertNotNull(got)
assertEquals("r1", got.id)
assertEquals(listOf("weak1", "weak2"), got.weakSpots)
}
@Test
fun testGetReturnsNullForMissing() = runTest {
assertNull(store.get("nope"))
}
@Test
fun testListRecentSortedByCreatedAtDesc() = runTest {
val t0 = Instant.parse("2026-09-15T10:00:00Z")
store.insert(refl("r1", ts = t0))
store.insert(refl("r2", ts = t0 + kotlin.time.Duration.parse("PT60S")))
store.insert(refl("r3", ts = t0 + kotlin.time.Duration.parse("PT120S")))
val recent = store.listRecent(limit = 3)
assertEquals(listOf("r3", "r2", "r1"), recent.map { it.id })
}
@Test
fun testListForConversationFilters() = runTest {
val t = Instant.parse("2026-09-15T10:00:00Z")
store.insert(refl("r1", conv = "c-1"))
store.insert(refl("r2", conv = "c-2"))
store.insert(refl("r3", conv = "c-1"))
val forC1 = store.listForConversation("c-1", limit = 10)
assertEquals(2, forC1.size)
}
@Test
fun testDeleteOlderThanRemoves() = runTest {
val t0 = Instant.parse("2026-09-15T10:00:00Z")
store.insert(refl("r1", ts = t0))
store.insert(refl("r2", ts = t0 + kotlin.time.Duration.parse("PT1H")))
store.deleteOlderThan(t0 + kotlin.time.Duration.parse("PT30M"))
assertNull(store.get("r1"))
assertNotNull(store.get("r2"))
}
@Test
fun testCountReturnsTotal() = runTest {
assertEquals(0, store.count())
store.insert(refl("r1"))
store.insert(refl("r2"))
assertEquals(2, store.count())
}
}
@@ -0,0 +1,64 @@
package pw.binom.agentik.reflection.ksqlite
import kotlinx.coroutines.test.runTest
import pw.binom.db.ksqlite.SQLiteConnection
import kotlin.test.Test
import kotlin.test.assertTrue
/**
* Тесты на Schema.migrate() в `:reflection-ksqlite`:
* - fresh DB → создаются `reflection` + индексы;
* - уже мигрированная БД → migrate() идемпотентен (no-op).
*/
class SchemaMigrationTest {
@Test
fun `fresh DB gets reflection table and indexes`() = runTest {
val conn = SQLiteConnection.memory("refl-mig-fresh-${kotlin.random.Random.nextLong()}")
try {
Schema.migrate(conn)
assertTrue(tableExists(conn, Schema.TABLE_REFLECTION), "table '${Schema.TABLE_REFLECTION}' should exist after migrate()")
for (index in listOf(
Schema.IDX_REFLECTION_CREATED,
Schema.IDX_REFLECTION_CONV,
)) {
assertTrue(indexExists(conn, index), "index '$index' should exist after migrate()")
}
} finally {
conn.close()
}
}
@Test
fun `migrate is idempotent on already-migrated DB`() = runTest {
val conn = SQLiteConnection.memory("refl-mig-idem-${kotlin.random.Random.nextLong()}")
try {
Schema.migrate(conn)
// повторный вызов не должен ни упасть, ни пересоздать таблицы
Schema.migrate(conn)
Schema.migrate(conn)
assertTrue(tableExists(conn, Schema.TABLE_REFLECTION))
} finally {
conn.close()
}
}
private fun tableExists(conn: SQLiteConnection, name: String): Boolean {
conn.prepare(
"SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = ?"
).use { stmt ->
stmt.bindText(1, name)
stmt.executeQuery().use { rs -> return rs.next() }
}
}
private fun indexExists(conn: SQLiteConnection, name: String): Boolean {
conn.prepare(
"SELECT 1 FROM sqlite_master WHERE type = 'index' AND name = ?"
).use { stmt ->
stmt.bindText(1, name)
stmt.executeQuery().use { rs -> return rs.next() }
}
}
}