Десктоп-клиент: Compose-UI, сессии, SQLite-кэш, голос, живые тесты

- settings/: settings.json (агенты, тема, clientId), атомарная запись
- persistence/: ConversationMetaRepository — группы, lastSeen, превью,
  watermark синхронизации (synced_at через PRAGMA table_info)
- session/: AgentRegistry (параллельное подключение, деградация в офлайн),
  AgentConnection (журнал + мета + живой список диалогов conversationsFlow,
  refresh по SSE-потоку), ChatSession (history + backfill + live + preview,
  дедуп по id, 1 мс шаг watermark, syncMutex), LiveTurn
- UI: рейка агентов, список диалогов (живое обновление), чат со стримингом,
  композер (Enter отправляет, Shift+Enter переносит, Ctrl+M — голос),
  окна «Агенты»/«Группы»/«Новый диалог», тёмная и светлая темы
- voice/: mic-api → Silero VAD → Vosk на одном выделенном потоке;
  VoiceNative.warmUp() до skia (иначе SIGSEGV), модель в артефакте
- markdown/: вендоренный рендер из ai/assistent
- ошибки в баннер через AppState.reportError (CancellationException игнорируется)
- тесты 103 (реальная SQLite + фейковый транспорт); против живого агента —
  RealServerIT, RealAppStateIT, RealSpeechIT
- LIVE-RUN-NOTES.md: живой прогон в GUI и найденные/починенные баги
This commit is contained in:
2026-09-28 10:08:10 +03:00
parent 27f349550c
commit 2c71542b13
53 changed files with 9147 additions and 0 deletions
@@ -0,0 +1,325 @@
package pw.binom.agentik.desktop.session
import io.ktor.server.cio.CIO as ServerCIO
import io.ktor.server.engine.embeddedServer
import io.ktor.server.routing.routing
import kotlinx.coroutines.delay
import kotlinx.coroutines.launch
import kotlinx.coroutines.runBlocking
import pw.binom.agentik.desktop.model.UiContent
import pw.binom.agentik.desktop.model.UiMessage
import pw.binom.agentik.desktop.settings.AgentConfig
import pw.binom.agentik.journal.ConversationRecord
import pw.binom.agentik.journal.ConversationStore
import pw.binom.agentik.journal.MessageRecord
import pw.binom.agentik.journal.MutableConversationStore
import pw.binom.agentik.journal.MutableJournalStore
import pw.binom.agentik.journal.inmemory.InMemoryJournalStore
import pw.binom.agentik.journal.inmemory.InMemoryMutableConversationStore
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.inmemory.InMemoryOutboxStore
import pw.binom.agentik.proto.Agent
import pw.binom.agentik.proto.AgentInfo
import pw.binom.agentik.proto.Content
import pw.binom.agentik.proto.Conversation
import pw.binom.agentik.proto.Message
import pw.binom.agentik.proto.MessageContext
import pw.binom.agentik.server.agentikAgent
import kotlin.test.AfterTest
import kotlin.test.Test
import io.ktor.client.request.get
import io.ktor.client.statement.bodyAsText
import kotlin.test.assertEquals
import kotlin.test.assertNotNull
import kotlin.test.assertTrue
import kotlin.time.Clock
import kotlin.time.Duration.Companion.milliseconds
import kotlin.time.Instant
/**
* Сквозной тест: настоящий `:server` (HTTP/JSON/SSE) поверх in-memory-агента,
* а с другой стороны — наш `AgentConnection` + `ChatSession`. Проверяем, что
* десктопный стек работает по живому протоколу: создать диалог → отправить →
* увидеть стриминг → после `End` история в локальном SQLite → повторное
* открытие не плодит дублей.
*
* Никаких моков протокола: `AgentikAgent` (HTTP) ↔ `embeddedServer` ↔
* `agentikAgent(fake)`. Именно это ловит расхождения JSON/SSE-контракта,
* которые юнит-тесты с ручными фейками не видят.
*/
class AgentConnectionE2ETest {
private lateinit var server: io.ktor.server.engine.EmbeddedServer<*, *>
private var port: Int = 0
@AfterTest
fun tearDown() {
if (::server.isInitialized) server.stop(100, 300)
}
private fun startServer(agent: Agent) {
server = embeddedServer(ServerCIO, port = 0, host = "127.0.0.1") {
routing { agentikAgent(agent, path = "/agentik", token = null) }
}.start(wait = false)
port = runBlocking { server.engine.resolvedConnectors().first().port }
}
@Test
fun `send streams through real server and persists into local sqlite`() = runBlocking {
val fake = EchoAgent()
startServer(fake)
val conn = AgentConnection.open(
config = AgentConfig.create(baseUrl = "http://127.0.0.1:$port/agentik"),
clientId = "e2e-client",
scope = this,
)
try {
val conv = assertNotNull(conn.createConversation(temp = false), "createConversation")
val session = assertNotNull(conn.openSession(conv.id), "openSession")
try {
session.start()
val beforeTurn = Clock.System.now() - 50.milliseconds
session.send("привет")
await(timeoutMs = 15_000) {
session.history.value.size >= 2 && !session.isStreaming
}
val history = session.history.value
assertEquals(2, history.size, "user + assistant: $history")
val user = history[0] as UiMessage.User
assertEquals("привет", (user.content.single() as UiContent.Text).body)
val bot = history[1] as UiMessage.Assistant
assertEquals("эхо: привет", (bot.content.single() as UiContent.Text).body)
// Локальный кэш — на диске (agentik.db), а не в HTTP.
val cached = conn.journalCache.list(conv.id, Instant.DISTANT_PAST, 0, Int.MAX_VALUE)
assertEquals(2, cached.size, "cache records")
// Превью денормализовано из ответа ассистента.
assertEquals("эхо: привет", conn.meta.meta(conv.id)?.preview)
assertNotNull(conn.meta.meta(conv.id)?.syncedAt)
// Счётчик непрочитанных считается сервером от lastSeen.
// Сначала дожидаемся, пока `:client`-кэш списка диалогов
// подтянет созданный диалог (seed/live-refresh асинхронны).
await(timeoutMs = 10_000) {
conn.conversationStore.list(offset = 0, limit = 500).any { it.id == conv.id }
}
conn.meta.markSeen(conv.id, beforeTurn)
val ui = conn.uiConversations()
val row = ui.first { it.id == conv.id }
assertEquals(2, row.unread, "unread since beforeTurn")
// Повторное открытие читает из кэша и НЕ дублирует бэкфиллом.
val again = assertNotNull(conn.openSession(conv.id))
try {
again.start()
assertEquals(2, again.history.value.size, "reopen must not duplicate")
} finally {
again.close()
}
} finally {
session.close()
}
} finally {
conn.close()
}
}
@Test
fun `interrupt emits interrupted and commits nothing`() = runBlocking {
val fake = EchoAgent(interruptible = true)
startServer(fake)
val conn = AgentConnection.open(
config = AgentConfig.create(baseUrl = "http://127.0.0.1:$port/agentik"),
clientId = "e2e-client",
scope = this,
)
try {
val conv = assertNotNull(conn.createConversation(temp = false))
val session = assertNotNull(conn.openSession(conv.id))
try {
session.start()
val sendJob = launch { session.send("прервись") }
// Ждём, пока сервер реально вошёл в send (user-сообщение в журнале).
await(timeoutMs = 10_000) { fake.journal.count(conv.id) >= 1 }
session.interrupt()
await(timeoutMs = 10_000) { !session.isStreaming }
sendJob.join()
assertTrue(
session.history.value.none { it is UiMessage.Assistant },
"interrupted turn must not commit an assistant message",
)
} finally {
session.close()
}
} finally {
conn.close()
}
}
private suspend fun await(timeoutMs: Long = 5_000, condition: suspend () -> Boolean) {
val deadline = System.currentTimeMillis() + timeoutMs
while (System.currentTimeMillis() < deadline) {
if (condition()) return
delay(25)
}
assertTrue(condition(), "condition not met within ${timeoutMs}ms")
}
}
/**
* Эхо-агент: на любой `send` отвечает `эхо: <текст>`, эмитит полный набор
* live-событий и пишет пару user/assistant в журнал — ровно то, что делает
* настоящий ChatAgent, но без LLM.
*
* @param interruptible если true — между `Working` и ответом ждёт
* `Interrupted`; ход закрывается прерыванием, ответ НЕ пишется.
*/
private class EchoAgent(private val interruptible: Boolean = false) : Agent {
override val id: String = "e2e-echo"
override val info: AgentInfo = AgentInfo(
name = "Эхо",
description = "Тестовый эхо-агент для сквозного теста",
)
override val journal: MutableJournalStore = InMemoryJournalStore()
override val outbox = InMemoryOutboxStore(maxMessages = null, ttl = null)
private val store: MutableConversationStore = InMemoryMutableConversationStore()
override val conversationStore: ConversationStore get() = store
private val conversations = mutableMapOf<String, EchoConversation>()
private var seq = 0
override fun createConversation(temp: Boolean): Conversation {
val id = "conv-${++seq}"
val now = Clock.System.now()
val conv = EchoConversation(id, journal, outbox, store, now, interruptible)
conversations[id] = conv
runBlocking {
store.upsert(ConversationRecord(id, null, temp, now, now))
outbox.append(CommonEvent.Agent(now, AgentEvent.Created(now, id)))
}
return conv
}
override suspend fun getConversation(id: String): Conversation? = conversations[id]
override suspend fun deleteConversation(id: String): Boolean {
val conv = conversations.remove(id) ?: return false
conv.close()
store.delete(id)
journal.clear(id)
val now = Clock.System.now()
outbox.append(CommonEvent.Agent(now, AgentEvent.Deleted(now, id)))
return true
}
override suspend fun renameConversation(id: String, title: String?): Instant? =
store.rename(id, title)
override fun close() {}
}
private class EchoConversation(
override val id: String,
private val journal: MutableJournalStore,
private val outbox: MutableOutboxStore,
private val store: MutableConversationStore,
createdAt: Instant,
private val interruptible: Boolean,
) : Conversation {
override val isSupportImageInput: Boolean = false
override val isSupportImageOutput: Boolean = false
override val isTemporal: Boolean = false
override var title: String? = null
private var updated: Instant = createdAt
override val updatedAt: Instant get() = updated
private var msgSeq = 0
/**
* id сообщения с zero-padded счётчиком В НАЧАЛЕ: журнал сортирует по
* `createdAt ASC, id ASC`, а фейк штампует user- и assistant-сообщение
* в одну и ту же миллисекунду — без padding'а лексикографически
* `...-a2` < `...-u1` и ассистент «уезжает» в начало истории.
*/
private fun nextMsgId(convId: String, role: String): String =
"$convId-${(++msgSeq).toString().padStart(6, '0')}-$role"
private val interrupted = java.util.concurrent.atomic.AtomicBoolean(false)
override suspend fun rename(title: String) {
this.title = title
}
override suspend fun interrupt() {
interrupted.set(true)
val now = Clock.System.now()
emit(now, Event.Interrupted(now))
}
override suspend fun getMessages(after: Instant, offset: Int, limit: Int): List<Message> = emptyList()
override suspend fun send(content: List<Content>, context: MessageContext?) {
val now = Clock.System.now()
journal.append(
MessageRecord.UserMessage(nextMsgId(id, "user"), id, content.toJournal(), now, context = null),
)
emit(now, Event.Working(now))
if (interruptible) {
// Имитация долгого LLM-цикла: крутимся, пока не придёт interrupt().
while (!interrupted.get()) delay(20)
return
}
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))
val endAt = Clock.System.now()
journal.append(
MessageRecord.AssistantMessage(
nextMsgId(id, "assistant"),
id,
listOf(pw.binom.agentik.journal.Content.Text(reply)),
endAt,
),
)
emit(endAt, Event.End(endAt))
updated = endAt
store.touch(id, endAt)
outbox.append(CommonEvent.Agent(endAt, AgentEvent.Touched(endAt, id, endAt)))
}
private suspend fun emit(at: Instant, event: Event) {
outbox.append(CommonEvent.Conversation(at, id, event))
}
override fun close() {}
}
/** `:proto.Content` → `:journal.Content` (в проде это делает storage-impl). */
private fun List<Content>.toJournal(): List<pw.binom.agentik.journal.Content> = map { c ->
when (c) {
is Content.Text -> pw.binom.agentik.journal.Content.Text(c.body)
is Content.Image -> pw.binom.agentik.journal.Content.Image(c.data, c.mime)
}
}
@@ -0,0 +1,128 @@
package pw.binom.agentik.desktop.session
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.cancel
import kotlinx.coroutines.runBlocking
import org.junit.jupiter.api.io.TempDir
import pw.binom.agentik.desktop.settings.AgentConfig
import pw.binom.agentik.desktop.settings.SettingsRepository
import pw.binom.agentik.desktop.settings.ThemeId
import java.nio.file.Path
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertNull
import kotlin.test.assertTrue
/**
* Тесты [AgentRegistry]: persist настроек и переходы статусов. Подключение
* подменено [AgentConnector], который всегда падает — так проверяем
* деградацию в [AgentStatus.Offline] без сети. Успешное [AgentStatus.Online]
* потребовало бы настоящего [AgentConnection] (его конструктор приватный),
* поэтому проверяется косвенно — см. подключение к реальному серверу вручную.
*/
class AgentRegistryTest {
@TempDir
lateinit var dir: Path
private fun registry(scope: CoroutineScope) = AgentRegistry(
settingsRepo = SettingsRepository.inDirectory(dir),
scope = scope,
connector = { _, _ -> throw RuntimeException("нет сети") },
)
@Test
fun `add agent persists config and goes offline on failure`() = runBlocking {
val scope = CoroutineScope(SupervisorJob() + Dispatchers.Default)
val reg = registry(scope)
try {
reg.load()
val cfg = AgentConfig.create("http://192.168.1.5:8080/agentik", nameOverride = "Тестовый")
val status = reg.addAgent(cfg)
assertTrue(status is AgentStatus.Offline)
assertEquals(1, reg.settings.value.agents.size)
assertNull(reg.connection(cfg.id))
// persist переживает переоткрытие
val reopened = SettingsRepository.inDirectory(dir).load()
assertEquals(listOf(cfg), reopened.agents)
} finally {
reg.close()
scope.cancel()
}
}
@Test
fun `update agent persists new baseUrl`() = runBlocking {
val scope = CoroutineScope(SupervisorJob() + Dispatchers.Default)
val reg = registry(scope)
try {
reg.load()
val cfg = AgentConfig.create("http://old:8080/agentik")
reg.addAgent(cfg)
reg.updateAgent(cfg.copy(baseUrl = "http://new:9090/agentik"))
assertEquals("http://new:9090/agentik", reg.settings.value.agents.single().baseUrl)
assertEquals("http://new:9090/agentik", SettingsRepository.inDirectory(dir).load().agents.single().baseUrl)
} finally {
reg.close()
scope.cancel()
}
}
@Test
fun `remove agent drops it from settings and statuses`() = runBlocking {
val scope = CoroutineScope(SupervisorJob() + Dispatchers.Default)
val reg = registry(scope)
try {
reg.load()
val cfg = AgentConfig.create("http://host:8080/agentik")
reg.addAgent(cfg)
reg.removeAgent(cfg.id)
assertTrue(reg.settings.value.agents.isEmpty())
assertNull(reg.statuses.value[cfg.id])
assertTrue(SettingsRepository.inDirectory(dir).load().agents.isEmpty())
} finally {
reg.close()
scope.cancel()
}
}
@Test
fun `updateSettings persists theme`() = runBlocking {
val scope = CoroutineScope(SupervisorJob() + Dispatchers.Default)
val reg = registry(scope)
try {
reg.load()
reg.updateSettings { it.copy(themeId = ThemeId.LIGHT) }
assertEquals(ThemeId.LIGHT, reg.settings.value.themeId)
assertEquals(ThemeId.LIGHT, SettingsRepository.inDirectory(dir).load().themeId)
} finally {
reg.close()
scope.cancel()
}
}
@Test
fun `update nonexistent agent is offline and does not add it`() = runBlocking {
val scope = CoroutineScope(SupervisorJob() + Dispatchers.Default)
val reg = registry(scope)
try {
reg.load()
val status = reg.updateAgent(AgentConfig.create("http://host:8080/agentik"))
assertTrue(status is AgentStatus.Offline)
assertTrue(reg.settings.value.agents.isEmpty())
} finally {
reg.close()
scope.cancel()
}
}
}
@@ -0,0 +1,325 @@
package pw.binom.agentik.desktop.session
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.model.UiMessage
import pw.binom.agentik.desktop.persistence.ConversationMetaRepository
import pw.binom.agentik.journal.Content as JournalContent
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.agentik.proto.Content
import pw.binom.agentik.proto.Conversation
import pw.binom.agentik.proto.Message
import pw.binom.agentik.proto.MessageContext
import pw.binom.db.ksqlite.SQLiteConnection
import kotlin.test.assertEquals
import kotlin.test.assertTrue
import kotlin.time.Instant
/**
* Тесты [ChatSession] на фейковом транспорте: реальные `KsqliteJournalStore`
* и `ConversationMetaRepository` (in-memory SQLite), но `JournalStore` /
* `OutboxStore` / `Conversation` — подделки. Проверяют склейку
* «бэкфилл → кэш → live-ход → End → история» без сети.
*
* Используют `runBlocking` (как `ConversationMetaRepositoryIT`), а не
* `runTest`: репозитории под капотом уходят в `Dispatchers.Default`, и
* виртуальное время `runTest` их не дожидается. Вместо этого — [await]
* с реальным `delay`.
*/
class ChatSessionTest {
private val convId = "conv-1"
private fun at(ms: Long) = Instant.fromEpochMilliseconds(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 user(id: String, ms: Long, body: String) = MessageRecord.UserMessage(
id = id,
conversationId = convId,
content = listOf<JournalContent>(JournalContent.Text(body)),
createdAt = at(ms),
)
/** Ждёт выполнения условия или падает по таймауту. */
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 = FakeOutboxStore()
var remote = FakeJournalStore()
val conversation = FakeConversation("conv-1")
val scope = CoroutineScope(SupervisorJob() + Dispatchers.Default)
fun session() = ChatSession(
conversationId = "conv-1",
conversation = conversation,
remoteJournal = remote,
journalCache = cache,
meta = meta,
outbox = outbox,
scope = scope,
)
fun close() {
scope.cancel()
cache.close()
meta.close()
connection.close()
}
}
@Test
fun `start backfills remote journal into local cache`() = runBlocking {
val f = Fixture()
f.remote.records = listOf(user("u1", 1, "привет"), assistant("a1", 2, "здравствуй"))
val session = f.session()
try {
session.start()
await { session.history.value.size == 2 }
assertEquals(2, session.history.value.size)
// watermark = createdAt последней записи минус 1 мс (см. backfill:
// на миллисекундной границе assistant-сообщение делит timestamp с
// user-сообщением, и строгий `createdAt > after` его бы потерял).
assertEquals(at(1), f.meta.meta(convId)?.syncedAt)
} finally {
f.close()
}
}
@Test
fun `second start does not duplicate history`() = runBlocking {
val f = Fixture()
f.remote.records = listOf(user("u1", 1, "привет"), assistant("a1", 2, "здравствуй"))
val session = f.session()
try {
session.start()
await { session.history.value.size == 2 }
session.start()
delay(50)
assertEquals(2, session.history.value.size)
} finally {
f.close()
}
}
@Test
fun `legacy cache without watermark tolerates duplicate backfill`() = runBlocking {
val f = Fixture()
// Кэш уже содержит сообщение, но meta.syncedAt пуст (БД до миграции).
f.cache.append(user("u1", 1, "старое"))
f.remote.records = listOf(user("u1", 1, "старое"), assistant("a1", 2, "новое"))
val session = f.session()
try {
session.start()
await { session.history.value.size == 2 }
// Дубль u1 не уронил вставку, обе записи на месте.
assertEquals(2, session.history.value.size)
} finally {
f.close()
}
}
@Test
fun `live turn streams then collapses into history on End`() = runBlocking {
val f = Fixture()
val session = f.session()
try {
session.start()
f.outbox.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"))
await { session.live.value.size == 1 }
assertTrue(session.live.value.single() is UiMessage.Assistant)
assertTrue(session.isStreaming)
assertEquals("Hello", (session.live.value.single() as UiMessage.Assistant)
.let { (it.content.single() as pw.binom.agentik.desktop.model.UiContent.Text).body })
// Сервер зафиксировал ход в журнале и прислал End.
f.remote.records = listOf(user("u1", 9, "вопрос"), assistant("a1", 14, "Hello"))
f.outbox.emit(Event.End(at(14)))
await { session.live.value.isEmpty() && session.history.value.size == 2 }
assertTrue(session.live.value.isEmpty())
assertTrue(!session.isStreaming)
} finally {
f.close()
}
}
@Test
fun `backfill keeps a reply that shares its millisecond with the watermark`() = runBlocking {
val f = Fixture()
// Кэшируем user-сообщение; watermark встанет на его же миллисекунду.
f.remote.records = listOf(user("u1", 5, "вопрос"))
val session = f.session()
try {
session.start()
await { session.history.value.size == 1 }
// Ответ пришёл в ТУ ЖЕ миллисекунду (быстрый локальный агент).
// Фильтр `createdAt > after` потерял бы его, если бы watermark
// стоял ровно на 5 мс — отсюда отступ назад на 1 мс + отсечение
// известных id вместо опоры на время.
f.remote.records = listOf(user("u1", 5, "вопрос"), assistant("a1", 5, "ответ"))
f.outbox.emit(Event.End(at(5)))
await { session.history.value.size == 2 }
// Порядок при равных timestamp'ах задаёт id (контракт SQL
// `createdAt ASC, id ASC`) — здесь он не важен, важно что оба
// сообщения доехали: ответ не потерян на границе watermark'а.
assertEquals(1, session.history.value.count { it is UiMessage.User })
assertEquals(1, session.history.value.count { it is UiMessage.Assistant })
} finally {
f.close()
}
}
@Test
fun `interrupt clears live turn without committing`() = runBlocking {
val f = Fixture()
val session = f.session()
try {
session.start()
f.outbox.awaitSubscriber()
f.outbox.emit(Event.StartResponse(at(10), Event.ResponseType.TEXT))
f.outbox.emit(Event.AppendText(at(11), "частичный"))
await { session.live.value.size == 1 }
f.outbox.emit(Event.Interrupted(at(12)))
await { session.live.value.isEmpty() }
assertTrue(session.live.value.isEmpty())
assertTrue(session.history.value.isEmpty())
} finally {
f.close()
}
}
@Test
fun `preview is denormalized from last assistant message`() = runBlocking {
val f = Fixture()
f.remote.records = listOf(
user("u1", 1, "привет"),
assistant("a1", 2, "Большой ответ\nв несколько строк"),
)
val session = f.session()
try {
session.start()
await { f.meta.meta(convId)?.preview != null }
assertEquals("Большой ответ в несколько строк", f.meta.meta(convId)?.preview)
} finally {
f.close()
}
}
}
// ───────────────────────── fakes ─────────────────────────
private class FakeJournalStore(
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 FakeOutboxStore : 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() {}
/** Ждёт, пока у потока появится подписчик (ChatSession подписался). */
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))
}
}
private class FakeConversation(
override val id: String,
override val isSupportImageInput: Boolean = true,
override val isSupportImageOutput: Boolean = true,
override val isTemporal: Boolean = false,
override var title: String? = null,
override var updatedAt: Instant = Instant.fromEpochMilliseconds(0),
) : Conversation {
val sent = mutableListOf<List<Content>>()
var interruptCount = 0
override suspend fun rename(title: String) {
this.title = title
}
override suspend fun send(content: List<Content>, context: MessageContext?) {
sent += content
}
override suspend fun interrupt() {
interruptCount++
}
override suspend fun getMessages(after: Instant, offset: Int, limit: Int): List<Message> = emptyList()
override fun close() {}
}
@@ -0,0 +1,209 @@
package pw.binom.agentik.desktop.session
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertIs
import kotlin.test.assertTrue
import pw.binom.agentik.outbox.Event
import kotlin.time.Clock
import kotlin.time.Instant
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) }
}
@Test
fun `full happy path accumulates reasoning and answer`() {
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(), "мир!"),
Event.End(at()),
)
assertEquals(TurnPhase.DONE, turn.phase)
val blocks = turn.blocks
assertEquals(2, blocks.size)
val reasoning = assertIs<LiveBlock.Reasoning>(blocks[0])
assertEquals("сначала подумаем", reasoning.text.toString())
val answer = assertIs<LiveBlock.Answer>(blocks[1])
assertEquals("Привет, мир!", answer.text.toString())
assertEquals("Привет, мир!", turn.answerText())
}
@Test
fun `tool call and result are stitched into one block`() {
val turn = LiveTurn()
turn.feed(
Event.Working(at()),
Event.StartResponse(at(), Event.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()),
)
val tools = turn.blocks.filterIsInstance<LiveBlock.Tool>()
assertEquals(1, tools.size)
assertEquals("t1", tools[0].callId)
assertEquals("grep", tools[0].toolName)
assertEquals("{\"q\":\"x\"}", tools[0].argsJson)
assertEquals("found 3", tools[0].result)
assertEquals(null, tools[0].failed)
}
@Test
fun `orphan tool result creates a block`() {
val turn = LiveTurn()
turn.feed(
Event.Working(at()),
Event.StartResponse(at(), Event.ResponseType.TEXT),
Event.ToolResult(at(), toolCallId = "missing", toolName = null, result = "ok"),
)
val tool = turn.blocks.filterIsInstance<LiveBlock.Tool>().single()
assertEquals("missing", tool.callId)
assertEquals("(unknown)", tool.toolName)
assertEquals("ok", tool.result)
}
@Test
fun `tool failed marks the block`() {
val turn = LiveTurn()
turn.feed(
Event.Working(at()),
Event.ToolCall(at(), id = "t9", title = null, toolName = "bash", toolArgs = "{}"),
Event.ToolFailed(at(), toolCallId = "t9", toolName = "bash", message = "boom", durationMs = 5),
)
val tool = turn.blocks.filterIsInstance<LiveBlock.Tool>().single()
assertEquals("boom", tool.failed)
}
@Test
fun `error appends failure block and fails the turn`() {
val turn = LiveTurn()
turn.feed(
Event.Working(at()),
Event.StartResponse(at(), Event.ResponseType.TEXT),
Event.AppendText(at(), "partial"),
Event.Error(at(), message = "llm down", code = "E500"),
)
assertEquals(TurnPhase.FAILED, turn.phase)
val failure = turn.blocks.filterIsInstance<LiveBlock.Failure>().single()
assertEquals("llm down", failure.message)
assertEquals("E500", failure.code)
// Частичный ответ остаётся в блоках — UI решает, показывать ли его.
assertEquals("partial", turn.answerText())
}
@Test
fun `interrupted turn is terminal and keeps partial text`() {
val turn = LiveTurn()
turn.feed(
Event.Working(at()),
Event.StartResponse(at(), Event.ResponseType.TEXT),
Event.AppendText(at(), "не до конца"),
Event.Interrupted(at()),
)
assertEquals(TurnPhase.INTERRUPTED, turn.phase)
assertTrue(turn.phase.isTerminal)
assertEquals("не до конца", turn.answerText())
}
@Test
fun `text after tool call starts a new answer block`() {
val turn = LiveTurn()
turn.feed(
Event.Working(at()),
Event.StartResponse(at(), Event.ResponseType.TEXT),
Event.AppendText(at(), "до"),
Event.ToolCall(at(), id = "t1", title = null, toolName = "grep", toolArgs = "{}"),
Event.AppendText(at(), "после"),
)
val answers = turn.blocks.filterIsInstance<LiveBlock.Answer>()
assertEquals(2, answers.size)
assertEquals("до", answers[0].text.toString())
assertEquals("после", answers[1].text.toString())
assertEquals("допосле", turn.answerText())
}
@Test
fun `image response produces image block`() {
val bytes = byteArrayOf(1, 2, 3)
val turn = LiveTurn()
turn.feed(
Event.Working(at()),
Event.StartResponse(at(), Event.ResponseType.IMAGE),
Event.AppendImage(at(), body = bytes, mime = "image/png"),
Event.End(at()),
)
assertTrue(turn.blocks.filterIsInstance<LiveBlock.Answer>().isEmpty())
val image = turn.blocks.filterIsInstance<LiveBlock.Image>().single()
assertEquals("image/png", image.mime)
assertEquals(3, image.data.size)
}
@Test
fun `working resets the previous turn`() {
val turn = LiveTurn()
turn.feed(
Event.Working(at()),
Event.StartResponse(at(), Event.ResponseType.TEXT),
Event.AppendText(at(), "старое"),
Event.End(at()),
)
turn.feed(Event.Working(at()), Event.AppendText(at(), "новое"))
assertEquals(1, turn.blocks.size)
assertEquals("новое", turn.answerText())
}
@Test
fun `append without explicit start response still streams as answer`() {
val turn = LiveTurn()
turn.feed(Event.Working(at()), Event.AppendText(at(), "без старта"))
assertEquals(TurnPhase.RESPONDING, turn.phase)
assertEquals("без старта", turn.answerText())
}
@Test
fun `onChange fires on mutation`() {
val turn = LiveTurn()
var count = 0
turn.onChange = { count++ }
turn.feed(Event.Working(at()), Event.AppendText(at(), "x"))
assertEquals(2, count)
}
@Test
fun `preview text collapses whitespace and truncates`() {
val turn = LiveTurn()
turn.feed(
Event.Working(at()),
Event.AppendText(at(), "Первая строка\n\nвторая строка "),
)
assertEquals("Первая строка вторая строка", turn.previewText())
val long = LiveTurn()
long.feed(Event.Working(at()), Event.AppendText(at(), "слово ".repeat(60)))
val preview = long.previewText(maxLength = 30)
assertTrue(preview.length <= 31)
assertTrue(preview.endsWith("…"))
}
}
@@ -0,0 +1,195 @@
package pw.binom.agentik.desktop.session
import kotlinx.coroutines.delay
import kotlinx.coroutines.runBlocking
import org.junit.jupiter.api.Assumptions.assumeTrue
import pw.binom.agentik.desktop.model.UiContent
import pw.binom.agentik.desktop.model.UiMessage
import pw.binom.agentik.desktop.settings.AgentConfig
import kotlin.test.AfterTest
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertNotNull
import kotlin.test.assertTrue
import kotlin.time.Instant
import java.nio.file.Files
/**
* Реальная интеграция с ЖИВЫМ `:standalone`-сервером (Qwen через vLLM),
* а не с in-memory-эхо-агентом из [AgentConnectionE2ETest].
*
* Адрес берётся из env `AGENTIK_E2E_URL`, по умолчанию
* `http://127.0.0.1:8080/agentik`. Если сервер не поднят — тест
* **пропускается** (`assumeTrue`), поэтому `./gradlew jvmTest` остаётся
* зелёным без запущенного агента, а с ним — проверяет настоящий контракт:
* создать диалог → отправить → дождаться живого ответа LLM → история в
* локальном SQLite → повторное открытие не плодит дублей.
*
* БД уводится в temp-папку, пользовательский `~/.agentik-desktop` не
* трогается.
*/
class RealServerIT {
private val connection = mutableListOf<AgentConnection>()
@AfterTest
fun tearDown() {
connection.forEach { runCatching { it.close() } }
connection.clear()
}
private fun url(): String =
System.getenv("AGENTIK_E2E_URL")?.takeIf { it.isNotBlank() }
?: "http://127.0.0.1:8080/agentik"
private fun connect(scope: kotlinx.coroutines.CoroutineScope, db: String): AgentConnection? =
runCatching {
AgentConnection.open(
config = AgentConfig.create(baseUrl = url()),
clientId = "agentik-desktop-real-e2e",
scope = scope,
dbPath = db,
)
}.getOrNull()
@Test
fun `live server answers and the turn lands in the local cache`() = runBlocking {
val dir = Files.createTempDirectory("agentik-real-e2e")
val db = dir.resolve("agentik.db").toString()
val conn = connect(this, db)
assumeTrue(conn != null, "standalone-сервер недоступен на ${url()} — тест пропущен")
conn!!
connection += conn
val conv = conn.createConversation(temp = false)
assertNotNull(conv, "createConversation")
val session = assertNotNull(conn.openSession(conv.id), "openSession")
try {
session.start()
session.send("Ответь ровно одним словом: привет")
await(timeoutMs = 120_000) {
session.history.value.size >= 2 && !session.isStreaming
}
val history = session.history.value
assertEquals(2, history.size, "user + assistant: $history")
val user = history[0] as UiMessage.User
assertEquals("Ответь ровно одним словом: привет", (user.content.single() as UiContent.Text).body)
val bot = history[1] as UiMessage.Assistant
val answer = (bot.content.single() as UiContent.Text).body
assertTrue(answer.isNotBlank(), "LLM вернул пустой ответ")
// Локальный кэш (offline-история) — на диске.
val cached = conn.journalCache.list(conv.id, Instant.DISTANT_PAST, 0, Int.MAX_VALUE)
assertEquals(2, cached.size, "cache records")
// Превью + watermark денормализованы из ответа.
assertNotNull(conn.meta.meta(conv.id)?.preview)
assertNotNull(conn.meta.meta(conv.id)?.syncedAt)
// Список диалогов от сервера содержит созданный.
await(timeoutMs = 20_000) {
conn.uiConversations().any { it.id == conv.id }
}
// Переоткрытие читает из кэша и НЕ дублирует бэкфиллом.
val again = assertNotNull(conn.openSession(conv.id))
try {
again.start()
assertEquals(2, again.history.value.size, "reopen must not duplicate")
} finally {
again.close()
}
} finally {
session.close()
}
}
/**
* Сценарий «перезапуск клиента»: закрываем соединение и открываем
* заново на той же БД — история должна читаться из локального кэша,
* а не теряться.
*/
@Test
fun `history survives a reconnect on the same database`() = runBlocking {
val dir = Files.createTempDirectory("agentik-real-e2e-reconnect")
val db = dir.resolve("agentik.db").toString()
val convId: String
run {
val conn = connect(this, db)
assumeTrue(conn != null, "standalone-сервер недоступен на ${url()} — тест пропущен")
conn!!
connection += conn
val conv = assertNotNull(conn.createConversation(temp = false))
val session = assertNotNull(conn.openSession(conv.id))
session.start()
session.send("Скажи 'ок'")
await(timeoutMs = 120_000) { session.history.value.size >= 2 && !session.isStreaming }
convId = conv.id
session.close()
conn.close()
connection -= conn
}
val conn2 = assertNotNull(connect(this, db), "повторное подключение")
connection += conn2
val session2 = assertNotNull(conn2.openSession(convId), "openSession после перезапуска")
try {
session2.start()
await(timeoutMs = 30_000) { session2.history.value.size >= 2 }
assertEquals(2, session2.history.value.size, "история из локального кэша после реконнекта")
} finally {
session2.close()
}
}
/**
* Регрессия на реальный баг: `POST /conversations/{id}/messages` на
* сервере отвечает не сразу, а по завершении всего хода агента. Длинный
* ответ LLM (десятки секунд) обрывался дефолтным request-timeout'ом
* клиента (15000 мс) — ход умирал на клиенте, продолжая жить на сервере.
* Клиент обязан ставить на этот запрос infinite-timeout
* (`ConversationClient.send` → `noReadTimeout()`).
*/
@Test
fun `long llm answer outlives the default client request timeout`() = runBlocking {
val dir = Files.createTempDirectory("agentik-real-e2e-long")
val conn = connect(this, dir.resolve("agentik.db").toString())
assumeTrue(conn != null, "standalone-сервер недоступен на ${url()} — тест пропущен")
conn!!
connection += conn
val conv = assertNotNull(conn.createConversation(temp = false))
val session = assertNotNull(conn.openSession(conv.id))
try {
session.start()
session.send("Напиши связный текст примерно на 300 слов про историю SQLite. Без списков и заголовков.")
await(timeoutMs = 300_000) {
session.history.value.size >= 2 && !session.isStreaming
}
val answer = (session.history.value[1] as UiMessage.Assistant)
.content.filterIsInstance<UiContent.Text>().joinToString("\n") { it.body }
assertTrue(
answer.length > 500,
"ожидался длинный ответ (иначе тест не проверяет таймаут), а пришло ${answer.length} символов: $answer",
)
} finally {
session.close()
}
}
private suspend fun await(timeoutMs: Long = 5_000, condition: suspend () -> Boolean) {
val deadline = System.currentTimeMillis() + timeoutMs
while (System.currentTimeMillis() < deadline) {
if (condition()) return
delay(50)
}
assertTrue(condition(), "condition not met within ${timeoutMs}ms")
}
}