From 26f94ea83d5fb09c57f5601937dbfe5038325f1f Mon Sep 17 00:00:00 2001 From: subochev Date: Wed, 30 Sep 2026 11:43:49 +0300 Subject: [PATCH] =?UTF-8?q?=D0=98=D0=B7=D0=BC=D0=B5=D0=BD=D0=B5=D0=BD?= =?UTF-8?q?=D0=B8=D1=8F=20=D0=B2=D0=B2=D0=BE=D0=B4=D1=8F=D1=82=20live-?= =?UTF-8?q?=D0=BA=D0=B0=D0=BD=D0=B0=D0=BB=20OnlineOutbox=20=D0=B8=20=D1=84?= =?UTF-8?q?=D0=BE=D0=BD=D0=BE=D0=B2=D1=83=D1=8E=20=D1=81=D0=B8=D0=BD=D1=85?= =?UTF-8?q?=D1=80=D0=BE=D0=BD=D0=B8=D0=B7=D0=B0=D1=86=D0=B8=D1=8E=20=D0=B4?= =?UTF-8?q?=D0=B8=D0=B0=D0=BB=D0=BE=D0=B3=D0=BE=D0=B2.?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- REQUIREMENTS.md | 26 +- .../persistence/ConversationMetaRepository.kt | 2 +- .../desktop/session/AgentConnection.kt | 37 +++ .../desktop/session/AgentSyncEngine.kt | 110 ++++++++ .../agentik/desktop/session/ChatSession.kt | 76 ++--- .../desktop/session/ConversationSync.kt | 101 +++++++ .../binom/agentik/desktop/session/LiveTurn.kt | 87 +++--- .../desktop/net/ConnectionCheckerTest.kt | 3 + .../desktop/session/AgentConnectionE2ETest.kt | 12 +- .../desktop/session/AgentSyncEngineTest.kt | 260 ++++++++++++++++++ .../desktop/session/ChatSessionTest.kt | 35 ++- .../agentik/desktop/session/LiveTurnTest.kt | 59 ++-- 12 files changed, 699 insertions(+), 109 deletions(-) create mode 100644 src/jvmMain/kotlin/pw/binom/agentik/desktop/session/AgentSyncEngine.kt create mode 100644 src/jvmMain/kotlin/pw/binom/agentik/desktop/session/ConversationSync.kt create mode 100644 src/jvmTest/kotlin/pw/binom/agentik/desktop/session/AgentSyncEngineTest.kt diff --git a/REQUIREMENTS.md b/REQUIREMENTS.md index 572b1db..b6a7584 100644 --- a/REQUIREMENTS.md +++ b/REQUIREMENTS.md @@ -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 понятие «домашний каталог» своё. diff --git a/src/jvmMain/kotlin/pw/binom/agentik/desktop/persistence/ConversationMetaRepository.kt b/src/jvmMain/kotlin/pw/binom/agentik/desktop/persistence/ConversationMetaRepository.kt index c3ec149..5460304 100644 --- a/src/jvmMain/kotlin/pw/binom/agentik/desktop/persistence/ConversationMetaRepository.kt +++ b/src/jvmMain/kotlin/pw/binom/agentik/desktop/persistence/ConversationMetaRepository.kt @@ -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`) и разные контракты diff --git a/src/jvmMain/kotlin/pw/binom/agentik/desktop/session/AgentConnection.kt b/src/jvmMain/kotlin/pw/binom/agentik/desktop/session/AgentConnection.kt index abc7bab..7aefd04 100644 --- a/src/jvmMain/kotlin/pw/binom/agentik/desktop/session/AgentConnection.kt +++ b/src/jvmMain/kotlin/pw/binom/agentik/desktop/session/AgentConnection.kt @@ -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() + 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() } diff --git a/src/jvmMain/kotlin/pw/binom/agentik/desktop/session/AgentSyncEngine.kt b/src/jvmMain/kotlin/pw/binom/agentik/desktop/session/AgentSyncEngine.kt new file mode 100644 index 0000000..5518461 --- /dev/null +++ b/src/jvmMain/kotlin/pw/binom/agentik/desktop/session/AgentSyncEngine.kt @@ -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 = emptySet() + + private var job: Job? = null + + /** + * Запустить ingest. Идемпотентно. + * + * @param conversations список диалогов агента на момент вызова (для + * стартового догона). + * @param onChanged дёргается, когда локальный кэш/превью изменились — + * чтобы UI перечитал список. + */ + fun start( + conversations: suspend () -> List, + 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 + } +} diff --git a/src/jvmMain/kotlin/pw/binom/agentik/desktop/session/ChatSession.kt b/src/jvmMain/kotlin/pw/binom/agentik/desktop/session/ChatSession.kt index 983ee20..f27c9c0 100644 --- a/src/jvmMain/kotlin/pw/binom/agentik/desktop/session/ChatSession.kt +++ b/src/jvmMain/kotlin/pw/binom/agentik/desktop/session/ChatSession.kt @@ -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>(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() { diff --git a/src/jvmMain/kotlin/pw/binom/agentik/desktop/session/ConversationSync.kt b/src/jvmMain/kotlin/pw/binom/agentik/desktop/session/ConversationSync.kt new file mode 100644 index 0000000..f8f5ca1 --- /dev/null +++ b/src/jvmMain/kotlin/pw/binom/agentik/desktop/session/ConversationSync.kt @@ -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() + + 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) +} diff --git a/src/jvmMain/kotlin/pw/binom/agentik/desktop/session/LiveTurn.kt b/src/jvmMain/kotlin/pw/binom/agentik/desktop/session/LiveTurn.kt index 49a824d..efcd31a 100644 --- a/src/jvmMain/kotlin/pw/binom/agentik/desktop/session/LiveTurn.kt +++ b/src/jvmMain/kotlin/pw/binom/agentik/desktop/session/LiveTurn.kt @@ -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() + } + } + } + /** * Дописывает текстовый чанк в последний блок нужного типа. Если * последний блок — другого типа (например, между двумя кусками текста diff --git a/src/jvmTest/kotlin/pw/binom/agentik/desktop/net/ConnectionCheckerTest.kt b/src/jvmTest/kotlin/pw/binom/agentik/desktop/net/ConnectionCheckerTest.kt index 2de7953..34ceebc 100644 --- a/src/jvmTest/kotlin/pw/binom/agentik/desktop/net/ConnectionCheckerTest.kt +++ b/src/jvmTest/kotlin/pw/binom/agentik/desktop/net/ConnectionCheckerTest.kt @@ -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 diff --git a/src/jvmTest/kotlin/pw/binom/agentik/desktop/session/AgentConnectionE2ETest.kt b/src/jvmTest/kotlin/pw/binom/agentik/desktop/session/AgentConnectionE2ETest.kt index 9557ca0..03aafe0 100644 --- a/src/jvmTest/kotlin/pw/binom/agentik/desktop/session/AgentConnectionE2ETest.kt +++ b/src/jvmTest/kotlin/pw/binom/agentik/desktop/session/AgentConnectionE2ETest.kt @@ -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().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( diff --git a/src/jvmTest/kotlin/pw/binom/agentik/desktop/session/AgentSyncEngineTest.kt b/src/jvmTest/kotlin/pw/binom/agentik/desktop/session/AgentSyncEngineTest.kt new file mode 100644 index 0000000..2ec98ec --- /dev/null +++ b/src/jvmTest/kotlin/pw/binom/agentik/desktop/session/AgentSyncEngineTest.kt @@ -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.Text(body)), + createdAt = at(ms), + ) + + private fun assistant(id: String, ms: Long, body: String) = MessageRecord.AssistantMessage( + id = id, + conversationId = convId, + content = listOf(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 = emptyList(), +) : JournalStore { + override suspend fun list( + conversationId: String, + after: Instant, + offset: Int, + limit: Int, + ): List = 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(replay = 256, extraBufferCapacity = 256) + + override fun events(after: Instant?): Flow = 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)) + } +} diff --git a/src/jvmTest/kotlin/pw/binom/agentik/desktop/session/ChatSessionTest.kt b/src/jvmTest/kotlin/pw/binom/agentik/desktop/session/ChatSessionTest.kt index e0bf75f..29b6e8e 100644 --- a/src/jvmTest/kotlin/pw/binom/agentik/desktop/session/ChatSessionTest.kt +++ b/src/jvmTest/kotlin/pw/binom/agentik/desktop/session/ChatSessionTest.kt @@ -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(replay = 256, extraBufferCapacity = 256) + + override fun onlineEvents(conversationId: String): Flow = 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, diff --git a/src/jvmTest/kotlin/pw/binom/agentik/desktop/session/LiveTurnTest.kt b/src/jvmTest/kotlin/pw/binom/agentik/desktop/session/LiveTurnTest.kt index 238cd09..9ad3ef8 100644 --- a/src/jvmTest/kotlin/pw/binom/agentik/desktop/session/LiveTurnTest.kt +++ b/src/jvmTest/kotlin/pw/binom/agentik/desktop/session/LiveTurnTest.kt @@ -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() @@ -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("…"))