" !in out, "html tags must be stripped, got: $out")
+ assertTrue("**" !in out && "`" !in out, "md markers must be stripped, got: $out")
+ }
+
+ @Test
+ fun `line that looks like table header but has no separator falls through as plain text`() {
+ // Строка с `|`, но без `| --- |` после неё — обычный текст,
+ // таблица не строится.
+ val md = "hello | world"
+ val out = MarkdownToTelegram.convert(md)
+ assertTrue("" !in out, "must not wrap non-table in , got: $out")
+ assertTrue("│" !in out, "must not emit box-drawing for non-table, got: $out")
+ }
+
+ @Test
+ fun `REPRO table from LLM now renders as boxed grid`() {
+ // Реальный случай из standalone-лога: агент выдал markdown-таблицу
+ // с HTML внутри ячеек. До фикса она уходила в Telegram как сырой
+ // `| Вещдоказательство | Состояние | ...`; теперь — boxed grid.
+ val input = """
+Вот расклад:
+
+| Вещдоказательство | Состояние | Примечание |
+|---|---|---|
+| Стол | Протрезвевший | Наконец-то |
+| Купе | Отменено | Не доехали |
+
+Готово.
+""".trimIndent()
+ val out = MarkdownToTelegram.convert(input)
+ // Рамка рисуется.
+ assertTrue("┌" in out, "missing top border, got: $out")
+ assertTrue("Стол" in out, "row must survive, got: $out")
+ assertTrue("Протрезвевший" in out, "cell text must survive, got: $out")
+ // Текст «Вот расклад:» и «Готово.» — обычными строками.
+ assertTrue("Вот расклад:" in out, "intro line missing, got: $out")
+ assertTrue("Готово." in out, "outro line missing, got: $out")
+ // Никакого raw-разделителя `|---|---|` в выводе.
+ assertTrue("|---|" !in out, "raw separator leaked, got: $out")
+ }
+}
\ No newline at end of file
diff --git a/integrations-telegram/src/test/kotlin/pw/binom/agentik/integrations/telegram/TelegramBridgeComponentTest.kt b/integrations-telegram/src/test/kotlin/pw/binom/agentik/integrations/telegram/TelegramBridgeComponentTest.kt
index 6c82a4c..2757d64 100644
--- a/integrations-telegram/src/test/kotlin/pw/binom/agentik/integrations/telegram/TelegramBridgeComponentTest.kt
+++ b/integrations-telegram/src/test/kotlin/pw/binom/agentik/integrations/telegram/TelegramBridgeComponentTest.kt
@@ -6,26 +6,25 @@ import kotlinx.coroutines.cancel
import kotlinx.coroutines.channels.Channel
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.MutableSharedFlow
-import kotlinx.coroutines.flow.collect
import kotlinx.coroutines.flow.asSharedFlow
import kotlinx.coroutines.flow.filter
+import kotlinx.coroutines.flow.filterIsInstance
import kotlinx.coroutines.flow.first
import kotlinx.coroutines.launch
import kotlinx.coroutines.runBlocking
-import kotlinx.coroutines.test.runTest
-import kotlinx.coroutines.withTimeout
import kotlinx.coroutines.withTimeout
import pw.binom.agentik.agent.Component
import pw.binom.agentik.agent.MutableAgent
import pw.binom.agentik.content.Content
-import pw.binom.agentik.outbox.AgentEvent
-import pw.binom.agentik.outbox.CommonEvent
-import pw.binom.agentik.outbox.OnlineEvent
-import pw.binom.agentik.outbox.OutboxStore
-import pw.binom.agentik.outbox.OnlineOutbox
import pw.binom.agentik.content.MessageContext
import pw.binom.agentik.journal.ConversationStore
import pw.binom.agentik.journal.JournalStore
+import pw.binom.agentik.outbox.AgentEvent
+import pw.binom.agentik.outbox.CommonEvent
+import pw.binom.agentik.outbox.DurableEvent
+import pw.binom.agentik.outbox.OnlineEvent
+import pw.binom.agentik.outbox.OnlineOutbox
+import pw.binom.agentik.outbox.OutboxStore
import pw.binom.agentik.proto.AgentInfo
import pw.binom.agentik.proto.ChatSnapshot
import pw.binom.agentik.proto.ConversationsSnapshot
@@ -33,21 +32,14 @@ import pw.binom.agentik.proto.Conversation
import pw.binom.db.ksqlite.SQLiteConnection
import pw.binom.telegram.TelegramClient
import pw.binom.telegram.dto.BotCommand
-import pw.binom.telegram.dto.CallbackQuery
import pw.binom.telegram.dto.Chat
import pw.binom.telegram.dto.ChatType
-import pw.binom.telegram.dto.ChosenInlineResult
import pw.binom.telegram.dto.EditMessageResult
import pw.binom.telegram.dto.EditTextRequest
import pw.binom.telegram.dto.File
-import pw.binom.telegram.dto.InlineQuery
import pw.binom.telegram.dto.Message
import pw.binom.telegram.dto.ParseMode
-import pw.binom.telegram.dto.Poll
-import pw.binom.telegram.dto.PollAnswer
-import pw.binom.telegram.dto.PreCheckoutQuery
import pw.binom.telegram.dto.SendChatEvent
-import pw.binom.telegram.dto.ShippingQuery
import pw.binom.telegram.dto.TextMessage
import pw.binom.telegram.dto.Update
import pw.binom.telegram.dto.User
@@ -64,39 +56,34 @@ import kotlin.test.fail
import kotlin.time.Instant
/**
- * Тесты моста [TelegramBridgeComponent].
+ * Тесты [TelegramBridgeComponent].
*
* Контракт:
* - входящий `Update.message` с текстом → создан/найден `Conversation` (по
- * [TelegramChatMap]), `conv.send(Content.Text(...))` дёрнут;
- * - «typing…» в чат через `sendChatAction` шлётся на каждое входящее;
- * - стрим ответа (StartResponse → AppendText → End) рендерится в TG через
- * `sendMessage("...")` + `editMessage(...)` с накопленным текстом;
- * - ошибка в `conv.send` падает в чат как «⚠️ …» (одна строка);
- * - `uninstall(agent)` отменяет polling-джобу.
+ * [TelegramChatMap]), `conv.send(Content.Text(...))` дёрнут, `sendChatAction(TYPING)`
+ * в чат уходит ДО `conv.send`;
+ * - на каждом online-событии для conv'а (Working / StartResponse /
+ * AppendText / End) мост освежает `sendChatAction(TYPING)` в чат;
+ * - на `DurableEvent.AssistantMessage` (целое сообщение) мост шлёт
+ * `sendMessage(...)` с полным текстом (стриминг через `editMessage`
+ * НЕ используется);
+ * - на `DurableEvent.Error` шлётся короткое сообщение об ошибке;
+ * - `uninstall(agent)` отменяет polling-джобу и подписчиков.
*/
-/**
- * Подписаться на [flow] в фоне, накапливать события в список и вернуть его.
- * Подписка стартует до теста (subscribe-before-act), потом тест шлёт события и ассертит.
- */
-private fun subscribe(scope: CoroutineScope, flow: kotlinx.coroutines.flow.Flow): MutableList {
+private fun subscribe(scope: CoroutineScope, flow: Flow): MutableList {
val list = mutableListOf()
- scope.launch {
- flow.collect { list.add(it) }
- }
+ scope.launch { flow.collect { list.add(it) } }
return list
}
-/** Ждём, пока в [list] появится элемент, удовлетворяющий [match], или таймаут. */
-private suspend fun List.awaitOne(match: (T) -> Boolean, timeoutMs: Long = 2_000): T {
- return kotlinx.coroutines.withTimeout(timeoutMs) {
+private suspend fun List.awaitOne(match: (T) -> Boolean, timeoutMs: Long = 2_000): T =
+ withTimeout(timeoutMs) {
while (true) {
firstOrNull(match)?.let { return@withTimeout it }
kotlinx.coroutines.delay(10)
}
@Suppress("UNREACHABLE_CODE") error("unreachable")
}
-}
class TelegramBridgeComponentTest {
@@ -112,14 +99,12 @@ class TelegramBridgeComponentTest {
fakeTg = FakeTelegramClient()
connection = SQLiteConnection.memory("tg-test-${Random.nextLong()}")
chatMap = TelegramChatMap(connection)
- fakeAgent = FakeMutableAgent(onlineOutbox = FakeOnlineOutbox())
+ fakeAgent = FakeMutableAgent()
scope = CoroutineScope(SupervisorJob() + kotlinx.coroutines.Dispatchers.Default)
}
@AfterTest
fun tearDown() {
- // Defensive: инициализация могла не дойти (например, SQLiteException
- // в @BeforeTest — connection создан, chatMap — нет).
runCatching { if (::bridge.isInitialized) bridge.uninstall(fakeAgent) }
runCatching { if (::scope.isInitialized) scope.cancel() }
runCatching { if (::chatMap.isInitialized) chatMap.close() }
@@ -152,22 +137,22 @@ class TelegramBridgeComponentTest {
),
)
-
+ // ============================================================================
+ // Inbound: маппинг chat ↔ conv, conv.send, typing ДО send
+ // ============================================================================
@Test
fun `first incoming text creates a persistent conversation and forwards the text`() = runBlocking {
- // persistent=true (default): первый Update.message → createConversation(temp=false),
- // chat_map сохраняет chat_id → conv_id, conv.send(Content.Text(text)) дёрнут.
bridge = newBridge(scope)
bridge.install(fakeAgent)
val incoming = subscribe(scope, bridge.incomingFlow())
try {
+ val typingBefore = fakeTg.chatActions.size
fakeTg.push(makeUpdate(updateId = 1, chatId = 100, text = "hi"))
- val ev = incoming.awaitOne({ e: TelegramBridgeComponent.IncomingEvent -> e.text == "hi" })
+ val ev = incoming.awaitOne({ it.text == "hi" })
assertEquals("hi", ev.text)
assertEquals("100", ev.chatId)
- // Та же convId выдаётся для одного чата.
val convId = ev.convId
// Маппинг зафиксирован.
@@ -179,6 +164,13 @@ class TelegramBridgeComponentTest {
val sent = fakeAgent.conversations.getValue(convId).sent
assertEquals(1, sent.size)
assertEquals("hi", (sent[0][0] as Content.Text).body)
+
+ // Typing в чат ушёл ДО conv.send (точнее, ДО того, как send
+ // завершится — но в тесте мост его шлёт сразу в handleUpdate
+ // синхронно, а send — fire-and-forget launch'ом).
+ val newTypings = fakeTg.chatActions.drop(typingBefore)
+ .filter { it.first == "100" && it.second == SendChatEvent.Action.TYPING }
+ assertEquals(1, newTypings.size, "expected exactly one TYPING for chat 100, got $newTypings")
} finally {
bridge.uninstall(fakeAgent)
}
@@ -186,18 +178,17 @@ class TelegramBridgeComponentTest {
@Test
fun `incoming text from same chat reuses the same conversation`() = runBlocking {
- // Два сообщения от одного chatId — маппинг один и тот же, новых createConversation нет.
bridge = newBridge(scope)
bridge.install(fakeAgent)
val incoming = subscribe(scope, bridge.incomingFlow())
try {
fakeTg.push(makeUpdate(updateId = 1, chatId = 200, text = "first"))
- val first = incoming.awaitOne({ e: TelegramBridgeComponent.IncomingEvent -> e.text == "first" })
+ val first = incoming.awaitOne({ it.text == "first" })
val firstConv = first.convId
- // Ждём, пока bridge обработает (poll вернётся пустым после Update).
fakeTg.drain()
fakeTg.push(makeUpdate(updateId = 2, chatId = 200, text = "second"))
- incoming.awaitOne({ e: TelegramBridgeComponent.IncomingEvent -> e.text == "second" })
+ incoming.awaitOne({ it.text == "second" })
+
// createConversation вызвался ровно один раз.
assertEquals(1, fakeAgent.created.size)
// Оба сообщения ушли в один conv.
@@ -212,8 +203,6 @@ class TelegramBridgeComponentTest {
@Test
fun `non-persistent config creates a fresh conversation per message`() = runBlocking {
- // persistentConversations=false: каждое сообщение = новая (temp) беседа,
- // маппинг всё равно обновляется на последнюю.
bridge = newBridge(
scope = scope,
config = TelegramConfig(token = "test", pollingTimeoutSec = 0, persistentConversations = false),
@@ -222,96 +211,207 @@ class TelegramBridgeComponentTest {
val incoming = subscribe(scope, bridge.incomingFlow())
try {
fakeTg.push(makeUpdate(updateId = 1, chatId = 300, text = "a"))
- val first = incoming.awaitOne({ e: TelegramBridgeComponent.IncomingEvent -> e.text == "a" })
+ val first = incoming.awaitOne({ it.text == "a" })
fakeTg.drain()
fakeTg.push(makeUpdate(updateId = 2, chatId = 300, text = "b"))
- val second = incoming.awaitOne({ e: TelegramBridgeComponent.IncomingEvent -> e.text == "b" })
+ val second = incoming.awaitOne({ it.text == "b" })
- // Два разных conv (temp).
assertEquals(2, fakeAgent.created.size)
assertEquals(true, fakeAgent.created.all { it.temp })
- // Маппинг указывает на последний.
assertEquals(second.convId, chatMap.getConversationId("300"))
- // Первая convId больше не текущая.
assertTrue(first.convId != second.convId)
} finally {
bridge.uninstall(fakeAgent)
}
}
+ // ============================================================================
+ // Typing: refresh на каждом online-событии
+ // ============================================================================
+
@Test
- fun `typing chat action is sent on every incoming text`() = runBlocking {
+ fun `typing chat action is throttled and only one TYPING per chat in burst`() = runBlocking {
+ // Telegram режет rate-limit на sendChatAction — мост троттлит
+ // до ≤1 запроса в 4 секунды на chatId. Здесь шлём 2 входящих
+ // и 5 online-событий подряд — должны получить ровно 1 TYPING
+ // (сразу на первое входящее), всё остальное подавлено.
bridge = newBridge(scope)
bridge.install(fakeAgent)
- val incoming: MutableList = subscribe(scope, bridge.incomingFlow())
+ val incoming = subscribe(scope, bridge.incomingFlow())
try {
fakeTg.push(makeUpdate(updateId = 1, chatId = 400, text = "ping"))
+ val ev = incoming.awaitOne({ it.text == "ping" })
+ val convId = ev.convId
fakeTg.drain()
fakeTg.push(makeUpdate(updateId = 2, chatId = 400, text = "pong"))
- incoming.awaitOne({ e: TelegramBridgeComponent.IncomingEvent -> e.text == "pong" })
+ incoming.awaitOne({ it.text == "pong" })
fakeTg.drain()
- val typing = fakeTg.chatActions.filter { it.second == SendChatEvent.Action.TYPING }
- assertTrue(typing.size >= 2, "expected ≥2 TYPING actions, got ${fakeTg.chatActions}")
- assertTrue(typing.all { it.first == "400" }, "expected all TYPING to chat 400, got $typing")
+ val typings = fakeTg.chatActions.filter { it.first == "400" && it.second == SendChatEvent.Action.TYPING }
+ assertEquals(1, typings.size, "burst of two messages must collapse to one TYPING, got $typings")
+
+ // Online-события подряд — все подавлены троттлом.
+ val baseline = typings.size
+ val now = Instant.fromEpochMilliseconds(1_700_000_000_000)
+ val out = fakeAgent.onlineOutbox
+ out.emit(OnlineEvent.Working(now, convId))
+ out.emit(OnlineEvent.StartResponse(now, convId, OnlineEvent.ResponseType.TEXT))
+ out.emit(OnlineEvent.AppendText(now, convId, "a"))
+ out.emit(OnlineEvent.AppendText(now, convId, "b"))
+ out.emit(OnlineEvent.End(now, convId))
+ // Даём корутинам время отработать burst, но без ожидания
+ // конкретного счётчика — ожидаем ровно 0 новых TYPING.
+ kotlinx.coroutines.delay(100)
+ val refreshed = fakeTg.chatActions
+ .drop(baseline)
+ .filter { it.first == "400" && it.second == SendChatEvent.Action.TYPING }
+ assertEquals(0, refreshed.size, "online burst within throttle window must not send more TYPING, got $refreshed")
} finally {
bridge.uninstall(fakeAgent)
}
}
@Test
- fun `outbound streaming renders answer via editMessage edits`() = runBlocking {
+ fun `typing chat action is refreshed after throttle interval expires`() = runBlocking {
+ // Мост шлёт TYPING сразу на первое входящее, потом пропускает
+ // burst online-событий; следующее входящее ПОСЛЕ 4s троттла
+ // должно освежить TYPING.
+ bridge = newBridge(scope)
+ bridge.install(fakeAgent)
+ val incoming = subscribe(scope, bridge.incomingFlow())
+ try {
+ fakeTg.push(makeUpdate(updateId = 1, chatId = 500, text = "hi"))
+ incoming.awaitOne({ it.text == "hi" })
+ fakeTg.drain()
+ val first = fakeTg.chatActions.count { it.first == "500" && it.second == SendChatEvent.Action.TYPING }
+ assertEquals(1, first)
+
+ // Провалидируем мост: между двумя входящими должен пройти
+ // TYPING_REFRESH_MIN_INTERVAL_MS — иначе второе входящее тоже
+ // схлопнется в первый TYPING. Здесь «время» управляется через
+ // ручное продвижение в `sendTypingThrottled` — мы проверяем,
+ // что МОСТ после интервала снова шлёт TYPING (эмулируем,
+ // подменив `lastTypingSentAt` через отражение/private API не
+ // делаем — это покрывается end-to-end сценарием ниже).
+ // Реальный прогон: меняем internal-флаг через рефлексию и
+ // шлём второе входящее.
+ val field = bridge.javaClass.getDeclaredField("lastTypingSentAt").apply { isAccessible = true }
+ @Suppress("UNCHECKED_CAST")
+ val map = field.get(bridge) as MutableMap
+ map["500"] = 0L // сбрасываем «последний раз» — троттл сразу разрешает
+
+ fakeTg.push(makeUpdate(updateId = 2, chatId = 500, text = "again"))
+ incoming.awaitOne({ it.text == "again" })
+ fakeTg.drain()
+ val second = fakeTg.chatActions.count { it.first == "500" && it.second == SendChatEvent.Action.TYPING }
+ assertEquals(2, second, "after throttle resets, next incoming must send fresh TYPING")
+ } finally {
+ bridge.uninstall(fakeAgent)
+ }
+ }
+
+ // ============================================================================
+ // Outbound: DurableEvent.AssistantMessage → sendMessage(полный текст)
+ // ============================================================================
+
+ @Test
+ fun `DurableEvent AssistantMessage sends full text as a single message`() = runBlocking {
bridge = newBridge(scope)
bridge.install(fakeAgent)
val incoming = subscribe(scope, bridge.incomingFlow())
try {
fakeTg.push(makeUpdate(updateId = 1, chatId = 500, text = "ask"))
- val ev = incoming.awaitOne({ e: TelegramBridgeComponent.IncomingEvent -> e.text == "ask" })
+ val ev = incoming.awaitOne({ it.text == "ask" })
val convId = ev.convId
-
val now = Instant.fromEpochMilliseconds(1_700_000_000_000)
- val out = fakeAgent.onlineOutbox
- out.emit(OnlineEvent.StartResponse(now, convId, OnlineEvent.ResponseType.TEXT))
- out.emit(OnlineEvent.AppendText(now, convId, "Hello"))
- out.emit(OnlineEvent.AppendText(now, conversationId = convId, body = ", "))
- out.emit(OnlineEvent.AppendText(now, conversationId = convId, body = "world!"))
- out.emit(OnlineEvent.End(now, convId))
- fakeTg.waitForEdits(3)
- val drafts = fakeTg.sentMessages.filter { it.chatId == "500" && it.text == "..." }
- assertEquals(1, drafts.size, "expected exactly one draft '...' sent, got ${fakeTg.sentMessages}")
- val edits = fakeTg.edited.filter { it.chatId == "500" }
- assertTrue(edits.size >= 3, "expected ≥3 edits (one per AppendText), got ${edits.size}")
- val lastText = edits.last().text
- assertTrue("Hello" in lastText && "world!" in lastText, "final edit must contain full text: $lastText")
+ // Полный текст приходит в DurableEvent.AssistantMessage.
+ fakeAgent.outbox.emitDurable(
+ convId = convId,
+ event = DurableEvent.AssistantMessage(
+ date = now,
+ id = "msg-1",
+ content = listOf(Content.Text("Hello, world!")),
+ ),
+ )
+ fakeTg.waitForMessageContaining("Hello, world!")
+
+ // Ровно ОДНО сообщение, никаких drafts / edits.
+ val final = fakeTg.sentMessages.filter { it.chatId == "500" && it.text == "Hello, world!" }
+ assertEquals(1, final.size, "expected one sendMessage with full text, got ${fakeTg.sentMessages}")
+ // Никаких editMessage'й не было.
+ assertEquals(0, fakeTg.edited.size, "editMessage must not be used, got ${fakeTg.edited}")
} finally {
bridge.uninstall(fakeAgent)
}
}
@Test
- fun `End with empty buffer replaces draft with (empty response)`() = runBlocking {
+ fun `AssistantMessage with markdown is sent as Telegram HTML with parseMode`() = runBlocking {
+ bridge = newBridge(scope)
+ bridge.install(fakeAgent)
+ val incoming = subscribe(scope, bridge.incomingFlow())
+ try {
+ fakeTg.push(makeUpdate(updateId = 1, chatId = 550, text = "ask"))
+ val ev = incoming.awaitOne({ it.text == "ask" })
+ val convId = ev.convId
+ val now = Instant.fromEpochMilliseconds(1_700_000_000_000)
+ fakeAgent.outbox.emitDurable(
+ convId = convId,
+ event = DurableEvent.AssistantMessage(
+ date = now,
+ id = "msg-md",
+ content = listOf(Content.Text("Hello **world** — see [docs](https://e.com)")),
+ ),
+ )
+ // Ждём конкретный HTML, а не plain.
+ fakeTg.waitForHtmlContaining("550", "world")
+
+ val sent = fakeTg.sentMessages.filter { it.chatId == "550" }
+ assertEquals(1, sent.size, "expected one sendMessage, got ${fakeTg.sentMessages}")
+ assertEquals("HTML", sent[0].parseMode?.code, "parseMode must be HTML, got ${sent[0].parseMode}")
+ assertEquals(
+ "Hello world — see docs",
+ sent[0].text,
+ )
+ } finally {
+ bridge.uninstall(fakeAgent)
+ }
+ }
+
+ @Test
+ fun `AssistantMessage with empty content falls back to (empty response)`() = runBlocking {
bridge = newBridge(scope)
bridge.install(fakeAgent)
val incoming = subscribe(scope, bridge.incomingFlow())
try {
fakeTg.push(makeUpdate(updateId = 1, chatId = 600, text = "ask"))
- val ev = incoming.awaitOne({ e: TelegramBridgeComponent.IncomingEvent -> e.text == "ask" })
+ val ev = incoming.awaitOne({ it.text == "ask" })
val convId = ev.convId
val now = Instant.fromEpochMilliseconds(1_700_000_000_000)
- val out = fakeAgent.onlineOutbox
- // StartResponse → draft "..." создан. Сразу End без AppendText → buffer пустой.
- out.emit(OnlineEvent.StartResponse(now, convId, OnlineEvent.ResponseType.TEXT))
- out.emit(OnlineEvent.End(now, convId))
- fakeTg.waitForEditContaining("(empty response)")
+ fakeAgent.outbox.emitDurable(
+ convId = convId,
+ event = DurableEvent.AssistantMessage(
+ date = now,
+ id = "msg-empty",
+ content = emptyList(),
+ ),
+ )
+ fakeTg.waitForMessageContaining("(empty response)")
- val emptyEdits = fakeTg.edited.filter { it.chatId == "600" && it.text == "(empty response)" }
- assertEquals(1, emptyEdits.size, "expected one (empty response) edit, got ${fakeTg.edited}")
+ val fallback = fakeTg.sentMessages.filter {
+ it.chatId == "600" && it.text == "(empty response)"
+ }
+ assertEquals(1, fallback.size, "expected one (empty response), got ${fakeTg.sentMessages}")
} finally {
bridge.uninstall(fakeAgent)
}
}
+ // ============================================================================
+ // Errors / lifecycle
+ // ============================================================================
+
@Test
fun `conv send error is reported to the chat as a single line`() = runBlocking {
bridge = newBridge(scope)
@@ -319,10 +419,10 @@ class TelegramBridgeComponentTest {
val incoming = subscribe(scope, bridge.incomingFlow())
try {
fakeTg.push(makeUpdate(updateId = 1, chatId = 700, text = "warmup"))
- val warmup = incoming.awaitOne({ e: TelegramBridgeComponent.IncomingEvent -> e.text == "warmup" })
+ val warmup = incoming.awaitOne({ it.text == "warmup" })
fakeAgent.conversations.getValue(warmup.convId).sendError = IllegalStateException("model offline")
fakeTg.push(makeUpdate(updateId = 2, chatId = 700, text = "boom"))
- incoming.awaitOne({ e: TelegramBridgeComponent.IncomingEvent -> e.text == "boom" })
+ incoming.awaitOne({ it.text == "boom" })
fakeTg.waitForMessageContaining("⚠️ agent error: model offline")
val errors = fakeTg.sentMessages.filter {
@@ -334,6 +434,31 @@ class TelegramBridgeComponentTest {
}
}
+ @Test
+ fun `DurableEvent Error is forwarded to the chat as a warning line`() = runBlocking {
+ bridge = newBridge(scope)
+ bridge.install(fakeAgent)
+ val incoming = subscribe(scope, bridge.incomingFlow())
+ try {
+ fakeTg.push(makeUpdate(updateId = 1, chatId = 750, text = "ask"))
+ val ev = incoming.awaitOne({ it.text == "ask" })
+ val convId = ev.convId
+ val now = Instant.fromEpochMilliseconds(1_700_000_000_000)
+ fakeAgent.outbox.emitDurable(
+ convId = convId,
+ event = DurableEvent.Error(date = now, message = "context overflow", code = "ctx"),
+ )
+ fakeTg.waitForMessageContaining("⚠️ context overflow")
+
+ val errs = fakeTg.sentMessages.filter {
+ it.chatId == "750" && it.text == "⚠️ context overflow"
+ }
+ assertEquals(1, errs.size, "expected one ⚠️ error line, got ${fakeTg.sentMessages}")
+ } finally {
+ bridge.uninstall(fakeAgent)
+ }
+ }
+
@Test
fun `uninstall cancels polling and no further updates are processed`() = runBlocking {
bridge = newBridge(scope)
@@ -341,7 +466,7 @@ class TelegramBridgeComponentTest {
val incoming = subscribe(scope, bridge.incomingFlow())
try {
fakeTg.push(makeUpdate(updateId = 1, chatId = 800, text = "alive"))
- incoming.awaitOne({ e: TelegramBridgeComponent.IncomingEvent -> e.text == "alive" })
+ incoming.awaitOne({ it.text == "alive" })
bridge.uninstall(fakeAgent)
val sentBefore = fakeAgent.conversations.values.sumOf { it.sent.size }
@@ -354,9 +479,61 @@ class TelegramBridgeComponentTest {
}
}
-// Fakes
-// ============================================================================
+ @Test
+ fun `Telegram rejection of AssistantMessage is swallowed (polling survives)`() = runBlocking {
+ bridge = newBridge(scope)
+ bridge.install(fakeAgent)
+ val incoming = subscribe(scope, bridge.incomingFlow())
+ try {
+ fakeTg.push(makeUpdate(updateId = 1, chatId = 700, text = "ask"))
+ val ev = incoming.awaitOne({ it.text == "ask" })
+ val convId = ev.convId
+ val now = Instant.fromEpochMilliseconds(1_700_000_000_000)
+ // Имитируем TelegramException (400 Bad Request: can't parse entities).
+ // Мост должен:
+ // - не упасть (polling выживает)
+ // - залогировать ошибку через slf4j (log.error("sendMessage(HTML) failed ..."))
+ // - не отправить сообщение в чат
+ fakeTg.failNextSend = pw.binom.telegram.TelegramException(
+ code = 400,
+ description = "Bad Request: can't parse entities",
+ )
+ fakeAgent.outbox.emitDurable(
+ convId = convId,
+ event = DurableEvent.AssistantMessage(
+ date = now,
+ id = "msg-bad",
+ content = listOf(Content.Text("**boom**")),
+ ),
+ )
+ // Дать корутине шанс обработать исключение.
+ kotlinx.coroutines.delay(150)
+
+ // Сообщение не дошло (Telegram отверг).
+ assertTrue(
+ fakeTg.sentMessages.none { it.chatId == "700" && "boom" in it.text },
+ "Telegram отверг сообщение — в sentMessages его быть не должно",
+ )
+
+ // Polling выжил: следующий AssistantMessage уходит нормально.
+ fakeAgent.outbox.emitDurable(
+ convId = convId,
+ event = DurableEvent.AssistantMessage(
+ date = now,
+ id = "msg-ok",
+ content = listOf(Content.Text("hi again")),
+ ),
+ )
+ fakeTg.waitForMessageContaining("hi again")
+ assertTrue(
+ fakeTg.sentMessages.any { it.chatId == "700" && "hi again" in it.text },
+ "после ошибки мост продолжил работать",
+ )
+ } finally {
+ bridge.uninstall(fakeAgent)
+ }
+ }
}
/**
@@ -364,15 +541,10 @@ class TelegramBridgeComponentTest {
* методы пишут в лог-структуры.
*/
class FakeTelegramClient : TelegramClient {
- /** Запросы `getUpdate`, ожидающие обработки polling'ом. */
private val queue = Channel>(capacity = Channel.UNLIMITED)
- /** «Опубликовать» пачку апдейтов — следующий getUpdate её заберёт. */
- fun push(vararg updates: Update) {
- queue.trySend(updates.toList())
- }
+ fun push(vararg updates: Update) { queue.trySend(updates.toList()) }
- /** Слить все текущие «опубликованные» апдейты (после теста). */
suspend fun drain() {
while (true) {
val r = queue.tryReceive()
@@ -385,6 +557,9 @@ class FakeTelegramClient : TelegramClient {
val chatActions: MutableList> = mutableListOf()
val deleted: MutableList> = mutableListOf()
+ /** Если задано — следующий `sendMessage` бросит [RuntimeException] (имитация 400/403/429). */
+ var failNextSend: RuntimeException? = null
+
private val messageIdSeq = java.util.concurrent.atomic.AtomicLong(1000L)
override suspend fun getUpdate(
@@ -395,6 +570,11 @@ class FakeTelegramClient : TelegramClient {
): List = queue.receive()
override suspend fun sendMessage(message: TextMessage): Message {
+ failNextSend?.let {
+ val e = it
+ failNextSend = null
+ throw e
+ }
sentMessages.add(message)
return Message(
messageId = messageIdSeq.getAndIncrement(),
@@ -422,27 +602,15 @@ class FakeTelegramClient : TelegramClient {
return true
}
- suspend fun waitForEdits(minCount: Int) {
- withTimeout(2_000) {
- while (edited.size < minCount) {
- kotlinx.coroutines.delay(10)
- }
- }
- }
-
- suspend fun waitForEditContaining(text: String) {
- withTimeout(2_000) {
- while (edited.none { text in it.text }) {
- kotlinx.coroutines.delay(10)
- }
- }
- }
-
suspend fun waitForMessageContaining(text: String) {
withTimeout(2_000) {
- while (sentMessages.none { text in it.text }) {
- kotlinx.coroutines.delay(10)
- }
+ while (sentMessages.none { text in it.text }) kotlinx.coroutines.delay(10)
+ }
+ }
+
+ suspend fun waitForHtmlContaining(chatId: String, html: String) {
+ withTimeout(2_000) {
+ while (sentMessages.none { it.chatId == chatId && html in it.text }) kotlinx.coroutines.delay(10)
}
}
@@ -467,10 +635,18 @@ class FakeTelegramClient : TelegramClient {
}
/**
- * Минимальный `MutableAgent`: создаёт `FakeConversation`, помнит их по id.
- * `onlineOutbox` управляется из теста (`FakeOnlineOutbox`).
+ * Минимальный `MutableAgent` для тестов моста.
+ *
+ * - `onlineOutbox` / `outbox` управляются из теста (тест пушит события
+ * через `emit(...)` / `emitDurable(...)`, мост слушает `onlineEvents(convId)`
+ * и `conversationEvents(convId)`).
+ * - `createConversation` возвращает `FakeConversation` с управляемой
+ * отправкой (`sendError`).
*/
-class FakeMutableAgent(override val onlineOutbox: FakeOnlineOutbox) : MutableAgent {
+class FakeMutableAgent : MutableAgent {
+ override val onlineOutbox: FakeOnlineOutbox = FakeOnlineOutbox()
+ override val outbox: FakeOutboxStore = FakeOutboxStore()
+
override val systemProviders: MutableList = mutableListOf()
override val toolProviders: MutableList = mutableListOf()
val conversations: MutableMap = mutableMapOf()
@@ -483,7 +659,6 @@ class FakeMutableAgent(override val onlineOutbox: FakeOnlineOutbox) : MutableAge
override val id: String = "test-agent"
override val info: AgentInfo = AgentInfo(name = "test", description = "", usefulness = "")
override val journal: JournalStore by lazy { error("journal not used in bridge test") }
- override val outbox: OutboxStore by lazy { error("outbox not used in bridge test") }
override val conversationStore: ConversationStore by lazy { error("conversationStore not used in bridge test") }
override fun createConversation(temp: Boolean): Conversation {
@@ -510,7 +685,8 @@ class FakeMutableAgent(override val onlineOutbox: FakeOnlineOutbox) : MutableAge
}
/**
- * Управляемый `OnlineOutbox` — тест пушит OnlineEvent'ы, мост подписан на `onlineEvents(convId)`.
+ * Управляемый `OnlineOutbox`: тест эмитит `OnlineEvent` через [emit],
+ * мост слушает `onlineEvents(convId)`.
*/
class FakeOnlineOutbox : OnlineOutbox {
private val bus = MutableSharedFlow(replay = 0, extraBufferCapacity = 64)
@@ -524,13 +700,37 @@ class FakeOnlineOutbox : OnlineOutbox {
override fun close() {}
}
+/**
+ * Управляемый `OutboxStore`: тест эмитит durable-события для конкретного
+ * conv'а через [emitDurable], мост слушает `conversationEvents(convId)`.
+ */
+class FakeOutboxStore : OutboxStore {
+ private val bus = MutableSharedFlow(replay = 0, extraBufferCapacity = 64)
+
+ suspend fun emitDurable(convId: String, event: DurableEvent) {
+ bus.emit(CommonEvent.Conversation(date = event.date, conversationId = convId, event = event))
+ }
+
+ override fun events(): Flow = bus.asSharedFlow()
+ override fun conversationEvents(conversationId: String?): Flow =
+ bus.asSharedFlow()
+ .filterIsInstance()
+ .let { filtered ->
+ if (conversationId == null) filtered
+ else filtered.filter { it.conversationId == conversationId }
+ }
+ override fun agentEvents(): Flow =
+ bus.asSharedFlow().filterIsInstance()
+
+ override fun close() {}
+}
+
class FakeConversation(
override val id: String,
override val isTemporal: Boolean = false,
) : Conversation {
val sent: MutableList> = mutableListOf()
var interrupted = false
- /** Если задано, следующий `send` кинет это исключение. */
var sendError: Throwable? = null
override val isSupportImageInput: Boolean = false
@@ -549,23 +749,7 @@ class FakeConversation(
interrupted = true
}
- override suspend fun getMessages(after: Instant, offset: Int, limit: Int) = emptyList()
+ override suspend fun getMessages(after: Instant, offset: Int, limit: Int) =
+ emptyList()
override fun close() {}
}
-
-// ============================================================================
-// Flow helpers
-// ============================================================================
-
-private suspend fun Flow.firstAsync(predicate: (T) -> Boolean = { true }): T =
- kotlinx.coroutines.withTimeout(2_000) {
- this@firstAsync.first { predicate(it) }
- }
-
-private suspend fun Flow.firstAsyncOrNull(predicate: (T) -> Boolean): T? = try {
- kotlinx.coroutines.withTimeout(1_000) {
- this@firstAsyncOrNull.first { predicate(it) }
- }
-} catch (_: kotlinx.coroutines.TimeoutCancellationException) {
- null
-}
diff --git a/standalone/README.md b/standalone/README.md
index c9b9db6..507c828 100644
--- a/standalone/README.md
+++ b/standalone/README.md
@@ -22,6 +22,9 @@
## Как запустить
+### Переменные среды
+`AGENTIK_SKILLS_DIR=./agentik/skills;AGENTIK_SYSTEM_PROMPT=Ты полезный ассистент.;OPENAI_BASE_URL=http://192.168.88.135:8001/v1;AGENTIK_DB_PATH=./agentik/db.db;OPENAI_CONTEXT_WINDOW=115000;AGENTIK_MEMORY_DIR=./agentik/memory;AGENTIK_LLM_BACKEND=openai;AGENTIK_PORT=8080;OPENAI_API_KEY=sk-76R2p5nxQxflPkIROr6r2xGiuSYzUCYM;OPENAI_MODEL=/root/.cache/huggingface/Qwen3.8-27B-NVFP4-RTX5090;AGENTIK_SOUL=./agentik/SOUL.md;AGENTIK_TELEGRAM_TOKEN=8857094360:AAFeVYMlutSaHptbhpO4mUryR6Je7OUJsm4`
+
### Требования
- JVM 21+.
@@ -107,14 +110,32 @@ AGENTIK_GOOGLE_MODEL_PATH=/root/gemma-4-E2B-it.litertlm \
- Long-polling через Bot API (`getUpdates?timeout=N`). Не webhook —
не нужен публичный URL, всё работает за NAT.
- Маппинг `chatId ↔ ConversationId` в общей SQLite-БД standalone'а
- (таблица `tg_chat_map`). Один Telegram-чат = один диалог с агентом;
- сообщения из чата стримят ответы обратно через `editMessageText`.
+ (таблица `tg_chat_map`). Один Telegram-чат = один диалог с агентом.
+- **«печатает…» в чате** через `sendChatAction(TYPING)`:
+ - одно сразу при получении апдейта (ДО вызова LLM — пользователь
+ видит статус ровно в момент отправки текста);
+ - рефреш на каждом online-событии для диалога (Working /
+ StartReasoning / StartResponse / AppendText / AppendImage / End) —
+ Telegram гасит статус через ~5 с, поэтому для длинных tool-call'ов
+ (когда дельт текста нет) мост продолжает рефрешить typing по
+ `Working`.
+- **Финальное сообщение** приходит одним `sendMessage(...)` по
+ `DurableEvent.AssistantMessage` (целиковый текст из `Content.Text`).
+ Никакого стриминга через `editMessage` — UX проще, и Telegram
+ нормально индексирует одиночное сообщение.
+- **Markdown → Telegram HTML**: мост конвертирует CommonMark-подобный
+ markdown агента в `ParseMode.HTML` (``, ``, ``, ``,
+ ``); остальные символы экранируются. Сообщения длиннее
+ 4096 символов режутся по переводам строк (вне HTML-тегов).
+- **Markdown-таблицы** (`| col | col |` + `| --- | --- |`) рендерятся
+ через Unicode box-drawing (`┌─┬─┐ │ │ │ ├─┼─┤ └─┴─┘`) внутри `` —
+ Telegram Bot API не поддерживает `
`, поэтому табличный вид
+ даёт только моноширинная рамка. Markdown/HTML-разметка внутри
+ ячеек снимается (в `` inline-форматирование не работает).
- Group-чаты (`group`, `supergroup`) и команда `/new` стартуют **новый
- временный** диалог (`ConversationRecord.isTemporal = true`) — не
- пишут в мапу, удаляются при следующем `/new`.
-- Tool-call показывает «⚙️ обрабатываю…» через `sendChatAction(typing)` —
- включается на `onlineOutbox.onlineEvents` пока обрабатывается ход.
-- Ошибки агента стримятся в чат одной строкой, polling не падает.
+ временный** диалог (`ConversationRecord.isTemporal = true`).
+- Ошибки агента (`DurableEvent.Error`, исключение из `conv.send`)
+ приходят одной строкой, polling не падает.
### Как включить
@@ -134,8 +155,9 @@ AGENTIK_GOOGLE_MODEL_PATH=/root/gemma-4-E2B-it.litertlm \
telegram: enabled (polling timeout=30s, persistent=true)
```
-4. Ответы приходят стримингом — `editMessageText` обновляет одно и то
- же сообщение каждые ~1с пока агент генерирует.
+4. Ответ приходит одним сообщением (после завершения хода ассистента).
+ Пока агент думает / вызывает тулзы / стримит токены — в чате горит
+ «печатает…».
### Persistent vs temporary