5bdc517988
Введение монотонного offset'а как персистентного состояния агента:
offset'ы переживают рестарт standalone-агента, клиент продолжает
синхронизацию инкрементально, без полной re-sync с нуля.
outbox-api:
- Cursor (offset: Long) — курсор в журнале событий агента.
- OffsetSequencer — интерфейс резервирования уникального offset.
- CursorStore — персистентное хранилище текущего offset'а.
- PersistentOffsetSequencer — декоратор над любым OutboxStore,
обновляет CursorStore на каждом append (atomic transaction).
- OutboxGapException — клиент запросил after < earliestCursor() →
сервер не может удовлетворить, клиент обязан делать full resync.
- DurableEvent переименован из Event.kt → DurableEvent.kt (Event.kt
был общим sealed-типом, теперь это термин из спеки).
- MutableOutboxStore и OutboxStore теперь читают offset через
CursorStore вместо in-memory counter'а.
outbox-inmemory:
- InMemoryOffsetSequencer — для тестов и dev-режима.
- InMemoryOutboxStore теперь принимает OffsetSequencer в конструкторе.
outbox-ksqlite (новый модуль):
- KsqliteCursorStore — таблица outbox_cursor (agent_id TEXT PK,
offset INTEGER NOT NULL DEFAULT 0, updated_at INTEGER NOT NULL).
- KsqliteCursorStoreTest — 4 теста (set/get, monotonic, concurrent).
proto + server:
- Snapshot.proto — server-state snapshot endpoint для клиентов,
которым нужна полная материализация (использование TBD).
- Routes.kt + SnapshotRouteTest — endpoint /agentik/snapshot (GET).
journal-ksqlite:
- KsqliteJournalStore.listFlow/append — без изменений по API,
нотации минимальные (codecs).
standalone:
- DurableLog (бывший ChatAgent-orchestration) — атомарный commit
события в OutboxStore + PersistentOffsetSequencer + materialization
(через Reducer) одной транзакцией.
- SqliteStores — добавляет KsqliteCursorStore в bundle, единая
shared-connection для всех ksqlite-сторов standalone-агента.
- ChatAgent / ConversationLoop / ConversationEvents / ReflectionScheduler /
ToolDispatcher — переход на новые абстракции.
- standalone/build.gradle.kts — implementation(project(':outbox-ksqlite'))
включено (раньше было закомментировано — модуль только создавался).
client:
- AgentikAgent / AgentClient / HttpEventStore / HttpJournalStore /
ReconnectingOutbox — используют Cursor через transport API.
- client/README.md — синхронизирован с новым поведением (468 строк
diff — это в основном оформление и примеры).
kotlinx-io: 0.8.0 → 0.9.1 в libs.versions.toml (см. sync-core tests).
SYNC-SYSTEM.md (в корне) — спецификация, на которую ссылается и
:sync-core (эта сессия), и эта Cursor-абстракция в outbox-api.
Тесты: standalone 132, journal-ksqlite 25, outbox-inmemory 20,
outbox-ksqlite 4, client 10, sync-core 74 — все зелёные на jvm;
sync-core linuxX64 74 тоже зелёный.
sync2/ (заброшенный stub с одним build.gradle.kts) удалён.
526 lines
29 KiB
Markdown
526 lines
29 KiB
Markdown
# `:client` — Ktor-клиент к `:server` (KMP, jvm + native)
|
||
|
||
Тонкий HTTP-клиент к `:server`-фасаду + локальные примитивы, чтобы
|
||
собирать свои клиенты (UI, CLI, parent-агенты, A2A-bridge) без бойлерплейта
|
||
про HTTP, JSON, SSE и lifecycle `Conversation`.
|
||
|
||
## Что есть
|
||
|
||
- `AgentikAgent(id, baseUrl, engineFactory, token?)` — entry-point. Возвращает
|
||
`Agent` (тот же интерфейс, что в `:proto`). HttpClient создаётся внутри
|
||
из переданной `engineFactory` (`CIO`, `OkHttp`, `Darwin`).
|
||
- `Agent`: `createConversation` / `getConversation` / `getConversations` /
|
||
`deleteConversation` / `journal` / `outbox` / `close`.
|
||
- `Conversation`: `send(content, context?)` / `events(after)` (SSE `Flow<Event>`)
|
||
/ `getMessages(after, offset, limit)` / `rename` / `interrupt` / `close`.
|
||
- `HttpJournalStore` — `list(convId, afterSeq, upToSeq, limit)` /
|
||
`count(convId, afterSeq)` → `List<MessageRecord>` со всеми типами записей
|
||
(User/Assistant/ToolCall/ToolResult/Error + tokens). Адресация — по `seq`
|
||
(см. «Курсорный протокол»), не по датам.
|
||
- `HttpEventStore` — `events(after: Cursor?)` / `agentEvents` /
|
||
`conversationEvents` (SSE), `currentCursor()` / `oldestCursor()`.
|
||
- `Agent.conversationsSnapshot()` / `Agent.chatSnapshot(convId)` — состояние +
|
||
`Cursor`, на котором оно валидно. Точка входа resync'а.
|
||
- `ReconnectingOutbox(outbox, scope, policy)` — обёртка над `OutboxStore` с
|
||
авто-reconnect при обрыве стрима (exponential backoff). Два независимых
|
||
потока: `events(after: Cursor?)` (те же `CommonEvent`) и `connectionStatus()`
|
||
(`Connecting`/`Connected`/`Disconnected`/`Failed`/`Gap`) — статус НЕ мешается
|
||
с основным потоком событий. Мёртвый курсор даёт `Gap` (не ретраится). См. ниже.
|
||
|
||
`Agent` — `AutoCloseable`; `agent.close()` закрывает HttpClient. Не нужно
|
||
вручную создавать `HttpClient` и накатывать на него JSON/Bearer-плагины.
|
||
|
||
## Подключение
|
||
|
||
```kotlin
|
||
// build.gradle.kts
|
||
dependencies {
|
||
api("pw.binom.agentik:client:0.1.0")
|
||
// Движок — на твой выбор (один из):
|
||
implementation("io.ktor:ktor-client-cio:3.x") // JVM/Native
|
||
implementation("io.ktor:ktor-client-okhttp:3.x") // JVM
|
||
implementation("io.ktor:ktor-client-darwin:3.x") // iOS/macOS
|
||
// Опционально — только если будешь использовать `InMemoryJournalStore`
|
||
// как клиентский кэш. Свой `MutableJournalStore` — не нужен.
|
||
api("pw.binom.agentik:journal-inmemory:0.1.0")
|
||
}
|
||
```
|
||
|
||
## Что клиент хранит локально (persistence)
|
||
|
||
Либа **не** имеет `SettingsRepository` / `Config` — это намеренно: UI-фреймворки хранят настройки по-разному (JSON-файл, Keychain, Android DataStore, NSUserDefaults, ...). Либа не навязывает формат, но клиент должен сериализовать у себя минимум:
|
||
|
||
| Поле | Что это | Где взять |
|
||
|---|---|---|
|
||
| `clientId` (параметр `id` в `AgentikAgent`) | Идентичность клиента в логах сервера (X-Client-Id header). Не user-id в агенте, не device-id — это **произвольная строка клиента**, обычно `<app-name>-<installation-uuid>`. Сервер использует для log multiplexing и не интерпретирует. | Генерируется один раз при первом запуске (`UUID.randomUUID().toString()`) и сохраняется. Никогда не меняется. |
|
||
| `baseUrl` | URL сервера (`http://host:8080/agentik`). Должен включать path-prefix фасада, не только хост. | Из настроек пользователя / дефолт |
|
||
| `token` | Bearer-токен. `null` = анонимный доступ (если сервер разрешает). | Из настроек пользователя / secure-storage |
|
||
|
||
Опционально (для UX): `engineFactory` — обычно compile-time выбор по платформе (`CIO` JVM/Native, `OkHttp` JVM, `Darwin` iOS/macOS).
|
||
|
||
Минимальный JSON для UI, который хранит в файле:
|
||
|
||
```json
|
||
{
|
||
"clientId": "my-android-app-550e8400-e29b-41d4-a716-446655440000",
|
||
"baseUrl": "https://agent.example.com/agentik",
|
||
"token": "s3cret"
|
||
}
|
||
```
|
||
|
||
⚠️ `clientId` **генерируется один раз** при установке и больше не меняется — иначе сломается log multiplexing на сервере.
|
||
|
||
## Курсорный протокол (как получить гарантированно актуальное состояние)
|
||
|
||
Всё серьёзное в `:client` крутится вокруг одного понятия — **курсор события**
|
||
(`Cursor(epoch, offset)`), аналога Kafka-offset. Он монотонный, сквозной на
|
||
все события агента (одна общая нумерация для `AgentEvent` и `Conversation`
|
||
-событий) и лежит **над** двумя хранилищами:
|
||
|
||
- **`OutboxStore`** — короткий bounded-tail live-поток `CommonEvent`
|
||
(уведомления/дельты). Хранится ограниченно (cap/TTL), события вытесняются.
|
||
- **`journal` + `conversationStore`** — персистентный источник истины
|
||
(полное состояние). У каждой записи есть свой `seq` из того же счётчика.
|
||
|
||
Ключевое свойство: **`CommonEvent.offset` == `MessageRecord.seq` ==
|
||
`ConversationRecord.seq`**. Событие с `offset = N` — это ровно «строка состояния
|
||
с `seq = N` изменилась (или появилась/удалилась)». События **абсолютные**: в
|
||
`UserMessage`/`AssistantMessage` лежит целая запись, `Renamed` несёт новый
|
||
заголовок, `Deleted` — «строки больше нет». Поэтому накатывать их на состояние
|
||
можно повторно (идемпотентно по `id`) и в любом порядке относительно снапшота.
|
||
|
||
### Инвариант, на котором стоит гарантия
|
||
|
||
1. **Писатель** (сервер) сначала пишет строку состояния с `seq = N`, потом
|
||
кладёт событие с `offset = N` в outbox. `seq` и `offset` — один счётчик.
|
||
2. **Читатель** (клиент) читает **сначала курсор, потом состояние**:
|
||
`C = currentCursor()` → `state = listUpTo(C)`. Всё, что `≤ C`, уже в снапшоте;
|
||
всё, что `> C`, придёт потоком.
|
||
3. **Применение идемпотентно** (upsert/delete/rename по `id`), поэтому
|
||
перекрытие снапшота и дельт безвредно.
|
||
|
||
Ничего не блокируется. Снапшот — это **не** «заморозка таблицы на время
|
||
выгрузки»: это baseline на курсоре `C` плюс накат всех дельт `> C`.
|
||
|
||
### Правильная последовательность синхронизации
|
||
|
||
```
|
||
1. lastSeen = локально сохранённый курсор (или null при первом запуске)
|
||
2. попытка: outbox.events(after = lastSeen) ← если сервер ответил
|
||
OutboxGapException / ConnectionStatus.Gap → курсор мёртв, иди в п.3
|
||
3. ПОЛНЫЙ RESYNC:
|
||
a. очистить локальную БД (строки + курсор), пометить «resyncing»
|
||
b. C = agent.conversationsSnapshot().cursor (или .chatSnapshot(convId))
|
||
c. подписаться events(after = C) и СКОПИРОВАТЬ события в буфер (не применять!)
|
||
d. прочитать полное состояние: listUpTo(C) / snapshot.messages
|
||
e. применить снапшот целиком
|
||
f. применить буфер дельт в порядке offset
|
||
4. дальше: применение каждого события из потока (upsert by id)
|
||
5. сохранить последний offset как lastSeen
|
||
```
|
||
|
||
Порядок из шага 3 критичен: **сначала подписка, потом снапшот**. Если сделать
|
||
наоборот (снапшот, потом подписка) — события, пришедшие в промежуток, потеряются.
|
||
Буферизация (а не «применять на лету») закрывает delete-resurrection: событие
|
||
`Deleted(offset > C)` для строки, которая ещё лежит в необработанной странице
|
||
снапшота, при применении «на лету» было бы стёрто, а потом снапшот вставил бы
|
||
строку обратно.
|
||
|
||
### Курсор мёртв: `OutboxGapException`
|
||
|
||
Клиент давно не заходил, outbox вытеснил его события (`after.offset <
|
||
oldestCursor().offset`), либо сменилась `epoch` (БД сервера откатили/
|
||
восстановили/скопировали — счётчик начал считаться заново). Сервер отвечает
|
||
`410 Gone`. Клиент **не ретраит** — это сигнал «сделай полный resync»
|
||
(шаг 3 выше). С `ReconnectingOutbox` это приходит как
|
||
`ConnectionStatus.Gap`, поток закрывается, background-loop встаёт.
|
||
|
||
**Никогда не ретрай `OutboxGapException`** — ретрай никогда не пройдёт.
|
||
|
||
### Простой вариант: пересоздать outbox на resync
|
||
|
||
Если своя реализация шага 3 кажется тяжёлой — минимальный корректный путь
|
||
через `ReconnectingOutbox`:
|
||
|
||
```kotlin
|
||
var recon = ReconnectingOutbox(agent.outbox, scope)
|
||
scope.launch { recon.events(after = lastSeen).collect { applyEvent(it) } }
|
||
scope.launch {
|
||
recon.connectionStatus().collect { s ->
|
||
if (s is ConnectionStatus.Gap) {
|
||
recon.close()
|
||
val snap = agent.chatSnapshot(convId) // state + cursor
|
||
applySnapshot(snap.messages) // upsert by id
|
||
lastSeen = snap.cursor
|
||
recon = ReconnectingOutbox(agent.outbox, scope)
|
||
scope.launch { recon.events(after = lastSeen).collect { applyEvent(it) } }
|
||
}
|
||
}
|
||
}
|
||
```
|
||
|
||
## Быстрый старт: свой клиент за 5 минут
|
||
|
||
Один self-contained пример: создаём агента, открываем диалог,
|
||
отправляем сообщение, печатаем streaming-ответ.
|
||
|
||
```kotlin
|
||
import pw.binom.agentik.client.AgentikAgent
|
||
import pw.binom.agentik.content.Content
|
||
import pw.binom.agentik.outbox.DurableEvent
|
||
import pw.binom.agentik.outbox.OnlineEvent
|
||
import io.ktor.client.engine.cio.CIO
|
||
import kotlinx.coroutines.launch
|
||
import kotlinx.coroutines.runBlocking
|
||
|
||
fun main() = runBlocking {
|
||
// 1. Agent — обёртка над :server фасадом. HttpClient создаётся внутри.
|
||
val agent = AgentikAgent(
|
||
id = "my-client",
|
||
baseUrl = "http://localhost:8080/agentik",
|
||
engineFactory = CIO,
|
||
token = "s3cret", // или null, если не нужен
|
||
)
|
||
|
||
val conv = agent.createConversation(temp = false)
|
||
|
||
// Подписка «после текущего курсора» — событий строго после этой точки.
|
||
val cursor = agent.outbox.currentCursor()
|
||
|
||
// 2. Два независимых потока событий диалога:
|
||
// durable (outbox) — целые события, с курсором после переподключения;
|
||
// online (OnlineOutbox) — стриминг ответа, только live (без курсора).
|
||
launch {
|
||
agent.outbox.conversationEvents(after = cursor, conversationId = conv.id)
|
||
.collect { ce ->
|
||
when (val ev = ce.event) {
|
||
is DurableEvent.AssistantMessage -> println("[answer ready: ${ev.content}]")
|
||
is DurableEvent.Interrupted -> println("[interrupted]")
|
||
is DurableEvent.Error -> println("[error: ${ev.message}]")
|
||
else -> Unit
|
||
}
|
||
}
|
||
}
|
||
launch {
|
||
agent.onlineOutbox.onlineEvents(conv.id).collect { ev ->
|
||
when (ev) {
|
||
is OnlineEvent.Working -> println("[working]")
|
||
is OnlineEvent.AppendText -> print(ev.body)
|
||
is OnlineEvent.StartResponse -> println("[start]")
|
||
is OnlineEvent.End -> println("\n[end]")
|
||
else -> Unit
|
||
}
|
||
}
|
||
}
|
||
|
||
// 3. Отправить ход (fire-and-forget — ответ придёт по подпискам выше).
|
||
conv.send(listOf(Content.Text("Привет")))
|
||
|
||
// 4. Чистый shutdown.
|
||
conv.close()
|
||
agent.close()
|
||
}
|
||
```
|
||
|
||
**Это весь клиент.** `:server` сам хранит историю, контекст, события.
|
||
Ты только получаешь два типизированных `Flow` и рендеришь как хочешь.
|
||
|
||
> **Durable vs online.** `DurableEvent` (в `agent.outbox`) — «целые» события, их
|
||
> можно перезапросить по курсору `after`. `OnlineEvent` (в
|
||
> `agent.onlineOutbox`) — поток стриминга (`Working`/`End`/`AppendText`/
|
||
> `AppendImage`), **никогда не сохраняется** и не реплеится: потерянный при
|
||
> обрыве фрагмент невосстановим, но целый ответ всегда придёт durable-
|
||
> `DurableEvent.AssistantMessage` и/или ляжет в journal.
|
||
|
||
`HttpClient`, `applyAgentikDefaults`, выбор engine'а — всё скрыто
|
||
внутри `AgentikAgent`. Один вызов — один готовый `Agent`.
|
||
|
||
### Добавить локальный кэш истории (ещё 4 строки)
|
||
|
||
```kotlin
|
||
import pw.binom.agentik.journal.inmemory.InMemoryJournalStore
|
||
|
||
// Свой кэш. Хочешь SQLite/JSON/etc. — реализуй MutableJournalStore сам.
|
||
val cache = InMemoryJournalStore()
|
||
|
||
// Снапшот на курсоре + подписка ПОСЛЕ него — без потерь (см. «Курсорный протокол»).
|
||
val snap = agent.chatSnapshot(conv.id)
|
||
cache.appendAll(snap.messages)
|
||
launch {
|
||
agent.outbox.conversationEvents(after = snap.cursor, conversationId = conv.id)
|
||
.collect { ce -> applyDurable(ce.event, cache) } // upsert by id
|
||
}
|
||
|
||
// История — теперь из кэша, без HTTP:
|
||
val history = cache.list(conv.id, afterSeq = 0L, upToSeq = Long.MAX_VALUE, limit = Int.MAX_VALUE)
|
||
history.forEach { rec ->
|
||
when (rec) {
|
||
is pw.binom.agentik.journal.MessageRecord.UserMessage -> print("user> ${rec.content.text()}")
|
||
is pw.binom.agentik.journal.MessageRecord.AssistantMessage -> print("agent> ${rec.content.text()}")
|
||
is pw.binom.agentik.journal.MessageRecord.ToolCall -> print("[tool: ${rec.toolName}]")
|
||
is pw.binom.agentik.journal.MessageRecord.ToolResult -> print("[result]")
|
||
is pw.binom.agentik.journal.MessageRecord.Error -> print("[error: ${rec.message}]")
|
||
}
|
||
}
|
||
```
|
||
|
||
Шаблон «snapshot(курсор) → local.apply → live-дельты после курсора» работает с
|
||
любым `MutableJournalStore` (см. `:journal-api`). Это и есть кэширование
|
||
«без геморроя» с гарантией актуальности.
|
||
|
||
### Что вообще не нужно писать самому
|
||
|
||
- HTTP-сериализация `DurableEvent`/`Message` — `agentikHttpClient` регистрирует
|
||
`agentikJson` и `InstantSerializer`.
|
||
- SSE-парсер — `readSse()` внутри `:client`.
|
||
- Cursor-менеджмент — сервер ведёт единый монотонный `offset`/`seq`, клиент
|
||
лишь хранит `Cursor(epoch, offset)`. Никаких `Instant`-сравнений и
|
||
pagination-циклов вручную.
|
||
- Lifecycle подписок на `events()` — `Conversation.close()` отменяет SSE-job.
|
||
- HTTP-клиент и Bearer — `AgentikAgent` создаёт `HttpClient(engineFactory)`
|
||
с Bearer'ом из `token=` под капотом; `agent.close()` его закрывает.
|
||
- Движковые настройки (requestTimeout и пр.) — `HttpClient(engineFactory) { ... }`
|
||
создаётся здесь; для нестандартных движковых настроек используй
|
||
`agentikHttpClient(engineFactory, token)` напрямую (он экспортирован).
|
||
|
||
### Что нужно написать самому
|
||
|
||
- UI-рендеринг `DurableEvent`'ов — это твоё (Compose/HTML/CLI).
|
||
- Диалог с пользователем — ввод текста, отображение кнопок и т.п.
|
||
- Persist кэша между запусками (если нужно) — замени `InMemoryJournalStore`
|
||
на свой `MutableJournalStore` (см. `:journal-ksqlite` как пример).
|
||
|
||
|
||
## Базовый пример: send + collect events
|
||
|
||
```kotlin
|
||
import pw.binom.agentik.client.AgentikAgent
|
||
import pw.binom.agentik.content.Content
|
||
import pw.binom.agentik.outbox.DurableEvent
|
||
import pw.binom.agentik.outbox.OnlineEvent
|
||
import io.ktor.client.engine.cio.CIO
|
||
import kotlinx.coroutines.launch
|
||
|
||
val agent = AgentikAgent(
|
||
id = "agentik",
|
||
baseUrl = "http://localhost:8080/agentik",
|
||
engineFactory = CIO,
|
||
)
|
||
|
||
val conv = agent.createConversation(temp = false)
|
||
|
||
// durable-поток (с курсором): terminal-события хода.
|
||
val cursor = agent.outbox.currentCursor()
|
||
launch {
|
||
agent.outbox.conversationEvents(after = cursor, conversationId = conv.id)
|
||
.collect { ce ->
|
||
when (ce.event) {
|
||
is DurableEvent.AssistantMessage -> println("\n--- answer ready ---")
|
||
is DurableEvent.Error -> error("agent error: ${(ce.event as DurableEvent.Error).message}")
|
||
else -> Unit
|
||
}
|
||
}
|
||
}
|
||
// online-поток (live-only): стриминг ответа.
|
||
launch {
|
||
agent.onlineOutbox.onlineEvents(conv.id).collect { ev ->
|
||
when (ev) {
|
||
is OnlineEvent.AppendText -> print(ev.body) // streaming чанки
|
||
is OnlineEvent.End -> println("\n--- end ---")
|
||
else -> Unit
|
||
}
|
||
}
|
||
}
|
||
|
||
conv.send(listOf(Content.Text("Привет, расскажи про себя")))
|
||
```
|
||
## Локальный кэш истории (правильный паттерн)
|
||
|
||
Клиент держит свой `MutableJournalStore` и наполняет его **снапшотом на
|
||
курсоре + дельтами после курсора** (см. «Курсорный протокол»). Чтение истории —
|
||
из локального кэша, без HTTP.
|
||
|
||
```kotlin
|
||
import pw.binom.agentik.journal.inmemory.InMemoryJournalStore
|
||
import pw.binom.agentik.journal.MessageRecord
|
||
import pw.binom.agentik.outbox.DurableEvent
|
||
|
||
// Кэш. Для диска — свой MutableJournalStore (KsqliteJournalStore в :journal-ksqlite).
|
||
val cache = InMemoryJournalStore()
|
||
|
||
class ChatSession(
|
||
private val agent: pw.binom.agentik.proto.Agent,
|
||
val conversationId: String,
|
||
) : AutoCloseable {
|
||
private val scope = kotlinx.coroutines.CoroutineScope(
|
||
kotlinx.coroutines.SupervisorJob() + kotlinx.coroutines.Dispatchers.Default,
|
||
)
|
||
var lastSeen: pw.binom.agentik.outbox.Cursor? = null
|
||
|
||
init {
|
||
scope.launch {
|
||
// 1. Снапшот: состояние + курсор, на котором оно валидно.
|
||
val snap = agent.chatSnapshot(conversationId)
|
||
snap.messages.forEach { cache.append(it) }
|
||
lastSeen = snap.cursor
|
||
// 2. Дельты строго после курсора снапшота.
|
||
agent.outbox.conversationEvents(after = snap.cursor, conversationId = conversationId)
|
||
.collect { ce ->
|
||
applyToCache(ce.event)
|
||
lastSeen = ce.cursor
|
||
}
|
||
}
|
||
}
|
||
|
||
private suspend fun applyToCache(e: DurableEvent) {
|
||
when (e) {
|
||
is DurableEvent.UserMessage -> cache.append(e.toRecord())
|
||
is DurableEvent.AssistantMessage -> cache.append(e.toRecord())
|
||
is DurableEvent.ToolCall -> cache.append(e.toRecord())
|
||
is DurableEvent.ToolResult -> cache.append(e.toRecord())
|
||
is DurableEvent.Error -> cache.append(e.toRecord())
|
||
is DurableEvent.Interrupted -> Unit
|
||
}
|
||
}
|
||
|
||
fun history() = kotlinx.coroutines.runBlocking {
|
||
cache.list(conversationId, afterSeq = 0L, upToSeq = Long.MAX_VALUE, limit = Int.MAX_VALUE)
|
||
}
|
||
|
||
override fun close() { scope.cancel() }
|
||
}
|
||
```
|
||
|
||
> `applyToCache` через `cache.append` даёт upsert по `id` (append-only store
|
||
> отбрасывает дубликаты `id`), поэтому перекрытие снапшота и дельт безвредно.
|
||
> Замените `InMemoryJournalStore` на `KsqliteJournalStore` — код не меняется.
|
||
|
||
### Когда курсор мёртв
|
||
|
||
Если `conversationEvents(after = ...)` бросает `OutboxGapException` (или
|
||
`ReconnectingOutbox` эмитит `ConnectionStatus.Gap`) — клиент был оффлайн дольше
|
||
retention'а. Повторите всю последовательность с шага 1 (снапшот), **предварительно
|
||
очистив локальную БД** (`cache.clear(conversationId)`), иначе воскреснут
|
||
удалённые строки. Полный алгоритм — в «Курсорный протокол» выше.
|
||
|
||
## Кэш списка бесед
|
||
|
||
`agent.conversationStore` — read-only projection поверх таблицы `conversation`
|
||
на сервере (`ConversationRecord` = id / title / isTemporal / createdAt /
|
||
updatedAt). `AgentikAgent` оборачивает его в локальный кэш
|
||
(`wrapWithLocalConversationCache`) по тому же протоколу, что и историю:
|
||
снапшот на курсоре + live-дельты.
|
||
|
||
```kotlin
|
||
import pw.binom.agentik.client.AgentikAgent
|
||
import pw.binom.agentik.outbox.AgentEvent
|
||
import io.ktor.client.engine.cio.CIO
|
||
|
||
val agent = AgentikAgent(id = "agentik", baseUrl = "http://localhost:8080/agentik", engineFactory = CIO)
|
||
|
||
// Снапшот списка бесед + его курсор (глобальный для агента).
|
||
val snap = agent.conversationsSnapshot()
|
||
snap.conversations.forEach { println("${it.id} ${it.title ?: "(no title)"} ${it.updatedAt}") }
|
||
|
||
// Дельты после курсора снапшота: Created / Deleted / Renamed / Touched.
|
||
agent.outbox.agentEvents(after = snap.cursor).collect { ce ->
|
||
when (val ev = ce.event) {
|
||
is AgentEvent.Created -> println("+ ${ev.conversationId}")
|
||
is AgentEvent.Renamed -> println("~ ${ev.id} -> ${ev.title}")
|
||
is AgentEvent.Touched -> println("~ ${ev.id} (${ev.updatedAt})")
|
||
is AgentEvent.Deleted -> println("- ${ev.id}")
|
||
}
|
||
}
|
||
```
|
||
|
||
**Команды** (создать / переименовать / удалить) идут через `agent`; сервер сам
|
||
эмитит соответствующее `AgentEvent` в outbox, клиент применяет его к кэшу:
|
||
|
||
```kotlin
|
||
val conv = agent.createConversation(temp = false) // POST /conversations -> Created
|
||
agent.renameConversation(conv.id, "Новый заголовок") // PATCH /conversations/{id} -> Renamed
|
||
agent.deleteConversation(conv.id) // DELETE /conversations/{id} -> Deleted
|
||
```
|
||
|
||
**Никогда не пиши в `conversationStore` напрямую.** Для активной работы
|
||
(send / interrupt) — handle через `agent.getConversation(id)`.
|
||
|
||
### Если хочется своего cache-импла
|
||
|
||
`InMemoryMutableConversationStore` подходит для большинства случаев. Для диска —
|
||
свой `MutableConversationStore` (см. `KsqliteMutableConversationStore` в
|
||
`:journal-ksqlite`). Методы `rename`/`touch` принимают `seq` из общего счётчика:
|
||
override suspend fun list(offset: Int, limit: Int): List<ConversationRecord> { /* SELECT ORDER BY updatedAt DESC */ }
|
||
override suspend fun delete(id: String): Boolean { /* DELETE */ }
|
||
override suspend fun rename(id: String, title: String?): Instant? { /* UPDATE + bump updatedAt */ }
|
||
override suspend fun touch(id: String, now: Instant) { /* UPDATE updatedAt */ }
|
||
override fun close() {}
|
||
}
|
||
```
|
||
|
||
## Тесты
|
||
|
||
```
|
||
./gradlew :client:jvmTest
|
||
```
|
||
|
||
Покрывают: JSON-парсинг `DurableEvent`-ов, SSE-стрим, recovery после разрыва,
|
||
401/404, reconnect-cycle `ReconnectingOutbox` (5 кейсов: успех / обрыв +
|
||
reconnect / exhausted attempts → Failed / мёртвый курсор → Gap (без ретрая) /
|
||
close → cancel).
|
||
|
||
## Auto-reconnect для живого outbox
|
||
|
||
Базовый `OutboxStore.events(after: Cursor?)` — cold SSE-стрим; при обрыве
|
||
(мобильная сеть, рестарт сервера) клиент должен сам реконнектиться с курсором
|
||
последнего увиденного события. Это повторяется в каждом клиенте, поэтому
|
||
`ReconnectingOutbox` берёт это на себя:
|
||
|
||
```kotlin
|
||
val recon = ReconnectingOutbox(
|
||
outbox = agent.outbox, // или HttpEventStore
|
||
scope = myScreenScope,
|
||
policy = BackoffPolicy.Default, // 1s → 2s → ... → 30s, ±20% jitter
|
||
)
|
||
|
||
// lastSeen — курсор из последнего снапшота / последнего события.
|
||
scope.launch { recon.events(after = lastSeen).collect { handle(it) } }
|
||
scope.launch {
|
||
recon.connectionStatus().collect { status ->
|
||
when (status) {
|
||
is Connecting -> ui.showBanner("connecting...")
|
||
is Connected -> ui.hideBanner()
|
||
is Disconnected -> ui.showBanner("reconnecting in ${status.willRetryIn}…")
|
||
is Failed -> ui.showError(status.cause)
|
||
is Gap -> resync(status.cause) // курсор мёртв — полный resync
|
||
}
|
||
}
|
||
}
|
||
|
||
// На выходе (например, navigation back):
|
||
recon.close() // отменяет background-loop, потоки терминируются
|
||
```
|
||
|
||
`Gap` — единственный статус, который **не** ретраится: курсор старше
|
||
retention'а или чужая эпоха. Обработка — полный resync (см. «Курсорный
|
||
протокол»). Если не передать `after`, при старте берётся
|
||
`outbox.currentCursor()` (live-only семантика).
|
||
|
||
Два потока **независимы** — `events()` содержит только `CommonEvent`,
|
||
`connectionStatus()` содержит только `ConnectionStatus`. Никакого
|
||
"мешающего" `Connecting`/`Disconnected` в потоке событий.
|
||
|
||
Параметры backoff (см. `BackoffPolicy`):
|
||
- `initial` / `max` — границы задержки
|
||
- `multiplier` — множитель на каждом шаге
|
||
- `jitter` — рандом-разброс (по умолчанию 20%)
|
||
- `maxAttempts` — лимит попыток; после — `Failed` + закрытие потока
|
||
|
||
Если нужен фиксированный delay для тестов — `BackoffPolicy.Fixed(10.milliseconds, attempts = 3)`.
|
||
|
||
## Известное ограничение
|
||
|
||
SSE event-stream в не-TTY ssh-сессии (без `-tt`) закрывается на
|
||
default-таймауте Ktor. Используйте либо `ssh -tt`, либо нативный
|
||
terminal (TTY). Это upstream-особенность Ktor SSE.
|