storage-inmemory: in-memory импл 4 store'ов
Новый KMP-модуль :storage-inmemory с тред-безопасными (Mutex) in-memory имплами для всех 4 store'ов из :storage-core: - InMemoryConversationStore (Map по id, sortedByDescending(updatedAt)) - InMemoryMessageStore (List per conversation, с tokenStats) - InMemoryWorkingMemoryStore (List с order_idx, atomic compact + summary) - InMemoryReflectionStore (List по conversationId, FIFO для listRecent) Фабрика InMemoryStorage.create() возвращает готовый StorageBundle. Семантика 1:1 с SQLite-имплами — параллельные suspend-вызовы атомарны через kotlinx.coroutines.sync.Mutex (lost-update race невозможен, как в SQLite-driver-locked версии). Тесты: 35 новых, проверяют round-trip всех CRUD-операций, тред-безопасность, tokenStats агрегацию, compact+summary, events-flow. Планируемое использование: - :standalone тесты (вместо SqliteStores.inMemory() с JDBC) - Android ART-сборка (commit 7+; SQLite требует JDBC драйвера, недоступного в Android base classes) - embedded/cold-start сценарии без SQLite-инициализации Tests: 299/299 green (264 ранее + 35 новых в :storage-inmemory)
This commit is contained in:
+64
@@ -0,0 +1,64 @@
|
||||
package pw.binom.agentik.storage.inmemory
|
||||
|
||||
import kotlinx.coroutines.sync.Mutex
|
||||
import kotlinx.coroutines.sync.withLock
|
||||
import kotlin.time.Clock
|
||||
import pw.binom.agentik.storage.ConversationRecord
|
||||
import pw.binom.agentik.storage.ConversationStore
|
||||
import kotlin.time.Instant
|
||||
|
||||
/**
|
||||
* Thread-safe Map-импл [ConversationStore].
|
||||
*
|
||||
* Использует `Mutex` для атомарности read-modify-write операций
|
||||
* (rename, touch) — иначе два параллельных `rename` могут потерять обновления
|
||||
* (lost-update race), что в SQLite невозможно из-за driver-level locking.
|
||||
*/
|
||||
class InMemoryConversationStore(
|
||||
private val clock: Clock = Clock.System,
|
||||
) : ConversationStore {
|
||||
|
||||
private val byId: MutableMap<String, ConversationRecord> = mutableMapOf()
|
||||
private val mutex = Mutex()
|
||||
|
||||
override suspend fun upsert(record: ConversationRecord) {
|
||||
mutex.withLock { byId[record.id] = record }
|
||||
}
|
||||
|
||||
override suspend fun get(id: String): ConversationRecord? {
|
||||
mutex.withLock { return byId[id] }
|
||||
}
|
||||
|
||||
override suspend fun delete(id: String): Boolean {
|
||||
mutex.withLock { return byId.remove(id) != null }
|
||||
}
|
||||
|
||||
override suspend fun list(offset: Int, limit: Int): List<ConversationRecord> {
|
||||
mutex.withLock {
|
||||
return byId.values
|
||||
.sortedByDescending { it.updatedAt }
|
||||
.drop(offset)
|
||||
.take(limit)
|
||||
}
|
||||
}
|
||||
|
||||
override suspend fun rename(id: String, title: String?): Instant? {
|
||||
mutex.withLock {
|
||||
val existing = byId[id] ?: return null
|
||||
val now = clock.now()
|
||||
byId[id] = existing.copy(title = title, updatedAt = now)
|
||||
return now
|
||||
}
|
||||
}
|
||||
|
||||
override suspend fun touch(id: String, now: Instant) {
|
||||
mutex.withLock {
|
||||
val existing = byId[id] ?: return
|
||||
byId[id] = existing.copy(updatedAt = now)
|
||||
}
|
||||
}
|
||||
|
||||
override fun close() {
|
||||
// no-op
|
||||
}
|
||||
}
|
||||
+75
@@ -0,0 +1,75 @@
|
||||
package pw.binom.agentik.storage.inmemory
|
||||
|
||||
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 pw.binom.agentik.storage.MessageEvent
|
||||
import pw.binom.agentik.storage.MessageRecord
|
||||
import pw.binom.agentik.storage.MessageStore
|
||||
import pw.binom.agentik.storage.TokenStats
|
||||
import kotlin.time.Instant
|
||||
|
||||
/**
|
||||
* Thread-safe append-only лог сообщений в памяти.
|
||||
*
|
||||
* Хранит все записи в одном `List<MessageRecord>`, индексированном
|
||||
* `conversationId`. `tokenStats` обходит только assistant-записи и
|
||||
* суммирует `TurnTokens`.
|
||||
*
|
||||
* В отличие от SQLite-импла, не нуждается в SQLDelight и работает в
|
||||
* любом KMP-таргете (включая iOS/native, где SQLite через Android driver
|
||||
* недоступен).
|
||||
*/
|
||||
class InMemoryMessageStore : MessageStore {
|
||||
|
||||
private val byConv: MutableMap<String, MutableList<MessageRecord>> = mutableMapOf()
|
||||
private val events = MutableSharedFlow<MessageEvent>(extraBufferCapacity = 64)
|
||||
private val mutex = Mutex()
|
||||
|
||||
override suspend fun append(record: MessageRecord) {
|
||||
mutex.withLock {
|
||||
val list = byConv.getOrPut(record.conversationId) { mutableListOf() }
|
||||
list.add(record)
|
||||
}
|
||||
events.tryEmit(MessageEvent.Appended(record.conversationId, record))
|
||||
}
|
||||
|
||||
override suspend fun list(conversationId: String, after: Instant, offset: Int, limit: Int): List<MessageRecord> {
|
||||
mutex.withLock {
|
||||
val all = byConv[conversationId].orEmpty()
|
||||
val filtered = all.filter { it.createdAt > after }
|
||||
.sortedBy { it.createdAt }
|
||||
return filtered.drop(offset).take(limit)
|
||||
}
|
||||
}
|
||||
|
||||
override suspend fun listAll(conversationId: String): List<MessageRecord> {
|
||||
mutex.withLock {
|
||||
return byConv[conversationId].orEmpty().sortedBy { it.createdAt }.toList()
|
||||
}
|
||||
}
|
||||
|
||||
override fun events(): Flow<MessageEvent> = events.asSharedFlow()
|
||||
|
||||
override suspend fun tokenStats(conversationId: String): TokenStats {
|
||||
mutex.withLock {
|
||||
var input = 0L
|
||||
var output = 0L
|
||||
var turns = 0
|
||||
for (rec in byConv[conversationId].orEmpty()) {
|
||||
if (rec !is MessageRecord.AssistantMessage) continue
|
||||
val t = rec.tokens ?: continue
|
||||
input += t.input
|
||||
output += t.output
|
||||
turns += 1
|
||||
}
|
||||
return TokenStats(turns = turns, inputTokens = input, outputTokens = output)
|
||||
}
|
||||
}
|
||||
|
||||
override fun close() {
|
||||
// no-op
|
||||
}
|
||||
}
|
||||
+79
@@ -0,0 +1,79 @@
|
||||
package pw.binom.agentik.storage.inmemory
|
||||
|
||||
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 pw.binom.agentik.storage.Reflection
|
||||
import pw.binom.agentik.storage.ReflectionEvent
|
||||
import pw.binom.agentik.storage.ReflectionStore
|
||||
import kotlin.time.Instant
|
||||
|
||||
/**
|
||||
* Thread-safe List-импл [ReflectionStore].
|
||||
*
|
||||
* `weakSpots` хранятся как List<String> в самой структуре — JSON-сериализация
|
||||
* делается на уровне SQLDelight в SQLite-импле; здесь просто держим в памяти.
|
||||
*/
|
||||
class InMemoryReflectionStore : ReflectionStore {
|
||||
|
||||
private val byId: MutableMap<String, Reflection> = mutableMapOf()
|
||||
private val byConv: MutableMap<String, MutableList<String>> = mutableMapOf()
|
||||
private val all: MutableList<String> = mutableListOf() // FIFO порядок insert (для listRecent)
|
||||
private val events = MutableSharedFlow<ReflectionEvent>(extraBufferCapacity = 16)
|
||||
private val mutex = Mutex()
|
||||
|
||||
override suspend fun insert(reflection: Reflection) {
|
||||
mutex.withLock {
|
||||
byId[reflection.id] = reflection
|
||||
byConv.getOrPut(reflection.conversationId ?: "") { mutableListOf() }.add(reflection.id)
|
||||
all.add(reflection.id)
|
||||
}
|
||||
events.tryEmit(ReflectionEvent.Created(reflection))
|
||||
}
|
||||
|
||||
override suspend fun get(id: String): Reflection? {
|
||||
mutex.withLock { return byId[id] }
|
||||
}
|
||||
|
||||
override suspend fun listRecent(limit: Int): List<Reflection> {
|
||||
mutex.withLock {
|
||||
return all.asReversed().asSequence()
|
||||
.mapNotNull { byId[it] }
|
||||
.take(limit)
|
||||
.toList()
|
||||
}
|
||||
}
|
||||
|
||||
override suspend fun listForConversation(conversationId: String, limit: Int): List<Reflection> {
|
||||
mutex.withLock {
|
||||
val ids = byConv[conversationId].orEmpty()
|
||||
return ids.asReversed().asSequence()
|
||||
.mapNotNull { byId[it] }
|
||||
.take(limit)
|
||||
.toList()
|
||||
}
|
||||
}
|
||||
|
||||
override suspend fun deleteOlderThan(cutoff: Instant) {
|
||||
mutex.withLock {
|
||||
val toRemove = byId.values.filter { it.createdAt < cutoff }.map { it.id }
|
||||
for (id in toRemove) {
|
||||
byId.remove(id)
|
||||
all.remove(id)
|
||||
byConv.values.forEach { it.remove(id) }
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
override suspend fun count(): Int {
|
||||
mutex.withLock { return byId.size }
|
||||
}
|
||||
|
||||
override fun events(): Flow<ReflectionEvent> = events.asSharedFlow()
|
||||
|
||||
override fun close() {
|
||||
// no-op
|
||||
}
|
||||
}
|
||||
+23
@@ -0,0 +1,23 @@
|
||||
package pw.binom.agentik.storage.inmemory
|
||||
|
||||
import pw.binom.agentik.storage.StorageBundle
|
||||
import kotlin.time.Clock
|
||||
|
||||
/**
|
||||
* Фабрика готового [StorageBundle] на базе in-memory имплов.
|
||||
*
|
||||
* Удобно для:
|
||||
* - тестов (быстрая инициализация, не нужен JDBC driver);
|
||||
* - embedded-сценариев (Android ART, edge-узлы);
|
||||
* - dry-run / preview, где SQLite не нужен.
|
||||
*
|
||||
* Семантика полностью совпадает с SQLite-имплами (`:storage-sqlite`).
|
||||
*/
|
||||
object InMemoryStorage {
|
||||
fun create(clock: Clock = Clock.System): StorageBundle = StorageBundle(
|
||||
conversationStore = InMemoryConversationStore(clock),
|
||||
messageStore = InMemoryMessageStore(),
|
||||
workingMemoryStore = InMemoryWorkingMemoryStore(),
|
||||
reflectionStore = InMemoryReflectionStore(),
|
||||
)
|
||||
}
|
||||
+90
@@ -0,0 +1,90 @@
|
||||
package pw.binom.agentik.storage.inmemory
|
||||
|
||||
import kotlinx.coroutines.sync.Mutex
|
||||
import kotlinx.coroutines.sync.withLock
|
||||
import pw.binom.agentik.storage.Ids
|
||||
import pw.binom.agentik.storage.WorkingMemoryEntry
|
||||
import pw.binom.agentik.storage.WorkingMemoryRow
|
||||
import pw.binom.agentik.storage.WorkingMemoryStore
|
||||
import kotlin.time.Instant
|
||||
|
||||
/**
|
||||
* Thread-safe List-импл [WorkingMemoryStore].
|
||||
*
|
||||
* Append дает `order_idx = max(existing.order_idx) + 1` (или 0 если пусто).
|
||||
* `compact` атомарно удаляет строки в диапазоне `[dropFromOrderIdx, +∞)`
|
||||
* и опционально вставляет [WorkingMemoryEntry.Summary] в хвост.
|
||||
*
|
||||
* Семантика 1:1 с SQLite-имплом (см. `:storage-sqlite` после commit 3).
|
||||
*/
|
||||
class InMemoryWorkingMemoryStore : WorkingMemoryStore {
|
||||
|
||||
private val byConv: MutableMap<String, MutableList<WorkingMemoryRow>> = mutableMapOf()
|
||||
private val mutex = Mutex()
|
||||
|
||||
private suspend fun nextOrderIdx(conversationId: String): Long = mutex.withLock {
|
||||
(byConv[conversationId]?.maxOfOrNull { it.orderIdx } ?: -1L) + 1L
|
||||
}
|
||||
|
||||
override suspend fun append(conversationId: String, entry: WorkingMemoryEntry, now: Instant) {
|
||||
mutex.withLock {
|
||||
val rows = byConv.getOrPut(conversationId) { mutableListOf() }
|
||||
val nextIdx = (rows.maxOfOrNull { it.orderIdx } ?: -1L) + 1L
|
||||
rows.add(
|
||||
WorkingMemoryRow(
|
||||
id = Ids.new("wm"),
|
||||
conversationId = conversationId,
|
||||
orderIdx = nextIdx,
|
||||
sourceMessageId = entry.sourceMessageId,
|
||||
entry = entry,
|
||||
createdAt = now,
|
||||
)
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
override suspend fun list(conversationId: String): List<WorkingMemoryRow> {
|
||||
mutex.withLock {
|
||||
return byConv[conversationId].orEmpty()
|
||||
.sortedBy { it.orderIdx }
|
||||
.toList()
|
||||
}
|
||||
}
|
||||
|
||||
override suspend fun clear(conversationId: String) {
|
||||
mutex.withLock {
|
||||
byConv.remove(conversationId)
|
||||
}
|
||||
}
|
||||
|
||||
override suspend fun compact(
|
||||
dropFromOrderIdx: Long,
|
||||
conversationId: String,
|
||||
summaryText: String?,
|
||||
): Long = mutex.withLock {
|
||||
val rows = byConv.getOrPut(conversationId) { mutableListOf() }
|
||||
rows.removeAll { it.orderIdx >= dropFromOrderIdx }
|
||||
val newMax = rows.maxOfOrNull { it.orderIdx } ?: -1L
|
||||
return if (summaryText.isNullOrBlank()) {
|
||||
newMax
|
||||
} else {
|
||||
val summary = WorkingMemoryEntry.Summary(text = summaryText)
|
||||
val newIdx = newMax + 1L
|
||||
rows.add(
|
||||
WorkingMemoryRow(
|
||||
id = Ids.new("wm"),
|
||||
conversationId = conversationId,
|
||||
orderIdx = newIdx,
|
||||
sourceMessageId = null,
|
||||
entry = summary,
|
||||
createdAt = kotlin.time.Clock.System.now(),
|
||||
)
|
||||
)
|
||||
newIdx
|
||||
}
|
||||
}
|
||||
|
||||
override fun close() {
|
||||
// no-op
|
||||
}
|
||||
}
|
||||
+119
@@ -0,0 +1,119 @@
|
||||
package pw.binom.agentik.storage.inmemory
|
||||
|
||||
import pw.binom.agentik.storage.ConversationRecord
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertNotNull
|
||||
import kotlin.test.assertNull
|
||||
import kotlin.test.assertTrue
|
||||
import kotlin.time.Instant
|
||||
import kotlinx.coroutines.test.runTest
|
||||
|
||||
class InMemoryConversationStoreTest {
|
||||
|
||||
@Test
|
||||
fun `upsert and get roundtrip preserves all fields`() = runTest {
|
||||
val store = InMemoryConversationStore()
|
||||
val rec = ConversationRecord(
|
||||
id = "c1",
|
||||
title = "test",
|
||||
isTemporal = false,
|
||||
createdAt = Instant.parse("2026-09-15T10:00:00Z"),
|
||||
updatedAt = Instant.parse("2026-09-15T10:00:00Z"),
|
||||
)
|
||||
store.upsert(rec)
|
||||
val got = store.get("c1")
|
||||
assertEquals(rec, got)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `get returns null for missing id`() = runTest {
|
||||
val store = InMemoryConversationStore()
|
||||
assertNull(store.get("nope"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `delete removes the record and returns true`() = runTest {
|
||||
val store = InMemoryConversationStore()
|
||||
store.upsert(
|
||||
ConversationRecord(
|
||||
"c1", null, false,
|
||||
Instant.parse("2026-09-15T10:00:00Z"),
|
||||
Instant.parse("2026-09-15T10:00:00Z"),
|
||||
)
|
||||
)
|
||||
assertTrue(store.delete("c1"))
|
||||
assertNull(store.get("c1"))
|
||||
// повторный delete — false
|
||||
assertEquals(false, store.delete("c1"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `list sorts by updatedAt DESC and respects offset+limit`() = runTest {
|
||||
val store = InMemoryConversationStore()
|
||||
val t0 = Instant.parse("2026-09-15T10:00:00Z")
|
||||
store.upsert(ConversationRecord("c1", null, false, t0, t0))
|
||||
store.upsert(ConversationRecord("c2", null, false, t0, t0.plus(kotlin.time.Duration.parse("PT60S"))))
|
||||
store.upsert(ConversationRecord("c3", null, false, t0, t0.plus(kotlin.time.Duration.parse("PT120S"))))
|
||||
store.upsert(ConversationRecord("c4", null, false, t0, t0.plus(kotlin.time.Duration.parse("PT180S"))))
|
||||
|
||||
val page0 = store.list(offset = 0, limit = 2)
|
||||
assertEquals(listOf("c4", "c3"), page0.map { it.id })
|
||||
|
||||
val page1 = store.list(offset = 2, limit = 2)
|
||||
assertEquals(listOf("c2", "c1"), page1.map { it.id })
|
||||
|
||||
val page2 = store.list(offset = 4, limit = 2)
|
||||
assertEquals(emptyList(), page2)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `rename updates title and updatedAt, returns new updatedAt`() = runTest {
|
||||
val store = InMemoryConversationStore()
|
||||
val t0 = Instant.parse("2026-09-15T10:00:00Z")
|
||||
store.upsert(ConversationRecord("c1", null, false, t0, t0))
|
||||
|
||||
val newUpdated = store.rename("c1", "new title")
|
||||
assertNotNull(newUpdated)
|
||||
assertTrue(newUpdated > t0)
|
||||
|
||||
val got = store.get("c1")
|
||||
assertEquals("new title", got?.title)
|
||||
assertEquals(newUpdated, got?.updatedAt)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `rename with null title clears it`() = runTest {
|
||||
val store = InMemoryConversationStore()
|
||||
val t0 = Instant.parse("2026-09-15T10:00:00Z")
|
||||
store.upsert(ConversationRecord("c1", "old", false, t0, t0))
|
||||
store.rename("c1", null)
|
||||
assertNull(store.get("c1")?.title)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `rename returns null for missing conversation`() = runTest {
|
||||
val store = InMemoryConversationStore()
|
||||
assertNull(store.rename("nope", "x"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `touch bumps updatedAt without changing other fields`() = runTest {
|
||||
val store = InMemoryConversationStore()
|
||||
val t0 = Instant.parse("2026-09-15T10:00:00Z")
|
||||
val t1 = Instant.parse("2026-09-15T10:01:00Z")
|
||||
store.upsert(ConversationRecord("c1", "title", false, t0, t0))
|
||||
store.touch("c1", t1)
|
||||
val got = store.get("c1")
|
||||
assertEquals("title", got?.title)
|
||||
assertEquals(t1, got?.updatedAt)
|
||||
assertEquals(t0, got?.createdAt)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `close is idempotent and does nothing`() {
|
||||
val store = InMemoryConversationStore()
|
||||
store.close()
|
||||
store.close() // должно быть no-op
|
||||
}
|
||||
}
|
||||
+154
@@ -0,0 +1,154 @@
|
||||
package pw.binom.agentik.storage.inmemory
|
||||
|
||||
import pw.binom.agentik.storage.Content
|
||||
import pw.binom.agentik.storage.MessageRecord
|
||||
import pw.binom.agentik.storage.TurnTokens
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertNull
|
||||
import kotlin.test.assertTrue
|
||||
import kotlin.time.Instant
|
||||
import kotlinx.coroutines.async
|
||||
import kotlinx.coroutines.flow.first
|
||||
import kotlinx.coroutines.yield
|
||||
import kotlinx.coroutines.test.runTest
|
||||
|
||||
class InMemoryMessageStoreTest {
|
||||
|
||||
@Test
|
||||
fun `append and listAll returns inserted records in createdAt order`() = runTest {
|
||||
val store = InMemoryMessageStore()
|
||||
val t0 = Instant.parse("2026-09-15T10:00:00Z")
|
||||
val u = MessageRecord.UserMessage("u1", "c1", listOf(Content.Text("hi")), t0)
|
||||
val a = MessageRecord.AssistantMessage("a1", "c1", listOf(Content.Text("hello")), t0.plus(kotlin.time.Duration.parse("PT1S")), null)
|
||||
store.append(u)
|
||||
store.append(a)
|
||||
val all = store.listAll("c1")
|
||||
assertEquals(listOf("u1", "a1"), all.map { it.id })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `list filters by after and supports offset+limit`() = runTest {
|
||||
val store = InMemoryMessageStore()
|
||||
val t0 = Instant.parse("2026-09-15T10:00:00Z")
|
||||
for (i in 0 until 5) {
|
||||
store.append(
|
||||
MessageRecord.UserMessage(
|
||||
"u$i",
|
||||
"c1",
|
||||
listOf(Content.Text("msg-$i")),
|
||||
t0.plus(kotlin.time.Duration.parse("PT${i}S")),
|
||||
)
|
||||
)
|
||||
}
|
||||
// after=t0+1s должны видеть только msg-2..4 (т.е. u2,u3,u4)
|
||||
val after = t0.plus(kotlin.time.Duration.parse("PT1S"))
|
||||
val page = store.list("c1", after = after, offset = 0, limit = 10)
|
||||
assertEquals(listOf("u2", "u3", "u4"), page.map { it.id })
|
||||
|
||||
val page2 = store.list("c1", after = after, offset = 1, limit = 10)
|
||||
assertEquals(listOf("u3", "u4"), page2.map { it.id })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `list returns empty for unknown conversation`() = runTest {
|
||||
val store = InMemoryMessageStore()
|
||||
assertEquals(emptyList(), store.listAll("none"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `tokenStats sums across assistant messages with tokens`() = runTest {
|
||||
val store = InMemoryMessageStore()
|
||||
val t0 = Instant.parse("2026-09-15T10:00:00Z")
|
||||
store.append(
|
||||
MessageRecord.AssistantMessage(
|
||||
"a1", "c1", listOf(Content.Text("r")), t0.plus(kotlin.time.Duration.parse("PT1S")),
|
||||
TurnTokens(input = 100, output = 50),
|
||||
)
|
||||
)
|
||||
store.append(
|
||||
MessageRecord.AssistantMessage(
|
||||
"a2", "c1", listOf(Content.Text("r")), t0.plus(kotlin.time.Duration.parse("PT2S")),
|
||||
TurnTokens(input = 200, output = 80),
|
||||
)
|
||||
)
|
||||
// user без tokens
|
||||
store.append(
|
||||
MessageRecord.UserMessage("u1", "c1", listOf(Content.Text("hi")), t0)
|
||||
)
|
||||
val stats = store.tokenStats("c1")
|
||||
assertEquals(2, stats.turns)
|
||||
assertEquals(300L, stats.inputTokens)
|
||||
assertEquals(130L, stats.outputTokens)
|
||||
assertEquals(430L, stats.totalTokens)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `tokenStats skips assistant messages without tokens`() = runTest {
|
||||
val store = InMemoryMessageStore()
|
||||
val t0 = Instant.parse("2026-09-15T10:00:00Z")
|
||||
store.append(
|
||||
MessageRecord.AssistantMessage("a1", "c1", listOf(Content.Text("r")), t0, tokens = null)
|
||||
)
|
||||
val stats = store.tokenStats("c1")
|
||||
assertEquals(0, stats.turns)
|
||||
assertEquals(0L, stats.inputTokens)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `tokenStats returns zeros for unknown conversation`() = runTest {
|
||||
val store = InMemoryMessageStore()
|
||||
val stats = store.tokenStats("none")
|
||||
assertEquals(0, stats.turns)
|
||||
assertEquals(0L, stats.inputTokens)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `events flow emits Appended on append`() = runTest {
|
||||
val store = InMemoryMessageStore()
|
||||
val t0 = Instant.parse("2026-09-15T10:00:00Z")
|
||||
// SharedFlow не реплеит — запускаем коллектор ДО append, чтобы не потерять эвент.
|
||||
val events = store.events()
|
||||
val deferred = async { events.first() }
|
||||
yield() // даём коллектору подписаться ДО append — иначе SharedFlow без replay потеряет эвент
|
||||
store.append(MessageRecord.UserMessage("u1", "c1", listOf(Content.Text("hi")), t0))
|
||||
val ev = deferred.await()
|
||||
assertTrue(ev is pw.binom.agentik.storage.MessageEvent.Appended)
|
||||
val appended = ev as pw.binom.agentik.storage.MessageEvent.Appended
|
||||
assertEquals("c1", appended.conversationId)
|
||||
assertEquals("u1", appended.record.id)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `isolates conversations - list returns only requested conv`() = runTest {
|
||||
val store = InMemoryMessageStore()
|
||||
val t0 = Instant.parse("2026-09-15T10:00:00Z")
|
||||
store.append(MessageRecord.UserMessage("u1", "c1", listOf(Content.Text("hi")), t0))
|
||||
store.append(MessageRecord.UserMessage("u2", "c2", listOf(Content.Text("hello")), t0))
|
||||
assertEquals(listOf("u1"), store.listAll("c1").map { it.id })
|
||||
assertEquals(listOf("u2"), store.listAll("c2").map { it.id })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `close is idempotent and does nothing`() {
|
||||
val store = InMemoryMessageStore()
|
||||
store.close()
|
||||
store.close()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `append with null context and null tokens is supported`() = runTest {
|
||||
val store = InMemoryMessageStore()
|
||||
val t0 = Instant.parse("2026-09-15T10:00:00Z")
|
||||
val msg = MessageRecord.UserMessage(
|
||||
id = "u1",
|
||||
conversationId = "c1",
|
||||
content = listOf(Content.Text("hi")),
|
||||
createdAt = t0,
|
||||
context = null,
|
||||
)
|
||||
store.append(msg)
|
||||
val got = store.listAll("c1").first() as MessageRecord.UserMessage
|
||||
assertNull(got.context)
|
||||
}
|
||||
}
|
||||
+118
@@ -0,0 +1,118 @@
|
||||
package pw.binom.agentik.storage.inmemory
|
||||
|
||||
import pw.binom.agentik.storage.Reflection
|
||||
import pw.binom.agentik.storage.ReflectionEvent
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertNotNull
|
||||
import kotlin.test.assertNull
|
||||
import kotlin.test.assertTrue
|
||||
import kotlin.time.Instant
|
||||
import kotlinx.coroutines.async
|
||||
import kotlinx.coroutines.flow.first
|
||||
import kotlinx.coroutines.yield
|
||||
import kotlinx.coroutines.test.runTest
|
||||
|
||||
class InMemoryReflectionStoreTest {
|
||||
|
||||
private fun sample(id: String, convId: String?, at: Instant, score: Int = 4) = Reflection(
|
||||
id = id,
|
||||
conversationId = convId,
|
||||
createdAt = at,
|
||||
turnsAnalyzed = 5,
|
||||
score = score,
|
||||
summary = "ok",
|
||||
weakSpots = listOf("weakness-1"),
|
||||
)
|
||||
|
||||
@Test
|
||||
fun `insert and get roundtrip preserves all fields`() = runTest {
|
||||
val store = InMemoryReflectionStore()
|
||||
val r = sample("r1", "c1", Instant.parse("2026-09-15T10:00:00Z"))
|
||||
store.insert(r)
|
||||
assertEquals(r, store.get("r1"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `get returns null for missing id`() = runTest {
|
||||
val store = InMemoryReflectionStore()
|
||||
assertNull(store.get("nope"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `listRecent returns most-recent first up to limit`() = runTest {
|
||||
val store = InMemoryReflectionStore()
|
||||
val t0 = Instant.parse("2026-09-15T10:00:00Z")
|
||||
for (i in 0 until 5) {
|
||||
store.insert(sample("r$i", "c$i", t0.plus(kotlin.time.Duration.parse("PT${i}S"))))
|
||||
}
|
||||
val top3 = store.listRecent(limit = 3)
|
||||
assertEquals(listOf("r4", "r3", "r2"), top3.map { it.id })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `listForConversation returns only entries for that conversation most-recent first`() = runTest {
|
||||
val store = InMemoryReflectionStore()
|
||||
val t0 = Instant.parse("2026-09-15T10:00:00Z")
|
||||
store.insert(sample("r1", "c1", t0))
|
||||
store.insert(sample("r2", "c2", t0.plus(kotlin.time.Duration.parse("PT1S"))))
|
||||
store.insert(sample("r3", "c1", t0.plus(kotlin.time.Duration.parse("PT2S"))))
|
||||
store.insert(sample("r4", "c1", t0.plus(kotlin.time.Duration.parse("PT3S"))))
|
||||
val c1List = store.listForConversation("c1", limit = 10)
|
||||
assertEquals(listOf("r4", "r3", "r1"), c1List.map { it.id })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `listForConversation with null conversationId returns global ones`() = runTest {
|
||||
val store = InMemoryReflectionStore()
|
||||
val t0 = Instant.parse("2026-09-15T10:00:00Z")
|
||||
store.insert(sample("r1", null, t0))
|
||||
store.insert(sample("r2", "c1", t0.plus(kotlin.time.Duration.parse("PT1S"))))
|
||||
val global = store.listForConversation("", limit = 10)
|
||||
assertEquals(listOf("r1"), global.map { it.id })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `deleteOlderThan removes entries created before cutoff`() = runTest {
|
||||
val store = InMemoryReflectionStore()
|
||||
val t0 = Instant.parse("2026-09-15T10:00:00Z")
|
||||
store.insert(sample("r-old", "c1", t0))
|
||||
store.insert(sample("r-old2", "c1", t0.plus(kotlin.time.Duration.parse("PT10S"))))
|
||||
store.insert(sample("r-new", "c1", t0.plus(kotlin.time.Duration.parse("PT60S"))))
|
||||
store.deleteOlderThan(cutoff = t0.plus(kotlin.time.Duration.parse("PT30S")))
|
||||
assertEquals(1, store.count())
|
||||
assertNotNull(store.get("r-new"))
|
||||
assertNull(store.get("r-old"))
|
||||
assertNull(store.get("r-old2"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `count reflects inserts`() = runTest {
|
||||
val store = InMemoryReflectionStore()
|
||||
assertEquals(0, store.count())
|
||||
store.insert(sample("r1", "c1", Instant.parse("2026-09-15T10:00:00Z")))
|
||||
store.insert(sample("r2", "c1", Instant.parse("2026-09-15T10:01:00Z")))
|
||||
assertEquals(2, store.count())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `events flow emits Created on insert`() = runTest {
|
||||
val store = InMemoryReflectionStore()
|
||||
val r = sample("r1", "c1", Instant.parse("2026-09-15T10:00:00Z"))
|
||||
// SharedFlow не реплеит — запускаем коллектор ДО insert.
|
||||
val events = store.events()
|
||||
val deferred = async { events.first() }
|
||||
yield() // даём коллектору подписаться ДО insert
|
||||
store.insert(r)
|
||||
val ev = deferred.await()
|
||||
assertTrue(ev is ReflectionEvent.Created)
|
||||
assertEquals("r1", (ev as ReflectionEvent.Created).reflection.id)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `close is idempotent and does nothing`() {
|
||||
val store = InMemoryReflectionStore()
|
||||
store.close()
|
||||
store.close()
|
||||
}
|
||||
}
|
||||
+111
@@ -0,0 +1,111 @@
|
||||
package pw.binom.agentik.storage.inmemory
|
||||
|
||||
import pw.binom.agentik.storage.Content
|
||||
import pw.binom.agentik.storage.WorkingMemoryEntry
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertNull
|
||||
import kotlin.test.assertTrue
|
||||
import kotlin.time.Instant
|
||||
import kotlinx.coroutines.test.runTest
|
||||
|
||||
class InMemoryWorkingMemoryStoreTest {
|
||||
|
||||
@Test
|
||||
fun `append assigns sequential order_idx starting from 0`() = runTest {
|
||||
val store = InMemoryWorkingMemoryStore()
|
||||
val t0 = Instant.parse("2026-09-15T10:00:00Z")
|
||||
store.append("c1", WorkingMemoryEntry.System("you are brief"), t0)
|
||||
store.append("c1", WorkingMemoryEntry.User("u1", listOf(Content.Text("hi")), null), t0.plus(kotlin.time.Duration.parse("PT1S")))
|
||||
store.append("c1", WorkingMemoryEntry.Assistant("a1", listOf(Content.Text("hello"))), t0.plus(kotlin.time.Duration.parse("PT2S")))
|
||||
val rows = store.list("c1")
|
||||
assertEquals(3, rows.size)
|
||||
assertEquals(listOf(0L, 1L, 2L), rows.map { it.orderIdx })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `order_idx continues across conversations independently`() = runTest {
|
||||
val store = InMemoryWorkingMemoryStore()
|
||||
val t0 = Instant.parse("2026-09-15T10:00:00Z")
|
||||
store.append("c1", WorkingMemoryEntry.System("a"), t0)
|
||||
store.append("c1", WorkingMemoryEntry.User("u1", listOf(Content.Text("hi")), null), t0)
|
||||
store.append("c2", WorkingMemoryEntry.System("b"), t0)
|
||||
// c2 должен начать с 0, не продолжать c1
|
||||
val rows2 = store.list("c2")
|
||||
assertEquals(1, rows2.size)
|
||||
assertEquals(0L, rows2[0].orderIdx)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `list returns sorted by order_idx ASC`() = runTest {
|
||||
val store = InMemoryWorkingMemoryStore()
|
||||
val t0 = Instant.parse("2026-09-15T10:00:00Z")
|
||||
// store назначает order_idx = max(existing) + 1, поэтому порядок вставки
|
||||
// определяет позицию в списке. Сортировка по order_idx даёт ровно
|
||||
// порядок append'ов.
|
||||
store.append("c1", WorkingMemoryEntry.User("u0", listOf(Content.Text("a")), null), t0)
|
||||
store.append("c1", WorkingMemoryEntry.User("u1", listOf(Content.Text("b")), null), t0.plus(kotlin.time.Duration.parse("PT1S")))
|
||||
store.append("c1", WorkingMemoryEntry.User("u2", listOf(Content.Text("c")), null), t0.plus(kotlin.time.Duration.parse("PT2S")))
|
||||
val rows = store.list("c1")
|
||||
assertEquals(listOf("u0", "u1", "u2"), rows.map { it.entry.sourceMessageId })
|
||||
// а createdAt — это переданный параметр, не пересчитывается
|
||||
assertEquals(listOf(t0, t0.plus(kotlin.time.Duration.parse("PT1S")), t0.plus(kotlin.time.Duration.parse("PT2S"))), rows.map { it.createdAt })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `clear removes all rows for a conversation`() = runTest {
|
||||
val store = InMemoryWorkingMemoryStore()
|
||||
val t0 = Instant.parse("2026-09-15T10:00:00Z")
|
||||
store.append("c1", WorkingMemoryEntry.System("a"), t0)
|
||||
store.append("c2", WorkingMemoryEntry.System("b"), t0)
|
||||
store.clear("c1")
|
||||
assertEquals(emptyList(), store.list("c1"))
|
||||
assertEquals(1, store.list("c2").size)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `compact without summary drops tail and returns new max`() = runTest {
|
||||
val store = InMemoryWorkingMemoryStore()
|
||||
val t0 = Instant.parse("2026-09-15T10:00:00Z")
|
||||
store.append("c1", WorkingMemoryEntry.System("a"), t0)
|
||||
store.append("c1", WorkingMemoryEntry.User("u1", listOf(Content.Text("hi")), null), t0.plus(kotlin.time.Duration.parse("PT1S")))
|
||||
store.append("c1", WorkingMemoryEntry.Assistant("a1", listOf(Content.Text("hello"))), t0.plus(kotlin.time.Duration.parse("PT2S")))
|
||||
// dropFromOrderIdx=2 → удаляет всё >= 2 (то есть только Assistant "a1")
|
||||
val newMax = store.compact(dropFromOrderIdx = 2, conversationId = "c1", summaryText = null)
|
||||
assertEquals(1L, newMax)
|
||||
val rows = store.list("c1")
|
||||
assertEquals(2, rows.size)
|
||||
assertEquals("a", (rows[0].entry as WorkingMemoryEntry.System).text)
|
||||
assertEquals("u1", rows[1].entry.sourceMessageId)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `compact with summary replaces tail with synthetic Summary row`() = runTest {
|
||||
val store = InMemoryWorkingMemoryStore()
|
||||
val t0 = Instant.parse("2026-09-15T10:00:00Z")
|
||||
store.append("c1", WorkingMemoryEntry.System("you are brief"), t0)
|
||||
store.append("c1", WorkingMemoryEntry.User("u1", listOf(Content.Text("hi")), null), t0.plus(kotlin.time.Duration.parse("PT1S")))
|
||||
store.append("c1", WorkingMemoryEntry.Assistant("a1", listOf(Content.Text("hello"))), t0.plus(kotlin.time.Duration.parse("PT2S")))
|
||||
store.append("c1", WorkingMemoryEntry.User("u2", listOf(Content.Text("how are you")), null), t0.plus(kotlin.time.Duration.parse("PT3S")))
|
||||
store.append("c1", WorkingMemoryEntry.Assistant("a2", listOf(Content.Text("fine, thanks"))), t0.plus(kotlin.time.Duration.parse("PT4S")))
|
||||
// drop tail from idx 2, insert summary
|
||||
val newMax = store.compact(
|
||||
dropFromOrderIdx = 2,
|
||||
conversationId = "c1",
|
||||
summaryText = "user asked hi and how-are-you, assistant replied",
|
||||
)
|
||||
assertEquals(2L, newMax)
|
||||
val rows = store.list("c1")
|
||||
assertEquals(3, rows.size)
|
||||
val summary = rows.last().entry as WorkingMemoryEntry.Summary
|
||||
assertEquals("user asked hi and how-are-you, assistant replied", summary.text)
|
||||
assertNull(rows.last().sourceMessageId)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `close is idempotent and does nothing`() {
|
||||
val store = InMemoryWorkingMemoryStore()
|
||||
store.close()
|
||||
store.close()
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user