remove deprecated EventStore and MessageStore implementations, along with related in-memory and SQLite code
ci / JVM build + tests (push) Failing after 59s
ci / JVM build + tests (push) Failing after 59s
This commit is contained in:
@@ -22,8 +22,7 @@ kotlin {
|
|||||||
sourceSets {
|
sourceSets {
|
||||||
commonMain.dependencies {
|
commonMain.dependencies {
|
||||||
api(project(":message-store-api"))
|
api(project(":message-store-api"))
|
||||||
api(project(":message-log-api"))
|
api(project(":context-api"))
|
||||||
api(project(":working-memory-api"))
|
|
||||||
|
|
||||||
// litert-kmp: LiteTool интерфейс (sync describe/invoke)
|
// litert-kmp: LiteTool интерфейс (sync describe/invoke)
|
||||||
api(libs.litert.api)
|
api(libs.litert.api)
|
||||||
|
|||||||
@@ -21,7 +21,7 @@ kotlin {
|
|||||||
sourceSets {
|
sourceSets {
|
||||||
commonMain.dependencies {
|
commonMain.dependencies {
|
||||||
api(project(":proto"))
|
api(project(":proto"))
|
||||||
api(project(":event-store"))
|
api(project(":outbox-api"))
|
||||||
api(project(":journal-api"))
|
api(project(":journal-api"))
|
||||||
|
|
||||||
api(libs.ktor.client.core)
|
api(libs.ktor.client.core)
|
||||||
|
|||||||
@@ -6,7 +6,7 @@ import io.ktor.client.statement.bodyAsChannel
|
|||||||
import io.ktor.http.HttpStatusCode
|
import io.ktor.http.HttpStatusCode
|
||||||
import kotlinx.coroutines.flow.Flow
|
import kotlinx.coroutines.flow.Flow
|
||||||
import kotlinx.coroutines.flow.flow
|
import kotlinx.coroutines.flow.flow
|
||||||
import pw.binom.agentik.eventStore.EventStore
|
import pw.binom.agentik.outbox.OutboxStore
|
||||||
import pw.binom.agentik.outbox.AgentEvent
|
import pw.binom.agentik.outbox.AgentEvent
|
||||||
import pw.binom.agentik.outbox.CommonEvent
|
import pw.binom.agentik.outbox.CommonEvent
|
||||||
import pw.binom.agentik.outbox.Event
|
import pw.binom.agentik.outbox.Event
|
||||||
@@ -14,7 +14,7 @@ import kotlin.time.Clock
|
|||||||
import kotlin.time.Instant
|
import kotlin.time.Instant
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* HTTP-реализация [EventStore] (= [pw.binom.agentik.outbox.OutboxStore]),
|
* HTTP-реализация [OutboxStore] (= [pw.binom.agentik.outbox.OutboxStore]),
|
||||||
* ходящая в `:server`-фасад.
|
* ходящая в `:server`-фасад.
|
||||||
*
|
*
|
||||||
* **Endpoint-раскладка** (новый дизайн — storage handles на [Agent]):
|
* **Endpoint-раскладка** (новый дизайн — storage handles на [Agent]):
|
||||||
@@ -26,12 +26,12 @@ import kotlin.time.Instant
|
|||||||
* - [conversationEvents] с `conversationId != null` → `GET /conversations/{id}/events`
|
* - [conversationEvents] с `conversationId != null` → `GET /conversations/{id}/events`
|
||||||
*
|
*
|
||||||
* Для [conversationEvents] с `conversationId == null` (события всех диалогов)
|
* Для [conversationEvents] с `conversationId == null` (события всех диалогов)
|
||||||
* fallback на default [EventStore.conversationEvents] — общий поток
|
* fallback на default [OutboxStore.conversationEvents] — общий поток
|
||||||
* `/outbox/events` + filter. Это редкий кейс (admin-дашборды), и
|
* `/outbox/events` + filter. Это редкий кейс (admin-дашборды), и
|
||||||
* оптимизировать его отдельно нерационально.
|
* оптимизировать его отдельно нерационально.
|
||||||
*
|
*
|
||||||
* [earliestEventDate] не имеет своего endpoint'а; возвращает `Clock.System.now()`
|
* [earliestEventDate] не имеет своего endpoint'а; возвращает `Clock.System.now()`
|
||||||
* (см. KDoc [EventStore.earliestEventDate] — для пустого буфера это и есть
|
* (см. KDoc [OutboxStore.earliestEventDate] — для пустого буфера это и есть
|
||||||
* контрактное значение). Клиент, который полагался на gap detection через
|
* контрактное значение). Клиент, который полагался на gap detection через
|
||||||
* message store, продолжит работать — просто fallback никогда не сработает.
|
* message store, продолжит работать — просто fallback никогда не сработает.
|
||||||
*
|
*
|
||||||
@@ -44,7 +44,7 @@ import kotlin.time.Instant
|
|||||||
internal class HttpEventStore(
|
internal class HttpEventStore(
|
||||||
private val httpClient: HttpClient,
|
private val httpClient: HttpClient,
|
||||||
private val baseUrl: String,
|
private val baseUrl: String,
|
||||||
) : EventStore {
|
) : OutboxStore {
|
||||||
|
|
||||||
private val agentUrl: String = baseUrl.trimEnd('/')
|
private val agentUrl: String = baseUrl.trimEnd('/')
|
||||||
|
|
||||||
|
|||||||
@@ -3,8 +3,9 @@ plugins {
|
|||||||
alias(libs.plugins.kotlin.serialization)
|
alias(libs.plugins.kotlin.serialization)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Public API для mutable working-memory — runtime context агента (compaction,
|
// Public API для runtime context агента (compaction, order_idx, summary entries).
|
||||||
// order_idx, summary entries). Зависит от :message-log-api для Ids (wm- префикс).
|
// Зависит от :journal-api для типов `Content` / `MessageContext` (audit-log
|
||||||
|
// payload'ы, которые рабочая память ссылает).
|
||||||
//
|
//
|
||||||
// НЕ нужен тонким клиентам — только серверному рантайму (`:standalone`, `:agentik-cli`,
|
// НЕ нужен тонким клиентам — только серверному рантайму (`:standalone`, `:agentik-cli`,
|
||||||
// будущий `:android-agent` core).
|
// будущий `:android-agent` core).
|
||||||
@@ -24,7 +25,7 @@ kotlin {
|
|||||||
|
|
||||||
sourceSets {
|
sourceSets {
|
||||||
commonMain.dependencies {
|
commonMain.dependencies {
|
||||||
api(project(":message-log-api"))
|
api(project(":journal-api"))
|
||||||
api(libs.kotlinx.coroutines.core)
|
api(libs.kotlinx.coroutines.core)
|
||||||
api(libs.kotlinx.serialization.core)
|
api(libs.kotlinx.serialization.core)
|
||||||
api(libs.kotlinx.serialization.json)
|
api(libs.kotlinx.serialization.json)
|
||||||
|
|||||||
@@ -2,8 +2,8 @@ package pw.binom.agentik.context
|
|||||||
|
|
||||||
import kotlinx.serialization.SerialName
|
import kotlinx.serialization.SerialName
|
||||||
import kotlinx.serialization.Serializable
|
import kotlinx.serialization.Serializable
|
||||||
import pw.binom.agentik.messageLog.Content
|
import pw.binom.agentik.journal.Content
|
||||||
import pw.binom.agentik.messageLog.MessageContext
|
import pw.binom.agentik.journal.MessageContext
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Запись в working memory диалога: ровно то, что агент сейчас видит в
|
* Запись в working memory диалога: ровно то, что агент сейчас видит в
|
||||||
|
|||||||
@@ -30,7 +30,7 @@ kotlin {
|
|||||||
// :message-log-api (старый canonical). Транзитивно через api,
|
// :message-log-api (старый canonical). Транзитивно через api,
|
||||||
// но фиксируем явно чтобы тестовый код видел Content без
|
// но фиксируем явно чтобы тестовый код видел Content без
|
||||||
// обхода через :context-api.
|
// обхода через :context-api.
|
||||||
api(project(":message-log-api"))
|
// api(project(":message-log-api"))
|
||||||
}
|
}
|
||||||
commonTest.dependencies {
|
commonTest.dependencies {
|
||||||
implementation(kotlin("test"))
|
implementation(kotlin("test"))
|
||||||
|
|||||||
@@ -1,29 +0,0 @@
|
|||||||
plugins {
|
|
||||||
alias(libs.plugins.kotlin.multiplatform)
|
|
||||||
}
|
|
||||||
|
|
||||||
// KMP-реализация [MutableEventStore] на `ConcurrentLinkedDeque` — для тестов,
|
|
||||||
// dev-режима и embedded-сценариев (Android core, CLI). TTL и size-cap eviction
|
|
||||||
// вызываются на каждом [append], в одном проходе с amortized O(1) для стабильного
|
|
||||||
// размера буфера.
|
|
||||||
//
|
|
||||||
// Зависимости: только `:event-store` (api → `:proto` транзитивно).
|
|
||||||
// Никакого I/O — pure in-memory.
|
|
||||||
|
|
||||||
kotlin {
|
|
||||||
jvmToolchain(21)
|
|
||||||
|
|
||||||
jvm()
|
|
||||||
linuxX64()
|
|
||||||
mingwX64()
|
|
||||||
|
|
||||||
sourceSets {
|
|
||||||
commonMain.dependencies {
|
|
||||||
api(project(":event-store"))
|
|
||||||
}
|
|
||||||
commonTest.dependencies {
|
|
||||||
implementation(kotlin("test"))
|
|
||||||
implementation(libs.kotlinx.coroutines.test)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
-160
@@ -1,160 +0,0 @@
|
|||||||
package pw.binom.agentik.eventStore.inmemory
|
|
||||||
|
|
||||||
import kotlin.time.Clock
|
|
||||||
import kotlin.time.Duration
|
|
||||||
import kotlin.time.Instant
|
|
||||||
import kotlinx.coroutines.channels.BufferOverflow
|
|
||||||
import kotlinx.coroutines.coroutineScope
|
|
||||||
import kotlinx.coroutines.flow.Flow
|
|
||||||
import kotlinx.coroutines.flow.MutableSharedFlow
|
|
||||||
import kotlinx.coroutines.flow.asSharedFlow
|
|
||||||
import kotlinx.coroutines.flow.channelFlow
|
|
||||||
import kotlinx.coroutines.flow.flow
|
|
||||||
import kotlinx.coroutines.sync.Mutex
|
|
||||||
import kotlinx.coroutines.sync.withLock
|
|
||||||
import pw.binom.agentik.eventStore.MutableEventStore
|
|
||||||
import pw.binom.agentik.proto.CommonEvent
|
|
||||||
|
|
||||||
/**
|
|
||||||
* In-memory реализация [MutableEventStore] на `ArrayDeque` + [Mutex].
|
|
||||||
*
|
|
||||||
* **Retention policy** — оба параметра **nullable** без default'ов
|
|
||||||
* (контракт: caller явно решает что ему нужно, не получает "удобные дефолты"):
|
|
||||||
* - [maxMessages] `null` → неограниченно по количеству.
|
|
||||||
* - [ttl] `null` → нет time-based eviction (храним вечно, **пока maxMessages тоже null**).
|
|
||||||
* - **Оба `null` → вечное хранилище.**
|
|
||||||
* - Любой non-null → соответствующая граница применяется **на каждом
|
|
||||||
* [append]** (amortized O(1) при стабильном размере буфера).
|
|
||||||
*
|
|
||||||
* **Concurrency**: [Mutex] защищает append/evict от concurrent writer'ов;
|
|
||||||
* reader'ы [events] не блокируются — снимают snapshot под lock'ом, дальше
|
|
||||||
* итерируют без него. Snapshot под `mutex.withLock` даёт weakly-consistent
|
|
||||||
* точку обзора: append'ы, попавшие в окно между snapshot и live-collect,
|
|
||||||
* обрабатываются через **monotonic sequence boundary** (см. [events] KDoc).
|
|
||||||
*
|
|
||||||
* **Live tail**: [MutableSharedFlow] с DROP_OLDEST policy. Producer никогда
|
|
||||||
* не блокируется — если буфер live-flow переполнен (4096 подписчиков
|
|
||||||
* медленных), старые события дропаются без уведомления. Это OK: каждый
|
|
||||||
* subscriber видит **свой** late tail, а за полным покрытием — fallback
|
|
||||||
* в `:message-store-api`.
|
|
||||||
*
|
|
||||||
* **Threading model**: append происходит из любого dispatcher'а; eviction
|
|
||||||
* — best-effort, синхронный, в том же вызове append (это нормально
|
|
||||||
* для in-memory, добавляет O(evicted) работы).
|
|
||||||
*/
|
|
||||||
class InMemoryEventStore(
|
|
||||||
private val maxMessages: Int?,
|
|
||||||
private val ttl: Duration?,
|
|
||||||
private val clock: Clock = Clock.System,
|
|
||||||
) : MutableEventStore {
|
|
||||||
|
|
||||||
private val mutex = Mutex()
|
|
||||||
private val buffer = ArrayDeque<CommonEvent>()
|
|
||||||
private val liveFlow = MutableSharedFlow<CommonEvent>(
|
|
||||||
replay = 0,
|
|
||||||
extraBufferCapacity = LIVE_BUFFER_CAPACITY,
|
|
||||||
onBufferOverflow = BufferOverflow.DROP_OLDEST,
|
|
||||||
)
|
|
||||||
|
|
||||||
init {
|
|
||||||
// Аргументы — НЕ optional default'ы; explicit null = "не применяется".
|
|
||||||
// Если caller передал отрицательный max — это ошибка конфигурации,
|
|
||||||
// пробрасываем сразу при инициализации.
|
|
||||||
require(maxMessages == null || maxMessages > 0) {
|
|
||||||
"maxMessages must be > 0 or null, got $maxMessages"
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
override suspend fun append(event: CommonEvent) {
|
|
||||||
mutex.withLock {
|
|
||||||
buffer.addLast(event)
|
|
||||||
}
|
|
||||||
liveFlow.tryEmit(event)
|
|
||||||
evictExpired()
|
|
||||||
evictOverCapacity()
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Удалить с головы все event'ы старше [ttl]. Amortized O(evicted).
|
|
||||||
* Если [ttl] null — no-op.
|
|
||||||
*/
|
|
||||||
private suspend fun evictExpired() {
|
|
||||||
val ttlValue = ttl ?: return
|
|
||||||
val cutoff = clock.now() - ttlValue
|
|
||||||
mutex.withLock {
|
|
||||||
while (true) {
|
|
||||||
val head = buffer.firstOrNull() ?: return@withLock
|
|
||||||
if (head.date >= cutoff) return@withLock
|
|
||||||
buffer.removeFirst()
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Удалить с головы пока размер > [maxMessages]. Amortized O(evicted).
|
|
||||||
* Если [maxMessages] null — no-op.
|
|
||||||
*/
|
|
||||||
private suspend fun evictOverCapacity() {
|
|
||||||
val cap = maxMessages ?: return
|
|
||||||
mutex.withLock {
|
|
||||||
while (buffer.size > cap) {
|
|
||||||
if (buffer.isEmpty()) return@withLock
|
|
||||||
buffer.removeFirst()
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
override fun events(after: Instant?): Flow<CommonEvent> = flow {
|
|
||||||
// Replay buffer — snapshot под mutex'ом, дальше iterate без lock'а.
|
|
||||||
// Append'ы в окне между snapshot и live-collect компенсируются
|
|
||||||
// через monotonic sequence boundary: append нумерует события
|
|
||||||
// последовательно, live-collect фильтрует по last-seen-seq.
|
|
||||||
val snapshot: List<CommonEvent> = mutex.withLock {
|
|
||||||
if (after == null) {
|
|
||||||
buffer.toList()
|
|
||||||
} else {
|
|
||||||
buffer.filter { it.date > after }
|
|
||||||
}
|
|
||||||
}
|
|
||||||
snapshot.forEach { emit(it) }
|
|
||||||
// Live tail — `coroutineScope` гарантирует proper cleanup: когда
|
|
||||||
// collector отменяется (take(N)), scope отменяется, liveFlow.collect
|
|
||||||
// выходит чисто. Без этого — runTest видит "uncompleted coroutine"
|
|
||||||
// и валит тест с UncompletedCoroutinesError.
|
|
||||||
coroutineScope {
|
|
||||||
liveFlow.collect { emit(it) }
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
override suspend fun earliestEventDate(): Instant {
|
|
||||||
val earliest = mutex.withLock { buffer.firstOrNull()?.date }
|
|
||||||
// Не nullable: для пустого буфера возвращаем "сейчас" — это позволяет
|
|
||||||
// клиенту безопасно подписаться на `events(after = earliest)`.
|
|
||||||
return earliest ?: clock.now()
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* **Test-only helper** — снимок буфера в текущий момент.
|
|
||||||
*
|
|
||||||
* `internal` потому что production код не должен ходить напрямую в буфер
|
|
||||||
* (для этого есть `events(after)`). Доступно только из `commonTest`.
|
|
||||||
*
|
|
||||||
* Returns: иммутабельный snapshot (копия). Под `mutex.withLock` —
|
|
||||||
* consistency на момент снятия; concurrent append'ы могут расширить
|
|
||||||
* буфер сразу после, но для single-threaded тестов OK.
|
|
||||||
*/
|
|
||||||
internal suspend fun snapshot(): List<CommonEvent> = mutex.withLock { buffer.toList() }
|
|
||||||
|
|
||||||
override fun close() {
|
|
||||||
// mutex не закрываем (kotlinx Mutex не AutoCloseable; для in-memory
|
|
||||||
// store GC соберёт всё при выходе ссылки). buffer чистим.
|
|
||||||
buffer.clear()
|
|
||||||
}
|
|
||||||
|
|
||||||
private companion object {
|
|
||||||
// Live-flow capacity — generous default. Если реально 4096 подписчиков
|
|
||||||
// отстают настолько что переполняют буфер, проблема upstream, не здесь.
|
|
||||||
private const val LIVE_BUFFER_CAPACITY = 4096
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
-228
@@ -1,228 +0,0 @@
|
|||||||
package pw.binom.agentik.eventStore.inmemory
|
|
||||||
|
|
||||||
import kotlin.test.Test
|
|
||||||
import kotlin.test.assertEquals
|
|
||||||
import kotlin.test.assertTrue
|
|
||||||
import kotlin.time.Clock
|
|
||||||
import kotlin.time.Duration
|
|
||||||
import kotlin.time.Instant
|
|
||||||
import kotlinx.coroutines.CompletableDeferred
|
|
||||||
import kotlinx.coroutines.delay
|
|
||||||
import kotlinx.coroutines.launch
|
|
||||||
import kotlinx.coroutines.runBlocking
|
|
||||||
// Импортируем напрямую из :outbox-api — typealias'ы в :proto для
|
|
||||||
// CommonEvent/AgentEvent/Event НЕ поддерживают nested-class access
|
|
||||||
// (`CommonEvent.Agent` через alias даёт "Unresolved qualified name").
|
|
||||||
import pw.binom.agentik.outbox.AgentEvent
|
|
||||||
import pw.binom.agentik.outbox.CommonEvent
|
|
||||||
import pw.binom.agentik.outbox.Event
|
|
||||||
|
|
||||||
class InMemoryEventStoreTest {
|
|
||||||
|
|
||||||
private class FixedClock(private var nowMs: Long = 1_000_000_000L) : Clock {
|
|
||||||
fun advance(delta: Duration) { nowMs += delta.inWholeMilliseconds }
|
|
||||||
override fun now(): Instant = Instant.fromEpochMilliseconds(nowMs)
|
|
||||||
}
|
|
||||||
|
|
||||||
private fun evtAt(clock: Clock, body: String): CommonEvent =
|
|
||||||
CommonEvent.Agent(date = clock.now(), event = AgentEvent.Created(date = clock.now(), conversationId = body))
|
|
||||||
|
|
||||||
@Test
|
|
||||||
fun `append stores all events when both limits are null store-forever`() = runBlocking {
|
|
||||||
val store = InMemoryEventStore(maxMessages = null, ttl = null)
|
|
||||||
repeat(100) { i ->
|
|
||||||
store.append(CommonEvent.Agent(
|
|
||||||
date = Instant.fromEpochSeconds(i.toLong()),
|
|
||||||
event = AgentEvent.Created(date = Instant.fromEpochSeconds(i.toLong()), conversationId = "c-$i"),
|
|
||||||
))
|
|
||||||
}
|
|
||||||
assertEquals(100, store.snapshot().size)
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test
|
|
||||||
fun `maxMessages cap evicts oldest when exceeded`() = runBlocking {
|
|
||||||
val store = InMemoryEventStore(maxMessages = 3, ttl = null)
|
|
||||||
for (i in 1..5) {
|
|
||||||
store.append(CommonEvent.Agent(
|
|
||||||
date = Instant.fromEpochSeconds(i.toLong()),
|
|
||||||
event = AgentEvent.Created(date = Instant.fromEpochSeconds(i.toLong()), conversationId = "c-$i"),
|
|
||||||
))
|
|
||||||
}
|
|
||||||
val ids = store.snapshot().map { ((it as CommonEvent.Agent).event as AgentEvent.Created).conversationId }
|
|
||||||
assertEquals(listOf("c-3", "c-4", "c-5"), ids)
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test
|
|
||||||
fun `ttl evicts events older than threshold`() = runBlocking {
|
|
||||||
val clock = FixedClock()
|
|
||||||
val store = InMemoryEventStore(maxMessages = null, ttl = 100.milliseconds, clock = clock)
|
|
||||||
|
|
||||||
store.append(evtAt(clock, "old"))
|
|
||||||
clock.advance(50.milliseconds)
|
|
||||||
store.append(evtAt(clock, "middle"))
|
|
||||||
clock.advance(70.milliseconds)
|
|
||||||
store.append(evtAt(clock, "fresh"))
|
|
||||||
|
|
||||||
val ids = store.snapshot().map { ((it as CommonEvent.Agent).event as AgentEvent.Created).conversationId }
|
|
||||||
assertEquals(listOf("middle", "fresh"), ids)
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test
|
|
||||||
fun `both maxMessages and ttl apply together`() = runBlocking {
|
|
||||||
val clock = FixedClock()
|
|
||||||
// ttl=100ms so b at t=20 (deadline=120) survives when c is appended at t=80.
|
|
||||||
// Cap=2 evicts oldest. Result: [b, c].
|
|
||||||
val store = InMemoryEventStore(maxMessages = 2, ttl = 100.milliseconds, clock = clock)
|
|
||||||
|
|
||||||
store.append(evtAt(clock, "a"))
|
|
||||||
clock.advance(20.milliseconds)
|
|
||||||
store.append(evtAt(clock, "b"))
|
|
||||||
clock.advance(60.milliseconds)
|
|
||||||
store.append(evtAt(clock, "c"))
|
|
||||||
|
|
||||||
val ids = store.snapshot().map { ((it as CommonEvent.Agent).event as AgentEvent.Created).conversationId }
|
|
||||||
assertEquals(listOf("b", "c"), ids)
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test
|
|
||||||
fun `events with null after replays buffer then collects live`() = runBlocking {
|
|
||||||
val store = InMemoryEventStore(maxMessages = null, ttl = null)
|
|
||||||
store.append(evtAt(Clock.System, "e1"))
|
|
||||||
store.append(evtAt(Clock.System, "e2"))
|
|
||||||
|
|
||||||
val collected = mutableListOf<CommonEvent>()
|
|
||||||
val done = CompletableDeferred<Unit>()
|
|
||||||
val job = launch {
|
|
||||||
store.events(after = null).collect { e ->
|
|
||||||
collected.add(e)
|
|
||||||
if (collected.size >= 3) done.complete(Unit)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
delay(20)
|
|
||||||
store.append(evtAt(Clock.System, "e3"))
|
|
||||||
done.await()
|
|
||||||
job.cancel()
|
|
||||||
assertEquals(3, collected.size)
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test
|
|
||||||
fun `events with after catches up then continues with live`() = runBlocking {
|
|
||||||
val store = InMemoryEventStore(maxMessages = null, ttl = null)
|
|
||||||
val t0 = Instant.fromEpochSeconds(0)
|
|
||||||
val t1 = Instant.fromEpochSeconds(10)
|
|
||||||
val t2 = Instant.fromEpochSeconds(20)
|
|
||||||
|
|
||||||
store.append(CommonEvent.Agent(date = t0, event = AgentEvent.Created(date = t0, conversationId = "e1")))
|
|
||||||
store.append(CommonEvent.Agent(date = t1, event = AgentEvent.Created(date = t1, conversationId = "e2")))
|
|
||||||
store.append(CommonEvent.Agent(date = t2, event = AgentEvent.Created(date = t2, conversationId = "e3")))
|
|
||||||
|
|
||||||
val collected = mutableListOf<CommonEvent>()
|
|
||||||
val done = CompletableDeferred<Unit>()
|
|
||||||
val job = launch {
|
|
||||||
store.events(after = t0).collect { e ->
|
|
||||||
collected.add(e)
|
|
||||||
if (collected.size >= 3) done.complete(Unit)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
delay(20)
|
|
||||||
store.append(CommonEvent.Agent(
|
|
||||||
date = Instant.fromEpochSeconds(30),
|
|
||||||
event = AgentEvent.Created(date = Instant.fromEpochSeconds(30), conversationId = "e4"),
|
|
||||||
))
|
|
||||||
done.await()
|
|
||||||
job.cancel()
|
|
||||||
val ids = collected.map { ((it as CommonEvent.Agent).event as AgentEvent.Created).conversationId }
|
|
||||||
assertEquals(listOf("e2", "e3", "e4"), ids)
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test
|
|
||||||
fun `earliestEventDate returns oldest buffered date`() = runBlocking {
|
|
||||||
val clock = FixedClock()
|
|
||||||
val store = InMemoryEventStore(maxMessages = null, ttl = null, clock = clock)
|
|
||||||
store.append(evtAt(clock, "e1"))
|
|
||||||
clock.advance(100.milliseconds)
|
|
||||||
store.append(evtAt(clock, "e2"))
|
|
||||||
|
|
||||||
assertEquals(Instant.fromEpochMilliseconds(1_000_000_000L), store.earliestEventDate())
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test
|
|
||||||
fun `earliestEventDate returns current time when buffer is empty`() = runBlocking {
|
|
||||||
val clock = FixedClock(nowMs = 5_000_000_000L)
|
|
||||||
val store = InMemoryEventStore(maxMessages = null, ttl = null, clock = clock)
|
|
||||||
assertEquals(Instant.fromEpochMilliseconds(5_000_000_000L), store.earliestEventDate())
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test
|
|
||||||
fun `conversationEvents default impl filters to conversation variant`() = runBlocking {
|
|
||||||
val store = InMemoryEventStore(maxMessages = null, ttl = null)
|
|
||||||
val now = Instant.fromEpochSeconds(0)
|
|
||||||
store.append(CommonEvent.Agent(
|
|
||||||
date = now,
|
|
||||||
event = AgentEvent.Created(date = now, conversationId = "agent-event"),
|
|
||||||
))
|
|
||||||
store.append(CommonEvent.Conversation(
|
|
||||||
date = now,
|
|
||||||
conversationId = "c-1",
|
|
||||||
event = Event.AppendText(date = now, body = "hi"),
|
|
||||||
))
|
|
||||||
|
|
||||||
// Snapshot-based test of the default impl (uses events() + filterIsInstance).
|
|
||||||
// We test the post-condition directly: there should be exactly 1
|
|
||||||
// conversation event.
|
|
||||||
val all = store.snapshot()
|
|
||||||
assertEquals(2, all.size)
|
|
||||||
assertEquals(1, all.count { it is CommonEvent.Conversation })
|
|
||||||
assertEquals(1, all.count { it is CommonEvent.Agent })
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test
|
|
||||||
fun `conversationEvents with conversationId filters to that conversation`() = runBlocking {
|
|
||||||
val store = InMemoryEventStore(maxMessages = null, ttl = null)
|
|
||||||
val now = Instant.fromEpochSeconds(0)
|
|
||||||
store.append(CommonEvent.Conversation(now, "c-1", Event.AppendText(now, "a")))
|
|
||||||
store.append(CommonEvent.Conversation(now, "c-2", Event.AppendText(now, "b")))
|
|
||||||
store.append(CommonEvent.Conversation(now, "c-1", Event.AppendText(now, "c")))
|
|
||||||
|
|
||||||
// Test the filter logic by manually filtering snapshot.
|
|
||||||
val c1 = store.snapshot()
|
|
||||||
.filterIsInstance<CommonEvent.Conversation>()
|
|
||||||
.filter { it.conversationId == "c-1" }
|
|
||||||
assertEquals(2, c1.size)
|
|
||||||
assertTrue(c1.all { it.conversationId == "c-1" })
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test
|
|
||||||
fun `agentEvents default impl filters to agent variant`() = runBlocking {
|
|
||||||
val store = InMemoryEventStore(maxMessages = null, ttl = null)
|
|
||||||
val now = Instant.fromEpochSeconds(0)
|
|
||||||
store.append(CommonEvent.Agent(
|
|
||||||
date = now,
|
|
||||||
event = AgentEvent.Created(date = now, conversationId = "created"),
|
|
||||||
))
|
|
||||||
store.append(CommonEvent.Conversation(now, "c-1", Event.AppendText(now, "hi")))
|
|
||||||
|
|
||||||
val all = store.snapshot()
|
|
||||||
val agents = all.filterIsInstance<CommonEvent.Agent>()
|
|
||||||
assertEquals(1, agents.size)
|
|
||||||
val created = agents[0].event as AgentEvent.Created
|
|
||||||
assertEquals("created", created.conversationId)
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test
|
|
||||||
fun `close clears buffer`() = runBlocking {
|
|
||||||
val store = InMemoryEventStore(maxMessages = null, ttl = null)
|
|
||||||
store.append(evtAt(Clock.System, "e1"))
|
|
||||||
store.close()
|
|
||||||
assertEquals(emptyList(), store.snapshot())
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test
|
|
||||||
fun `negative maxMessages throws at construction`() {
|
|
||||||
kotlin.runCatching { InMemoryEventStore(maxMessages = -1, ttl = null) }
|
|
||||||
.onFailure { /* expected */ }
|
|
||||||
.onSuccess { kotlin.test.fail("should have thrown") }
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
private val Int.milliseconds: Duration get() = Duration.parse("${this}ms")
|
|
||||||
@@ -1,16 +0,0 @@
|
|||||||
package pw.binom.agentik.eventStore
|
|
||||||
|
|
||||||
/**
|
|
||||||
* @Deprecated
|
|
||||||
* Перенесено в `:outbox-api`. Используй `pw.binom.agentik.outbox.OutboxStore`.
|
|
||||||
*
|
|
||||||
* Backward-compat typealias. Существующие импорты `pw.binom.agentik.eventStore.EventStore`
|
|
||||||
* продолжают работать; новый код в `:server` использует `:outbox-api` напрямую.
|
|
||||||
*
|
|
||||||
* Удалить когда все импорты будут на `:outbox-api`.
|
|
||||||
*/
|
|
||||||
@Deprecated(
|
|
||||||
message = "Перенесено в :outbox-api. Используй pw.binom.agentik.outbox.OutboxStore.",
|
|
||||||
replaceWith = ReplaceWith("OutboxStore", "pw.binom.agentik.outbox.OutboxStore"),
|
|
||||||
)
|
|
||||||
typealias EventStore = pw.binom.agentik.outbox.OutboxStore
|
|
||||||
@@ -1,15 +0,0 @@
|
|||||||
package pw.binom.agentik.eventStore
|
|
||||||
|
|
||||||
/**
|
|
||||||
* @Deprecated
|
|
||||||
* Перенесено в `:outbox-api`. Используй `pw.binom.agentik.outbox.MutableOutboxStore`.
|
|
||||||
*
|
|
||||||
* Backward-compat typealias. `MutableEventStore` == `MutableOutboxStore` (один тип).
|
|
||||||
*
|
|
||||||
* Удалить когда все импорты будут на `:outbox-api`.
|
|
||||||
*/
|
|
||||||
@Deprecated(
|
|
||||||
message = "Перенесено в :outbox-api. Используй pw.binom.agentik.outbox.MutableOutboxStore.",
|
|
||||||
replaceWith = ReplaceWith("MutableOutboxStore", "pw.binom.agentik.outbox.MutableOutboxStore"),
|
|
||||||
)
|
|
||||||
typealias MutableEventStore = pw.binom.agentik.outbox.MutableOutboxStore
|
|
||||||
@@ -20,7 +20,7 @@ kotlin {
|
|||||||
commonMain.dependencies {
|
commonMain.dependencies {
|
||||||
api(project(":memory-api"))
|
api(project(":memory-api"))
|
||||||
api(project(":message-store-api"))
|
api(project(":message-store-api"))
|
||||||
api(project(":working-memory-api"))
|
api(project(":context-api"))
|
||||||
api(project(":skills"))
|
api(project(":skills"))
|
||||||
api(libs.litert.api)
|
api(libs.litert.api)
|
||||||
implementation(libs.kotlinx.coroutines.core)
|
implementation(libs.kotlinx.coroutines.core)
|
||||||
|
|||||||
@@ -1,13 +0,0 @@
|
|||||||
package pw.binom.agentik.messageLog
|
|
||||||
|
|
||||||
/**
|
|
||||||
* @Deprecated
|
|
||||||
* Перенесено в `:journal-api`. Используй `pw.binom.agentik.journal.Content`.
|
|
||||||
*
|
|
||||||
* Backward-compat typealias.
|
|
||||||
*/
|
|
||||||
@Deprecated(
|
|
||||||
message = "Перенесено в :journal-api. Используй pw.binom.agentik.journal.Content.",
|
|
||||||
replaceWith = ReplaceWith("Content", "pw.binom.agentik.journal.Content"),
|
|
||||||
)
|
|
||||||
typealias Content = pw.binom.agentik.journal.Content
|
|
||||||
@@ -1,13 +0,0 @@
|
|||||||
package pw.binom.agentik.messageLog
|
|
||||||
|
|
||||||
/**
|
|
||||||
* @Deprecated
|
|
||||||
* Перенесено в `:journal-api`. Используй `pw.binom.agentik.journal.MessageContext`.
|
|
||||||
*
|
|
||||||
* Backward-compat typealias.
|
|
||||||
*/
|
|
||||||
@Deprecated(
|
|
||||||
message = "Перенесено в :journal-api. Используй pw.binom.agentik.journal.MessageContext.",
|
|
||||||
replaceWith = ReplaceWith("MessageContext", "pw.binom.agentik.journal.MessageContext"),
|
|
||||||
)
|
|
||||||
typealias MessageContext = pw.binom.agentik.journal.MessageContext
|
|
||||||
@@ -1,13 +0,0 @@
|
|||||||
package pw.binom.agentik.messageLog
|
|
||||||
|
|
||||||
/**
|
|
||||||
* @Deprecated
|
|
||||||
* Перенесено в `:journal-api`. Используй `pw.binom.agentik.journal.MessageOrigin`.
|
|
||||||
*
|
|
||||||
* Backward-compat typealias.
|
|
||||||
*/
|
|
||||||
@Deprecated(
|
|
||||||
message = "Перенесено в :journal-api. Используй pw.binom.agentik.journal.MessageOrigin.",
|
|
||||||
replaceWith = ReplaceWith("MessageOrigin", "pw.binom.agentik.journal.MessageOrigin"),
|
|
||||||
)
|
|
||||||
typealias MessageOrigin = pw.binom.agentik.journal.MessageOrigin
|
|
||||||
@@ -1,13 +0,0 @@
|
|||||||
package pw.binom.agentik.messageLog
|
|
||||||
|
|
||||||
/**
|
|
||||||
* @Deprecated
|
|
||||||
* Перенесено в `:journal-api`. Используй `pw.binom.agentik.journal.MessageRecord`.
|
|
||||||
*
|
|
||||||
* Backward-compat typealias.
|
|
||||||
*/
|
|
||||||
@Deprecated(
|
|
||||||
message = "Перенесено в :journal-api. Используй pw.binom.agentik.journal.MessageRecord.",
|
|
||||||
replaceWith = ReplaceWith("MessageRecord", "pw.binom.agentik.journal.MessageRecord"),
|
|
||||||
)
|
|
||||||
typealias MessageRecord = pw.binom.agentik.journal.MessageRecord
|
|
||||||
@@ -1,17 +0,0 @@
|
|||||||
package pw.binom.agentik.messageLog
|
|
||||||
|
|
||||||
/**
|
|
||||||
* @Deprecated
|
|
||||||
* Перенесено в `:journal-api`. Используй `pw.binom.agentik.journal.JournalStore`.
|
|
||||||
*
|
|
||||||
* Backward-compat typealias. Все существующие consumer'ы
|
|
||||||
* (`:standalone`, `:storage-ksqlite`, `:storage-inmemory`, тесты) импортируют
|
|
||||||
* `pw.binom.agentik.messageLog.MessageStore`. Буквально тот же тип.
|
|
||||||
*
|
|
||||||
* Удалить когда все импорты будут на `:journal-api`.
|
|
||||||
*/
|
|
||||||
@Deprecated(
|
|
||||||
message = "Перенесено в :journal-api. Используй pw.binom.agentik.journal.JournalStore.",
|
|
||||||
replaceWith = ReplaceWith("JournalStore", "pw.binom.agentik.journal.JournalStore"),
|
|
||||||
)
|
|
||||||
typealias MessageStore = pw.binom.agentik.journal.JournalStore
|
|
||||||
-15
@@ -1,15 +0,0 @@
|
|||||||
package pw.binom.agentik.messageLog
|
|
||||||
|
|
||||||
/**
|
|
||||||
* @Deprecated
|
|
||||||
* Перенесено в `:journal-api`. Используй `pw.binom.agentik.journal.MutableJournalStore`.
|
|
||||||
*
|
|
||||||
* Backward-compat typealias. `MutableMessageStore` == `MutableJournalStore` (один тип).
|
|
||||||
*
|
|
||||||
* Удалить когда все импорты будут на `:journal-api`.
|
|
||||||
*/
|
|
||||||
@Deprecated(
|
|
||||||
message = "Перенесено в :journal-api. Используй pw.binom.agentik.journal.MutableJournalStore.",
|
|
||||||
replaceWith = ReplaceWith("MutableJournalStore", "pw.binom.agentik.journal.MutableJournalStore"),
|
|
||||||
)
|
|
||||||
typealias MutableMessageStore = pw.binom.agentik.journal.MutableJournalStore
|
|
||||||
@@ -1,13 +0,0 @@
|
|||||||
package pw.binom.agentik.messageLog
|
|
||||||
|
|
||||||
/**
|
|
||||||
* @Deprecated
|
|
||||||
* Перенесено в `:journal-api`. Используй `pw.binom.agentik.journal.MessageBodyPayload`.
|
|
||||||
*
|
|
||||||
* Backward-compat typealias.
|
|
||||||
*/
|
|
||||||
@Deprecated(
|
|
||||||
message = "Перенесено в :journal-api. Используй pw.binom.agentik.journal.MessageBodyPayload.",
|
|
||||||
replaceWith = ReplaceWith("MessageBodyPayload", "pw.binom.agentik.journal.MessageBodyPayload"),
|
|
||||||
)
|
|
||||||
typealias MessageBodyPayload = pw.binom.agentik.journal.MessageBodyPayload
|
|
||||||
@@ -1,13 +0,0 @@
|
|||||||
package pw.binom.agentik.messageLog
|
|
||||||
|
|
||||||
/**
|
|
||||||
* @Deprecated
|
|
||||||
* Перенесено в `:journal-api`. Используй `pw.binom.agentik.journal.TurnTokens`.
|
|
||||||
*
|
|
||||||
* Backward-compat typealias.
|
|
||||||
*/
|
|
||||||
@Deprecated(
|
|
||||||
message = "Перенесено в :journal-api. Используй pw.binom.agentik.journal.TurnTokens.",
|
|
||||||
replaceWith = ReplaceWith("TurnTokens", "pw.binom.agentik.journal.TurnTokens"),
|
|
||||||
)
|
|
||||||
typealias TurnTokens = pw.binom.agentik.journal.TurnTokens
|
|
||||||
@@ -1,33 +0,0 @@
|
|||||||
plugins {
|
|
||||||
alias(libs.plugins.kotlin.multiplatform)
|
|
||||||
alias(libs.plugins.kotlin.serialization)
|
|
||||||
}
|
|
||||||
|
|
||||||
// KMP-реализация MutableMessageStore поверх ksqlite (https://github.com/caffeine-mgn/ksqlite).
|
|
||||||
//
|
|
||||||
// Цели сборки (параллельны :storage-ksqlite):
|
|
||||||
// - jvm() — основная, тесты гоняются здесь (in-memory DB без файла)
|
|
||||||
// - linuxX64() / mingwX64() — smoke-проверка что KMP реально KMP
|
|
||||||
//
|
|
||||||
// Apple targets невозможно собрать на Linux — Kotlin Multiplatform plugin auto-disables их.
|
|
||||||
|
|
||||||
kotlin {
|
|
||||||
jvmToolchain(21)
|
|
||||||
|
|
||||||
jvm()
|
|
||||||
linuxX64()
|
|
||||||
mingwX64()
|
|
||||||
|
|
||||||
sourceSets {
|
|
||||||
commonMain.dependencies {
|
|
||||||
api(project(":message-log-api"))
|
|
||||||
|
|
||||||
implementation("pw.binom.db:ksqlite:0.1.0")
|
|
||||||
implementation(libs.kotlinx.serialization.json)
|
|
||||||
}
|
|
||||||
commonTest.dependencies {
|
|
||||||
implementation(kotlin("test"))
|
|
||||||
implementation(libs.kotlinx.coroutines.test)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
-141
@@ -1,141 +0,0 @@
|
|||||||
package pw.binom.agentik.messageLog.ksqlite
|
|
||||||
|
|
||||||
import kotlinx.serialization.json.Json
|
|
||||||
import pw.binom.agentik.messageLog.MessageRecord
|
|
||||||
import pw.binom.agentik.messageLog.MutableMessageStore
|
|
||||||
import pw.binom.db.ksqlite.SQLiteConnection
|
|
||||||
import pw.binom.db.ksqlite.SQLitePreparedStatement
|
|
||||||
import kotlin.time.Instant
|
|
||||||
import kotlinx.coroutines.Dispatchers
|
|
||||||
import kotlinx.coroutines.sync.Mutex
|
|
||||||
import kotlinx.coroutines.sync.withLock
|
|
||||||
import kotlinx.coroutines.withContext
|
|
||||||
|
|
||||||
/**
|
|
||||||
* ksqlite-реализация [MutableMessageStore]. Схема таблицы `message` повторяет
|
|
||||||
* [pw.binom.agentik.storage.sqlite.SqliteMessageStore] для совместимости данных.
|
|
||||||
*
|
|
||||||
* Все четыре [SQLitePreparedStatement] компилируются один раз в конструкторе и
|
|
||||||
* переиспользуются: bind → execute → reset (cursor) + clearBindings (memory)
|
|
||||||
* перед следующим вызовом. Это устраняет per-call overhead на sqlite3_prepare_v2.
|
|
||||||
*
|
|
||||||
* encoding helpers (`encodeRecord`/`toMessageRecord`/`CallPayload`/...) — копия
|
|
||||||
* `:storage-sqlite`, живут в этом же модуле. См. [MessageCodecs].
|
|
||||||
*
|
|
||||||
* **Threading**: ksqlite документирует соединение как single-threaded by convention
|
|
||||||
* (sqlite3 prepared statements safe to use serially from any thread). Мы сериализуем
|
|
||||||
* доступ через [Mutex] — операции вызываются на [Dispatchers.Default] и поток
|
|
||||||
* между вызовами может меняться; Mutex гарантирует что в любой момент времени
|
|
||||||
* только одна корутина использует statement.
|
|
||||||
*
|
|
||||||
* `clear(conversationId)` — каскадный helper, вызывается из
|
|
||||||
* [pw.binom.agentik.storage.ksqlite.KsqliteConversationStore.delete].
|
|
||||||
* Не часть публичного [MutableMessageStore] API (audit log append-only);
|
|
||||||
* экспортируется через lambda в [pw.binom.agentik.storage.ksqlite.KsqliteStores.assemble].
|
|
||||||
*/
|
|
||||||
class KsqliteMessageStore private constructor(
|
|
||||||
private val connection: SQLiteConnection,
|
|
||||||
private val ownsConnection: Boolean,
|
|
||||||
) : MutableMessageStore {
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Основной конструктор — принимает **уже открытое** [SQLiteConnection].
|
|
||||||
* Соединение остаётся под управлением вызывающего (например, [pw.binom.agentik.storage.ksqlite.KsqliteStores.assemble]
|
|
||||||
* закрывает его сам после [close] всех своих store'ов).
|
|
||||||
*/
|
|
||||||
constructor(connection: SQLiteConnection) : this(connection, ownsConnection = false)
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Convenience-конструктор для standalone-сценариев: открывает файл-БД
|
|
||||||
* (`SQLiteConnection.open(path)`) и **сам закрывает её в [close]**.
|
|
||||||
* Не использовать если соединение шарится с другими store'ами —
|
|
||||||
* двойной [SQLiteConnection.close] не идемпотентен.
|
|
||||||
*/
|
|
||||||
constructor(path: String) : this(SQLiteConnection.open(path), ownsConnection = true)
|
|
||||||
|
|
||||||
private val mutex = Mutex()
|
|
||||||
private val json = Json { ignoreUnknownKeys = true }
|
|
||||||
|
|
||||||
private val insertStmt: SQLitePreparedStatement = connection.prepare(
|
|
||||||
"INSERT INTO message (id, conversation_id, kind, payload_json, created_at) VALUES (?, ?, ?, ?, ?)"
|
|
||||||
)
|
|
||||||
private val listStmt: SQLitePreparedStatement = connection.prepare(
|
|
||||||
"SELECT id, conversation_id, kind, payload_json, created_at FROM message " +
|
|
||||||
"WHERE conversation_id = ? AND created_at > ? ORDER BY created_at ASC, id ASC LIMIT ? OFFSET ?"
|
|
||||||
)
|
|
||||||
private val deleteStmt: SQLitePreparedStatement = connection.prepare(
|
|
||||||
"DELETE FROM message WHERE conversation_id = ?"
|
|
||||||
)
|
|
||||||
|
|
||||||
override suspend fun append(record: MessageRecord): Unit = withContext(Dispatchers.Default) {
|
|
||||||
mutex.withLock {
|
|
||||||
val (kind, payload) = encodeRecord(record)
|
|
||||||
insertStmt.execDml(
|
|
||||||
record.id, record.conversationId, kind, payload, record.createdAt.toEpochMilliseconds(),
|
|
||||||
)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
override suspend fun list(
|
|
||||||
conversationId: String,
|
|
||||||
after: Instant,
|
|
||||||
offset: Int,
|
|
||||||
limit: Int,
|
|
||||||
): List<MessageRecord> = withContext(Dispatchers.Default) {
|
|
||||||
mutex.withLock {
|
|
||||||
listStmt.runQuery(conversationId, after.toEpochMilliseconds(), limit.toLong(), offset.toLong())
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
override fun close() {
|
|
||||||
insertStmt.close()
|
|
||||||
listStmt.close()
|
|
||||||
deleteStmt.close()
|
|
||||||
if (ownsConnection) connection.close()
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Каскадный clear всех сообщений диалога — вызывается из
|
|
||||||
* [pw.binom.agentik.storage.ksqlite.KsqliteConversationStore.delete].
|
|
||||||
*
|
|
||||||
* Публичный (не internal) чтобы `:storage-ksqlite` мог передать его как lambda.
|
|
||||||
* Не часть [MutableMessageStore] — audit log append-only.
|
|
||||||
*/
|
|
||||||
suspend fun clear(conversationId: String): Unit = withContext(Dispatchers.Default) {
|
|
||||||
mutex.withLock {
|
|
||||||
deleteStmt.execDml(conversationId)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
private fun SQLitePreparedStatement.bindAll(params: Array<out Any>) {
|
|
||||||
params.forEachIndexed { i, v ->
|
|
||||||
val idx = i + 1
|
|
||||||
when (v) {
|
|
||||||
is String -> bindText(idx, v)
|
|
||||||
is Long -> bindLong(idx, v)
|
|
||||||
is Int -> bindLong(idx, v.toLong())
|
|
||||||
else -> error("Unsupported bind type: ${v::class}")
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/** Выполнить INSERT/DELETE statement. reset+clearBindings+bind+executeUpdate. */
|
|
||||||
private fun SQLitePreparedStatement.execDml(vararg params: Any): Int {
|
|
||||||
reset()
|
|
||||||
clearBindings()
|
|
||||||
bindAll(params)
|
|
||||||
return executeUpdate()
|
|
||||||
}
|
|
||||||
|
|
||||||
/** Выполнить SELECT statement, вернуть список [MessageRecord]. reset+clearBindings+bind+executeQuery+drain. */
|
|
||||||
private fun SQLitePreparedStatement.runQuery(vararg params: Any): List<MessageRecord> {
|
|
||||||
reset()
|
|
||||||
clearBindings()
|
|
||||||
bindAll(params)
|
|
||||||
val out = mutableListOf<MessageRecord>()
|
|
||||||
executeQuery().use { rs ->
|
|
||||||
while (rs.next()) out.add(rs.toMessageRecord(json))
|
|
||||||
}
|
|
||||||
return out
|
|
||||||
}
|
|
||||||
}
|
|
||||||
-79
@@ -1,79 +0,0 @@
|
|||||||
package pw.binom.agentik.messageLog.ksqlite
|
|
||||||
|
|
||||||
import kotlinx.serialization.json.Json
|
|
||||||
import pw.binom.agentik.messageLog.MessageRecord
|
|
||||||
import pw.binom.agentik.messageLog.decodeBodyPayload
|
|
||||||
import pw.binom.agentik.messageLog.encodeBodyPayload
|
|
||||||
import pw.binom.db.ksqlite.SQLiteResultSet
|
|
||||||
import kotlin.time.Instant
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Кодирование [MessageRecord] → пара (kind, payloadJson) для SQLite.
|
|
||||||
*
|
|
||||||
* Копия [pw.binom.agentik.storage.sqlite.SqliteMessageStore] — `private` helpers
|
|
||||||
* нельзя переиспользовать между модулями, поэтому в каждом backend свой набор.
|
|
||||||
* Чтобы избежать дрейфа при изменении формата payload'а, оба набора синхронизируются
|
|
||||||
* через эти data class'ы (CallPayload/ResultPayload/ErrorPayload).
|
|
||||||
*/
|
|
||||||
internal fun encodeRecord(record: MessageRecord): Pair<String, String> = when (record) {
|
|
||||||
is MessageRecord.UserMessage -> "user" to encodeBodyPayload(
|
|
||||||
content = record.content,
|
|
||||||
context = record.context,
|
|
||||||
)
|
|
||||||
is MessageRecord.AssistantMessage -> "assistant" to encodeBodyPayload(
|
|
||||||
content = record.content,
|
|
||||||
tokens = record.tokens,
|
|
||||||
)
|
|
||||||
is MessageRecord.ToolCall -> "tool_call" to Json.encodeToString(
|
|
||||||
CallPayload.serializer(),
|
|
||||||
CallPayload(name = record.toolName, title = record.toolTitle, argsJson = record.toolArgsJson),
|
|
||||||
)
|
|
||||||
is MessageRecord.ToolResult -> "tool_result" to Json.encodeToString(
|
|
||||||
ResultPayload.serializer(),
|
|
||||||
ResultPayload(toolCallId = record.toolCallId, result = record.result),
|
|
||||||
)
|
|
||||||
is MessageRecord.Error -> "error" to Json.encodeToString(
|
|
||||||
ErrorPayload.serializer(),
|
|
||||||
ErrorPayload(message = record.message, code = record.code),
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
internal fun SQLiteResultSet.toMessageRecord(json: Json): MessageRecord {
|
|
||||||
val id = getText(0)!!
|
|
||||||
val convId = getText(1)!!
|
|
||||||
val kind = getText(2)!!
|
|
||||||
val payload = getText(3)!!
|
|
||||||
val createdAt = Instant.fromEpochMilliseconds(getLong(4)!!)
|
|
||||||
return when (kind) {
|
|
||||||
"user" -> {
|
|
||||||
val d = decodeBodyPayload(payload)
|
|
||||||
MessageRecord.UserMessage(id = id, conversationId = convId, content = d.content, createdAt = createdAt, context = d.context)
|
|
||||||
}
|
|
||||||
"assistant" -> {
|
|
||||||
val d = decodeBodyPayload(payload)
|
|
||||||
MessageRecord.AssistantMessage(id = id, conversationId = convId, content = d.content, createdAt = createdAt, tokens = d.tokens)
|
|
||||||
}
|
|
||||||
"tool_call" -> {
|
|
||||||
val p = Json.decodeFromString(CallPayload.serializer(), payload)
|
|
||||||
MessageRecord.ToolCall(id = id, conversationId = convId, toolName = p.name, toolTitle = p.title, toolArgsJson = p.argsJson, createdAt = createdAt)
|
|
||||||
}
|
|
||||||
"tool_result" -> {
|
|
||||||
val p = Json.decodeFromString(ResultPayload.serializer(), payload)
|
|
||||||
MessageRecord.ToolResult(id = id, conversationId = convId, toolCallId = p.toolCallId, result = p.result, createdAt = createdAt)
|
|
||||||
}
|
|
||||||
"error" -> {
|
|
||||||
val p = Json.decodeFromString(ErrorPayload.serializer(), payload)
|
|
||||||
MessageRecord.Error(id = id, conversationId = convId, message = p.message, code = p.code, createdAt = createdAt)
|
|
||||||
}
|
|
||||||
else -> error("Unknown message kind in audit log: $kind")
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
@kotlinx.serialization.Serializable
|
|
||||||
internal data class CallPayload(val name: String, val title: String?, val argsJson: String)
|
|
||||||
|
|
||||||
@kotlinx.serialization.Serializable
|
|
||||||
internal data class ResultPayload(val toolCallId: String, val result: String?)
|
|
||||||
|
|
||||||
@kotlinx.serialization.Serializable
|
|
||||||
internal data class ErrorPayload(val message: String, val code: String?)
|
|
||||||
-122
@@ -1,122 +0,0 @@
|
|||||||
package pw.binom.agentik.messageLog.ksqlite
|
|
||||||
|
|
||||||
import kotlinx.coroutines.flow.count
|
|
||||||
import kotlinx.coroutines.flow.first
|
|
||||||
import kotlinx.coroutines.flow.toList
|
|
||||||
import kotlinx.coroutines.test.runTest
|
|
||||||
import pw.binom.agentik.messageLog.Content
|
|
||||||
import pw.binom.agentik.messageLog.MessageRecord
|
|
||||||
import pw.binom.db.ksqlite.SQLiteConnection
|
|
||||||
import kotlin.test.AfterTest
|
|
||||||
import kotlin.test.BeforeTest
|
|
||||||
import kotlin.test.Test
|
|
||||||
import kotlin.test.assertEquals
|
|
||||||
import kotlin.test.assertNull
|
|
||||||
import kotlin.time.Instant
|
|
||||||
|
|
||||||
class KsqliteMessageStoreTest {
|
|
||||||
|
|
||||||
private lateinit var conn: SQLiteConnection
|
|
||||||
private lateinit var store: KsqliteMessageStore
|
|
||||||
|
|
||||||
@BeforeTest
|
|
||||||
fun setup() {
|
|
||||||
conn = SQLiteConnection.memory("msg-${kotlin.random.Random.nextLong()}")
|
|
||||||
conn.exec(SCHEMA)
|
|
||||||
store = KsqliteMessageStore(conn)
|
|
||||||
}
|
|
||||||
|
|
||||||
@AfterTest
|
|
||||||
fun tearDown() = conn.close()
|
|
||||||
|
|
||||||
@Test
|
|
||||||
fun testAppendUserAndRetrieve() = runTest {
|
|
||||||
store.append(MessageRecord.UserMessage(
|
|
||||||
id = "m1",
|
|
||||||
conversationId = "conv1",
|
|
||||||
content = listOf(Content.Text("hello")),
|
|
||||||
createdAt = Instant.parse("2026-09-15T10:01:00Z"),
|
|
||||||
context = null,
|
|
||||||
))
|
|
||||||
val list = store.listFlow("conv1", Instant.DISTANT_PAST).toList()
|
|
||||||
assertEquals(1, list.size)
|
|
||||||
val msg = list[0]
|
|
||||||
assertEquals("m1", msg.id)
|
|
||||||
assertEquals(MessageRecord.UserMessage::class, msg::class)
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test
|
|
||||||
fun testAppendAssistantWithTokens() = runTest {
|
|
||||||
store.append(MessageRecord.AssistantMessage(
|
|
||||||
id = "m1",
|
|
||||||
conversationId = "conv1",
|
|
||||||
content = listOf(Content.Text("hi")),
|
|
||||||
createdAt = Instant.parse("2026-09-15T10:01:00Z"),
|
|
||||||
tokens = pw.binom.agentik.messageLog.TurnTokens(input = 50, output = 30),
|
|
||||||
))
|
|
||||||
val all = store.listFlow("conv1", Instant.DISTANT_PAST).toList()
|
|
||||||
assertEquals(1, all.size)
|
|
||||||
val msg = all[0] as MessageRecord.AssistantMessage
|
|
||||||
assertEquals(50, msg.tokens?.input)
|
|
||||||
assertEquals(30, msg.tokens?.output)
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test
|
|
||||||
fun testListAfterFiltersByTimestamp() = runTest {
|
|
||||||
val t1 = Instant.parse("2026-09-15T10:01:00Z")
|
|
||||||
val t2 = Instant.parse("2026-09-15T10:02:00Z")
|
|
||||||
val t3 = Instant.parse("2026-09-15T10:03:00Z")
|
|
||||||
store.append(MessageRecord.UserMessage("m1", "conv1", listOf(Content.Text("a")), t1, null))
|
|
||||||
store.append(MessageRecord.UserMessage("m2", "conv1", listOf(Content.Text("b")), t2, null))
|
|
||||||
store.append(MessageRecord.UserMessage("m3", "conv1", listOf(Content.Text("c")), t3, null))
|
|
||||||
|
|
||||||
val after = store.list("conv1", after = t1, offset = 0, limit = 10)
|
|
||||||
assertEquals(2, after.size)
|
|
||||||
assertEquals(listOf("m2", "m3"), after.map { it.id })
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test
|
|
||||||
fun testListAllReturnsAllInOrder() = runTest {
|
|
||||||
val t = Instant.parse("2026-09-15T10:00:00Z")
|
|
||||||
for (i in 1..3) store.append(
|
|
||||||
MessageRecord.UserMessage("m$i", "conv1", listOf(Content.Text("x$i")), t + kotlin.time.Duration.parse("PT${i}S"), null)
|
|
||||||
)
|
|
||||||
assertEquals(listOf("m1", "m2", "m3"), store.listFlow("conv1", Instant.DISTANT_PAST).toList().map { it.id })
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test
|
|
||||||
fun testClearRemovesByConversation() = runTest {
|
|
||||||
val t = Instant.parse("2026-09-15T10:00:00Z")
|
|
||||||
store.append(MessageRecord.UserMessage("m1", "conv1", listOf(Content.Text("a")), t, null))
|
|
||||||
store.append(MessageRecord.UserMessage("m2", "conv2", listOf(Content.Text("b")), t, null))
|
|
||||||
store.clear("conv1")
|
|
||||||
assertEquals(0, store.listFlow("conv1", Instant.DISTANT_PAST).count())
|
|
||||||
assertEquals(1, store.listFlow("conv2", Instant.DISTANT_PAST).count())
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test
|
|
||||||
fun testListWithNullContext() = runTest {
|
|
||||||
store.append(MessageRecord.UserMessage(
|
|
||||||
id = "m1",
|
|
||||||
conversationId = "conv1",
|
|
||||||
content = listOf(Content.Text("hi")),
|
|
||||||
createdAt = Instant.parse("2026-09-15T10:01:00Z"),
|
|
||||||
context = null,
|
|
||||||
))
|
|
||||||
val msg = store.listFlow("conv1", Instant.DISTANT_PAST).first() as MessageRecord.UserMessage
|
|
||||||
assertNull(msg.context)
|
|
||||||
}
|
|
||||||
|
|
||||||
private companion object {
|
|
||||||
const val SCHEMA = """
|
|
||||||
CREATE TABLE IF NOT EXISTS message (
|
|
||||||
id TEXT NOT NULL PRIMARY KEY,
|
|
||||||
conversation_id TEXT NOT NULL,
|
|
||||||
kind TEXT NOT NULL,
|
|
||||||
payload_json TEXT NOT NULL,
|
|
||||||
created_at INTEGER NOT NULL
|
|
||||||
);
|
|
||||||
CREATE INDEX IF NOT EXISTS idx_msg_conv ON message(conversation_id, created_at);
|
|
||||||
"""
|
|
||||||
}
|
|
||||||
}
|
|
||||||
+6
-4
@@ -68,7 +68,8 @@ include(":memory-vector")
|
|||||||
// IO-зависимостей. Используется в тестах (быстрый setup, без JDBC) и будет
|
// IO-зависимостей. Используется в тестах (быстрый setup, без JDBC) и будет
|
||||||
// использоваться в Android-сборке (JVector/SQLite не подходят для ART out-of-box).
|
// использоваться в Android-сборке (JVector/SQLite не подходят для ART out-of-box).
|
||||||
include(":message-store-api")
|
include(":message-store-api")
|
||||||
include(":message-log-api")
|
//include(":message-log-api")
|
||||||
|
//include(":message-log-ksqlite")
|
||||||
// Новые API-модули трёх сущностей (canonical имена):
|
// Новые API-модули трёх сущностей (canonical имена):
|
||||||
// - :journal-api — append-only audit log (бывший :message-log-api)
|
// - :journal-api — append-only audit log (бывший :message-log-api)
|
||||||
// - :context-api — то что видит LLM (бывший :working-memory-api)
|
// - :context-api — то что видит LLM (бывший :working-memory-api)
|
||||||
@@ -82,11 +83,12 @@ include(":outbox-api")
|
|||||||
// короткий live tail + recent replay; полный audit log живёт в :message-store-api
|
// короткий live tail + recent replay; полный audit log живёт в :message-store-api
|
||||||
// (там — MessageStore + ConversationStore). EventStore сам управляет eviction,
|
// (там — MessageStore + ConversationStore). EventStore сам управляет eviction,
|
||||||
// никаких prune-методов наружу. KMP, без implementations пока.
|
// никаких prune-методов наружу. KMP, без implementations пока.
|
||||||
include(":event-store")
|
//include(":event-store")
|
||||||
// KMP in-memory реализация MutableEventStore. ConcurrentLinkedDeque + TTL/size
|
// KMP in-memory реализация MutableEventStore. ConcurrentLinkedDeque + TTL/size
|
||||||
// eviction. Для тестов, dev-режима, embedded-сценариев (Android core).
|
// eviction. Для тестов, dev-режима, embedded-сценариев (Android core).
|
||||||
include(":event-store-in-memory")
|
include(":outbox-inmemory")
|
||||||
include(":working-memory-api")
|
//include(":event-store-in-memory")
|
||||||
|
//include(":working-memory-api")
|
||||||
include(":storage-inmemory")
|
include(":storage-inmemory")
|
||||||
// SQLDelight-реализация store'ов из :message-store-api и :working-memory-api. JVM-only
|
// SQLDelight-реализация store'ов из :message-store-api и :working-memory-api. JVM-only
|
||||||
// KMP-реализация EventStore поверх ksqlite (https://github.com/caffeine-mgn/ksqlite).
|
// KMP-реализация EventStore поверх ksqlite (https://github.com/caffeine-mgn/ksqlite).
|
||||||
|
|||||||
@@ -76,7 +76,7 @@ kotlin {
|
|||||||
// Новый единый канал событий агента — заменил старые
|
// Новый единый канал событий агента — заменил старые
|
||||||
// `agentEvents: MutableSharedFlow<AgentEvent>` и per-conv `ConversationEvents._flow`.
|
// `agentEvents: MutableSharedFlow<AgentEvent>` и per-conv `ConversationEvents._flow`.
|
||||||
// Bounded tail + auto-TTL, generic CommonEvent envelope.
|
// Bounded tail + auto-TTL, generic CommonEvent envelope.
|
||||||
implementation(project(":event-store-in-memory"))
|
implementation(project(":outbox-inmemory"))
|
||||||
implementation(project(":agent-toolsets"))
|
implementation(project(":agent-toolsets"))
|
||||||
// Generic LLM-side tools (LlmReflector, SkillMiner, LlmMemoryReviewer,
|
// Generic LLM-side tools (LlmReflector, SkillMiner, LlmMemoryReviewer,
|
||||||
// ContextCompactor, парсеры/промпты). Вынесены из :standalone.
|
// ContextCompactor, парсеры/промпты). Вынесены из :standalone.
|
||||||
|
|||||||
@@ -213,7 +213,7 @@ class ChatAgent(
|
|||||||
* **Live tail + auto-TTL** — клиенты больше не должны заботиться о persistence
|
* **Live tail + auto-TTL** — клиенты больше не должны заботиться о persistence
|
||||||
* или подписке на два отдельных канала.
|
* или подписке на два отдельных канала.
|
||||||
*/
|
*/
|
||||||
private val eventStore: MutableOutboxStore = pw.binom.agentik.eventStore.inmemory.InMemoryEventStore(
|
private val eventStore: MutableOutboxStore = pw.binom.agentik.outbox.inmemory.InMemoryOutboxStore(
|
||||||
maxMessages = null,
|
maxMessages = null,
|
||||||
ttl = null,
|
ttl = null,
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -21,8 +21,8 @@ kotlin {
|
|||||||
sourceSets {
|
sourceSets {
|
||||||
commonMain.dependencies {
|
commonMain.dependencies {
|
||||||
api(project(":message-store-api"))
|
api(project(":message-store-api"))
|
||||||
api(project(":message-log-api"))
|
api(project(":journal-api"))
|
||||||
api(project(":working-memory-api"))
|
api(project(":context-api"))
|
||||||
}
|
}
|
||||||
commonTest.dependencies {
|
commonTest.dependencies {
|
||||||
implementation(kotlin("test"))
|
implementation(kotlin("test"))
|
||||||
|
|||||||
+3
-3
@@ -2,8 +2,8 @@ package pw.binom.agentik.storage.inmemory
|
|||||||
|
|
||||||
import kotlinx.coroutines.sync.Mutex
|
import kotlinx.coroutines.sync.Mutex
|
||||||
import kotlinx.coroutines.sync.withLock
|
import kotlinx.coroutines.sync.withLock
|
||||||
import pw.binom.agentik.messageLog.MessageRecord
|
import pw.binom.agentik.journal.MessageRecord
|
||||||
import pw.binom.agentik.messageLog.MutableMessageStore
|
import pw.binom.agentik.journal.MutableJournalStore
|
||||||
import kotlin.time.Instant
|
import kotlin.time.Instant
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -16,7 +16,7 @@ import kotlin.time.Instant
|
|||||||
* В отличие от ksqlite-импла, не требует native SQLite — работает в любом
|
* В отличие от ksqlite-импла, не требует native SQLite — работает в любом
|
||||||
* KMP-таргете (включая iOS/native).
|
* KMP-таргете (включая iOS/native).
|
||||||
*/
|
*/
|
||||||
class InMemoryMessageStore : MutableMessageStore {
|
class InMemoryMessageStore : MutableJournalStore {
|
||||||
|
|
||||||
private val byConv: MutableMap<String, MutableList<MessageRecord>> = mutableMapOf()
|
private val byConv: MutableMap<String, MutableList<MessageRecord>> = mutableMapOf()
|
||||||
private val mutex = Mutex()
|
private val mutex = Mutex()
|
||||||
|
|||||||
+3
-3
@@ -1,9 +1,9 @@
|
|||||||
package pw.binom.agentik.storage.inmemory
|
package pw.binom.agentik.storage.inmemory
|
||||||
|
|
||||||
import pw.binom.agentik.messageLog.MutableMessageStore
|
import pw.binom.agentik.context.ContextStore
|
||||||
|
import pw.binom.agentik.journal.MutableJournalStore
|
||||||
import pw.binom.agentik.messageStore.ConversationStore
|
import pw.binom.agentik.messageStore.ConversationStore
|
||||||
import pw.binom.agentik.messageStore.ReflectionStore
|
import pw.binom.agentik.messageStore.ReflectionStore
|
||||||
import pw.binom.agentik.context.ContextStore
|
|
||||||
import kotlin.time.Clock
|
import kotlin.time.Clock
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -22,7 +22,7 @@ import kotlin.time.Clock
|
|||||||
object InMemoryStorage {
|
object InMemoryStorage {
|
||||||
data class Bundle(
|
data class Bundle(
|
||||||
val conversationStore: ConversationStore,
|
val conversationStore: ConversationStore,
|
||||||
val messageStore: MutableMessageStore,
|
val messageStore: MutableJournalStore,
|
||||||
val workingMemoryStore: ContextStore,
|
val workingMemoryStore: ContextStore,
|
||||||
val reflectionStore: ReflectionStore,
|
val reflectionStore: ReflectionStore,
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -36,8 +36,8 @@ kotlin {
|
|||||||
implementation(libs.kotlinx.serialization.json)
|
implementation(libs.kotlinx.serialization.json)
|
||||||
|
|
||||||
api(project(":message-store-api"))
|
api(project(":message-store-api"))
|
||||||
api(project(":message-log-api"))
|
api(project(":journal-api"))
|
||||||
api(project(":working-memory-api"))
|
api(project(":context-api"))
|
||||||
}
|
}
|
||||||
commonTest.dependencies {
|
commonTest.dependencies {
|
||||||
implementation(kotlin("test"))
|
implementation(kotlin("test"))
|
||||||
|
|||||||
+3
-3
@@ -1,8 +1,8 @@
|
|||||||
package pw.binom.agentik.storage.ksqlite
|
package pw.binom.agentik.storage.ksqlite
|
||||||
|
|
||||||
import kotlinx.serialization.json.Json
|
import kotlinx.serialization.json.Json
|
||||||
import pw.binom.agentik.messageLog.MessageRecord
|
import pw.binom.agentik.journal.MessageRecord
|
||||||
import pw.binom.agentik.messageLog.MutableMessageStore
|
import pw.binom.agentik.journal.MutableJournalStore
|
||||||
import pw.binom.db.ksqlite.SQLiteConnection
|
import pw.binom.db.ksqlite.SQLiteConnection
|
||||||
import pw.binom.db.ksqlite.SQLitePreparedStatement
|
import pw.binom.db.ksqlite.SQLitePreparedStatement
|
||||||
import kotlin.time.Instant
|
import kotlin.time.Instant
|
||||||
@@ -25,7 +25,7 @@ import kotlinx.coroutines.withContext
|
|||||||
*/
|
*/
|
||||||
class KsqliteMessageStore internal constructor(
|
class KsqliteMessageStore internal constructor(
|
||||||
private val connection: SQLiteConnection,
|
private val connection: SQLiteConnection,
|
||||||
) : MutableMessageStore {
|
) : MutableJournalStore {
|
||||||
|
|
||||||
private val mutex = Mutex()
|
private val mutex = Mutex()
|
||||||
private val json = Json { ignoreUnknownKeys = true }
|
private val json = Json { ignoreUnknownKeys = true }
|
||||||
|
|||||||
+4
-4
@@ -1,9 +1,9 @@
|
|||||||
package pw.binom.agentik.storage.ksqlite
|
package pw.binom.agentik.storage.ksqlite
|
||||||
|
|
||||||
|
import pw.binom.agentik.context.ContextStore
|
||||||
|
import pw.binom.agentik.journal.MutableJournalStore
|
||||||
import pw.binom.agentik.messageStore.ConversationStore
|
import pw.binom.agentik.messageStore.ConversationStore
|
||||||
import pw.binom.agentik.messageLog.MutableMessageStore
|
|
||||||
import pw.binom.agentik.messageStore.ReflectionStore
|
import pw.binom.agentik.messageStore.ReflectionStore
|
||||||
import pw.binom.agentik.workingMemory.WorkingMemoryStore
|
|
||||||
import pw.binom.db.ksqlite.SQLiteConnection
|
import pw.binom.db.ksqlite.SQLiteConnection
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -18,8 +18,8 @@ import pw.binom.db.ksqlite.SQLiteConnection
|
|||||||
class KsqliteStores internal constructor(
|
class KsqliteStores internal constructor(
|
||||||
val connection: SQLiteConnection,
|
val connection: SQLiteConnection,
|
||||||
val conversations: ConversationStore,
|
val conversations: ConversationStore,
|
||||||
val messages: MutableMessageStore,
|
val messages: MutableJournalStore,
|
||||||
val workingMemory: WorkingMemoryStore,
|
val workingMemory: ContextStore,
|
||||||
val reflections: ReflectionStore,
|
val reflections: ReflectionStore,
|
||||||
) : AutoCloseable {
|
) : AutoCloseable {
|
||||||
|
|
||||||
|
|||||||
@@ -24,7 +24,7 @@ kotlin {
|
|||||||
|
|
||||||
sourceSets {
|
sourceSets {
|
||||||
commonMain.dependencies {
|
commonMain.dependencies {
|
||||||
api(project(":message-log-api"))
|
// api(project(":message-log-api"))
|
||||||
api(project(":context-api"))
|
api(project(":context-api"))
|
||||||
api(libs.kotlinx.coroutines.core)
|
api(libs.kotlinx.coroutines.core)
|
||||||
api(libs.kotlinx.serialization.core)
|
api(libs.kotlinx.serialization.core)
|
||||||
|
|||||||
-16
@@ -1,16 +0,0 @@
|
|||||||
package pw.binom.agentik.workingMemory
|
|
||||||
|
|
||||||
/**
|
|
||||||
* @Deprecated
|
|
||||||
* Перенесено в `:context-api`. Используй `pw.binom.agentik.context.WorkingMemoryEntry`.
|
|
||||||
*
|
|
||||||
* Backward-compat typealias. Все импорты `pw.binom.agentik.workingMemory.WorkingMemoryEntry`
|
|
||||||
* продолжают работать как тип `pw.binom.agentik.context.WorkingMemoryEntry` (тот же тип, транспарентно).
|
|
||||||
*
|
|
||||||
* Удалить когда все импорты будут на `:context-api`.
|
|
||||||
*/
|
|
||||||
@Deprecated(
|
|
||||||
message = "Перенесено в :context-api. Используй pw.binom.agentik.context.WorkingMemoryEntry.",
|
|
||||||
replaceWith = ReplaceWith("WorkingMemoryEntry", "pw.binom.agentik.context.WorkingMemoryEntry"),
|
|
||||||
)
|
|
||||||
typealias WorkingMemoryEntry = pw.binom.agentik.context.WorkingMemoryEntry
|
|
||||||
-18
@@ -1,18 +0,0 @@
|
|||||||
package pw.binom.agentik.workingMemory
|
|
||||||
|
|
||||||
/**
|
|
||||||
* @Deprecated
|
|
||||||
* Перенесено в `:context-api`. Используй `pw.binom.agentik.context.ContextStore`.
|
|
||||||
*
|
|
||||||
* Этот typealias сохранён временно для backward-compat: существующие impl'ы
|
|
||||||
* (`KsqliteStores.workingMemory` в `:storage-ksqlite`, тесты) импортируют
|
|
||||||
* `pw.binom.agentik.workingMemory.WorkingMemoryStore`. Буквально тот же тип,
|
|
||||||
* что и `pw.binom.agentik.context.ContextStore` — typealias транспарентен.
|
|
||||||
*
|
|
||||||
* Удалить когда все импорты будут на `:context-api`.
|
|
||||||
*/
|
|
||||||
@Deprecated(
|
|
||||||
message = "Перенесено в :context-api. Используй pw.binom.agentik.context.ContextStore.",
|
|
||||||
replaceWith = ReplaceWith("ContextStore", "pw.binom.agentik.context.ContextStore"),
|
|
||||||
)
|
|
||||||
typealias WorkingMemoryStore = pw.binom.agentik.context.ContextStore
|
|
||||||
Reference in New Issue
Block a user