refactor(message-log-api): split MessageStore into read-only and mutable interfaces
ci / JVM build + tests (push) Successful in 6m14s
ci / JVM build + tests (push) Successful in 6m14s
- Introduced `MutableMessageStore` for producers with an `append` operation, separate from read-only `MessageStore`. - Updated all consumers and implementations to use the appropriate interface (`read-only` for observers, `mutable` for producers). - Improves modularity and ensures compile-time guarantees against unintended write operations in the audit log.
This commit is contained in:
@@ -4,21 +4,27 @@ import kotlinx.coroutines.flow.Flow
|
|||||||
import kotlin.time.Instant
|
import kotlin.time.Instant
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Append-only audit log сообщений.
|
* Append-only audit log сообщений — read-only представление.
|
||||||
*
|
*
|
||||||
* Только `append` и чтение. Никаких обновлений, никакого удаления
|
* Producer-операция [MutableMessageStore.append] находится на
|
||||||
* (кроме каскадного вместе с ConversationStore.delete).
|
* [MutableMessageStore] — этот интерфейс только для чтения, чтобы
|
||||||
|
* consumer'ы физически не могли писать в audit log.
|
||||||
|
*
|
||||||
|
* Никаких обновлений, никакого удаления (кроме каскадного вместе
|
||||||
|
* с ConversationStore.delete).
|
||||||
*/
|
*/
|
||||||
interface MessageStore : AutoCloseable {
|
interface MessageStore : AutoCloseable {
|
||||||
|
|
||||||
suspend fun append(record: MessageRecord)
|
|
||||||
|
|
||||||
suspend fun list(conversationId: String, after: Instant, offset: Int, limit: Int): List<MessageRecord>
|
suspend fun list(conversationId: String, after: Instant, offset: Int, limit: Int): List<MessageRecord>
|
||||||
|
|
||||||
suspend fun listAll(conversationId: String): List<MessageRecord>
|
suspend fun listAll(conversationId: String): List<MessageRecord>
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Subscribe на события append'ов. Default — `emptyFlow()` для store'ов
|
||||||
|
* без live-уведомлений (например, SQLite без триггеров); impl'ы с
|
||||||
|
* in-memory notification (`:storage-inmemory`) override'ят.
|
||||||
|
*/
|
||||||
fun events(): Flow<MessageEvent> = kotlinx.coroutines.flow.emptyFlow()
|
fun events(): Flow<MessageEvent> = kotlinx.coroutines.flow.emptyFlow()
|
||||||
|
|
||||||
suspend fun tokenStats(conversationId: String): TokenStats
|
suspend fun tokenStats(conversationId: String): TokenStats
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
+17
@@ -0,0 +1,17 @@
|
|||||||
|
package pw.binom.agentik.messageLog
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Mutable вариант [MessageStore] — добавляет producer-операцию [append].
|
||||||
|
*
|
||||||
|
* Этот интерфейс предназначен **только для producer'ов** (ChatAgent,
|
||||||
|
* ConversationLoop, ToolDispatcher, sub-agents, A2A-bridge).
|
||||||
|
* Consumer'ы (DebugRoutes, admin dashboards, parent agents) должны
|
||||||
|
* принимать **read-only** [MessageStore] — тогда невозможно случайно
|
||||||
|
* записать в audit log из observer'а.
|
||||||
|
*
|
||||||
|
* **Append семантика**: см. KDoc [MessageStore.append][MessageStore] —
|
||||||
|
* на этом интерфейсе (не дублируем).
|
||||||
|
*/
|
||||||
|
interface MutableMessageStore : MessageStore {
|
||||||
|
suspend fun append(record: MessageRecord)
|
||||||
|
}
|
||||||
@@ -30,7 +30,7 @@ import pw.binom.agentik.messageStore.ConversationStore
|
|||||||
import pw.binom.agentik.messageStore.Ids
|
import pw.binom.agentik.messageStore.Ids
|
||||||
import pw.binom.agentik.messageStore.Reflection
|
import pw.binom.agentik.messageStore.Reflection
|
||||||
import pw.binom.agentik.messageStore.ReflectionStore
|
import pw.binom.agentik.messageStore.ReflectionStore
|
||||||
import pw.binom.agentik.messageLog.MessageStore
|
import pw.binom.agentik.messageLog.MutableMessageStore
|
||||||
import pw.binom.agentik.workingMemory.WorkingMemoryStore
|
import pw.binom.agentik.workingMemory.WorkingMemoryStore
|
||||||
import pw.binom.agentik.toolsets.DisableToolsetTool
|
import pw.binom.agentik.toolsets.DisableToolsetTool
|
||||||
import pw.binom.agentik.toolsets.EnableToolsetTool
|
import pw.binom.agentik.toolsets.EnableToolsetTool
|
||||||
@@ -64,7 +64,7 @@ import pw.binom.litert.LiteLlm
|
|||||||
class ChatAgent(
|
class ChatAgent(
|
||||||
override val id: String,
|
override val id: String,
|
||||||
private val conversationStore: ConversationStore,
|
private val conversationStore: ConversationStore,
|
||||||
private val messageStore: MessageStore,
|
private val messageStore: MutableMessageStore,
|
||||||
private val workingMemoryStore: WorkingMemoryStore,
|
private val workingMemoryStore: WorkingMemoryStore,
|
||||||
private val reflectionStore: ReflectionStore,
|
private val reflectionStore: ReflectionStore,
|
||||||
private val llm: LiteLlm,
|
private val llm: LiteLlm,
|
||||||
|
|||||||
@@ -34,7 +34,7 @@ import pw.binom.agentik.messageStore.ConversationStore
|
|||||||
import pw.binom.agentik.messageLog.MessageContext
|
import pw.binom.agentik.messageLog.MessageContext
|
||||||
import pw.binom.agentik.messageLog.MessageOrigin
|
import pw.binom.agentik.messageLog.MessageOrigin
|
||||||
import pw.binom.agentik.messageLog.MessageRecord
|
import pw.binom.agentik.messageLog.MessageRecord
|
||||||
import pw.binom.agentik.messageLog.MessageStore
|
import pw.binom.agentik.messageLog.MutableMessageStore
|
||||||
import pw.binom.agentik.messageStore.ReflectionStore
|
import pw.binom.agentik.messageStore.ReflectionStore
|
||||||
import pw.binom.agentik.messageLog.TurnTokens
|
import pw.binom.agentik.messageLog.TurnTokens
|
||||||
import pw.binom.agentik.workingMemory.WorkingMemoryEntry
|
import pw.binom.agentik.workingMemory.WorkingMemoryEntry
|
||||||
@@ -55,7 +55,7 @@ import pw.binom.agentik.toolsets.NamedTool
|
|||||||
class ConversationLoop(
|
class ConversationLoop(
|
||||||
record: ConversationRecord,
|
record: ConversationRecord,
|
||||||
private val conversationStore: ConversationStore,
|
private val conversationStore: ConversationStore,
|
||||||
private val messageStore: MessageStore,
|
private val messageStore: MutableMessageStore,
|
||||||
private val workingMemoryStore: WorkingMemoryStore,
|
private val workingMemoryStore: WorkingMemoryStore,
|
||||||
private val reflectionStore: ReflectionStore?,
|
private val reflectionStore: ReflectionStore?,
|
||||||
private val eventStore: pw.binom.agentik.eventStore.MutableEventStore,
|
private val eventStore: pw.binom.agentik.eventStore.MutableEventStore,
|
||||||
|
|||||||
@@ -6,7 +6,7 @@ import kotlinx.coroutines.async
|
|||||||
import mu.KotlinLogging
|
import mu.KotlinLogging
|
||||||
import pw.binom.agentik.proto.Event as ProtoEvent
|
import pw.binom.agentik.proto.Event as ProtoEvent
|
||||||
import pw.binom.agentik.messageLog.MessageRecord
|
import pw.binom.agentik.messageLog.MessageRecord
|
||||||
import pw.binom.agentik.messageLog.MessageStore
|
import pw.binom.agentik.messageLog.MutableMessageStore
|
||||||
import pw.binom.agentik.workingMemory.WorkingMemoryEntry
|
import pw.binom.agentik.workingMemory.WorkingMemoryEntry
|
||||||
import pw.binom.agentik.toolsets.ToolsetDispatchPolicy
|
import pw.binom.agentik.toolsets.ToolsetDispatchPolicy
|
||||||
import pw.binom.litert.LiteToolCall
|
import pw.binom.litert.LiteToolCall
|
||||||
@@ -16,7 +16,7 @@ import pw.binom.agentik.toolsets.NamedTool
|
|||||||
|
|
||||||
internal class ToolDispatcher(
|
internal class ToolDispatcher(
|
||||||
private val state: ConversationState,
|
private val state: ConversationState,
|
||||||
private val messageStore: MessageStore,
|
private val messageStore: MutableMessageStore,
|
||||||
private val events: ConversationEvents,
|
private val events: ConversationEvents,
|
||||||
private val backgroundEvents: BackgroundEventBus,
|
private val backgroundEvents: BackgroundEventBus,
|
||||||
private val toolsByName: MutableMap<String, NamedTool>,
|
private val toolsByName: MutableMap<String, NamedTool>,
|
||||||
|
|||||||
+2
-2
@@ -7,7 +7,7 @@ import kotlinx.coroutines.sync.Mutex
|
|||||||
import kotlinx.coroutines.sync.withLock
|
import kotlinx.coroutines.sync.withLock
|
||||||
import pw.binom.agentik.messageLog.MessageEvent
|
import pw.binom.agentik.messageLog.MessageEvent
|
||||||
import pw.binom.agentik.messageLog.MessageRecord
|
import pw.binom.agentik.messageLog.MessageRecord
|
||||||
import pw.binom.agentik.messageLog.MessageStore
|
import pw.binom.agentik.messageLog.MutableMessageStore
|
||||||
import pw.binom.agentik.messageLog.TokenStats
|
import pw.binom.agentik.messageLog.TokenStats
|
||||||
import kotlin.time.Instant
|
import kotlin.time.Instant
|
||||||
|
|
||||||
@@ -22,7 +22,7 @@ import kotlin.time.Instant
|
|||||||
* любом KMP-таргете (включая iOS/native, где SQLite через Android driver
|
* любом KMP-таргете (включая iOS/native, где SQLite через Android driver
|
||||||
* недоступен).
|
* недоступен).
|
||||||
*/
|
*/
|
||||||
class InMemoryMessageStore : MessageStore {
|
class InMemoryMessageStore : MutableMessageStore {
|
||||||
|
|
||||||
private val byConv: MutableMap<String, MutableList<MessageRecord>> = mutableMapOf()
|
private val byConv: MutableMap<String, MutableList<MessageRecord>> = mutableMapOf()
|
||||||
private val events = MutableSharedFlow<MessageEvent>(extraBufferCapacity = 64)
|
private val events = MutableSharedFlow<MessageEvent>(extraBufferCapacity = 64)
|
||||||
|
|||||||
+2
-2
@@ -1,6 +1,6 @@
|
|||||||
package pw.binom.agentik.storage.inmemory
|
package pw.binom.agentik.storage.inmemory
|
||||||
|
|
||||||
import pw.binom.agentik.messageLog.MessageStore
|
import pw.binom.agentik.messageLog.MutableMessageStore
|
||||||
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.workingMemory.WorkingMemoryStore
|
import pw.binom.agentik.workingMemory.WorkingMemoryStore
|
||||||
@@ -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: MessageStore,
|
val messageStore: MutableMessageStore,
|
||||||
val workingMemoryStore: WorkingMemoryStore,
|
val workingMemoryStore: WorkingMemoryStore,
|
||||||
val reflectionStore: ReflectionStore,
|
val reflectionStore: ReflectionStore,
|
||||||
)
|
)
|
||||||
|
|||||||
+2
-2
@@ -2,7 +2,7 @@ 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.messageLog.MessageRecord
|
||||||
import pw.binom.agentik.messageLog.MessageStore
|
import pw.binom.agentik.messageLog.MutableMessageStore
|
||||||
import pw.binom.agentik.messageLog.TokenStats
|
import pw.binom.agentik.messageLog.TokenStats
|
||||||
import pw.binom.agentik.messageLog.decodeBodyPayload
|
import pw.binom.agentik.messageLog.decodeBodyPayload
|
||||||
import pw.binom.agentik.messageLog.encodeBodyPayload
|
import pw.binom.agentik.messageLog.encodeBodyPayload
|
||||||
@@ -24,7 +24,7 @@ import kotlinx.coroutines.withContext
|
|||||||
*/
|
*/
|
||||||
class KsqliteMessageStore(
|
class KsqliteMessageStore(
|
||||||
private val connection: SQLiteConnection,
|
private val connection: SQLiteConnection,
|
||||||
) : MessageStore {
|
) : MutableMessageStore {
|
||||||
|
|
||||||
private val mutex = Mutex()
|
private val mutex = Mutex()
|
||||||
private val json = Json { ignoreUnknownKeys = true }
|
private val json = Json { ignoreUnknownKeys = true }
|
||||||
|
|||||||
+2
-2
@@ -1,7 +1,7 @@
|
|||||||
package pw.binom.agentik.storage.ksqlite
|
package pw.binom.agentik.storage.ksqlite
|
||||||
|
|
||||||
import pw.binom.agentik.messageStore.ConversationStore
|
import pw.binom.agentik.messageStore.ConversationStore
|
||||||
import pw.binom.agentik.messageLog.MessageStore
|
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.agentik.workingMemory.WorkingMemoryStore
|
||||||
import pw.binom.db.ksqlite.SQLiteConnection
|
import pw.binom.db.ksqlite.SQLiteConnection
|
||||||
@@ -18,7 +18,7 @@ import pw.binom.db.ksqlite.SQLiteConnection
|
|||||||
class KsqliteStores private constructor(
|
class KsqliteStores private constructor(
|
||||||
val connection: SQLiteConnection,
|
val connection: SQLiteConnection,
|
||||||
val conversations: ConversationStore,
|
val conversations: ConversationStore,
|
||||||
val messages: MessageStore,
|
val messages: MutableMessageStore,
|
||||||
val workingMemory: WorkingMemoryStore,
|
val workingMemory: WorkingMemoryStore,
|
||||||
val reflections: ReflectionStore,
|
val reflections: ReflectionStore,
|
||||||
) : AutoCloseable {
|
) : AutoCloseable {
|
||||||
|
|||||||
+2
-2
@@ -2,7 +2,7 @@ package pw.binom.agentik.storage.sqlite
|
|||||||
|
|
||||||
import kotlinx.serialization.json.Json
|
import kotlinx.serialization.json.Json
|
||||||
import pw.binom.agentik.messageLog.MessageRecord
|
import pw.binom.agentik.messageLog.MessageRecord
|
||||||
import pw.binom.agentik.messageLog.MessageStore
|
import pw.binom.agentik.messageLog.MutableMessageStore
|
||||||
import pw.binom.agentik.messageLog.TokenStats
|
import pw.binom.agentik.messageLog.TokenStats
|
||||||
import pw.binom.agentik.messageLog.decodeBodyPayload
|
import pw.binom.agentik.messageLog.decodeBodyPayload
|
||||||
import pw.binom.agentik.messageLog.encodeBodyPayload
|
import pw.binom.agentik.messageLog.encodeBodyPayload
|
||||||
@@ -19,7 +19,7 @@ import kotlin.time.Instant
|
|||||||
* SQLDelight сохраняет snake_case в сгенерированной data class (`Message`),
|
* SQLDelight сохраняет snake_case в сгенерированной data class (`Message`),
|
||||||
* поэтому обращаемся через `conversation_id`, `payload_json`, `created_at`.
|
* поэтому обращаемся через `conversation_id`, `payload_json`, `created_at`.
|
||||||
*/
|
*/
|
||||||
class SqliteMessageStore(private val db: AgentikDatabase) : MessageStore {
|
class SqliteMessageStore(private val db: AgentikDatabase) : MutableMessageStore {
|
||||||
|
|
||||||
private val q get() = db.messageQueries
|
private val q get() = db.messageQueries
|
||||||
|
|
||||||
|
|||||||
@@ -4,7 +4,7 @@ import app.cash.sqldelight.db.QueryResult
|
|||||||
import app.cash.sqldelight.db.SqlDriver
|
import app.cash.sqldelight.db.SqlDriver
|
||||||
import app.cash.sqldelight.driver.jdbc.sqlite.JdbcSqliteDriver
|
import app.cash.sqldelight.driver.jdbc.sqlite.JdbcSqliteDriver
|
||||||
import pw.binom.agentik.messageStore.ConversationStore
|
import pw.binom.agentik.messageStore.ConversationStore
|
||||||
import pw.binom.agentik.messageLog.MessageStore
|
import pw.binom.agentik.messageLog.MutableMessageStore
|
||||||
import pw.binom.agentik.messageStore.ReflectionStore
|
import pw.binom.agentik.messageStore.ReflectionStore
|
||||||
import pw.binom.agentik.storage.sqlite.SqliteReflectionStore
|
import pw.binom.agentik.storage.sqlite.SqliteReflectionStore
|
||||||
import pw.binom.agentik.workingMemory.WorkingMemoryStore
|
import pw.binom.agentik.workingMemory.WorkingMemoryStore
|
||||||
@@ -15,7 +15,7 @@ import pw.binom.agentik.workingMemory.WorkingMemoryStore
|
|||||||
class SqliteStores private constructor(
|
class SqliteStores private constructor(
|
||||||
val driver: SqlDriver,
|
val driver: SqlDriver,
|
||||||
val conversations: ConversationStore,
|
val conversations: ConversationStore,
|
||||||
val messages: MessageStore,
|
val messages: MutableMessageStore,
|
||||||
val workingMemory: WorkingMemoryStore,
|
val workingMemory: WorkingMemoryStore,
|
||||||
val reflections: ReflectionStore,
|
val reflections: ReflectionStore,
|
||||||
) : AutoCloseable {
|
) : AutoCloseable {
|
||||||
|
|||||||
Reference in New Issue
Block a user