Изменения вводят live-канал OnlineOutbox и фоновую синхронизацию диалогов.

This commit is contained in:
2026-09-30 11:43:49 +03:00
parent 2ad2d65f06
commit 26f94ea83d
12 changed files with 699 additions and 109 deletions
+24 -2
View File
@@ -302,6 +302,30 @@
- **R43.** Полезный признак, что что-то не вынесено: **в экране накопилась логика
длиннее пары десятков строк.** Значит, это место жить в экране не должно.
### 8.4. Кэш наполняется фоново, а не по факту открытия
**Задача:** превью и история должны быть у **всех** диалогов, а не только у тех,
что открывали в этом окне. Раньше «…» в списке стояло именно потому, что запись
появлялась лишь при открытии диалога.
**Схема (на каждого агента):**
1. Одна подписка на **все** события-сообщения
(`outbox.conversationEvents(after=null, conversationId=null)`).
2. На `End`/`Interrupted` любого диалога — догон журнала по watermark'у и
пересчёт превью; данные кладутся в локальный кэш.
3. На старте — догон диалогов, чей `updatedAt` с сервера **новее** нашего
watermark'а, включая **никогда не открытые**.
- **R44.** **Источник истины — журнал, не SSE.** Outbox — короткий bounded tail:
события за время, пока клиент был выключен, из него выпадают. Поэтому догон
идёт по журналу (`updatedAt` / `synced_at`), а SSE лишь триггерит его.
- **R45.** **Открытый диалог фоновый ingest не трогает.** Пока диалог открыт,
его сообщения ведёт сессия; engine пропускает этот диалог, чтобы не дублировать
работу. Общий мьютекс на диалог сериализует запись на транзишене
«открыли / закрыли».
- **R46.** **API менять не нужно** (снимает открытый вопрос из §9): превью
считается локально из журнала, число непрочитанных — `journal.count(after=lastSeen)`.
## 9. Решения, которые ещё не приняты
- Запись: по нажатию или на удержание.
@@ -315,8 +339,6 @@
- **Группы — только на устройстве или на сервере?** Сервер про них не знает
(`STORAGE.md`, §7). Если только у нас — раскладка не поедет между десктопом
и телефоном.
- **Меняем ли API: превью последнего сообщения и число непрочитанных.**
Взять неоткуда, в макете они нарисованы (`STORAGE.md`, §6).
- **Бейдж — число или точка.** Точка сервер не трогает вообще (`STORAGE.md`, §6.2).
- **Где каталог клиента** — `~/.agentik/` или системный («Документы»).
На Android понятие «домашний каталог» своё.
@@ -20,7 +20,7 @@ import kotlin.time.Instant
* `journal.count(convId, after = lastSeen)` и здесь не лежит.
* - **`preview`** — последний фрагмент текста хода, денормализованный
* в строку. Сервер такого поля не отдаёт (STORAGE.md §6.1), клиент
* поддерживает сам на `Event.AppendText` / `Event.End`.
* поддерживает сам на `OnlineEvent.AppendText` / `Event.End`.
*
* **Почему НЕ наследуется от `KsqliteJournalStore`:** разные таблицы
* (`conversation_meta` + `chat_group` против `message`) и разные контракты
@@ -2,6 +2,7 @@ package pw.binom.agentik.desktop.session
import io.ktor.client.engine.HttpClientEngineFactory
import io.ktor.client.engine.cio.CIO
import java.util.concurrent.ConcurrentHashMap
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.SupervisorJob
@@ -24,6 +25,7 @@ import pw.binom.agentik.journal.ConversationStore
import pw.binom.agentik.journal.MessageRecord
import pw.binom.agentik.journal.MutableJournalStore
import pw.binom.agentik.journal.ksqlite.KsqliteJournalStore
import pw.binom.agentik.outbox.OnlineOutbox
import pw.binom.agentik.outbox.OutboxStore
import pw.binom.agentik.proto.Agent
import pw.binom.agentik.proto.AgentInfo
@@ -107,6 +109,25 @@ class AgentConnection private constructor(
/** Live-события (SSE). */
val outbox: OutboxStore get() = agent.outbox
/** Live-стриминг ответа (дельты текста/картинок), без сохранения. */
val onlineOutbox: OnlineOutbox get() = agent.onlineOutbox
/** Общие мьютексы бэкфилла: сессия + фоновый ingest не пишут один диалог разом. */
private val locks = ConversationLocks()
/** Фоновый ingest ВСЕХ conversation-событий: локальный кэш сообщений + превью. */
private val syncEngine = AgentSyncEngine(
remoteJournal = agent.journal,
journalCache = journalCache,
meta = meta,
outbox = agent.outbox,
locks = locks,
scope = watchScope,
)
/** Диалоги с открытой сессией — их ведёт [ChatSession], engine их пропускает. */
private val openSessions = ConcurrentHashMap.newKeySet<String>()
init {
startConversationWatch()
}
@@ -120,6 +141,13 @@ class AgentConnection private constructor(
delay(SEED_RETRY_DELAY)
}
refreshConversations()
// Глобальный ingest: события всех диалогов → кэш + превью. Отдельной
// подписки на сообщения у UI-списка больше нет — превью неоткрытых
// диалогов догоняются журналом, а не только по факту открытия.
syncEngine.start(
conversations = { agent.conversationStore.list(offset = 0, limit = 500) },
onChanged = { refreshConversations() },
)
runCatching {
agent.outbox.agentEvents(Instant.DISTANT_PAST).collect { refreshConversations() }
}
@@ -205,6 +233,8 @@ class AgentConnection private constructor(
*/
suspend fun openSession(id: String): ChatSession? {
val conv = agent.getConversation(id) ?: return null
openSessions.add(id)
syncEngine.openConversationIds = openSessions.toSet()
return ChatSession(
conversationId = id,
conversation = conv,
@@ -212,11 +242,18 @@ class AgentConnection private constructor(
journalCache = journalCache,
meta = meta,
outbox = agent.outbox,
onlineOutbox = agent.onlineOutbox,
scope = scope,
locks = locks,
onClosed = {
openSessions.remove(id)
syncEngine.openConversationIds = openSessions.toSet()
},
)
}
override fun close() {
syncEngine.close()
watchScope.cancel()
runCatching { agent.close() }
runCatching { journalCache.close() }
@@ -0,0 +1,110 @@
package pw.binom.agentik.desktop.session
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Job
import kotlinx.coroutines.launch
import pw.binom.agentik.desktop.persistence.ConversationMetaRepository
import pw.binom.agentik.journal.ConversationRecord
import pw.binom.agentik.journal.JournalStore
import pw.binom.agentik.journal.MutableJournalStore
import pw.binom.agentik.outbox.Event
import pw.binom.agentik.outbox.OutboxStore
/**
* Фоновый ingest одного агента: **одна** подписка на все conversation-события
* (`conversationEvents(after = null, conversationId = null)`) держит локальный
* кэш сообщений и превью списка в актуальном состоянии.
*
* **Зачем.** Раньше сообщения диалога попадали в кэш только когда диалог
* открывали ([ChatSession]) — у неоткрытых диалогов не было ни истории, ни
* превью («…»). Теперь:
* - на старте engine догоняет журнал для диалогов, изменившихся с прошлого
* раза (`updatedAt > syncedAt`), включая никогда не открытые — кэш и превью
* заполняются сразу;
* - на каждый `End`/`Interrupted` любого диалога — инкрементальный догон
* журнала + пересчёт превью;
* - `updatedAt`/`createdAt` берутся у сервера, дубли отсекаются по id.
*
* **Почему не полагаемся только на SSE.** Outbox — bounded TTL tail: события
* старше буфера (пока приложение было выключено) из него выпадают. Поэтому
* источник истины — журнал, а SSE лишь триггерит догон.
*
* **Открытый диалог ведёт [ChatSession]** — engine его пропускает
* ([openConversationIds]); общий [ConversationLocks] сериализует бэкфилл, если
* транзишен «открыли/закрыли» совпал с фоновым догоном.
*/
internal class AgentSyncEngine(
private val remoteJournal: JournalStore,
private val journalCache: MutableJournalStore,
private val meta: ConversationMetaRepository,
private val outbox: OutboxStore,
private val locks: ConversationLocks,
private val scope: CoroutineScope,
) {
/** Диалоги, которые прямо сейчас держит открытая [ChatSession]. */
@Volatile
var openConversationIds: Set<String> = emptySet()
private var job: Job? = null
/**
* Запустить ingest. Идемпотентно.
*
* @param conversations список диалогов агента на момент вызова (для
* стартового догона).
* @param onChanged дёргается, когда локальный кэш/превью изменились —
* чтобы UI перечитал список.
*/
fun start(
conversations: suspend () -> List<ConversationRecord>,
onChanged: suspend () -> Unit,
) {
if (job != null) return
job = scope.launch {
// Подписка — первой: стартовый догон её не задерживает.
launch {
outbox.conversationEvents(after = null, conversationId = null).collect { ev ->
if (ev.conversationId in openConversationIds) return@collect
when (ev.event) {
is Event.End, is Event.Interrupted -> {
runCatching { syncConversation(ev.conversationId) }
onChanged()
}
else -> Unit
}
}
}
// Стартовый догон изменённых диалогов (включая никогда не открытые).
runCatching {
var changed = false
for (rec in conversations()) {
if (rec.id in openConversationIds) continue
val syncedAt = meta.meta(rec.id)?.syncedAt
if (syncedAt == null || rec.updatedAt > syncedAt) {
syncConversation(rec.id, fallbackWatermark = rec.updatedAt)
changed = true
}
}
if (changed) onChanged()
}
}
}
/** Догнать один диалог (бэкфилл + превью) под общим мьютексом. */
suspend fun syncConversation(conversationId: String, fallbackWatermark: kotlin.time.Instant? = null) =
syncConversationLocked(
remoteJournal = remoteJournal,
journalCache = journalCache,
meta = meta,
locks = locks,
conversationId = conversationId,
fallbackWatermark = fallbackWatermark,
)
fun close() {
job?.cancel()
job = null
}
}
@@ -16,11 +16,11 @@ import pw.binom.agentik.desktop.persistence.ConversationMetaRepository
import pw.binom.agentik.journal.JournalStore
import pw.binom.agentik.journal.MutableJournalStore
import pw.binom.agentik.outbox.Event
import pw.binom.agentik.outbox.OnlineOutbox
import pw.binom.agentik.outbox.OutboxStore
import pw.binom.agentik.proto.Content
import pw.binom.agentik.proto.Conversation
import kotlin.time.Clock
import kotlin.time.Duration.Companion.milliseconds
import kotlin.time.Instant
/**
@@ -32,8 +32,9 @@ import kotlin.time.Instant
* старте догоняем то, что появилось с прошлого раза, и кладём в кэш.
* Watermark — `conversation_meta.synced_at`, чтобы не вставлять дубли
* (`message.id` — PRIMARY KEY, повторный INSERT бросает).
* - **Live-ход** — `OutboxStore.conversationEvents` (SSE): события
* `Working → StartReasoning → AppendText → … → End` накапливаются в
* - **Live-ход** — два SSE-потока: durable
* `OutboxStore.conversationEvents` (`Working → … → End`) и онлайн
* `OnlineOutbox.onlineEvents` (стриминг ответа). Оба копятся в
* [LiveTurn] и рисуются поверх истории, пока ход не закрыт.
* - **Наши поля** — [ConversationMetaRepository]: `lastSeen` (бейдж),
* `preview` (строка списка).
@@ -61,7 +62,15 @@ class ChatSession(
private val journalCache: MutableJournalStore,
private val meta: ConversationMetaRepository,
private val outbox: OutboxStore,
private val onlineOutbox: OnlineOutbox,
private val scope: CoroutineScope,
/**
* Общие с [AgentSyncEngine] мьютексы по диалогам: сериализуют бэкфилл,
* если фоновый ingest и открытая сессия синхронят один диалог.
*/
private val locks: ConversationLocks = ConversationLocks(),
/** Колбэк на [close] — engine снимает диалог из своего open-набора. */
private val onClosed: () -> Unit = {},
) : AutoCloseable {
private val _history = MutableStateFlow<List<UiMessage>>(emptyList())
@@ -88,16 +97,17 @@ class ChatSession(
private val liveTurn = LiveTurn()
private var eventsJob: Job? = null
private var onlineJob: Job? = null
private var started = false
/**
* Сериализует [sync]. Без него два параллельных триггера (оптимистичный
* `sync()` из [sendContent] и `sync()` из обработчика SSE `End`, который
* приходит почти сразу на быстром ответе) читают один и тот же watermark
* `null` → `after = DISTANT_PAST` и оба вставляют user-сообщение →
* Сериализует [sync] через общий [locks]. Без него два параллельных
* триггера (оптимистичный `sync()` из [sendContent] и `sync()` из
* обработчика SSE `End`) читают один и тот же watermark `null` →
* `after = DISTANT_PAST` и оба вставляют user-сообщение →
* `UNIQUE constraint failed: message.id`.
*/
private val syncMutex = Mutex()
private fun syncLock(): Mutex = locks.lock(conversationId)
/**
* Догоняет историю и подписывается на live-события. Идемпотентно —
@@ -112,6 +122,7 @@ class ChatSession(
}
sync()
subscribeEvents()
subscribeOnline()
}
/** Отправить текстовый ход. Пустой/пробельный текст игнорируется. */
@@ -147,7 +158,10 @@ class ChatSession(
override fun close() {
eventsJob?.cancel()
eventsJob = null
onlineJob?.cancel()
onlineJob = null
runCatching { conversation.close() }
onClosed()
}
// ───── internals ─────
@@ -168,44 +182,36 @@ class ChatSession(
}
}
/**
* Подписка на **онлайн-поток** — стриминг ответа (`StartReasoning` /
* `StartResponse` / `AppendText` / `AppendImage`). Отдельный live-канал
* без курсора и без catchup: пропущенный при обрыве фрагмент
* невосстановим, но целый результат придёт durable-`End` и/или ляжет
* в журнал (см. [sync]).
*/
private fun subscribeOnline() {
onlineJob = scope.launch {
onlineOutbox.onlineEvents(conversationId).collect { event ->
liveTurn.applyOnline(event)
}
}
}
private suspend fun currentSyncedAt(): Instant? = meta.meta(conversationId)?.syncedAt
/** Инкрементальный бэкфилл + перечитка кэша + превью. Сериализован. */
private suspend fun sync() = syncMutex.withLock {
private suspend fun sync() = syncLock().withLock {
backfill()
reloadHistory()
updatePreview()
}
/**
* Дотягивает сообщения, появившиеся с прошлого watermark'а.
*
* **Идемпотентность по id, а не по времени.** Watermark — миллисекундный:
* user- и assistant-сообщение быстрого ответа делят одну миллисекунду, и
* фильтр `createdAt > after` просто потерял бы второе. Поэтому watermark
* отступает на 1 мс назад, а повторно пришедшие записи отсекаются по
* множеству уже известных id (никакого `runCatching` — реальные ошибки
* записи должны быть видны).
* Дотягивает сообщения, появившиеся с прошлого watermark'а (см.
* [backfillConversation] — там же про идемпотентность по id).
*/
private suspend fun backfill() {
val row = meta.meta(conversationId)
val after = row?.syncedAt ?: Instant.DISTANT_PAST
val known = journalCache.list(
conversationId = conversationId,
after = Instant.DISTANT_PAST,
offset = 0,
limit = Int.MAX_VALUE,
).mapTo(mutableSetOf()) { it.id }
var newest: Instant? = null
remoteJournal.listFlow(conversationId = conversationId, after = after).collect { rec ->
if (known.add(rec.id)) {
journalCache.append(rec)
}
val n = newest
if (n == null || rec.createdAt > n) newest = rec.createdAt
}
newest?.let { meta.setSyncedAt(conversationId, it - 1.milliseconds) }
backfillConversation(remoteJournal, journalCache, meta, conversationId)
}
private suspend fun reloadHistory() {
@@ -0,0 +1,101 @@
package pw.binom.agentik.desktop.session
import java.util.concurrent.ConcurrentHashMap
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
import pw.binom.agentik.desktop.model.UiMessage
import pw.binom.agentik.desktop.model.toUiMessages
import pw.binom.agentik.desktop.persistence.ConversationMetaRepository
import pw.binom.agentik.journal.JournalStore
import pw.binom.agentik.journal.MutableJournalStore
import kotlin.time.Duration.Companion.milliseconds
import kotlin.time.Instant
/**
* Мьютексы по диалогам. Сериализуют инкрементальный бэкфилл одного диалога
* между открытой [ChatSession] и фоновым [AgentSyncEngine]: при транзишене
* «диалог только что открыли / только что закрыли» оба могут дёрнуть sync
* одновременно, а двойной INSERT одного `message.id` бросает
* `UNIQUE constraint failed`.
*/
class ConversationLocks {
private val locks = ConcurrentHashMap<String, Mutex>()
fun lock(conversationId: String): Mutex = locks.getOrPut(conversationId) { Mutex() }
}
/**
* Инкрементальный бэкфилл одного диалога: дотягивает из удалённого журнала
* всё, что появилось после локального watermark'а, и кладёт в кэш.
*
* **Идемпотентность по id, а не по времени.** Watermark — миллисекундный:
* user- и assistant-сообщение быстрого ответа делят одну миллисекунду, и
* фильтр `createdAt > after` потерял бы второе. Поэтому watermark отступает
* на 1 мс назад, а повторно пришедшие записи отсекаются по множеству уже
* известных id (никакого `runCatching` — реальные ошибки записи должны быть
* видны).
*
* @param fallbackWatermark если в журнале не оказалось ни одной записи (пустой
* диалог), watermark ставится по нему, чтобы не перезапрашивать пустой
* диалог на каждом старте.
* @return `createdAt` самой свежей дотянутой записи, или `null`.
*/
internal suspend fun backfillConversation(
remoteJournal: JournalStore,
journalCache: MutableJournalStore,
meta: ConversationMetaRepository,
conversationId: String,
fallbackWatermark: Instant? = null,
): Instant? {
val after = meta.meta(conversationId)?.syncedAt ?: Instant.DISTANT_PAST
val known = journalCache.list(
conversationId = conversationId,
after = Instant.DISTANT_PAST,
offset = 0,
limit = Int.MAX_VALUE,
).mapTo(mutableSetOf()) { it.id }
var newest: Instant? = null
remoteJournal.listFlow(conversationId = conversationId, after = after).collect { rec ->
if (known.add(rec.id)) {
journalCache.append(rec)
}
val n = newest
if (n == null || rec.createdAt > n) newest = rec.createdAt
}
val watermark = newest ?: fallbackWatermark
watermark?.let { meta.setSyncedAt(conversationId, it - 1.milliseconds) }
return newest
}
/**
* Пересчитывает денормализованное превью строки списка: последнее
* assistant-сообщение, иначе последнее вообще. Читает из локального кэша.
*/
internal suspend fun updateConversationPreview(
journalCache: MutableJournalStore,
meta: ConversationMetaRepository,
conversationId: String,
) {
val history = journalCache.list(
conversationId = conversationId,
after = Instant.DISTANT_PAST,
offset = 0,
limit = Int.MAX_VALUE,
).toUiMessages()
val last = history.lastOrNull { it is UiMessage.Assistant } ?: history.lastOrNull()
meta.setPreview(conversationId, last?.previewText())
}
/** Бэкфилл + пересчёт превью одного диалога под общим мьютексом. */
internal suspend fun syncConversationLocked(
remoteJournal: JournalStore,
journalCache: MutableJournalStore,
meta: ConversationMetaRepository,
locks: ConversationLocks,
conversationId: String,
fallbackWatermark: Instant? = null,
) = locks.lock(conversationId).withLock {
backfillConversation(remoteJournal, journalCache, meta, conversationId, fallbackWatermark)
updateConversationPreview(journalCache, meta, conversationId)
}
@@ -1,6 +1,7 @@
package pw.binom.agentik.desktop.session
import pw.binom.agentik.outbox.Event
import pw.binom.agentik.outbox.OnlineEvent
/**
* Фаза текущего хода ассистента. Управляет спиннером и тем, что UI
@@ -13,10 +14,10 @@ enum class TurnPhase {
/** `send()` ушёл, ответа ещё нет (`Event.Working`). Спиннер «агент думает». */
WORKING,
/** Идут рассуждения (`StartReasoning` + `AppendText` в reasoning-блок). */
/** Идут рассуждения (`StartReasoning` + `AppendText` в reasoning-блок). Live-канал. */
REASONING,
/** Стримится ответ (`StartResponse` + `AppendText`/`AppendImage`/tool'ы). */
/** Стримится ответ (`StartResponse` + `AppendText`/`AppendImage`/tool'ы). Live-канал. */
RESPONDING,
/** Ход закрыт нормально (`End`). Live-блоки пора заменить историей из кэша. */
@@ -92,8 +93,8 @@ sealed interface LiveBlock {
}
/**
* Аккумулятор live-хода: применяет [Event] по мере SSE-стрима и держит
* текущий список [LiveBlock] + [phase].
* Аккумулятор live-хода: применяет durable-[Event] и онлайн-[OnlineEvent]
* по мере двух SSE-стримов и держит текущий список [LiveBlock] + [phase].
*
* Это **модель представления**, не хранилище: после `End` ход
* «схлопывается» в обычную `MessageRecord`, и UI переключается на историю
@@ -151,7 +152,9 @@ class LiveTurn {
}
/**
* Применяет одно событие. Неизвестные/нерелевантные события
* Применяет одно **durable** событие ([Event]). Стриминговые события
* (дельты текста/картинок, маркеры фаз) сюда НЕ приходят — они в
* [applyOnline] ([OnlineEvent]). Неизвестные/нерелевантные
* (`ConversationClosing`, `CompactionTriggered`) игнорируются.
*/
fun apply(event: Event) {
@@ -161,38 +164,6 @@ class LiveTurn {
start()
}
is Event.StartReasoning -> {
phase = TurnPhase.REASONING
_blocks.add(LiveBlock.Reasoning(nextKey("reasoning")))
fire()
}
is Event.StartResponse -> {
phase = TurnPhase.RESPONDING
if (event.responseType == Event.ResponseType.TEXT) {
_blocks.add(LiveBlock.Answer(nextKey("answer")))
}
fire()
}
is Event.AppendText -> {
when (phase) {
TurnPhase.REASONING -> appendText(event.body, reasoning = true)
// AppendText без предшествующего StartResponse —
// трактуем как ответ (forward-compat).
else -> {
if (phase != TurnPhase.RESPONDING) phase = TurnPhase.RESPONDING
appendText(event.body, reasoning = false)
}
}
}
is Event.AppendImage -> {
phase = TurnPhase.RESPONDING
_blocks.add(LiveBlock.Image(nextKey("image"), event.mime, event.body))
fire()
}
is Event.ToolCall -> {
_blocks.add(
LiveBlock.Tool(
@@ -253,6 +224,48 @@ class LiveTurn {
}
}
/**
* Применяет одно **онлайн** событие ([OnlineEvent]) — live-стриминг
* ответа. Онлайн-события нигде не сохраняются: при обрыве соединения
* потерянный фрагмент невосстановим (целый результат придёт
* durable-[Event.End] и/или ляжет в журнал).
*/
fun applyOnline(event: OnlineEvent) {
when (event) {
is OnlineEvent.StartReasoning -> {
phase = TurnPhase.REASONING
_blocks.add(LiveBlock.Reasoning(nextKey("reasoning")))
fire()
}
is OnlineEvent.StartResponse -> {
phase = TurnPhase.RESPONDING
if (event.responseType == OnlineEvent.ResponseType.TEXT) {
_blocks.add(LiveBlock.Answer(nextKey("answer")))
}
fire()
}
is OnlineEvent.AppendText -> {
when (phase) {
TurnPhase.REASONING -> appendText(event.body, reasoning = true)
// AppendText без предшествующего StartResponse —
// трактуем как ответ (forward-compat).
else -> {
if (phase != TurnPhase.RESPONDING) phase = TurnPhase.RESPONDING
appendText(event.body, reasoning = false)
}
}
}
is OnlineEvent.AppendImage -> {
phase = TurnPhase.RESPONDING
_blocks.add(LiveBlock.Image(nextKey("image"), event.mime, event.body))
fire()
}
}
}
/**
* Дописывает текстовый чанк в последний блок нужного типа. Если
* последний блок — другого типа (например, между двумя кусками текста
@@ -8,6 +8,7 @@ import pw.binom.agentik.desktop.settings.AgentConfig
import pw.binom.agentik.journal.ConversationStore
import pw.binom.agentik.journal.inmemory.InMemoryJournalStore
import pw.binom.agentik.journal.inmemory.InMemoryMutableConversationStore
import pw.binom.agentik.outbox.inmemory.InMemoryOnlineOutbox
import pw.binom.agentik.outbox.inmemory.InMemoryOutboxStore
import pw.binom.agentik.proto.Agent
import pw.binom.agentik.proto.AgentInfo
@@ -98,6 +99,8 @@ private class StubAgent : Agent {
override val outbox = InMemoryOutboxStore(maxMessages = null, ttl = null)
override val onlineOutbox = InMemoryOnlineOutbox()
private val store = InMemoryMutableConversationStore()
override val conversationStore: ConversationStore get() = store
@@ -20,6 +20,9 @@ import pw.binom.agentik.outbox.AgentEvent
import pw.binom.agentik.outbox.CommonEvent
import pw.binom.agentik.outbox.Event
import pw.binom.agentik.outbox.MutableOutboxStore
import pw.binom.agentik.outbox.MutableOnlineOutbox
import pw.binom.agentik.outbox.OnlineEvent
import pw.binom.agentik.outbox.inmemory.InMemoryOnlineOutbox
import pw.binom.agentik.outbox.inmemory.InMemoryOutboxStore
import pw.binom.agentik.proto.Agent
import pw.binom.agentik.proto.AgentInfo
@@ -198,6 +201,8 @@ private class EchoAgent(private val interruptible: Boolean = false) : Agent {
override val outbox = InMemoryOutboxStore(maxMessages = null, ttl = null)
override val onlineOutbox = InMemoryOnlineOutbox()
private val store: MutableConversationStore = InMemoryMutableConversationStore()
override val conversationStore: ConversationStore get() = store
@@ -208,7 +213,7 @@ private class EchoAgent(private val interruptible: Boolean = false) : Agent {
override fun createConversation(temp: Boolean): Conversation {
val id = "conv-${++seq}"
val now = Clock.System.now()
val conv = EchoConversation(id, journal, outbox, store, now, interruptible)
val conv = EchoConversation(id, journal, outbox, onlineOutbox, store, now, interruptible)
conversations[id] = conv
runBlocking {
store.upsert(ConversationRecord(id, null, temp, now, now))
@@ -239,6 +244,7 @@ private class EchoConversation(
override val id: String,
private val journal: MutableJournalStore,
private val outbox: MutableOutboxStore,
private val onlineOutbox: MutableOnlineOutbox,
private val store: MutableConversationStore,
createdAt: Instant,
private val interruptible: Boolean,
@@ -290,8 +296,8 @@ private class EchoConversation(
val text = content.filterIsInstance<Content.Text>().joinToString("") { it.body }
val reply = "эхо: $text"
emit(now, Event.StartResponse(now, Event.ResponseType.TEXT))
emit(now, Event.AppendText(now, reply))
onlineOutbox.tryAppendOnline(id, OnlineEvent.StartResponse(now, OnlineEvent.ResponseType.TEXT))
onlineOutbox.tryAppendOnline(id, OnlineEvent.AppendText(now, reply))
val endAt = Clock.System.now()
journal.append(
@@ -0,0 +1,260 @@
package pw.binom.agentik.desktop.session
import java.util.concurrent.atomic.AtomicInteger
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.cancel
import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.MutableSharedFlow
import kotlinx.coroutines.flow.asSharedFlow
import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.withTimeoutOrNull
import org.junit.jupiter.api.Test
import pw.binom.agentik.desktop.persistence.ConversationMetaRepository
import pw.binom.agentik.journal.Content as JournalContent
import pw.binom.agentik.journal.ConversationRecord
import pw.binom.agentik.journal.JournalStore
import pw.binom.agentik.journal.MessageRecord
import pw.binom.agentik.journal.ksqlite.KsqliteJournalStore
import pw.binom.agentik.outbox.CommonEvent
import pw.binom.agentik.outbox.Event
import pw.binom.agentik.outbox.OutboxStore
import pw.binom.db.ksqlite.SQLiteConnection
import kotlin.test.assertEquals
import kotlin.test.assertTrue
import kotlin.time.Instant
/**
* Тесты [AgentSyncEngine]: фоновый ingest всех диалогов агента. Реальные
* `KsqliteJournalStore` + `ConversationMetaRepository` (in-memory SQLite),
* `JournalStore` / `OutboxStore` — подделки. Сеть не нужна.
*
* `runBlocking` + [await] (не `runTest`): репозиторий под капотом уходит в
* `Dispatchers.Default`, виртуальное время `runTest` его не дожидается.
*/
class AgentSyncEngineTest {
private val convId = "conv-1"
private fun at(ms: Long) = Instant.fromEpochMilliseconds(ms)
private fun user(id: String, ms: Long, body: String) = MessageRecord.UserMessage(
id = id,
conversationId = convId,
content = listOf<JournalContent>(JournalContent.Text(body)),
createdAt = at(ms),
)
private fun assistant(id: String, ms: Long, body: String) = MessageRecord.AssistantMessage(
id = id,
conversationId = convId,
content = listOf<JournalContent>(JournalContent.Text(body)),
createdAt = at(ms),
)
private fun rec(updatedAt: Instant) = ConversationRecord(
id = convId,
title = null,
isTemporal = false,
createdAt = at(0),
updatedAt = updatedAt,
)
private suspend fun await(timeoutMs: Long = 5_000, condition: suspend () -> Boolean) {
val ok = withTimeoutOrNull(timeoutMs) {
while (!condition()) delay(5)
true
}
check(ok == true) { "await: условие не выполнено за ${timeoutMs}мс" }
}
private class Fixture {
val connection: SQLiteConnection = SQLiteConnection.memory()
val cache = KsqliteJournalStore(connection)
val meta = ConversationMetaRepository(connection)
val outbox = EngineFakeOutboxStore()
val remote = EngineFakeJournalStore()
val scope = CoroutineScope(SupervisorJob() + Dispatchers.Default)
val engine = AgentSyncEngine(
remoteJournal = remote,
journalCache = cache,
meta = meta,
outbox = outbox,
locks = ConversationLocks(),
scope = scope,
)
fun close() {
engine.close()
scope.cancel()
cache.close()
meta.close()
connection.close()
}
}
@Test
fun `start backfills never-opened conversation into cache and preview`() = runBlocking {
val f = Fixture()
try {
f.remote.records = listOf(user("u1", 1, "привет"), assistant("a1", 2, "здравствуй"))
val changes = AtomicInteger(0)
f.engine.start(
conversations = { listOf(rec(updatedAt = at(2))) },
onChanged = { changes.incrementAndGet() },
)
await { f.cache.count(convId) == 2L }
await { f.meta.meta(convId)?.preview == "здравствуй" }
await { changes.get() >= 1 }
} finally {
f.close()
}
}
@Test
fun `start skips conversation already synced past its updatedAt`() = runBlocking {
val f = Fixture()
try {
f.meta.setSyncedAt(convId, at(100))
f.remote.records = listOf(user("u1", 1, "не должно попасть в кэш"))
f.engine.start(
conversations = { listOf(rec(updatedAt = at(50))) },
onChanged = {},
)
delay(150)
assertEquals(0L, f.cache.count(convId))
} finally {
f.close()
}
}
@Test
fun `start skips conversation held by an open session`() = runBlocking {
val f = Fixture()
try {
f.engine.openConversationIds = setOf(convId)
f.remote.records = listOf(user("u1", 1, "привет"))
f.engine.start(
conversations = { listOf(rec(updatedAt = at(2))) },
onChanged = {},
)
delay(150)
assertEquals(0L, f.cache.count(convId))
} finally {
f.close()
}
}
@Test
fun `End event for non-open conversation backfills journal and preview`() = runBlocking {
val f = Fixture()
try {
f.engine.start(conversations = { emptyList() }, onChanged = {})
f.outbox.awaitSubscriber()
f.remote.records = listOf(user("u1", 1, "привет"), assistant("a1", 2, "пока"))
f.outbox.emit(Event.End(date = at(3)))
await { f.cache.count(convId) == 2L }
await { f.meta.meta(convId)?.preview == "пока" }
} finally {
f.close()
}
}
@Test
fun `End event for open conversation is ignored by the engine`() = runBlocking {
val f = Fixture()
try {
f.engine.openConversationIds = setOf(convId)
f.engine.start(conversations = { emptyList() }, onChanged = {})
f.outbox.awaitSubscriber()
f.remote.records = listOf(user("u1", 1, "привет"))
f.outbox.emit(Event.End(date = at(3)))
delay(150)
assertEquals(0L, f.cache.count(convId))
} finally {
f.close()
}
}
@Test
fun `second start is a no-op`() = runBlocking {
val f = Fixture()
try {
f.remote.records = listOf(user("u1", 1, "привет"))
val changes = AtomicInteger(0)
f.engine.start(
conversations = { listOf(rec(updatedAt = at(1))) },
onChanged = { changes.incrementAndGet() },
)
await { f.cache.count(convId) == 1L }
await { changes.get() >= 1 }
val afterFirst = changes.get()
f.engine.start(
conversations = { listOf(rec(updatedAt = at(1))) },
onChanged = { changes.incrementAndGet() },
)
delay(100)
assertEquals(afterFirst, changes.get())
assertEquals(1L, f.cache.count(convId))
assertTrue(afterFirst >= 1)
} finally {
f.close()
}
}
}
private class EngineFakeJournalStore(
var records: List<MessageRecord> = emptyList(),
) : JournalStore {
override suspend fun list(
conversationId: String,
after: Instant,
offset: Int,
limit: Int,
): List<MessageRecord> = records
.filter { it.conversationId == conversationId && it.createdAt > after }
.sortedWith(compareBy({ it.createdAt }, { it.id }))
.drop(offset)
.take(limit)
override suspend fun count(conversationId: String): Long =
records.count { it.conversationId == conversationId }.toLong()
override suspend fun count(conversationId: String, after: Instant): Long =
records.count { it.conversationId == conversationId && it.createdAt > after }.toLong()
override fun close() {}
}
private class EngineFakeOutboxStore : OutboxStore {
private val shared = MutableSharedFlow<CommonEvent>(replay = 256, extraBufferCapacity = 256)
override fun events(after: Instant?): Flow<CommonEvent> = shared.asSharedFlow()
override suspend fun earliestEventDate(): Instant = Instant.DISTANT_PAST
override fun close() {}
suspend fun awaitSubscriber() {
withTimeoutOrNull(5_000) {
while (shared.subscriptionCount.value == 0) delay(5)
}
}
suspend fun emit(event: Event, conversationId: String = "conv-1") {
shared.emit(CommonEvent.Conversation(date = event.date, conversationId = conversationId, event = event))
}
}
@@ -19,6 +19,8 @@ import pw.binom.agentik.journal.MessageRecord
import pw.binom.agentik.journal.ksqlite.KsqliteJournalStore
import pw.binom.agentik.outbox.CommonEvent
import pw.binom.agentik.outbox.Event
import pw.binom.agentik.outbox.OnlineEvent
import pw.binom.agentik.outbox.OnlineOutbox
import pw.binom.agentik.outbox.OutboxStore
import pw.binom.agentik.proto.Content
import pw.binom.agentik.proto.Conversation
@@ -74,6 +76,7 @@ class ChatSessionTest {
val cache = KsqliteJournalStore(connection)
val meta = ConversationMetaRepository(connection)
val outbox = FakeOutboxStore()
val online = FakeOnlineOutbox()
var remote = FakeJournalStore()
val conversation = FakeConversation("conv-1")
val scope = CoroutineScope(SupervisorJob() + Dispatchers.Default)
@@ -85,6 +88,7 @@ class ChatSessionTest {
journalCache = cache,
meta = meta,
outbox = outbox,
onlineOutbox = online,
scope = scope,
)
@@ -157,11 +161,12 @@ class ChatSessionTest {
try {
session.start()
f.outbox.awaitSubscriber()
f.online.awaitSubscriber()
f.outbox.emit(Event.Working(at(10)))
f.outbox.emit(Event.StartResponse(at(11), Event.ResponseType.TEXT))
f.outbox.emit(Event.AppendText(at(12), "Hel"))
f.outbox.emit(Event.AppendText(at(13), "lo"))
f.online.emit(OnlineEvent.StartResponse(at(11), OnlineEvent.ResponseType.TEXT))
f.online.emit(OnlineEvent.AppendText(at(12), "Hel"))
f.online.emit(OnlineEvent.AppendText(at(13), "lo"))
await { session.live.value.size == 1 }
assertTrue(session.live.value.single() is UiMessage.Assistant)
@@ -216,9 +221,10 @@ class ChatSessionTest {
try {
session.start()
f.outbox.awaitSubscriber()
f.online.awaitSubscriber()
f.outbox.emit(Event.StartResponse(at(10), Event.ResponseType.TEXT))
f.outbox.emit(Event.AppendText(at(11), "частичный"))
f.online.emit(OnlineEvent.StartResponse(at(10), OnlineEvent.ResponseType.TEXT))
f.online.emit(OnlineEvent.AppendText(at(11), "частичный"))
await { session.live.value.size == 1 }
f.outbox.emit(Event.Interrupted(at(12)))
@@ -296,6 +302,25 @@ private class FakeOutboxStore : OutboxStore {
}
}
private class FakeOnlineOutbox : OnlineOutbox {
private val shared = MutableSharedFlow<OnlineEvent>(replay = 256, extraBufferCapacity = 256)
override fun onlineEvents(conversationId: String): Flow<OnlineEvent> = shared.asSharedFlow()
override fun close() {}
/** Ждёт, пока у потока появится подписчик (ChatSession подписался). */
suspend fun awaitSubscriber() {
withTimeoutOrNull(5_000) {
while (shared.subscriptionCount.value == 0) delay(5)
}
}
suspend fun emit(event: OnlineEvent) {
shared.emit(event)
}
}
private class FakeConversation(
override val id: String,
override val isSupportImageInput: Boolean = true,
@@ -5,6 +5,7 @@ import kotlin.test.assertEquals
import kotlin.test.assertIs
import kotlin.test.assertTrue
import pw.binom.agentik.outbox.Event
import pw.binom.agentik.outbox.OnlineEvent
import kotlin.time.Clock
import kotlin.time.Instant
@@ -13,8 +14,14 @@ class LiveTurnTest {
private val now = Clock.System.now()
private fun at() = Instant.fromEpochMilliseconds(now.toEpochMilliseconds())
private fun LiveTurn.feed(vararg events: Event) {
events.forEach { apply(it) }
private fun LiveTurn.feed(vararg events: Any) {
events.forEach {
when (it) {
is Event -> apply(it)
is OnlineEvent -> applyOnline(it)
else -> error("unexpected event: $it")
}
}
}
@Test
@@ -22,12 +29,12 @@ class LiveTurnTest {
val turn = LiveTurn()
turn.feed(
Event.Working(at()),
Event.StartReasoning(at()),
Event.AppendText(at(), "сначала "),
Event.AppendText(at(), "подумаем"),
Event.StartResponse(at(), Event.ResponseType.TEXT),
Event.AppendText(at(), "Привет, "),
Event.AppendText(at(), "мир!"),
OnlineEvent.StartReasoning(at()),
OnlineEvent.AppendText(at(), "сначала "),
OnlineEvent.AppendText(at(), "подумаем"),
OnlineEvent.StartResponse(at(), OnlineEvent.ResponseType.TEXT),
OnlineEvent.AppendText(at(), "Привет, "),
OnlineEvent.AppendText(at(), "мир!"),
Event.End(at()),
)
@@ -46,7 +53,7 @@ class LiveTurnTest {
val turn = LiveTurn()
turn.feed(
Event.Working(at()),
Event.StartResponse(at(), Event.ResponseType.TEXT),
OnlineEvent.StartResponse(at(), OnlineEvent.ResponseType.TEXT),
Event.ToolCall(at(), id = "t1", title = "grep", toolName = "grep", toolArgs = "{\"q\":\"x\"}"),
Event.ToolResult(at(), toolCallId = "t1", toolName = "grep", result = "found 3"),
Event.End(at()),
@@ -66,7 +73,7 @@ class LiveTurnTest {
val turn = LiveTurn()
turn.feed(
Event.Working(at()),
Event.StartResponse(at(), Event.ResponseType.TEXT),
OnlineEvent.StartResponse(at(), OnlineEvent.ResponseType.TEXT),
Event.ToolResult(at(), toolCallId = "missing", toolName = null, result = "ok"),
)
@@ -94,8 +101,8 @@ class LiveTurnTest {
val turn = LiveTurn()
turn.feed(
Event.Working(at()),
Event.StartResponse(at(), Event.ResponseType.TEXT),
Event.AppendText(at(), "partial"),
OnlineEvent.StartResponse(at(), OnlineEvent.ResponseType.TEXT),
OnlineEvent.AppendText(at(), "partial"),
Event.Error(at(), message = "llm down", code = "E500"),
)
@@ -112,8 +119,8 @@ class LiveTurnTest {
val turn = LiveTurn()
turn.feed(
Event.Working(at()),
Event.StartResponse(at(), Event.ResponseType.TEXT),
Event.AppendText(at(), "не до конца"),
OnlineEvent.StartResponse(at(), OnlineEvent.ResponseType.TEXT),
OnlineEvent.AppendText(at(), "не до конца"),
Event.Interrupted(at()),
)
@@ -127,10 +134,10 @@ class LiveTurnTest {
val turn = LiveTurn()
turn.feed(
Event.Working(at()),
Event.StartResponse(at(), Event.ResponseType.TEXT),
Event.AppendText(at(), "до"),
OnlineEvent.StartResponse(at(), OnlineEvent.ResponseType.TEXT),
OnlineEvent.AppendText(at(), "до"),
Event.ToolCall(at(), id = "t1", title = null, toolName = "grep", toolArgs = "{}"),
Event.AppendText(at(), "после"),
OnlineEvent.AppendText(at(), "после"),
)
val answers = turn.blocks.filterIsInstance<LiveBlock.Answer>()
@@ -146,8 +153,8 @@ class LiveTurnTest {
val turn = LiveTurn()
turn.feed(
Event.Working(at()),
Event.StartResponse(at(), Event.ResponseType.IMAGE),
Event.AppendImage(at(), body = bytes, mime = "image/png"),
OnlineEvent.StartResponse(at(), OnlineEvent.ResponseType.IMAGE),
OnlineEvent.AppendImage(at(), body = bytes, mime = "image/png"),
Event.End(at()),
)
@@ -162,11 +169,11 @@ class LiveTurnTest {
val turn = LiveTurn()
turn.feed(
Event.Working(at()),
Event.StartResponse(at(), Event.ResponseType.TEXT),
Event.AppendText(at(), "старое"),
OnlineEvent.StartResponse(at(), OnlineEvent.ResponseType.TEXT),
OnlineEvent.AppendText(at(), "старое"),
Event.End(at()),
)
turn.feed(Event.Working(at()), Event.AppendText(at(), "новое"))
turn.feed(Event.Working(at()), OnlineEvent.AppendText(at(), "новое"))
assertEquals(1, turn.blocks.size)
assertEquals("новое", turn.answerText())
@@ -175,7 +182,7 @@ class LiveTurnTest {
@Test
fun `append without explicit start response still streams as answer`() {
val turn = LiveTurn()
turn.feed(Event.Working(at()), Event.AppendText(at(), "без старта"))
turn.feed(Event.Working(at()), OnlineEvent.AppendText(at(), "без старта"))
assertEquals(TurnPhase.RESPONDING, turn.phase)
assertEquals("без старта", turn.answerText())
@@ -186,7 +193,7 @@ class LiveTurnTest {
val turn = LiveTurn()
var count = 0
turn.onChange = { count++ }
turn.feed(Event.Working(at()), Event.AppendText(at(), "x"))
turn.feed(Event.Working(at()), OnlineEvent.AppendText(at(), "x"))
assertEquals(2, count)
}
@@ -196,12 +203,12 @@ class LiveTurnTest {
val turn = LiveTurn()
turn.feed(
Event.Working(at()),
Event.AppendText(at(), "Первая строка\n\nвторая строка "),
OnlineEvent.AppendText(at(), "Первая строка\n\nвторая строка "),
)
assertEquals("Первая строка вторая строка", turn.previewText())
val long = LiveTurn()
long.feed(Event.Working(at()), Event.AppendText(at(), "слово ".repeat(60)))
long.feed(Event.Working(at()), OnlineEvent.AppendText(at(), "слово ".repeat(60)))
val preview = long.previewText(maxLength = 30)
assertTrue(preview.length <= 31)
assertTrue(preview.endsWith("…"))