standalone v1: dual-backend (openai + litert-google) with SQLDelight dual-log persistence

Replace EchoProtoAgent / EchoAgent / EchoA2aHandler placeholders with a real
stateful agent on top of SQLite (SQLDelight 2.3.2) and litert-api v6.

persistence (commonMain):
- ConversationStore / MessageStore / WorkingMemoryStore — three narrow
  interfaces, all operations suspend, AutoCloseable.
- MessageRecord sealed: UserMessage / AssistantMessage (Body subtype),
  ToolCall / ToolResult (audit-only), Summary / System (working-memory-only
  synthetic). Snake-case @SerialName discriminators.
- WorkingMemoryEntry sealed: System / User(sourceMessageId) /
  Assistant(sourceMessageId); sourceMessageId is null for System.
- Two-table dual-log model: append-only message audit + mutable
  working_memory with monotonic order_idx.

SQLite (jvmMain):
- SQLDelight schema + SqliteConversationStore / SqliteMessageStore /
  SqliteWorkingMemoryStore under src/jvmMain/sqldelight/.
- SqliteStores.open(path) / inMemory(); Schema.create gated on
  sqlite_master probe for idempotency.
- All payload_json is the MessageRecord encoded as JSON; subtype-specific
  fields avoid migrations.

agent (jvmMain):
- ChatAgent — stateful proto.Agent with live in-memory cache, lock-protected,
  AgentEvent bus (Created/Deleted).
- ChatConversation — long-lived LiteConversation handle; created lazily on
  first send from working_memory (system + initial messages), reused across
  all subsequent turns (REQUIRED for litert-google KV-cache).
- Per turn: append User to audit + WM → sendStreamContents (wrapped in
  transformWhile for litert-google-jvm 0.16.1 isDone workaround) → emit
  AppendText deltas → append Assistant to audit + WM + touch conversation.
- isClosed flag so getConversation reconstructs after close.

llm (jvmMain):
- LlmConfig data class with LlmBackend enum (OPENAI / GOOGLE); fromEnv
  parses AGENTIK_LLM_BACKEND and dispatches to backend-specific config.
- OpenAI: litert-openai, OpenAI-compatible endpoint, validated
  baseUrl/apiKey/model.
- Google: litert-google (reflection-resolved pw.binom.litert.google
  factory) on top of litertlm-jvm 0.16.1 native engine;
  visionBackend/audioBackend = null (LiteRT-LM 0.16.1 binds encoder
  graph even with null backend, but a model lacking encoder crashes;
  null is the correct "don't bind" signal).
- foldSystemIntoFirstUser (default true for GOOGLE) folds system prompt
  into the first user message to avoid chat template alternation issues.

build:
- Add sqldelight plugin + runtime + sqlite-driver + coroutines-extensions
  to gradle/libs.versions.toml.
- litert-openai: implementation; litert-google: runtimeOnly (resolved via
  reflection at runtime).
- KMP jvm executable via @OptIn(ExperimentalKotlinGradlePluginApi) +
  jvm { binaries { executable { mainClass.set("...MainKt") } } }.

tests (jvmTest): 30 passing
- PersistenceTest (11): conversation upsert/list/cascade-delete/rename/
  touch; message audit append/list; working-memory order preservation;
  image-content payload roundtrip.
- ChatAgentTest (14): system-prompt seeding; persistent vs temp
  persistence across SqliteStores reopen; multi-turn audit + WM growth;
  interrupt of in-flight slow send; agentEvents Created/Deleted flow;
  closed-conv reconstruct via getConversation.
- LlmConfigTest (6): env happy path, defaults, missing fields throw.

smoke tested e2e:
- openai backend against real llm.binom.pw/v1 (myopenai/local/codding)
  — multi-turn dialogue persisted, kill -9 + restart survives.
- google backend against gemma-4-E2B-it.litertlm — multi-turn
  ("Hello there!" → "2 + 2 = 4"), KV-cache survives across turns,
  SSE start→append_text*→end cleanly closes.

docs/STANDALONE.md updated for v1 architecture, dual-backend env table,
long-lived LiteConversation invariant, and litert-google-jvm 0.16.1
isDone-stream workaround.
This commit is contained in:
2026-09-13 12:39:36 +03:00
parent a3581abf84
commit 9e5d61707d
31 changed files with 2507 additions and 316 deletions
+296
View File
@@ -0,0 +1,296 @@
# Standalone — рантайм агента agentik
`standalone` — это исполняемое JVM-приложение (точка входа `pw.binom.agentik.standalone.MainKt`), которое поднимает реальный агент `ChatAgent` (stateful, SQLite-персистентный) и навешивает на него HTTP+SSE фасад `:server`. LLM-движок выбирается через `AGENTIK_LLM_BACKEND` — на v1 поддерживаются `litert-openai` (любой OpenAI-совместимый endpoint) и `litert-google` (on-device движок LiteRT-LM 0.16.1 через нативную `.so`-библиотеку).
Этот документ описывает, как `standalone` собран и как его расширять.
---
## 1. Что в коробке после `git clone`
```
:proto — единое ядро протокола (KMP, commonMain)
Agent / Conversation / Event / Message / Content / AgentEvent
stateful: агент сам хранит историю и working memory
:server — HTTP+SSE фасад :proto
public Route.agentikAgent(agent, path = "/agentik")
:standalone — JVM-рантайм с реальным LLM-агентом
ChatAgent + ChatConversation поверх SQLite и litert-* (openai/google)
default 8080:
GET /health health check
POST /agentik/conversations создать диалог (201)
GET /agentik/conversations список
GET /agentik/conversations/{id} один диалог
PATCH /agentik/conversations/{id} переименовать
DELETE /agentik/conversations/{id} удалить (204)
POST /agentik/conversations/{id}/messages отправить user-сообщение (202)
POST /agentik/conversations/{id}/interrupt прервать текущий ход (202)
GET /agentik/conversations/{id}/messages страница истории
GET /agentik/conversations/{id}/events SSE live-события хода
GET /agentik/events SSE live-события агента
```
Полная таблица эндпоинтов — в `docs/ARCHITECTURE.md` (раздел «:server»).
---
## 2. Архитектура слоёв
```
клиенты транспорт
┌───────────────┐ ┌─────────────────────────────┐
│ Web / CLI / │ ──HTTP──► │ Route.agentikAgent(agent) │
│ desktop │ ──SSE───► │ :server (Ktor + Netty) │
│ │ └──────────────┬──────────────┘
└───────────────┘ │
▼
pw.binom.agentik.proto.Agent
(ChatAgent)
│
┌───────────────┴───────────────┐
▼ ▼
ChatConversation.send(content) agent.events / agent.getConversations
│
▼
┌──────────────────────────────────────┐
│ 1. audit: append UserMessage │
│ 2. working_memory: append User │
│ 3. ensureLiteConversation: │
│ first turn → create from WM; │
│ next turns → reuse (KV-cache) │
│ 4. sendStreamContents → emit │
│ StartResponse / AppendText / │
│ End │
│ 5. audit + WM: append AssistantMessage│
└──────────────────────────────────────┘
│ │
▼ ▼
Conversation.events(after) Conversation.getMessages(after)
(live, no replay) (история)
```
Слои рантайма:
```
:standalone
├── persistence/ ← интерфейсы и records (commonMain, без зависимостей)
│ ConversationStore / MessageStore / WorkingMemoryStore
│ ConversationRecord / MessageRecord / WorkingMemoryEntry / Content
│ Payload.kt — JSON-сериализация
├── persistence/sqlite/ ← JVM: SQLDelight-схема + три SQLite-реализации
│ SqliteStores.open(path | inMemory)
│ src/jvmMain/sqldelight/.../*.sq
├── llm/ ← LlmConfig (env → OpenAI/Google), backend-agnostic
└── agent/ ← ChatAgent + ChatConversation (stateful, long-lived LiteConv)
```
Ключевой инвариант: **разговор живёт внутри агента, а не в клиенте и не в транспорте.** Транспорт — лишь сериализатор: HTTP пишет в/читает из `:server`-эндпоинтов, SSE шлёт события. У них нет своего состояния диалога.
---
## 3. Точка входа: `pw.binom.agentik.standalone.MainKt`
```kotlin
fun main() {
val port = System.getenv("AGENTIK_PORT")?.toIntOrNull() ?: 8080
val dbPath = System.getenv("AGENTIK_DB_PATH")?.takeIf { it.isNotBlank() } ?: "./agentik.db"
val llmConfig = LlmConfig.fromEnv()
val llm = llmConfig.createLlm()
val stores = SqliteStores.open(dbPath = dbPath)
val agent = ChatAgent(
id = "agentik",
stores = stores,
llm = llm,
llmConfig = llmConfig,
)
val server = embeddedServer(Netty, port = port) {
routing {
get("/health") { call.respondText("ok") }
agentikAgent(agent, path = "/agentik")
}
}
Runtime.getRuntime().addShutdownHook(Thread {
agent.close(); stores.close(); llm.close()
})
server.start(wait = true)
}
```
Один `ChatAgent` отвечает и за диалоги (`/agentik/conversations/...`), и за live-события (`/agentik/events`). Все три ресурса — БД, LLM, Netty — корректно закрываются в shutdown-хуке.
---
## 4. Контракт `Agent` (от `pw.binom.agentik.proto`)
Реализация **обязана** уметь:
| метод | смысл |
|---|---|
| `id: String` | идентификатор агента |
| `createConversation(temp: Boolean): Conversation` | новая сессия, `temp=true` — не персистить |
| `getConversation(id): Conversation?` | достать по id, `null` если нет |
| `deleteConversation(id): Boolean` | удалить |
| `getConversations(offset, limit)` | страница списка |
| `events(after: Instant): Flow<AgentEvent>` | live-события по множеству разговоров |
Реализация `Conversation`:
| метод | смысл |
|---|---|
| `id: String` | идентификатор диалога |
| `title: String?` | заголовок (может быть `null`) |
| `isTemporal: Boolean` | `true` = не персистить (`temp=true` при создании) |
| `isSupportImageInput/Output: Boolean` | мультимодальные возможности (для v1 оба `false`) |
| `updatedAt: Instant` | последний `send`/`rename` |
| `send(content: List<Content>)` | **fire-and-forget**: добавить user-сообщение, запустить ход, выйти |
| `interrupt()` | остановить текущий ход (best-effort) |
| `events(after): Flow<Event>` | live-события хода (StartReasoning, StartResponse, AppendText, End, Interrupted, Error) |
| `getMessages(after, offset, limit)` | страница истории |
| `rename(title)` | переименовать |
| `close()` | освободить ресурсы |
Главное: `send` ничего не возвращает. Чтобы получить события, нужно **отдельно** подписаться на `events(after)` ДО `send` либо сразу после — поток событий стартует с момента подписки, бэкфилл через `getMessages`.
---
## 5. Persistence — dual-log
`ChatAgent` хранит каждую сессию в двух логически разных таблицах:
| таблица | назначение | мутации |
|---|---|---|
| `message` | append-only audit log. Все user/assistant/tool-call/tool-result сообщения. Никогда не редактируется (кроме каскадного `DELETE` при удалении диалога). | только `INSERT` |
| `working_memory` | mutable LLM-контекст. System-prompt + текущая история + (в v2) суммаризации. | `INSERT`, `compact(dropFromIdx, summary)` |
Маппинг `:proto.Message ↔ MessageRecord` живёт в `ChatConversation.kt` (`toProto`/`toStorage`) — сами `MessageRecord` намеренно НЕ зависят от `:proto`, чтобы можно было сменить транспорт без миграции таблиц.
Подробный контракт — в комментариях к `MessageRecord.kt` и `WorkingMemoryEntry.kt`.
### `MessageStore`
```kotlin
suspend fun append(record: MessageRecord)
suspend fun list(conversationId: String, after: Instant, offset: Int, limit: Int): List<MessageRecord>
suspend fun listAll(conversationId: String): List<MessageRecord>
```
### `WorkingMemoryStore`
```kotlin
suspend fun append(conversationId: String, entry: WorkingMemoryEntry, now: Instant)
suspend fun list(conversationId: String): List<WorkingMemoryRow>
suspend fun clear(conversationId: String)
suspend fun compact(dropFromOrderIdx: Long, conversationId: String): Long
```
`compact` — атомарный «выбросить всё от `dropFromOrderIdx` и дальше, вставить новую синтетическую запись на следующий `order_idx`». Для v1 — просто `DELETE` от индекса (суммаризация появится в v2 вместе с LLM-вызовом для генерации текста).
### `ConversationStore`
```kotlin
suspend fun upsert(record: ConversationRecord)
suspend fun get(id: String): ConversationRecord?
suspend fun delete(id: String): Boolean // каскадно чистит message + working_memory
suspend fun list(offset: Int, limit: Int): List<ConversationRecord>
suspend fun rename(id: String, title: String?): Instant?
suspend fun touch(id: String, now: Instant)
```
Все три store — `AutoCloseable`; корневой ресурс `SqliteStores` закрывает их вместе с `SqlDriver`.
---
## 6. LLM-конфигурация (`LlmConfig`)
`LlmConfig.fromEnv()` парсит env, валидирует обязательные поля и выбирает бэкенд через `AGENTIK_LLM_BACKEND`:
- `openai` (default) — `litert-openai`, текст-онли чат против любого OpenAI-совместимого endpoint.
- `google` — `litert-google` (LiteRT-LM 0.16.1, on-device `.task`/`.litertlm` модель, требует нативной библиотеки через `litertlm-jvm`).
### Общие env
| env | смысл | default |
|---|---|---|
| `AGENTIK_PORT` | порт Netty (`/agentik`, `/health`) | `8080` |
| `AGENTIK_DB_PATH` | путь к SQLite-файлу | `./agentik.db` |
| `AGENTIK_LLM_BACKEND` | `openai` или `google` | `openai` |
| `AGENTIK_SYSTEM_PROMPT` | текст системного промпта | «Ты полезный ассистент. Отвечай кратко и по делу.» |
### Backend `openai`
| env | смысл |
|---|---|
| `OPENAI_BASE_URL` | endpoint (например, `https://api.openai.com/v1` или `http://localhost:11434/v1`) — обязательно |
| `OPENAI_API_KEY` | ключ модели — обязательно |
| `OPENAI_MODEL` | имя модели (например, `gpt-4o-mini`, `myopenai/local/codding`) — обязательно |
### Backend `google` (on-device LiteRT-LM)
| env | смысл | default |
|---|---|---|
| `AGENTIK_GOOGLE_MODEL_PATH` | путь к `.task` или `.litertlm` модели — обязательно |
| `AGENTIK_GOOGLE_CACHE_DIR` | каталог кеша скомпилированных graph'ов | пусто (системный tmp) |
| `AGENTIK_GOOGLE_THREADS` | число CPU-потоков для движка | `4` |
`AGENTIK_DB_PATH=:memory:` создаёт in-memory БД (только для тестов и интеграционных проверок).
### Long-lived LiteConversation — ОБЯЗАТЕЛЬНО для google
`LiteConversation` от любого litert-бэкенда — это **долгоживущая stateful ручка**: она держит историю сообщений и (для google) KV-cache/sampler-state. На v1 `ChatConversation` создаёт `LiteConversation` один раз (на первом `send`) и переиспользует на всех последующих turn'ах той же беседы. Пересоздание LiteConversation на каждый send ломает KV-cache для google (и в v1 приводит к ошибке chat template «roles must alternate» после двух user-сообщений).
`AGENTIK_GOOGLE_FOLD_SYSTEM_INTO_FIRST_USER` (default `true` для google) — фолдит системный промпт в первое user-сообщение, чтобы движок увидел только `[User+system, Assistant, User, ...]` и не падал на alternation. Для openai этот режим отключён.
### Workaround для litert-google-jvm 0.16.1
Баг в `sendStreamContents` — поток не закрывается после `isDone = true`. `ChatConversation` оборачивает стрим в `transformWhile { !delta.isDone }`, поэтому SSE получает корректный `End` event.
---
## 7. Расширение
### Подключить тул (v2)
```kotlin
val tools: List<LiteTool> = listOf(ReadFileTool(Path.of("/work")))
val cfg = LiteConversationConfig(systemInstruction = "…", tools = tools)
```
LiteConversation сообщит модели о доступных тулах. После появления `LiteContentPart.ToolResult` (v2) появится и `ChatConversation.addToolResult(...)` — ручной tool-loop на on-device движке без пересоздания беседы.
### Добавить ещё один транспорт
Каждый транспорт — отдельный модуль, который получает `Agent` и сериализует его под свой протокол:
* `:server` (HTTP+SSE) — готов, `Route.agentikAgent(agent, path = "/agentik")`
* `:client` (HTTP-клиент) — готов, `AgentikAgent(id, baseUrl, httpClient)`
* `:irc-server` — IRC-фасад, в планах
Транспорт **не имеет доступа к внутренностям `ChatAgent`** — он видит только интерфейс `Agent`. Это и есть «транспортно-агностичное ядро».
### Добавить ещё один LLM-бэкенд
1. Описать `Config` data class с нужными полями.
2. Реализовать `LiteLlm`/`LiteConversation` поверх движка (см. litert-kmp — там уже есть `litert-google`, `litert-openai`, `litert-koog`).
3. Расширить `LlmConfig.createLlm()` веткой `when`.
### Заменить SQLite на Postgres / MongoDB / etc
1. Реализовать три store-интерфейса поверх нового движка.
2. Передать их в `ChatAgent` вместо `SqliteStores`.
3. Удалить (или оставить за `:standalone`-флагом) `:persistence/sqlite/`.
---
## 8. Что НЕ делает `standalone` сегодня
* **Нет инструментов (тулов).** Поддержка ждёт `LiteContentPart.ToolResult` в litert-api (v2). До тех пор агент — text-only чат.
* **Нет суммаризации.** `WorkingMemoryStore.compact` уже есть, но без LLM-вызова для генерации текста суммаризации.
* **Нет MCP-клиента.**
* **Нет авторизации.** Все эндпоинты открыты.
* **Нет инкрементальной догрузки старых сообщений.** `getMessages(after)` работает с offset/limit, но без «схлопывания» (compaction в визуальной истории — задача клиента).
Каждый пункт закрывается отдельным коммитом; код логически разделён по слоям так, чтобы точечные изменения не требовали переделки соседей.
+16
View File
@@ -6,11 +6,14 @@ kotlinx-datetime = "0.8.0"
ktor = "3.1.3"
agui = "0.1.0"
a2a = "1.0.0-SNAPSHOT"
litert = "6"
sqldelight = "2.3.2"
[plugins]
kotlin-multiplatform = { id = "org.jetbrains.kotlin.multiplatform", version.ref = "kotlin" }
kotlin-jvm = { id = "org.jetbrains.kotlin.jvm", version.ref = "kotlin" }
kotlin-serialization = { id = "org.jetbrains.kotlin.plugin.serialization", version.ref = "kotlin" }
sqldelight = { id = "app.cash.sqldelight", version.ref = "sqldelight" }
[libraries]
# --- AG-UI (pw.binom.agui) — KMP: jvm + native ---
@@ -23,6 +26,16 @@ a2a-shared = { module = "pw.binom.a2a:shared", version.ref = "a2a" }
a2a-client = { module = "pw.binom.a2a:client", version.ref = "a2a" }
a2a-server = { module = "pw.binom.a2a:server", version.ref = "a2a" }
# --- litert-kmp (pw.binom.litert) — universal LLM wrapper ---
litert-api = { module = "pw.binom.litert:litert-api", version.ref = "litert" }
litert-openai = { module = "pw.binom.litert:litert-openai-jvm", version.ref = "litert" }
litert-google = { module = "pw.binom.litert:litert-google", version.ref = "litert" }
# --- SQLDelight (app.cash.sqldelight) — KMP SQLite, JDBC driver ---
sqldelight-runtime = { module = "app.cash.sqldelight:runtime", version.ref = "sqldelight" }
sqldelight-sqlite-driver = { module = "app.cash.sqldelight:sqlite-driver", version.ref = "sqldelight" }
sqldelight-coroutines = { module = "app.cash.sqldelight:coroutines-extensions", version.ref = "sqldelight" }
# --- Ktor (сервер, JVM) ---
ktor-server-core = { module = "io.ktor:ktor-server-core", version.ref = "ktor" }
ktor-server-sse = { module = "io.ktor:ktor-server-sse", version.ref = "ktor" }
@@ -32,6 +45,9 @@ ktor-serialization-kotlinx-json = { module = "io.ktor:ktor-serialization-kotlinx
ktor-client-core = { module = "io.ktor:ktor-client-core", version.ref = "ktor" }
ktor-client-cio = { module = "io.ktor:ktor-client-cio", version.ref = "ktor" }
ktor-client-content-negotiation = { module = "io.ktor:ktor-client-content-negotiation", version.ref = "ktor" }
ktor-server-test-host = { module = "io.ktor:ktor-server-test-host", version.ref = "ktor" }
kotlin-test = { module = "org.jetbrains.kotlin:kotlin-test", version.ref = "kotlin" }
# --- commons ---
kotlinx-coroutines-core = { module = "org.jetbrains.kotlinx:kotlinx-coroutines-core", version.ref = "kotlinx-coroutines" }
+54 -17
View File
@@ -1,14 +1,18 @@
@file:OptIn(org.jetbrains.kotlin.gradle.ExperimentalKotlinGradlePluginApi::class)
import org.jetbrains.kotlin.gradle.ExperimentalKotlinGradlePluginApi
import org.jetbrains.kotlin.gradle.targets.jvm.KotlinJvmTarget
plugins {
alias(libs.plugins.kotlin.multiplatform)
alias(libs.plugins.kotlin.serialization)
alias(libs.plugins.sqldelight)
}
kotlin {
jvmToolchain(21)
jvm {
@OptIn(ExperimentalKotlinGradlePluginApi::class)
binaries {
executable {
mainClass.set("pw.binom.agentik.standalone.MainKt")
@@ -18,30 +22,63 @@ kotlin {
sourceSets {
commonMain.dependencies {
// AG-UI: протокол (события, RunAgentInput, Agent) — KMP, jvm + native
implementation(libs.agui.api)
// proto: наш in-house протокол (KMP)
implementation(project(":proto"))
implementation(project(":server"))
// Commons
implementation(libs.kotlinx.coroutines.core)
implementation(libs.kotlinx.serialization.json)
implementation(libs.kotlinx.datetime)
// litert-kmp: контракт (commonMain)
api(libs.litert.api)
// SQLDelight runtime (commonMain)
api(libs.sqldelight.runtime)
api(libs.sqldelight.coroutines)
}
jvmMain.dependencies {
// AG-UI: Ktor-хелперы маршрута + движок Netty (JVM)
implementation(libs.agui.server)
// litert-openai: JVM-реализация
implementation(libs.litert.openai)
// litert-google: встроенный LiteRT-LM движок, нужен только на runtime
runtimeOnly(libs.litert.google)
// SQLDelight JDBC driver (JVM)
implementation(libs.sqldelight.sqlite.driver)
// Ktor server (для :server facade + a2aServer)
implementation(libs.ktor.server.core)
implementation(libs.ktor.server.sse)
implementation(libs.ktor.server.netty)
implementation(libs.ktor.server.content.negotiation)
implementation(libs.ktor.serialization.kotlinx.json)
// Транспортные фасады
implementation(libs.agui.server)
implementation(libs.a2a.server)
}
commonTest.dependencies {
implementation(libs.kotlinx.coroutines.core)
implementation(libs.kotlin.test)
}
jvmTest.dependencies {
implementation(libs.kotlinx.coroutines.core)
implementation(libs.kotlinx.serialization.json)
// A2A (pw.binom.a2a): shared — протокол, server + client — JVM
implementation(libs.a2a.shared)
implementation(libs.a2a.server)
implementation(libs.a2a.client)
implementation(libs.ktor.client.core)
implementation(libs.ktor.client.cio)
// server: Ktor-фасад нашего :proto (JVM)
implementation(project(":server"))
}
commonTest.dependencies {
implementation(kotlin("test"))
implementation(libs.kotlin.test)
// Ktor test engine для smoke-тестов HTTP
implementation(libs.ktor.server.test.host)
}
}
}
sqldelight {
databases {
create("AgentikDatabase") {
packageName.set("pw.binom.agentik.standalone.persistence.sqlite")
srcDirs.setFrom("src/jvmMain/sqldelight")
}
}
}
@@ -0,0 +1,29 @@
package pw.binom.agentik.standalone.persistence
import kotlinx.serialization.SerialName
import kotlinx.serialization.Serializable
/**
* Часть контента сообщения на уровне хранилища.
*
* Намеренно НЕ зависит от [pw.binom.agentik.proto.Content] — маппинг
* `:proto.Content ↔ Content` живёт в `Mapping.kt`. Структурно типы
* идентичны, но даёт возможность заменить transport-протокол без миграции
* таблиц.
*
* Image сериализуется в JSON через base64 (стандарт для kotlinx-serialization).
*/
@Serializable
sealed interface Content {
@Serializable
@SerialName("text")
data class Text(val body: String) : Content
@Serializable
@SerialName("image")
data class Image(val data: ByteArray, val mime: String) : Content {
override fun equals(other: Any?): Boolean =
this === other || (other is Image && mime == other.mime && data.contentEquals(other.data))
override fun hashCode(): Int = 31 * mime.hashCode() + data.contentHashCode()
}
}
@@ -0,0 +1,14 @@
package pw.binom.agentik.standalone.persistence
import kotlin.time.Instant
/**
* Snapshot диалога. В таблице `conversation` хранится как есть.
*/
data class ConversationRecord(
val id: String,
val title: String?,
val isTemporal: Boolean,
val createdAt: Instant,
val updatedAt: Instant,
)
@@ -0,0 +1,27 @@
package pw.binom.agentik.standalone.persistence
import kotlin.time.Instant
/**
* CRUD по таблице `conversation`.
*/
interface ConversationStore : AutoCloseable {
/** Создать или обновить snapshot диалога. */
suspend fun upsert(record: ConversationRecord)
/** Диалог по id, или `null`. */
suspend fun get(id: String): ConversationRecord?
/** Удалить диалог (вместе с его сообщениями и working memory). */
suspend fun delete(id: String): Boolean
/** Список диалогов, отсортированный по `updatedAt` DESC. */
suspend fun list(offset: Int, limit: Int): List<ConversationRecord>
/** Переименовать диалог; `null` для сброса заголовка. Возвращает новый `updatedAt` или `null`, если не найден. */
suspend fun rename(id: String, title: String?): Instant?
/** Обновить `updatedAt` диалога (например, после отправки сообщения). */
suspend fun touch(id: String, now: Instant)
}
@@ -0,0 +1,91 @@
package pw.binom.agentik.standalone.persistence
import kotlinx.serialization.SerialName
import kotlinx.serialization.Serializable
import kotlin.time.Instant
/**
* Запись в таблице `message` (append-only audit) и `working_memory` (mutable view).
*
* Используется sealed-иерархия: подтипы `User`/`Assistant`/`ToolCall`/`ToolResult`
* живут и там, и там. `Summary`/`System` — только в `working_memory`
* (синтетические строки, созданные при суммаризации или как system-prompt).
*
* Все подтипы несут [id] (UUID, стабильный между лайв-стримом Event и историей),
* [conversationId] и [createdAt].
*
* Поля, специфичные для подтипа, сериализуются в JSON в `payload_json`
* колонке SQLite — это даёт гибкость без миграций при добавлении полей.
*/
@Serializable
sealed interface MessageRecord {
val id: String
val conversationId: String
val createdAt: Instant
/** Подтип сообщения с телом из [Content]. */
@Serializable
sealed interface Body : MessageRecord {
val content: List<Content>
}
@Serializable
@SerialName("user")
data class UserMessage(
override val id: String,
override val conversationId: String,
override val content: List<Content>,
override val createdAt: Instant,
) : Body
@Serializable
@SerialName("assistant")
data class AssistantMessage(
override val id: String,
override val conversationId: String,
override val content: List<Content>,
override val createdAt: Instant,
) : Body
@Serializable
@SerialName("tool_call")
data class ToolCall(
override val id: String,
override val conversationId: String,
val toolName: String,
val toolTitle: String?,
val toolArgsJson: String,
override val createdAt: Instant,
) : MessageRecord
@Serializable
@SerialName("tool_result")
data class ToolResult(
override val id: String,
override val conversationId: String,
val toolCallId: String,
val result: String?,
override val createdAt: Instant,
) : MessageRecord
/** Синтетическое: суммаризация старого контекста. Только в working_memory. */
@Serializable
@SerialName("summary")
data class Summary(
override val id: String,
override val conversationId: String,
val text: String,
override val createdAt: Instant,
) : MessageRecord
/** Синтетическое: system-prompt, введённый при создании диалога. Только в working_memory. */
@Serializable
@SerialName("system")
data class System(
override val id: String,
override val conversationId: String,
val text: String,
override val createdAt: Instant,
) : MessageRecord
}
@@ -0,0 +1,24 @@
package pw.binom.agentik.standalone.persistence
import kotlin.time.Instant
/**
* Append-only audit log сообщений (`message` table).
*
* Только `insert` и чтение. Никаких обновлений, никакого удаления (кроме
* каскадного удаления вместе с [ConversationStore.delete]).
*/
interface MessageStore : AutoCloseable {
/** Добавить запись в audit log. `conversationId` берётся из [MessageRecord.conversationId]. */
suspend fun append(record: MessageRecord)
/**
* Страница audit-сообщений диалога после [after] (UTC), отсортированная
* по `createdAt ASC`. Для первоначальной загрузки передай `Instant.DISTANT_PAST`.
*/
suspend fun list(conversationId: String, after: Instant, offset: Int, limit: Int): List<MessageRecord>
/** Все сообщения диалога, отсортированные по `createdAt ASC` (для rebuild working memory). */
suspend fun listAll(conversationId: String): List<MessageRecord>
}
@@ -0,0 +1,26 @@
package pw.binom.agentik.standalone.persistence
import kotlinx.serialization.builtins.ListSerializer
import kotlinx.serialization.json.Json
/**
* JSON-формат для тел user/assistant сообщений: список [Content],
* сериализованный в строку (через kotlinx-serialization).
*/
private val bodyJson = Json {
ignoreUnknownKeys = true
encodeDefaults = true
explicitNulls = false
}
/**
* Сериализует список [Content] в JSON-строку для хранения в `payload_json`.
*/
fun encodeBodyPayload(content: List<Content>): String =
bodyJson.encodeToString(ListSerializer(Content.serializer()), content)
/**
* Десериализует список [Content] из JSON-строки `payload_json`.
*/
fun decodeBodyPayload(json: String): List<Content> =
bodyJson.decodeFromString(ListSerializer(Content.serializer()), json)
@@ -0,0 +1,42 @@
package pw.binom.agentik.standalone.persistence
import kotlinx.serialization.SerialName
import kotlinx.serialization.Serializable
/**
* Запись в working memory диалога: ровно то, что агент сейчас видит в
* LLM-контексте. Упорядочено по `order_idx` (заполняется в store при append).
*
* Sealed-иерархия: для v1 — `System` (синтетический system-prompt),
* `User`/`Assistant` (реплики с ссылкой на audit log через [sourceMessageId]).
* Суммаризация (для v2) добавит подтип `Summary`.
*/
@Serializable
sealed interface WorkingMemoryEntry {
/** Ссылка на исходное сообщение в audit log (`message.id`). `null` для синтетических строк. */
val sourceMessageId: String?
/** Синтетический system-prompt, добавляется при создании диалога. */
@Serializable
@SerialName("system")
data class System(val text: String) : WorkingMemoryEntry {
override val sourceMessageId: String? = null
}
/** Реплика пользователя. */
@Serializable
@SerialName("user")
data class User(
override val sourceMessageId: String,
val content: List<Content>,
) : WorkingMemoryEntry
/** Реплика ассистента. */
@Serializable
@SerialName("assistant")
data class Assistant(
override val sourceMessageId: String,
val content: List<Content>,
) : WorkingMemoryEntry
}
@@ -0,0 +1,50 @@
package pw.binom.agentik.standalone.persistence
import kotlin.time.Instant
/**
* Одна строка `working_memory` таблицы (внутреннее представление store).
*
* Используется для тестов и для перестроения [WorkingMemoryEntry] из row.
* Агент не должен с этим типом работать напрямую — он работает с
* [WorkingMemoryEntry] через [WorkingMemoryStore].
*/
data class WorkingMemoryRow(
val id: String,
val conversationId: String,
val orderIdx: Long,
val sourceMessageId: String?,
val entry: WorkingMemoryEntry,
val createdAt: Instant,
)
/**
* Мутируемое представление LLM-контекста диалога (`working_memory` table).
*
* Аудит-лог — [MessageStore], неизменный; здесь — ровно то, что агент сейчас
* «видит»: системный промпт + реплики + (опционально) суммаризации.
* Строки упорядочены по `order_idx` ASC; `source_message_id` NULL указывает
* на синтетические строки (System, а в v2 — Summary).
*
* Суммаризация / чистка — один атомарный вызов [compact].
*/
interface WorkingMemoryStore : AutoCloseable {
/** Добавить запись в конец working memory (новый максимальный `order_idx`). */
suspend fun append(conversationId: String, entry: WorkingMemoryEntry, now: Instant)
/** Все строки working memory диалога в порядке отправки. */
suspend fun list(conversationId: String): List<WorkingMemoryRow>
/** Очистить working memory диалога (используется при reset/rebuild). */
suspend fun clear(conversationId: String)
/**
* Атомарная суммаризация (v2): удаляет все строки с `order_idx` в диапазоне
* `[dropFromOrderIdx, +∞)` и вставляет вместо них новый `summary` с
* указанным текстом. Возвращает новый максимальный `order_idx`.
*
* Для v1 просто удаляет — суммаризация появится в v2.
*/
suspend fun compact(dropFromOrderIdx: Long, conversationId: String): Long
}
@@ -1,11 +0,0 @@
package pw.binom.agentik.standalone
import kotlin.test.Test
import kotlin.test.assertTrue
class PlaceholderTest {
@Test
fun placeholder() {
assertTrue(true)
}
}
@@ -1,48 +0,0 @@
package pw.binom.agentik.standalone
import pw.binom.a2a.client.A2AClient
import pw.binom.a2a.model.AgentCard
import pw.binom.a2a.model.Message
import pw.binom.a2a.model.Role
import pw.binom.a2a.model.Task
import pw.binom.a2a.model.TextPart
data class RemoteAgent(val name: String, val baseUrl: String, val token: String? = null)
object A2aOutbound {
private val clients = mutableMapOf<String, A2AClient>()
fun remoteAgents(): List<RemoteAgent> =
System.getenv("AGENTIK_A2A_AGENTS")
?.split(";")
?.mapNotNull { entry ->
if (!entry.contains("=")) return@mapNotNull null
val name = entry.substringBefore("=").trim()
val rest = entry.substringAfter("=").split(",", limit = 2)
val baseUrl = rest.getOrNull(0)?.trim().orEmpty()
if (name.isEmpty() || baseUrl.isEmpty()) return@mapNotNull null
val token = rest.getOrNull(1)?.trim()?.takeIf { it.isNotEmpty() }
RemoteAgent(name = name, baseUrl = baseUrl, token = token)
}
?: emptyList()
private fun clientFor(name: String): A2AClient =
clients.getOrPut(name) {
val remote =
remoteAgents().firstOrNull { it.name == name }
?: error("Remote agent $name is not configured (AGENTIK_A2A_AGENTS)")
A2AClient.create(baseUrl = remote.baseUrl, bearerToken = remote.token)
}
suspend fun send(name: String, text: String, contextId: String? = null): Task {
val message = Message(role = Role.USER, parts = listOf(TextPart(text = text)), contextId = contextId)
return clientFor(name).sendMessage(message, contextId)
}
suspend fun agentCard(name: String): AgentCard = clientFor(name).agentCard()
fun closeAll() {
clients.values.forEach { it.close() }
clients.clear()
}
}
@@ -1,21 +0,0 @@
package pw.binom.agentik.standalone
import pw.binom.a2a.model.Message
import pw.binom.a2a.model.Role
import pw.binom.a2a.model.TextPart
import pw.binom.a2a.server.AgentHandler
/**
* A2A-обработчик-заглушка: эхоит входящее сообщение.
* Здесь позже будет реальный агент (LLM / инструменты).
*/
object EchoA2aHandler : AgentHandler {
override suspend fun handle(request: Message, contextId: String?): Message {
val text = request.parts.filterIsInstance<TextPart>().joinToString("") { it.text }
return Message(
role = Role.AGENT,
parts = listOf(TextPart("echo: $text")),
contextId = contextId,
)
}
}
@@ -1,33 +0,0 @@
package pw.binom.agentik.standalone
import pw.binom.agui.api.agent.Agent
import pw.binom.agui.api.event.BaseEvent
import pw.binom.agui.api.event.RunFinishedEvent
import pw.binom.agui.api.event.RunStartedEvent
import pw.binom.agui.api.event.TextMessageContentEvent
import pw.binom.agui.api.event.TextMessageEndEvent
import pw.binom.agui.api.event.TextMessageStartEvent
import pw.binom.agui.api.message.MessageRole
import pw.binom.agui.api.run.RunAgentInput
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.flow
/**
* Заглушка AG-UI агента: отвечает последовательностью событий
* RUN_STARTED -> TEXT_MESSAGE_* -> RUN_FINISHED. Пока — эхо последнего
* пользовательского сообщения. Дальше сюда зайдёт реальный LLM-агент.
*/
object EchoAgent : Agent {
override fun run(input: RunAgentInput): Flow<BaseEvent> = flow {
emit(RunStartedEvent(threadId = input.threadId, runId = input.runId))
val userText = input.messages.lastOrNull { it.role == MessageRole.USER }?.content.orEmpty()
val messageId = "msg-${input.runId}"
emit(TextMessageStartEvent(messageId = messageId, role = MessageRole.ASSISTANT))
emit(TextMessageContentEvent(messageId = messageId, delta = "echo: $userText"))
emit(TextMessageEndEvent(messageId = messageId))
emit(RunFinishedEvent(threadId = input.threadId, runId = input.runId))
}
}
@@ -1,147 +0,0 @@
package pw.binom.agentik.standalone
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.Job
import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.MutableSharedFlow
import kotlinx.coroutines.flow.asSharedFlow
import kotlinx.coroutines.flow.flow
import kotlinx.coroutines.launch
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
import pw.binom.agentik.proto.Agent
import pw.binom.agentik.proto.AgentEvent
import pw.binom.agentik.proto.Content
import pw.binom.agentik.proto.Conversation
import pw.binom.agentik.proto.Event
import pw.binom.agentik.proto.Message
import java.util.UUID
import java.util.concurrent.ConcurrentHashMap
import kotlin.time.Instant
/**
* Минимальная in-memory реализация [Agent] из нашего :proto: эхо-агент,
* повторяет текст пользователя в [AppendText].
*
* Не persistent, без очередей ходов и без реальной отмены — stub для проверки
* что HTTP-фасад `:server` корректно подключён к агенту и прокидывает все
* события/историю согласно контракту.
*/
class EchoProtoAgent(override val id: String = "echo-proto") : Agent {
private val conversations = ConcurrentHashMap<String, EchoProtoConversation>()
private val agentEvents = MutableSharedFlow<AgentEvent>(replay = 0, extraBufferCapacity = 64)
override fun createConversation(temp: Boolean): Conversation {
val c = EchoProtoConversation(isTemporal = temp)
conversations[c.id] = c
agentEvents.tryEmit(
AgentEvent.Created(
date = c.updatedAt,
conversationId = c.id,
)
)
return c
}
override suspend fun getConversation(id: String): Conversation? = conversations[id]
override suspend fun deleteConversation(id: String): Boolean {
val removed = conversations.remove(id) ?: return false
removed.close()
agentEvents.tryEmit(AgentEvent.Deleted(date = Instant.fromEpochMilliseconds(System.currentTimeMillis()), id = id))
return true
}
override suspend fun getConversations(offset: Int, limit: Int): List<Conversation> =
conversations.values.sortedByDescending { it.updatedAt }.drop(offset).take(limit)
override fun events(after: Instant): Flow<AgentEvent> = flow {
agentEvents.asSharedFlow().collect { event ->
if (event.date > after) emit(event)
}
}
}
internal class EchoProtoConversation(
override val id: String = UUID.randomUUID().toString(),
override val isTemporal: Boolean,
) : Conversation {
override val isSupportImageInput: Boolean = false
override val isSupportImageOutput: Boolean = false
override var title: String? = null
private set
override var updatedAt: Instant = Instant.fromEpochMilliseconds(System.currentTimeMillis())
private set
private val messages = mutableListOf<Message>()
private val events = MutableSharedFlow<Event>(replay = 0, extraBufferCapacity = 64)
private val mutex = Mutex()
private var activeJob: Job? = null
private val scope = CoroutineScope(SupervisorJob() + Dispatchers.Default)
override suspend fun rename(title: String) = mutex.withLock {
this.title = title
updatedAt = Instant.fromEpochMilliseconds(System.currentTimeMillis())
}
override suspend fun send(content: List<Content>) {
val userMessage = Message.UserMessage(
id = "um-${UUID.randomUUID()}",
content = content,
date = Instant.fromEpochMilliseconds(System.currentTimeMillis()),
)
mutex.withLock {
messages.add(userMessage)
updatedAt = Instant.fromEpochMilliseconds(System.currentTimeMillis())
}
// stub-политика: прерываем предыдущий ход, запускаем новый.
activeJob?.cancel()
activeJob = scope.launch {
delay(20)
val now = Instant.fromEpochMilliseconds(System.currentTimeMillis())
val echo = content.joinToString(separator = " ") { c ->
when (c) {
is Content.Text -> c.body
is Content.Image -> "[image ${c.mime}, ${c.data.size}B]"
}
}
events.emit(Event.StartResponse(date = now, responseType = Event.ResponseType.TEXT))
val body = "echo: $echo"
events.emit(Event.AppendText(date = now, body = body))
val assistant = Message.AssistantMessage(
id = "am-${UUID.randomUUID()}",
content = listOf(Content.Text(body)),
date = Instant.fromEpochMilliseconds(System.currentTimeMillis()),
)
mutex.withLock { messages.add(assistant) }
events.emit(Event.End(date = Instant.fromEpochMilliseconds(System.currentTimeMillis())))
}
}
override suspend fun interrupt() {
activeJob?.cancel()
activeJob = null
events.emit(Event.Interrupted(date = Instant.fromEpochMilliseconds(System.currentTimeMillis())))
}
override fun events(after: Instant): Flow<Event> = flow {
events.asSharedFlow().collect { event ->
if (event.date > after) emit(event)
}
}
override suspend fun getMessages(after: Instant, offset: Int, limit: Int): List<Message> =
mutex.withLock { messages.asSequence().filter { it.date > after }.drop(offset).take(limit).toList() }
override fun close() {
activeJob?.cancel()
scope.coroutineContext[Job]?.cancel()
}
}
@@ -5,22 +5,14 @@ import io.ktor.server.netty.Netty
import io.ktor.server.response.respondText
import io.ktor.server.routing.get
import io.ktor.server.routing.routing
import kotlinx.coroutines.runBlocking
import pw.binom.a2a.server.A2AServer
import pw.binom.agentik.proto.Conversation
import pw.binom.agentik.proto.Message
import pw.binom.agentik.server.agentikAgent
import pw.binom.agui.server.aguiAgent
import pw.binom.agentik.standalone.agent.ChatAgent
import pw.binom.agentik.standalone.llm.LlmConfig
import pw.binom.agentik.standalone.persistence.sqlite.SqliteStores
/**
* standalone-контейнер agentik:
* - AG-UI: встраиваемый Ktor (Netty), порт AGENTIK_PORT (default 8080)
* POST /agui -> text/event-stream (протокол AG-UI, агент [EchoAgent])
* GET /health -> "ok"
* - A2A: встраиваемый Ktor (CIO), порт AGENTIK_A2A_PORT (default 8081)
* POST / -> JSON-RPC (message/send, tasks/get, tasks/cancel)
* GET /.well-known/agent-card.json
* - :server (proto): встраиваемый Ktor (Netty), порт AGENTIK_PORT (default 8080) — общий с AG-UI
* - :server (proto): встраиваемый Ktor (Netty), порт AGENTIK_PORT (default 8080)
* POST /agentik/conversations -> ConversationSnapshot (201)
* GET /agentik/conversations -> [ConversationSnapshot]
* GET /agentik/conversations/{id} -> ConversationSnapshot
@@ -31,36 +23,43 @@ import pw.binom.agui.server.aguiAgent
* GET /agentik/conversations/{id}/messages -> [Message]
* GET /agentik/conversations/{id}/events -> text/event-stream (SSE)
* GET /agentik/events -> text/event-stream (SSE, Agent-level)
* GET /health -> "ok"
*
* Для обращения к другим агентам: pw.binom.a2a.client.A2AClient.create(baseUrl, token).
* Для in-process вызова :server: pw.binom.agentik.client.AgentikAgent(id, baseUrl, httpClient).
* Хранилище — SQLite (env: AGENTIK_DB_PATH, default `./agentik.db`, `:memory:` для тестов).
* LLM — litert-openai (env: OPENAI_BASE_URL, OPENAI_API_KEY, OPENAI_MODEL).
*
* System prompt — AGENTIK_SYSTEM_PROMPT (default: встроенный `Ты полезный ассистент...`).
*/
fun main() {
val aguiPort = System.getenv("AGENTIK_PORT")?.toIntOrNull() ?: 8080
val a2aPort = System.getenv("AGENTIK_A2A_PORT")?.toIntOrNull() ?: 8081
val port = System.getenv("AGENTIK_PORT")?.toIntOrNull() ?: 8080
val dbPath = System.getenv("AGENTIK_DB_PATH")?.takeIf { it.isNotBlank() } ?: "./agentik.db"
// A2A: agent <-> agent (CIO). Нестреляющий старт — свой event-loop.
val a2a = A2AServer.create(port = a2aPort, agentName = "agentik", handler = EchoA2aHandler)
runBlocking { a2a.start() }
println("A2A server -> http://localhost:$a2aPort/ (JSON-RPC: message/send, tasks/get, tasks/cancel)")
val llmConfig = LlmConfig.fromEnv()
val llm = llmConfig.createLlm()
val stores = SqliteStores.open(dbPath = dbPath)
val agent = ChatAgent(
id = "agentik",
stores = stores,
llm = llm,
llmConfig = llmConfig,
)
// shared port: AG-UI (Netty) + :server (Netty) живут на 8080, A2A — на 8081.
val agui =
embeddedServer(Netty, port = aguiPort) {
val server = embeddedServer(Netty, port = port) {
routing {
get("/health") { call.respondText("ok") }
aguiAgent(EchoAgent, path = "/agui")
agentikAgent(EchoProtoAgent(), path = "/agentik")
agentikAgent(agent, path = "/agentik")
}
}
println("AG-UI -> http://localhost:$aguiPort/agui (SSE), /health")
println(":server proto -> http://localhost:$aguiPort/agentik/... (REST+SSE, агент [EchoProtoAgent])")
agui.start(wait = true)
println("agentik standalone listening on http://localhost:$port")
println(" GET /health")
println(" POST /agentik/conversations -> 201")
println(" GET /agentik/conversations/{id}/events -> SSE")
println(" storage: $dbPath")
println(" llm: ${llmConfig.backend} ${llmConfig.modelInfo()}")
Runtime.getRuntime().addShutdownHook(Thread {
agent.close()
stores.close()
llm.close()
})
server.start(wait = true)
}
// References for IDE noise suppression when source not auto-imported.
@Suppress("unused")
private val keepReferences: Array<Class<*>> = arrayOf(
Conversation::class.java,
Message::class.java,
)
@@ -0,0 +1,121 @@
package pw.binom.agentik.standalone.agent
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.MutableSharedFlow
import kotlinx.coroutines.flow.asSharedFlow
import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
import pw.binom.agentik.proto.Agent as ProtoAgent
import pw.binom.agentik.proto.AgentEvent
import pw.binom.agentik.proto.Conversation as ProtoConversation
import pw.binom.agentik.standalone.llm.LlmConfig
import pw.binom.agentik.standalone.persistence.ConversationRecord
import pw.binom.agentik.standalone.persistence.WorkingMemoryEntry
import pw.binom.agentik.standalone.persistence.sqlite.SqliteStores
import pw.binom.litert.LiteLlm
import kotlin.time.Instant
/**
* Stateful [ProtoAgent] на базе SQLite (история + working memory) и
* [LiteLlm] (универсальный LLM-контракт).
*
* Создаёт [ChatConversation] — те самые stateful диалоги, которые
* хранят свой собственный LiteConversation в ОЗУ (для скорости) и
* логируют в SQLite (для долговечности).
*
* Один [LiteLlm] шарится между всеми беседами агента.
*/
class ChatAgent(
override val id: String,
private val stores: SqliteStores,
private val llm: LiteLlm,
private val llmConfig: LlmConfig,
) : ProtoAgent, AutoCloseable {
private val agentEvents = MutableSharedFlow<AgentEvent>(
extraBufferCapacity = 64,
)
/** Защищает карту живых диалогов. */
private val liveLock = Mutex()
private val live: MutableMap<String, ChatConversation> = HashMap()
/** Live-подписка на события уровня агента (создание/удаление/переименование). */
override fun events(after: Instant): Flow<AgentEvent> {
// Реализация событийной шины упрощённая: возвращаем общий поток.
// Фильтр по `after` не делаем — для v1 после-семантика не нужна
// (см. Memory #3704: replay-free, бэкфилл через getConversations/getConversation).
return agentEvents.asSharedFlow()
}
override fun createConversation(temp: Boolean): ProtoConversation {
val now = now()
val id = "conv-${java.util.UUID.randomUUID()}"
val rec = ConversationRecord(
id = id,
title = null,
isTemporal = temp,
createdAt = now,
updatedAt = now,
)
// Temp-беседы не пишем в SQLite — они живут только в RAM-карте `live`
// и не переживают рестарт агента (см. Memory #3709).
if (!temp) {
runBlocking {
stores.conversations.upsert(rec)
stores.workingMemory.append(
conversationId = id,
entry = WorkingMemoryEntry.System(text = llmConfig.systemPrompt),
now = now,
)
}
}
val conv = ChatConversation(record = rec, stores = stores, llm = llm, systemPrompt = llmConfig.systemPrompt, foldSystemIntoFirstUser = llmConfig.foldSystemIntoFirstUser)
runBlocking {
liveLock.withLock { live[conv.id] = conv }
}
agentEvents.tryEmit(AgentEvent.Created(date = now(), conversationId = conv.id))
return conv
}
override suspend fun getConversation(id: String): ProtoConversation? {
liveLock.withLock { live[id] }?.let { if (!it.isClosed) return it }
val rec = stores.conversations.get(id) ?: return null
return ChatConversation(record = rec, stores = stores, llm = llm, systemPrompt = llmConfig.systemPrompt, foldSystemIntoFirstUser = llmConfig.foldSystemIntoFirstUser).also {
liveLock.withLock { live[id] = it }
}
}
override suspend fun deleteConversation(id: String): Boolean {
val conv = liveLock.withLock { live.remove(id) }
conv?.close()
val ok = stores.conversations.delete(id)
if (ok) agentEvents.tryEmit(AgentEvent.Deleted(date = now(), id = id))
return ok
}
override suspend fun getConversations(offset: Int, limit: Int): List<ProtoConversation> =
stores.conversations.list(offset = offset, limit = limit).map { rec ->
liveLock.withLock { live[rec.id] }
?: ChatConversation(record = rec, stores = stores, llm = llm, systemPrompt = llmConfig.systemPrompt).also {
liveLock.withLock { live[rec.id] = it }
}
}
override fun close() {
runBlocking {
liveLock.withLock {
live.values.forEach { it.close() }
live.clear()
}
}
runCatching { llm.close() }
}
internal fun unregister(id: String) {
runBlocking { liveLock.withLock { live.remove(id) } }
}
private fun now(): Instant = Instant.fromEpochMilliseconds(System.currentTimeMillis())
}
@@ -0,0 +1,353 @@
package pw.binom.agentik.standalone.agent
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.Job
import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.cancel
import kotlinx.coroutines.channels.BufferOverflow
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.MutableSharedFlow
import kotlinx.coroutines.flow.asSharedFlow
import kotlinx.coroutines.flow.transformWhile
import kotlinx.coroutines.launch
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
import pw.binom.agentik.proto.Content as ProtoContent
import pw.binom.agentik.proto.Conversation as ProtoConversation
import pw.binom.agentik.proto.Event as ProtoEvent
import pw.binom.agentik.proto.Message as ProtoMessage
import pw.binom.agentik.standalone.persistence.Content
import pw.binom.agentik.standalone.persistence.ConversationRecord
import pw.binom.agentik.standalone.persistence.ConversationStore
import pw.binom.agentik.standalone.persistence.MessageRecord
import pw.binom.agentik.standalone.persistence.MessageStore
import pw.binom.agentik.standalone.persistence.WorkingMemoryEntry
import pw.binom.agentik.standalone.persistence.WorkingMemoryStore
import pw.binom.agentik.standalone.persistence.sqlite.SqliteStores
import pw.binom.litert.LiteContentPart
import pw.binom.litert.LiteConversation
import pw.binom.litert.LiteConversationConfig
import pw.binom.litert.LiteLlm
import pw.binom.litert.LiteMessage
import pw.binom.litert.LiteRole
import kotlin.time.Instant
/**
* Stateful [pw.binom.agentik.proto.Conversation] поверх [LiteLlm].
*
* Контракт: один [LiteConversation] живёт столько же, сколько [ChatConversation].
* Это требование абстракции LiteLlm — у реализаций внутри `LiteConversation` хранится
* KV-cache движка (LiteRT-LM) либо инкрементальная история (litert-openai). Пересоздание
* на каждый send сломало бы и то, и другое.
*
* На каждый [send]:
* 1. Записывает user-сообщение в audit log (`message`) + working memory.
* 2. Создаёт [LiteConversation] **один раз** при первом send (initial messages =
* системный промпт + пары User/Assistant из working memory, **без** только что
* добавленного user-сообщения — мы передадим его через sendStreamContents).
* 3. Стримит ответ [LiteDelta] в [sendStreamContents] — эмитит [ProtoEvent.AppendText].
* Движок сам добавляет user/assistant к своей внутренней истории.
* 4. По завершении записывает AssistantMessage в audit + working memory.
*
* Если LiteConversation упал при инициализации или во время send — закрываем его,
* следующий send попробует создать заново. После [interrupt] LiteConversation жив;
* отменяется только текущий send.
*/
class ChatConversation(
record: ConversationRecord,
private val stores: SqliteStores,
private val llm: LiteLlm,
private val systemPrompt: String,
private val foldSystemIntoFirstUser: Boolean = false,
) : ProtoConversation, AutoCloseable {
private var record: ConversationRecord = record
override val id: String get() = record.id
override val isSupportImageInput: Boolean get() = false
override val isSupportImageOutput: Boolean get() = false
override val isTemporal: Boolean get() = record.isTemporal
override val title: String? get() = record.title
override val updatedAt: Instant get() = record.updatedAt
private val conversationStore: ConversationStore get() = stores.conversations
private val messageStore: MessageStore get() = stores.messages
private val workingMemory: WorkingMemoryStore get() = stores.workingMemory
private val events = MutableSharedFlow<ProtoEvent>(
replay = 0,
extraBufferCapacity = 128,
onBufferOverflow = BufferOverflow.DROP_OLDEST,
)
private val turnLock = Mutex()
private val scope = CoroutineScope(SupervisorJob() + Dispatchers.IO)
@Volatile
private var liteConv: LiteConversation? = null
@Volatile
private var activeTurn: Job? = null
@Volatile
private var closed = false
internal val isClosed: Boolean get() = closed
override suspend fun rename(title: String) {
val newRecord = conversationStore.rename(id, title)?.let { ts ->
record.copy(title = title, updatedAt = ts)
} ?: record.copy(title = title)
record = newRecord
}
override suspend fun send(content: List<ProtoContent>) {
check(!closed) { "Conversation closed: $id" }
val turnStarted = now()
val userMessageId = newId("msg")
val userRecord = MessageRecord.UserMessage(
id = userMessageId,
conversationId = id,
content = content.map { it.toStorage() },
createdAt = turnStarted,
)
if (!record.isTemporal) {
messageStore.append(userRecord)
workingMemory.append(
conversationId = id,
entry = WorkingMemoryEntry.User(
sourceMessageId = userMessageId,
content = userRecord.content,
),
now = turnStarted,
)
}
activeTurn = scope.launch {
turnLock.withLock {
runTurn(userRecord, turnStarted)
}
}
activeTurn?.join()
}
override suspend fun interrupt() {
runCatching { liteConv?.cancel() }
activeTurn?.cancel()
emitEvent(ProtoEvent.Interrupted(date = now()))
}
override fun events(after: Instant): Flow<ProtoEvent> =
events.asSharedFlow()
override suspend fun getMessages(after: Instant, offset: Int, limit: Int): List<ProtoMessage> =
messageStore.list(conversationId = id, after = after, offset = offset, limit = limit)
.map { it.toProto() }
override fun close() {
if (closed) return
closed = true
runCatching { liteConv?.close() }
runCatching { activeTurn?.cancel() }
scope.cancel()
}
private suspend fun runTurn(userRecord: MessageRecord.UserMessage, turnStarted: Instant) {
emitEvent(ProtoEvent.StartReasoning(date = turnStarted))
emitEvent(ProtoEvent.StartResponse(date = now(), responseType = ProtoEvent.ResponseType.TEXT))
val parts = userRecord.content.mapNotNull { c ->
when (c) {
is Content.Text -> LiteContentPart.Text(c.body)
is Content.Image -> {
System.err.println("[agentik] dropping image input (v1 text-only): mime=${c.mime}, ${c.data.size} bytes")
null
}
}
}
if (parts.isEmpty()) {
emitEvent(ProtoEvent.Error(date = now(), message = "Empty user input (no text content)"))
emitEvent(ProtoEvent.End(date = now()))
return
}
val liteConv = try {
getOrCreateLiteConversation(excludeUserSourceId = if (record.isTemporal) null else userRecord.id)
} catch (e: Throwable) {
this.liteConv = null
emitEvent(ProtoEvent.Error(date = now(), message = e.message ?: "LiteConversation init failed"))
emitEvent(ProtoEvent.End(date = now()))
return
}
val reply = StringBuilder()
try {
// Некоторые бэкенды (litert-google-jvm 0.16.1) не закрывают стрим после `isDone = true`,
// поэтому заворачиваем в transformWhile: видим isDone → отдаём дельту и завершаем flow.
liteConv.sendStreamContents(parts).transformWhile { delta ->
emit(delta)
!delta.isDone
}.collect { delta ->
if (delta.text.isNotEmpty()) {
reply.append(delta.text)
emitEvent(ProtoEvent.AppendText(date = now(), body = delta.text))
}
}
} catch (e: kotlinx.coroutines.CancellationException) {
throw e
} catch (e: Throwable) {
emitEvent(ProtoEvent.Error(date = now(), message = e.message ?: e.javaClass.simpleName))
emitEvent(ProtoEvent.End(date = now()))
return
}
val assistantId = newId("msg")
val assistantAt = now()
val assistantContent = listOf(Content.Text(reply.toString()))
val assistantRecord = MessageRecord.AssistantMessage(
id = assistantId,
conversationId = id,
content = assistantContent,
createdAt = assistantAt,
)
if (!record.isTemporal) {
messageStore.append(assistantRecord)
workingMemory.append(
conversationId = id,
entry = WorkingMemoryEntry.Assistant(
sourceMessageId = assistantId,
content = assistantContent,
),
now = assistantAt,
)
record = record.copy(updatedAt = assistantAt)
conversationStore.touch(id, assistantAt)
}
emitEvent(ProtoEvent.End(date = assistantAt))
}
/**
* Возвращает существующий [LiteConversation] или создаёт новый, инициализированный
* системным промптом и прошлыми User/Assistant из working memory.
*
* [excludeUserSourceId] — если задан, исключает одну строку (свежее user-сообщение,
* уже записанное в audit + working memory, но ещё не отправленное в LLM — мы отдадим
* его через [LiteConversation.sendStreamContents]). Это предотвращает дублирование
* "user → user" в LiteConversation history.
*
* Если [foldSystemIntoFirstUser] — система не передаётся как `systemInstruction`,
* а фолдится в первое user-сообщение (нужно для Gemma-3 шаблона LiteRT-LM 0.16.1).
*/
private suspend fun getOrCreateLiteConversation(excludeUserSourceId: String? = null): LiteConversation {
liteConv?.let { return it }
val wm = if (record.isTemporal) emptyList() else workingMemory.list(id)
val resolvedSystemPrompt = if (record.isTemporal) systemPrompt else wm
.firstOrNull { it.entry is WorkingMemoryEntry.System }
?.let { (it.entry as WorkingMemoryEntry.System).text }
?: systemPrompt
// Берём пары (User, Assistant) из прошлой истории, исключая свежее user-сообщение,
// которое отправим через sendStreamContents.
val pastTurns = if (record.isTemporal) emptyList() else wm
.filter { row ->
val isUserOrAssistant = row.entry is WorkingMemoryEntry.User || row.entry is WorkingMemoryEntry.Assistant
val isPendingUser = excludeUserSourceId != null && row.sourceMessageId == excludeUserSourceId
isUserOrAssistant && !isPendingUser
}
.map { row ->
when (val e = row.entry) {
is WorkingMemoryEntry.User -> LiteMessage(LiteRole.USER, e.content.toLiteContents())
is WorkingMemoryEntry.Assistant -> LiteMessage(LiteRole.MODEL, e.content.toLiteContents())
else -> error("unreachable")
}
}
val config = LiteConversationConfig(
systemInstruction = if (foldSystemIntoFirstUser) null else resolvedSystemPrompt.takeIf { it.isNotBlank() },
initialMessages = if (foldSystemIntoFirstUser) {
val firstUser = pastTurns.firstOrNull { it.role == LiteRole.USER }
if (firstUser != null && resolvedSystemPrompt.isNotBlank()) {
val folded = LiteMessage(
role = LiteRole.USER,
contents = listOf(LiteContentPart.Text(resolvedSystemPrompt + "\n\n")) + firstUser.contents,
)
listOf(folded) + pastTurns.drop(1)
} else {
pastTurns
}
} else {
pastTurns
},
)
return llm.createConversation(config).also { liteConv = it }
}
private fun emitEvent(event: ProtoEvent) {
events.tryEmit(event)
}
private fun now(): Instant =
Instant.fromEpochMilliseconds(System.currentTimeMillis())
private fun newId(prefix: String): String = "$prefix-${java.util.UUID.randomUUID()}"
}
internal fun List<Content>.toLiteContents(): List<LiteContentPart> = map { it.toLite() }
internal fun Content.toLite(): LiteContentPart = when (this) {
is Content.Text -> LiteContentPart.Text(body)
is Content.Image -> LiteContentPart.Image(data, mime)
}
private fun Content.toProto(): ProtoContent = when (this) {
is Content.Text -> ProtoContent.Text(body = body)
is Content.Image -> ProtoContent.Image(data = data, mime = mime)
}
internal fun ProtoContent.toStorage(): Content = when (this) {
is ProtoContent.Text -> Content.Text(body)
is ProtoContent.Image -> Content.Image(data, mime)
}
internal fun MessageRecord.toProto(): ProtoMessage = when (this) {
is MessageRecord.UserMessage -> ProtoMessage.UserMessage(
id = id,
date = createdAt,
content = content.map { it.toProto() },
)
is MessageRecord.AssistantMessage -> ProtoMessage.AssistantMessage(
id = id,
date = createdAt,
content = content.map { it.toProto() },
)
is MessageRecord.ToolCall -> ProtoMessage.ToolCall(
id = id,
date = createdAt,
title = toolTitle,
toolName = toolName,
toolArgs = toolArgsJson,
)
is MessageRecord.ToolResult -> ProtoMessage.ToolResult(
id = id,
date = createdAt,
result = result,
)
// Синтетические строки рабочей памяти: отдаём клиенту как текст ассистента
// (Summary — это результат суммаризации, по форме — ответ модели) или юзера (System).
is MessageRecord.Summary -> ProtoMessage.AssistantMessage(
id = id,
date = createdAt,
content = listOf(ProtoContent.Text(body = text)),
)
is MessageRecord.System -> ProtoMessage.UserMessage(
id = id,
date = createdAt,
content = listOf(ProtoContent.Text(body = text)),
)
}
@@ -0,0 +1,46 @@
package pw.binom.agentik.standalone.llm
import pw.binom.litert.LiteConfig
import pw.binom.litert.LiteLlm
/**
* Factory для движка Google LiteRT-LM на JVM.
*
* `litert-google` опубликован только как Android AAR (с .so внутри), поэтому на JVM
* приходится вручную подгружать JNI-библиотеки LiteRT-LM, прежде чем инстанциировать
* движок. Эта функция:
*
* 1. Резолвит `pw.binom.litert.google.GoogleLiteLlm` через reflection.
* 2. Если процесс уже загрузил нативные библиотеки (`-Djava.library.path`) — успешно
* создаёт движок.
* 3. Если нет — кидает `IllegalStateException` с инструкцией по настройке.
*/
fun googleLiteLlmJvm(config: LiteConfig): LiteLlm {
val cls = try {
Class.forName("pw.binom.litert.google.GoogleLiteLlm")
} catch (e: ClassNotFoundException) {
error(
"litert-google classes are not on the classpath. " +
"Add 'pw.binom.litert:litert-google-android:6' as a runtime dependency " +
"to use the GOOGLE backend on JVM."
)
}
val ctor = cls.constructors.firstOrNull { it.parameterCount == 1 }
?: error("pw.binom.litert.google.GoogleLiteLlm constructor not found")
val engine = try {
ctor.newInstance(config)
} catch (e: UnsatisfiedLinkError) {
throw IllegalStateException(
"LiteRT-LM native libraries are not loaded. " +
"Extract .so/.dylib/.dll from litertlm-android-0.16.1.aar and pass " +
"-Djava.library.path=<dir>, or build :standalone for the androidJvm target.",
e,
)
} catch (e: java.lang.reflect.InvocationTargetException) {
throw e.targetException ?: e
}
@Suppress("UNCHECKED_CAST")
return engine as LiteLlm
}
private fun error(message: String): Nothing = throw IllegalStateException(message)
@@ -0,0 +1,93 @@
package pw.binom.agentik.standalone.llm
import pw.binom.litert.LiteBackend
import pw.binom.litert.LiteConfig
import pw.binom.litert.LiteExperimental
import pw.binom.litert.LiteLlm
import pw.binom.litert.openai.OpenAiConfig as LitertOpenAiConfig
import pw.binom.litert.openai.openAiLiteLlm
data class LlmConfig(
val backend: LlmBackend,
val systemPrompt: String,
val openai: LitertOpenAiConfig? = null,
val google: GoogleConfig? = null,
val foldSystemIntoFirstUser: Boolean = backend == LlmBackend.GOOGLE,
) {
init {
when (backend) {
LlmBackend.OPENAI -> require(openai != null) { "OPENAI backend requires openai config" }
LlmBackend.GOOGLE -> require(google != null) { "GOOGLE backend requires google config" }
}
}
fun createLlm(): LiteLlm = when (backend) {
LlmBackend.OPENAI -> openAiLiteLlm(checkNotNull(openai))
LlmBackend.GOOGLE -> googleLiteLlmJvm(google!!.toLiteConfig())
}
fun modelInfo(): String = when (backend) {
LlmBackend.OPENAI -> "${checkNotNull(openai).model} @ ${checkNotNull(openai).baseUrl}"
LlmBackend.GOOGLE -> "${checkNotNull(google).modelPath}"
}
companion object {
const val DEFAULT_SYSTEM_PROMPT: String = "Ты полезный ассистент. Отвечай кратко и по делу."
fun fromEnv(env: (String) -> String? = System::getenv): LlmConfig {
val backend = LlmBackend.parse(env("AGENTIK_LLM_BACKEND"))
val systemPromptRaw = env("AGENTIK_SYSTEM_PROMPT")
val systemPrompt = if (systemPromptRaw.isNullOrBlank()) DEFAULT_SYSTEM_PROMPT else systemPromptRaw
return when (backend) {
LlmBackend.OPENAI -> {
val openai = LitertOpenAiConfig(
baseUrl = requireEnv(env, "OPENAI_BASE_URL"),
apiKey = requireEnv(env, "OPENAI_API_KEY"),
model = requireEnv(env, "OPENAI_MODEL"),
)
LlmConfig(backend, systemPrompt, openai = openai)
}
LlmBackend.GOOGLE -> {
val google = GoogleConfig(
modelPath = requireEnv(env, "AGENTIK_GOOGLE_MODEL_PATH"),
cacheDir = env("AGENTIK_GOOGLE_CACHE_DIR"),
threads = env("AGENTIK_GOOGLE_THREADS")?.toInt(),
)
LlmConfig(backend, systemPrompt, google = google)
}
}
}
private fun requireEnv(env: (String) -> String?, name: String): String =
env(name) ?: error("Required env var $name is not set")
}
}
enum class LlmBackend {
OPENAI,
GOOGLE;
companion object {
fun parse(raw: String?): LlmBackend = when (raw?.lowercase()) {
null, "", "openai" -> OPENAI
"google", "litert", "litert-google" -> GOOGLE
else -> error("Unknown LLM backend '$raw', expected 'openai' or 'google'")
}
}
}
data class GoogleConfig(
val modelPath: String,
val cacheDir: String? = null,
val threads: Int? = null,
) {
fun toLiteConfig(): LiteConfig = LiteConfig(
modelPath = modelPath,
cacheDir = cacheDir ?: "",
threads = threads,
backend = LiteBackend.CPU,
visionBackend = null,
audioBackend = null,
experimental = LiteExperimental(false, emptyMap()),
)
}
@@ -0,0 +1,80 @@
package pw.binom.agentik.standalone.persistence.sqlite
import kotlin.time.Instant
import pw.binom.agentik.standalone.persistence.ConversationRecord
import pw.binom.agentik.standalone.persistence.ConversationStore
/**
* SQLite-реализация [ConversationStore].
*/
class SqliteConversationStore(private val db: AgentikDatabase) : ConversationStore {
private val q get() = db.conversationQueries
override suspend fun upsert(record: ConversationRecord) {
// SQLite-конфликт по PRIMARY KEY → сначала пробуем insert, при ошибке → update
val existing = q.getById(record.id).executeAsOneOrNull()
if (existing == null) {
q.insert(
id = record.id,
title = record.title,
is_temporal = if (record.isTemporal) 1L else 0L,
created_at = record.createdAt.toEpochMilliseconds(),
updated_at = record.updatedAt.toEpochMilliseconds(),
)
} else {
q.update(
title = record.title,
is_temporal = if (record.isTemporal) 1L else 0L,
updated_at = record.updatedAt.toEpochMilliseconds(),
id = record.id,
)
}
}
override suspend fun get(id: String): ConversationRecord? {
val row = q.getById(id).executeAsOneOrNull() ?: return null
return row.toRecord()
}
override suspend fun delete(id: String): Boolean {
// Проверяем существование ДО удаления — иначе пустой `deleteById` вернёт
// «успех», и get() == null будет true, хотя диалога и не было.
val existed = q.getById(id).executeAsOneOrNull() != null
if (!existed) return false
db.transaction {
db.messageQueries.deleteByConversation(id)
db.workingMemoryQueries.clearByConversation(id)
q.deleteById(id)
}
return true
}
override suspend fun list(offset: Int, limit: Int): List<ConversationRecord> =
q.list(limit = limit.toLong(), offset = offset.toLong())
.executeAsList()
.map { it.toRecord() }
override suspend fun rename(id: String, title: String?): Instant? {
val now = Instant.fromEpochMilliseconds(System.currentTimeMillis())
q.rename(title = title, updated_at = now.toEpochMilliseconds(), id = id)
val ts = q.getUpdatedAt(id).executeAsOneOrNull() ?: return null
return Instant.fromEpochMilliseconds(ts)
}
override suspend fun touch(id: String, now: Instant) {
q.touch(updated_at = now.toEpochMilliseconds(), id = id)
}
override fun close() {
// driver закрывается во внешнем SqliteStores
}
}
private fun Conversation.toRecord(): ConversationRecord = ConversationRecord(
id = id,
title = title,
isTemporal = is_temporal != 0L,
createdAt = Instant.fromEpochMilliseconds(created_at),
updatedAt = Instant.fromEpochMilliseconds(updated_at),
)
@@ -0,0 +1,115 @@
package pw.binom.agentik.standalone.persistence.sqlite
import kotlinx.serialization.json.Json
import pw.binom.agentik.standalone.persistence.MessageRecord
import pw.binom.agentik.standalone.persistence.MessageStore
import pw.binom.agentik.standalone.persistence.decodeBodyPayload
import pw.binom.agentik.standalone.persistence.encodeBodyPayload
import kotlin.time.Instant
/**
* SQLite-реализация [MessageStore] (append-only audit log).
*
* `payloadJson` хранит JSON-сериализованные kind-specific поля. Для
* `user`/`assistant` это `List<Content>` (см. [encodeBodyPayload]).
* Для `tool_call`/`tool_result` payload хранит JSON-объект
* (см. [CallPayload]/[ResultPayload]).
*
* SQLDelight сохраняет snake_case в сгенерированной data class (`Message`),
* поэтому обращаемся через `conversation_id`, `payload_json`, `created_at`.
*/
class SqliteMessageStore(private val db: AgentikDatabase) : MessageStore {
private val q get() = db.messageQueries
override suspend fun append(record: MessageRecord) {
val (kind, payload) = encodeRecord(record)
q.insert(
conversation_id = record.conversationId,
kind = kind,
payload_json = payload,
created_at = record.createdAt.toEpochMilliseconds(),
id = record.id,
)
}
override suspend fun list(
conversationId: String,
after: Instant,
offset: Int,
limit: Int,
): List<MessageRecord> = q.listAfter(
conversation_id = conversationId,
created_at = after.toEpochMilliseconds(),
limit = limit.toLong(),
offset = offset.toLong(),
).executeAsList().map { it.toRecord() }
override suspend fun listAll(conversationId: String): List<MessageRecord> =
q.listByConversationAll(conversation_id = conversationId).executeAsList().map { it.toRecord() }
override fun close() {}
}
private fun encodeRecord(record: MessageRecord): Pair<String, String> = when (record) {
is MessageRecord.UserMessage -> "user" to encodeBodyPayload(record.content)
is MessageRecord.AssistantMessage -> "assistant" to encodeBodyPayload(record.content)
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.Summary,
is MessageRecord.System -> error("Summary/System — synthetic, cannot append to audit log")
}
@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?)
private fun Message.toRecord(): MessageRecord {
val id = id
val convId = conversation_id
val createdAt = Instant.fromEpochMilliseconds(created_at)
return when (kind) {
"user" -> MessageRecord.UserMessage(
id = id,
conversationId = convId,
content = decodeBodyPayload(payload_json),
createdAt = createdAt,
)
"assistant" -> MessageRecord.AssistantMessage(
id = id,
conversationId = convId,
content = decodeBodyPayload(payload_json),
createdAt = createdAt,
)
"tool_call" -> {
val p = Json.decodeFromString(CallPayload.serializer(), payload_json)
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_json)
MessageRecord.ToolResult(
id = id,
conversationId = convId,
toolCallId = p.toolCallId,
result = p.result,
createdAt = createdAt,
)
}
else -> error("Unknown message kind in audit log: $kind")
}
}
@@ -0,0 +1,73 @@
package pw.binom.agentik.standalone.persistence.sqlite
import app.cash.sqldelight.db.QueryResult
import app.cash.sqldelight.db.SqlDriver
import app.cash.sqldelight.driver.jdbc.sqlite.JdbcSqliteDriver
import pw.binom.agentik.standalone.persistence.ConversationStore
import pw.binom.agentik.standalone.persistence.MessageStore
import pw.binom.agentik.standalone.persistence.WorkingMemoryStore
/**
* Корневой объект SQLite-слоя: держит [SqlDriver] и три [WorkingMemoryStore]/[MessageStore]/[ConversationStore].
* Закрывается вместе с приложением.
*/
class SqliteStores private constructor(
val driver: SqlDriver,
val conversations: ConversationStore,
val messages: MessageStore,
val workingMemory: WorkingMemoryStore,
) : AutoCloseable {
override fun close() {
conversations.close()
messages.close()
workingMemory.close()
driver.close()
}
companion object {
/** Открыть/создать БД по пути `dbPath` (например, `"./agentik.db"` или абсолютный путь). */
fun open(dbPath: String): SqliteStores {
val driver = JdbcSqliteDriver("jdbc:sqlite:$dbPath")
createSchema(driver)
val db = AgentikDatabase(driver)
return SqliteStores(
driver = driver,
conversations = SqliteConversationStore(db),
messages = SqliteMessageStore(db),
workingMemory = SqliteWorkingMemoryStore(db),
)
}
/** Открыть/создать БД в памяти (для тестов). */
fun inMemory(): SqliteStores {
val driver = JdbcSqliteDriver(JdbcSqliteDriver.IN_MEMORY)
createSchema(driver)
val db = AgentikDatabase(driver)
return SqliteStores(
driver = driver,
conversations = SqliteConversationStore(db),
messages = SqliteMessageStore(db),
workingMemory = SqliteWorkingMemoryStore(db),
)
}
private fun createSchema(driver: SqlDriver) {
// Если таблица `conversation` уже есть — БД уже инициализирована,
// просто пропускаем create (иначе CREATE TABLE упадёт на дубликате).
val existing = driver.executeQuery(
identifier = null,
sql = "SELECT name FROM sqlite_master WHERE type='table' AND name='conversation'",
mapper = { cursor ->
QueryResult.Value(
if (cursor.next().value) cursor.getString(0) else null,
)
},
parameters = 0,
).value
if (existing != null) return
AgentikDatabase.Schema.create(driver)
}
}
}
@@ -0,0 +1,70 @@
package pw.binom.agentik.standalone.persistence.sqlite
import kotlinx.serialization.json.Json
import pw.binom.agentik.standalone.persistence.WorkingMemoryEntry
import pw.binom.agentik.standalone.persistence.WorkingMemoryRow
import pw.binom.agentik.standalone.persistence.WorkingMemoryStore
import kotlin.time.Instant
/**
* SQLite-реализация [WorkingMemoryStore] (мутируемый LLM-контекст).
*
* Строки упорядочены по `order_idx ASC`. Compact — атомарная операция:
* DELETE строк `>= dropFromOrderIdx` (без summary в v1).
*/
class SqliteWorkingMemoryStore(private val db: AgentikDatabase) : WorkingMemoryStore {
private val q get() = db.workingMemoryQueries
private val json = Json { ignoreUnknownKeys = true }
override suspend fun append(conversationId: String, entry: WorkingMemoryEntry, now: Instant) {
val newIdx = (q.maxOrderIdx(conversationId).executeAsOne()) + 1
q.insert(
id = newId(),
conversation_id = conversationId,
order_idx = newIdx,
source_message_id = entry.sourceMessageId,
kind = entryKind(entry),
payload_json = json.encodeToString(WorkingMemoryEntry.serializer(), entry),
created_at = now.toEpochMilliseconds(),
)
}
override suspend fun list(conversationId: String): List<WorkingMemoryRow> =
q.listByConversation(conversationId).executeAsList().map { it.toRow() }
override suspend fun clear(conversationId: String) {
q.clearByConversation(conversationId)
}
override suspend fun compact(dropFromOrderIdx: Long, conversationId: String): Long {
var newMax = 0L
db.transaction {
q.compactDelete(conversation_id = conversationId, order_idx = dropFromOrderIdx)
newMax = q.maxOrderIdx(conversationId).executeAsOne()
}
return newMax
}
override fun close() {}
}
private fun entryKind(e: WorkingMemoryEntry): String = when (e) {
is WorkingMemoryEntry.System -> "system"
is WorkingMemoryEntry.User -> "user"
is WorkingMemoryEntry.Assistant -> "assistant"
}
private fun Working_memory.toRow(): WorkingMemoryRow {
val entry = Json.decodeFromString(WorkingMemoryEntry.serializer(), payload_json)
return WorkingMemoryRow(
id = id,
conversationId = conversation_id,
orderIdx = order_idx,
sourceMessageId = source_message_id,
entry = entry,
createdAt = Instant.fromEpochMilliseconds(created_at),
)
}
private fun newId(): String = "wm-${java.util.UUID.randomUUID()}"
@@ -0,0 +1,42 @@
CREATE TABLE conversation (
id TEXT NOT NULL PRIMARY KEY,
title TEXT,
is_temporal INTEGER NOT NULL DEFAULT 0,
created_at INTEGER NOT NULL,
updated_at INTEGER NOT NULL
);
CREATE INDEX idx_conv_updated ON conversation(updated_at DESC);
insert:
INSERT INTO conversation (id, title, is_temporal, created_at, updated_at)
VALUES (?, ?, ?, ?, ?);
update:
UPDATE conversation SET title = ?, is_temporal = ?, updated_at = ?
WHERE id = ?;
getById:
SELECT * FROM conversation WHERE id = ?;
list:
SELECT * FROM conversation
WHERE is_temporal = 0
ORDER BY updated_at DESC
LIMIT :limit OFFSET :offset;
listAll:
SELECT * FROM conversation
ORDER BY updated_at DESC;
deleteById:
DELETE FROM conversation WHERE id = ?;
rename:
UPDATE conversation SET title = ?, updated_at = ? WHERE id = ?;
getUpdatedAt:
SELECT updated_at FROM conversation WHERE id = ?;
touch:
UPDATE conversation SET updated_at = ? WHERE id = ?;
@@ -0,0 +1,33 @@
CREATE TABLE 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 idx_msg_conv ON message(conversation_id, created_at);
insert:
INSERT INTO message (id, conversation_id, kind, payload_json, created_at)
VALUES (?, ?, ?, ?, ?);
listByConversation:
SELECT * FROM message
WHERE conversation_id = ?
ORDER BY created_at ASC, id ASC
LIMIT :limit OFFSET :offset;
listByConversationAll:
SELECT * FROM message
WHERE conversation_id = ?
ORDER BY created_at ASC, id ASC;
listAfter:
SELECT * FROM message
WHERE conversation_id = ? AND created_at > ?
ORDER BY created_at ASC, id ASC
LIMIT :limit OFFSET :offset;
deleteByConversation:
DELETE FROM message WHERE conversation_id = ?;
@@ -0,0 +1,34 @@
CREATE TABLE working_memory (
id TEXT NOT NULL PRIMARY KEY,
conversation_id TEXT NOT NULL,
order_idx INTEGER NOT NULL,
source_message_id TEXT,
kind TEXT NOT NULL,
payload_json TEXT NOT NULL,
created_at INTEGER NOT NULL
);
CREATE UNIQUE INDEX idx_wm_unique ON working_memory(conversation_id, order_idx);
CREATE INDEX idx_wm_conv ON working_memory(conversation_id, order_idx);
insert:
INSERT INTO working_memory (id, conversation_id, order_idx, source_message_id, kind, payload_json, created_at)
VALUES (?, ?, ?, ?, ?, ?, ?);
listByConversation:
SELECT * FROM working_memory
WHERE conversation_id = ?
ORDER BY order_idx ASC;
clearByConversation:
DELETE FROM working_memory WHERE conversation_id = ?;
maxOrderIdx:
SELECT COALESCE(MAX(order_idx), 0) FROM working_memory WHERE conversation_id = ?;
-- Atomic compact: delete rows >= dropFromOrderIdx and insert summary.
-- Caller supplies summaryId, summaryText, now epoch millis, and the new summary
-- gets order_idx = current max (after delete = before max).
compactDelete:
DELETE FROM working_memory
WHERE conversation_id = ? AND order_idx >= ?;
@@ -0,0 +1,433 @@
package pw.binom.agentik.standalone.agent
import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.collect
import kotlinx.coroutines.flow.flowOf
import kotlinx.coroutines.flow.toList
import kotlinx.coroutines.launch
import kotlinx.coroutines.test.runTest
import pw.binom.agentik.proto.AgentEvent
import pw.binom.agentik.proto.Content
import pw.binom.agentik.proto.Event as ProtoEvent
import pw.binom.agentik.standalone.llm.LlmConfig
import pw.binom.agentik.standalone.persistence.sqlite.SqliteStores
import pw.binom.litert.LiteContentPart
import pw.binom.litert.LiteConversation
import pw.binom.litert.LiteConversationConfig
import pw.binom.litert.LiteDelta
import pw.binom.litert.LiteLlm
import pw.binom.litert.LiteMessage
import pw.binom.litert.LiteRole
import pw.binom.litert.openai.OpenAiConfig
import kotlin.test.AfterTest
import kotlin.test.BeforeTest
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertIs
import kotlin.test.assertNotNull
import kotlin.test.assertNull
import kotlin.test.assertTrue
import kotlin.time.Instant
class ChatAgentTest {
private lateinit var stores: SqliteStores
private lateinit var fakeLlm: FakeLiteLlm
@BeforeTest
fun setup() {
stores = SqliteStores.inMemory()
fakeLlm = FakeLiteLlm()
}
@AfterTest
fun tearDown() {
stores.close()
}
private fun newAgent(): ChatAgent {
val cfg = LlmConfig(
backend = pw.binom.agentik.standalone.llm.LlmBackend.OPENAI,
systemPrompt = "be brief",
openai = OpenAiConfig(baseUrl = "http://test", apiKey = "test", model = "test"),
)
return ChatAgent(id = "agentik", stores = stores, llm = fakeLlm, llmConfig = cfg)
}
@Test
fun `createConversation seeds system prompt into working memory`() = runTest {
val agent = newAgent()
val conv = agent.createConversation(temp = false) as ChatConversation
val wm = stores.workingMemory.list(conv.id)
assertEquals(1, wm.size)
val first = wm[0]
assertIs<pw.binom.agentik.standalone.persistence.WorkingMemoryEntry.System>(first.entry)
assertEquals("be brief", first.entry.text)
}
@Test
fun `getConversation returns null for unknown id`() = runTest {
val agent = newAgent()
assertNull(agent.getConversation("nope"))
}
@Test
fun `getConversations returns all stored persistent conversations`() = runTest {
val agent = newAgent()
agent.createConversation(temp = false)
agent.createConversation(temp = true)
val list = agent.getConversations(0, 10)
// temp-беседы не персистятся, в списке только persistent
assertEquals(1, list.size)
}
@Test
fun `deleteConversation removes conversation and data`() = runTest {
val agent = newAgent()
val conv = agent.createConversation(temp = false)
val id = conv.id
// добавим сообщение, чтобы потом убедиться, что каскад сработал
stores.messages.append(
pw.binom.agentik.standalone.persistence.MessageRecord.UserMessage(
id = "m1",
conversationId = id,
content = listOf(pw.binom.agentik.standalone.persistence.Content.Text("hi")),
createdAt = Instant.fromEpochMilliseconds(1_700_000_000_000),
),
)
assertTrue(agent.deleteConversation(id))
assertNull(agent.getConversation(id))
assertNull(stores.conversations.get(id))
assertEquals(emptyList(), stores.messages.listAll(id))
}
@Test
fun `deleteConversation returns false for unknown id`() = runTest {
val agent = newAgent()
assertEquals(false, agent.deleteConversation("nope"))
}
@Test
fun `send emits start_reasoning, start_response, append_text, end`() = runTest {
val agent = newAgent()
fakeLlm.reply = "hello back"
val conv = agent.createConversation(temp = false)
conv.send(listOf(Content.Text("hi")))
// user message записан в audit + working memory
val msgs = stores.messages.listAll(conv.id)
assertEquals(2, msgs.size)
assertEquals("hi", (msgs[0] as pw.binom.agentik.standalone.persistence.MessageRecord.UserMessage).content.let {
(it[0] as pw.binom.agentik.standalone.persistence.Content.Text).body
})
assertEquals("hello back", (msgs[1] as pw.binom.agentik.standalone.persistence.MessageRecord.AssistantMessage).content.let {
(it[0] as pw.binom.agentik.standalone.persistence.Content.Text).body
})
}
@Test
fun `send reconstructs conversation history from working memory`() = runTest {
val agent = newAgent()
fakeLlm.rememberHistory = true
fakeLlm.reply = "first reply"
val conv1 = agent.createConversation(temp = false)
conv1.send(listOf(Content.Text("first user")))
// Новая беседа не должна видеть историю первой
val conv2 = agent.createConversation(temp = false)
fakeLlm.reply = "second reply"
conv2.send(listOf(Content.Text("second user")))
// Первая беседа должна иметь только свою систему + 1 user + 1 assistant
val wm1 = stores.workingMemory.list(conv1.id)
assertEquals(3, wm1.size)
// Вторая беседа — только своё
val wm2 = stores.workingMemory.list(conv2.id)
assertEquals(3, wm2.size)
}
@Test
fun `send passes system prompt and past history to LLM on first send`() = runTest {
val agent = newAgent()
fakeLlm.rememberHistory = true
val conv = agent.createConversation(temp = false)
fakeLlm.reply = "hi"
conv.send(listOf(Content.Text("hello")))
// Длинно-живущий LiteConversation: первый send создаёт его с systemInstruction
// и пустыми initialMessages (свежее user-сообщение пойдёт через sendStreamContents).
assertNotNull(fakeLlm.lastConfig)
assertEquals("be brief", fakeLlm.lastConfig!!.systemInstruction)
assertEquals(0, fakeLlm.lastConfig!!.initialMessages.size)
// Свежее user-сообщение отправлено через sendStreamContents
assertEquals(1, fakeLlm.conversations.size)
val sent = fakeLlm.lastContents
assertNotNull(sent)
assertEquals(1, sent!!.size)
assertEquals("hello", (sent[0] as LiteContentPart.Text).text)
}
@Test
fun `multi-turn conversation accumulates in single LiteConversation`() = runTest {
val agent = newAgent()
fakeLlm.rememberHistory = true
val conv = agent.createConversation(temp = false)
fakeLlm.reply = "first reply"
conv.send(listOf(Content.Text("first user")))
// первый turn: WM = [system, user, assistant]
assertEquals(3, stores.workingMemory.list(conv.id).size)
fakeLlm.reply = "second reply"
conv.send(listOf(Content.Text("second user")))
// второй turn: WM должен вырасти до [system, user, assistant, user, assistant]
val wm = stores.workingMemory.list(conv.id)
assertEquals(5, wm.size)
// Длинно-живущий LiteConversation: один на ChatConversation, история
// накапливается через sendStreamContents, без пересоздания.
assertEquals(1, fakeLlm.conversations.size)
val history = fakeLlm.conversations[0].history
assertEquals(4, history.size)
assertEquals("first user", history[0].text)
assertEquals(LiteRole.USER, history[0].role)
assertEquals("first reply", history[1].text)
assertEquals(LiteRole.MODEL, history[1].role)
assertEquals("second user", history[2].text)
assertEquals(LiteRole.USER, history[2].role)
assertEquals("second reply", history[3].text)
assertEquals(LiteRole.MODEL, history[3].role)
}
@Test
fun `reloaded conversation reconstructs LiteConversation from working memory`() = runTest {
val agent = newAgent()
fakeLlm.rememberHistory = true
val conv = agent.createConversation(temp = false)
fakeLlm.reply = "first reply"
conv.send(listOf(Content.Text("first user")))
val convId = conv.id
conv.close()
// Открываем новое ChatConversation с тем же id — LiteConversation должен
// быть создан заново из working memory (первый user+assistant как initial).
val reopened = agent.getConversation(convId)!!
fakeLlm.reply = "second reply"
reopened.send(listOf(Content.Text("second user")))
val allConvs = fakeLlm.conversations
assertEquals(2, allConvs.size) // original + reopened
val reopenedLite = allConvs.last()
// Initial messages: только прошлые user+assistant (НЕ включая текущий "second user")
assertEquals(2, reopenedLite.initialMessages.size)
assertEquals("first user", reopenedLite.initialMessages[0].text)
assertEquals(LiteRole.USER, reopenedLite.initialMessages[0].role)
assertEquals("first reply", reopenedLite.initialMessages[1].text)
assertEquals(LiteRole.MODEL, reopenedLite.initialMessages[1].role)
}
@Test
fun `interrupt cancels active send`() = runTest {
val agent = newAgent()
fakeLlm.slow = true
val conv = agent.createConversation(temp = false)
val sendJob = launch {
try {
conv.send(listOf(Content.Text("hi")))
} catch (_: kotlinx.coroutines.CancellationException) {
// ok
}
}
// ждём, пока корутина дойдёт до sendStreamContents и повиснет на slow-эмиссии
delay(200)
conv.interrupt()
sendJob.join()
// user сообщение в audit должно быть, assistant — нет (был отменён)
val msgs = stores.messages.listAll(conv.id)
assertEquals(1, msgs.size)
assertIs<pw.binom.agentik.standalone.persistence.MessageRecord.UserMessage>(msgs[0])
}
@Test
fun `temp conversation is not persisted across agent instances`() = runTest {
// Поднимаем file-backed БД, создаём temp-беседу
stores.close()
val dbPath = (System.getProperty("java.io.tmpdir") + "/agentik-test-${System.nanoTime()}.db")
stores = SqliteStores.open(dbPath)
val agent1 = ChatAgent(
id = "agentik",
stores = stores,
llm = FakeLiteLlm().also { fakeLlm = it },
llmConfig = LlmConfig(
backend = pw.binom.agentik.standalone.llm.LlmBackend.OPENAI,
systemPrompt = "be brief",
openai = OpenAiConfig(baseUrl = "http://test", apiKey = "test", model = "test"),
),
)
val tempConv = agent1.createConversation(temp = true)
val tempId = tempConv.id
assertNotNull(agent1.getConversation(tempId))
// Переоткрываем БД — temp-беседа не должна пережить рестарт
stores.close()
stores = SqliteStores.open(dbPath)
val agent2 = ChatAgent(
id = "agentik",
stores = stores,
llm = fakeLlm,
llmConfig = LlmConfig(
backend = pw.binom.agentik.standalone.llm.LlmBackend.OPENAI,
systemPrompt = "be brief",
openai = OpenAiConfig(baseUrl = "http://test", apiKey = "test", model = "test"),
),
)
assertNull(agent2.getConversation(tempId))
java.io.File(dbPath).delete()
}
@Test
fun `non-temp conversation persists across agent instances`() = runTest {
stores.close()
val dbPath = (System.getProperty("java.io.tmpdir") + "/agentik-test-${System.nanoTime()}.db")
stores = SqliteStores.open(dbPath)
val agent1 = ChatAgent(
id = "agentik",
stores = stores,
llm = FakeLiteLlm().also { fakeLlm = it },
llmConfig = LlmConfig(
backend = pw.binom.agentik.standalone.llm.LlmBackend.OPENAI,
systemPrompt = "be brief",
openai = OpenAiConfig(baseUrl = "http://test", apiKey = "test", model = "test"),
),
)
val conv = agent1.createConversation(temp = false)
val id = conv.id
stores.close()
stores = SqliteStores.open(dbPath)
val agent2 = ChatAgent(
id = "agentik",
stores = stores,
llm = fakeLlm,
llmConfig = LlmConfig(
backend = pw.binom.agentik.standalone.llm.LlmBackend.OPENAI,
systemPrompt = "be brief",
openai = OpenAiConfig(baseUrl = "http://test", apiKey = "test", model = "test"),
),
)
assertNotNull(agent2.getConversation(id))
java.io.File(dbPath).delete()
}
private fun fakeLiteLlmForReload(): LiteLlm = object : LiteLlm {
override val backendName: String = "fake"
override fun isInitialized(): Boolean = true
override fun createConversation(config: LiteConversationConfig): LiteConversation =
error("not used in reload test")
override fun infer(request: pw.binom.litert.LiteRequest): String = error("not used")
override fun inferStream(request: pw.binom.litert.LiteRequest): Flow<LiteDelta> = error("not used")
override fun close() {}
}
@Test
fun `agentEvents — Created + Deleted flow`() = runTest {
val agent = newAgent()
val events = mutableListOf<AgentEvent>()
val job = launch(start = kotlinx.coroutines.CoroutineStart.UNDISPATCHED) {
agent.events(Instant.DISTANT_PAST).collect { events.add(it) }
}
val conv = agent.createConversation(temp = false)
agent.deleteConversation(conv.id)
delay(50)
job.cancel()
assertEquals(2, events.size)
val created = events[0] as AgentEvent.Created
val deleted = events[1] as AgentEvent.Deleted
assertEquals(conv.id, created.conversationId)
assertEquals(conv.id, deleted.id)
}
}
/** Поддельный LiteLlm: возвращает fakeLlm.reply в sendStreamContents, опционально запоминает history. */
private class FakeLiteLlm : LiteLlm {
override val backendName: String = "fake"
var reply: String = ""
var rememberHistory: Boolean = false
var slow: Boolean = false
var lastConfig: LiteConversationConfig? = null
var lastContents: List<LiteContentPart>? = null
val conversations = mutableListOf<FakeLiteConversation>()
override fun isInitialized(): Boolean = true
override fun createConversation(config: LiteConversationConfig): LiteConversation {
lastConfig = config
val conv = FakeLiteConversation(this, config)
conversations.add(conv)
return conv
}
override fun infer(request: pw.binom.litert.LiteRequest): String {
throw UnsupportedOperationException("not used in test")
}
override fun inferStream(request: pw.binom.litert.LiteRequest): Flow<LiteDelta> {
throw UnsupportedOperationException("not used in test")
}
override fun close() {}
fun emit(text: String, sink: FakeLiteConversation): List<LiteDelta> {
// Эмулируем один-два фрагмента + done
return listOf(
LiteDelta(text = text.substring(0, text.length / 2), isDone = false),
LiteDelta(text = text.substring(text.length / 2), isDone = true),
)
}
}
private class FakeLiteConversation(
private val parent: FakeLiteLlm,
config: LiteConversationConfig,
) : LiteConversation {
val initialMessages: List<LiteMessage> = config.initialMessages
private val mutableHistory: MutableList<LiteMessage> = config.initialMessages.toMutableList()
override val history: List<LiteMessage>
get() = mutableHistory.toList()
override fun sendStream(prompt: String): Flow<LiteDelta> =
sendStreamContents(listOf(LiteContentPart.Text(prompt)))
override fun sendStreamContents(contents: List<LiteContentPart>): Flow<LiteDelta> {
parent.lastContents = contents
mutableHistory.add(LiteMessage(LiteRole.USER, contents))
if (parent.slow) {
return kotlinx.coroutines.flow.flow {
emit(LiteDelta(text = parent.reply.substring(0, parent.reply.length / 2)))
kotlinx.coroutines.delay(10_000)
emit(LiteDelta(text = parent.reply.substring(parent.reply.length / 2), isDone = true))
mutableHistory.add(LiteMessage.model(parent.reply))
}
}
val first = parent.reply.substring(0, parent.reply.length / 2)
val second = parent.reply.substring(parent.reply.length / 2)
return flowOf(
LiteDelta(text = first),
LiteDelta(text = second, isDone = true),
).also {
mutableHistory.add(LiteMessage.model(parent.reply))
}
}
override fun send(prompt: String): String = parent.reply
override fun sendContents(contents: List<LiteContentPart>): String = parent.reply
override fun cancel() {}
override fun tokenCount(): Int = history.size
override fun addToolResult(callId: String?, name: String, result: String) { error("not used") }
override fun close() {}
}
@@ -0,0 +1,90 @@
package pw.binom.agentik.standalone.llm
import pw.binom.litert.openai.OpenAiConfig
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFailsWith
class LlmConfigTest {
@Test
fun `fromEnv — happy path`() {
val cfg = LlmConfig.fromEnv { name ->
when (name) {
"OPENAI_BASE_URL" -> "https://api.openai.com/v1"
"OPENAI_API_KEY" -> "sk-test"
"OPENAI_MODEL" -> "gpt-4o-mini"
"AGENTIK_SYSTEM_PROMPT" -> "be brief"
else -> null
}
}
assertEquals("be brief", cfg.systemPrompt)
assertEquals(OpenAiConfig(baseUrl = "https://api.openai.com/v1", apiKey = "sk-test", model = "gpt-4o-mini"), cfg.openai)
}
@Test
fun `fromEnv — falls back to default system prompt`() {
val cfg = LlmConfig.fromEnv { name ->
when (name) {
"OPENAI_BASE_URL" -> "https://api.openai.com/v1"
"OPENAI_API_KEY" -> "sk-test"
"OPENAI_MODEL" -> "gpt-4o-mini"
else -> null
}
}
assertEquals(LlmConfig.DEFAULT_SYSTEM_PROMPT, cfg.systemPrompt)
}
@Test
fun `fromEnv — missing base url throws`() {
assertFailsWith<IllegalStateException> {
LlmConfig.fromEnv { name ->
when (name) {
"OPENAI_API_KEY" -> "sk-test"
"OPENAI_MODEL" -> "gpt-4o-mini"
else -> null
}
}
}
}
@Test
fun `fromEnv — missing api key throws`() {
assertFailsWith<IllegalStateException> {
LlmConfig.fromEnv { name ->
when (name) {
"OPENAI_BASE_URL" -> "https://api.openai.com/v1"
"OPENAI_MODEL" -> "gpt-4o-mini"
else -> null
}
}
}
}
@Test
fun `fromEnv — missing model throws`() {
assertFailsWith<IllegalStateException> {
LlmConfig.fromEnv { name ->
when (name) {
"OPENAI_BASE_URL" -> "https://api.openai.com/v1"
"OPENAI_API_KEY" -> "sk-test"
else -> null
}
}
}
}
@Test
fun `blank system prompt from env falls back to default`() {
val cfg = LlmConfig.fromEnv { name ->
when (name) {
"OPENAI_BASE_URL" -> "https://api.openai.com/v1"
"OPENAI_API_KEY" -> "sk-test"
"OPENAI_MODEL" -> "gpt-4o-mini"
"AGENTIK_SYSTEM_PROMPT" -> " "
else -> null
}
}
assertEquals(LlmConfig.DEFAULT_SYSTEM_PROMPT, cfg.systemPrompt)
}
}
@@ -0,0 +1,217 @@
package pw.binom.agentik.standalone.persistence
import kotlinx.coroutines.test.runTest
import pw.binom.agentik.standalone.persistence.sqlite.SqliteStores
import kotlin.test.AfterTest
import kotlin.test.BeforeTest
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFalse
import kotlin.test.assertNotNull
import kotlin.test.assertNull
import kotlin.test.assertTrue
import kotlin.time.Instant
class PersistenceTest {
private lateinit var stores: SqliteStores
@BeforeTest
fun setup() {
stores = SqliteStores.inMemory()
}
@AfterTest
fun tearDown() {
stores.close()
}
@Test
fun `upsert + get conversation — roundtrip`() = runTest {
val now = Instant.fromEpochMilliseconds(1_700_000_000_000)
val rec = ConversationRecord(
id = "c1",
title = "Hello",
isTemporal = false,
createdAt = now,
updatedAt = now,
)
stores.conversations.upsert(rec)
val got = stores.conversations.get("c1")
assertNotNull(got)
assertEquals(rec.id, got.id)
assertEquals(rec.title, got.title)
assertEquals(rec.isTemporal, got.isTemporal)
assertEquals(rec.createdAt, got.createdAt)
assertEquals(rec.updatedAt, got.updatedAt)
}
@Test
fun `upsert overwrites existing record`() = runTest {
val t0 = Instant.fromEpochMilliseconds(1_700_000_000_000)
stores.conversations.upsert(
ConversationRecord("c1", title = "A", isTemporal = false, createdAt = t0, updatedAt = t0),
)
val t1 = Instant.fromEpochMilliseconds(1_700_000_001_000)
stores.conversations.upsert(
ConversationRecord("c1", title = "B", isTemporal = true, createdAt = t0, updatedAt = t1),
)
val got = stores.conversations.get("c1")!!
assertEquals("B", got.title)
assertTrue(got.isTemporal)
assertEquals(t1, got.updatedAt)
}
@Test
fun `list returns conversations ordered by updated_at desc`() = runTest {
val t0 = Instant.fromEpochMilliseconds(1_700_000_000_000)
repeat(3) { i ->
stores.conversations.upsert(
ConversationRecord(
id = "c$i",
title = null,
isTemporal = false,
createdAt = t0,
updatedAt = Instant.fromEpochMilliseconds(1_700_000_000_000 + i * 1000),
),
)
}
val list = stores.conversations.list(offset = 0, limit = 10)
assertEquals(listOf("c2", "c1", "c0"), list.map { it.id })
}
@Test
fun `delete cascades messages and working_memory`() = runTest {
val t0 = Instant.fromEpochMilliseconds(1_700_000_000_000)
stores.conversations.upsert(
ConversationRecord("c1", null, false, t0, t0),
)
stores.messages.append(
MessageRecord.UserMessage(
id = "m1",
conversationId = "c1",
content = listOf(Content.Text("hello")),
createdAt = t0,
),
)
stores.workingMemory.append(
conversationId = "c1",
entry = WorkingMemoryEntry.User(
sourceMessageId = "m1",
content = listOf(Content.Text("hello")),
),
now = t0,
)
assertEquals(1, stores.messages.listAll("c1").size)
assertEquals(1, stores.workingMemory.list("c1").size)
val removed = stores.conversations.delete("c1")
assertTrue(removed)
assertNull(stores.conversations.get("c1"))
assertEquals(emptyList(), stores.messages.listAll("c1"))
assertEquals(emptyList(), stores.workingMemory.list("c1"))
}
@Test
fun `message audit log — append and read back`() = runTest {
val t0 = Instant.fromEpochMilliseconds(1_700_000_000_000)
stores.messages.append(
MessageRecord.UserMessage("m1", "c1", listOf(Content.Text("hi")), t0),
)
stores.messages.append(
MessageRecord.AssistantMessage("m2", "c1", listOf(Content.Text("yo")), t0),
)
val all = stores.messages.listAll("c1")
assertEquals(2, all.size)
assertEquals("m1", all[0].id)
assertEquals("m2", all[1].id)
assertTrue(all[0] is MessageRecord.UserMessage)
assertTrue(all[1] is MessageRecord.AssistantMessage)
assertEquals("hi", (all[0] as MessageRecord.UserMessage).content[0].let {
(it as Content.Text).body
})
}
@Test
fun `message after timestamp filter`() = runTest {
val t0 = Instant.fromEpochMilliseconds(1_700_000_000_000)
val t1 = Instant.fromEpochMilliseconds(1_700_000_001_000)
stores.messages.append(MessageRecord.UserMessage("m1", "c1", listOf(Content.Text("a")), t0))
stores.messages.append(MessageRecord.UserMessage("m2", "c1", listOf(Content.Text("b")), t1))
val after = stores.messages.list("c1", after = t0, offset = 0, limit = 10)
assertEquals(1, after.size)
assertEquals("m2", after[0].id)
}
@Test
fun `working memory — append + list preserves order`() = runTest {
val t0 = Instant.fromEpochMilliseconds(1_700_000_000_000)
val t1 = Instant.fromEpochMilliseconds(1_700_000_001_000)
stores.workingMemory.append(
conversationId = "c1",
entry = WorkingMemoryEntry.System(text = "you are a bot"),
now = t0,
)
stores.workingMemory.append(
conversationId = "c1",
entry = WorkingMemoryEntry.User(sourceMessageId = "m1", content = listOf(Content.Text("hi"))),
now = t0,
)
stores.workingMemory.append(
conversationId = "c1",
entry = WorkingMemoryEntry.Assistant(sourceMessageId = "m2", content = listOf(Content.Text("yo"))),
now = t1,
)
val list = stores.workingMemory.list("c1")
assertEquals(3, list.size)
assertTrue(list[0].entry is WorkingMemoryEntry.System)
assertTrue(list[1].entry is WorkingMemoryEntry.User)
assertTrue(list[2].entry is WorkingMemoryEntry.Assistant)
}
@Test
fun `rename updates title and bumps updated_at`() = runTest {
val t0 = Instant.fromEpochMilliseconds(1_700_000_000_000)
stores.conversations.upsert(ConversationRecord("c1", null, false, t0, t0))
val newTs = stores.conversations.rename("c1", "Renamed")
assertNotNull(newTs)
assertTrue(newTs > t0)
val got = stores.conversations.get("c1")!!
assertEquals("Renamed", got.title)
assertEquals(newTs, got.updatedAt)
}
@Test
fun `rename can clear title`() = runTest {
val t0 = Instant.fromEpochMilliseconds(1_700_000_000_000)
stores.conversations.upsert(ConversationRecord("c1", "Title", false, t0, t0))
stores.conversations.rename("c1", null)
val got = stores.conversations.get("c1")!!
assertNull(got.title)
}
@Test
fun `delete returns false when conversation does not exist`() = runTest {
assertFalse(stores.conversations.delete("nope"))
}
@Test
fun `image content roundtrip through message payload`() = runTest {
val t0 = Instant.fromEpochMilliseconds(1_700_000_000_000)
val bytes = byteArrayOf(0x89.toByte(), 0x50, 0x4E, 0x47) // PNG header
stores.messages.append(
MessageRecord.UserMessage(
id = "m1",
conversationId = "c1",
content = listOf(Content.Image(data = bytes, mime = "image/png")),
createdAt = t0,
),
)
val all = stores.messages.listAll("c1")
val image = (all[0] as MessageRecord.UserMessage).content[0] as Content.Image
assertEquals("image/png", image.mime)
assertTrue(bytes.contentEquals(image.data))
}
}