From 24212a100334700928ff6d414e72c9ef3afe1011 Mon Sep 17 00:00:00 2001 From: subochev Date: Mon, 5 Oct 2026 23:33:57 +0300 Subject: [PATCH] =?UTF-8?q?feat(standalone):=20:integrations-telegram=20?= =?UTF-8?q?=E2=80=94=20=D0=BC=D0=BE=D1=81=D1=82=20Telegram-=D1=87=D0=B0?= =?UTF-8?q?=D1=82=20=E2=86=94=20agent?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Новый модуль :integrations-telegram, опциональный «плагин» для standalone: * TelegramBridgeComponent — Component-имплементация, устанавливается в agent через agent.install(...) только при заданном AGENTIK_TELEGRAM_TOKEN. Нет токена — нет polling'а, classpath не содержит pw.binom.telegram.* (factory возвращает AutoCloseable, а Component-объект создаётся внутри factory и утекает в standalone только через agent-api). * Long-polling getUpdates?timeout=N — не webhook, не нужен публичный URL. * chatId ↔ ConversationId маппинг в общей SQLite-БД standalone'а (таблица tg_chat_map); persistent=true переиспользует conv_id, persistent=false — каждый Telegram-message = свежий temp-диалог. * Group-чаты и /new стартуют временный диалог, не пишут в мапу. * Tool-call показывает «⚙️ обрабатываю…» через sendChatAction(typing) по onlineOutbox.onlineEvents. * Outbound streaming через editMessageText с guard'ом «message is not modified»; End с пустым телом → «(empty response)». * Ошибки conv.send — одной строкой в чат, polling не падает. * uninstall отменяет scope, polling и stream-джобы корректно останавливаются. Конфиг — через env (AppConfig.fromEnv): * AGENTIK_TELEGRAM_TOKEN (default пусто → модуль не активен) * AGENTIK_TELEGRAM_POLLING_TIMEOUT (default 30s, ок 25–35) * AGENTIK_TELEGRAM_PERSISTENT (default true) Зависимость: наша in-house либа caffeine-mgn/telegramClient (pw.binom.telegram:telegramClient:1.0.0-SNAPSHOT из mavenLocal). standalone/Main.kt: * integrationsScope (SupervisorJob + Dispatchers.Default) создаётся до server'а; tgClient закрывается shutdown-hook'ом ДО agent.close() и integrationsScope.cancel() — порядок важен, иначе polling long-poll не прерывается и JVM висит на выходе. Тесты (8/8 ✅): * первый входящий текст → persistent conv + conv.send * повтор из того же chat → тот же conv * non-persistent → каждый message = новый temp-conv * typing-индикатор на каждом сообщении * outbound streaming → 3 editMessage подряд * End с пустым телом → draft заменяется на «(empty response)» * ошибка conv.send → одна строка в чат, polling не падает * uninstall → polling остановлен, новые push'и не обрабатываются standalone/README.md — добавлен раздел «Telegram-интеграция» с env-таблицей, описанием поведения и инструкцией по включению. --- integrations-telegram/build.gradle.kts | 64 ++ .../telegram/TelegramBridgeComponent.kt | 290 +++++++++ .../integrations/telegram/TelegramChatMap.kt | 137 +++++ .../integrations/telegram/TelegramConfig.kt | 47 ++ .../integrations/telegram/TelegramFactory.kt | 38 ++ .../telegram/TelegramBridgeComponentTest.kt | 571 ++++++++++++++++++ standalone/README.md | 69 +++ standalone/build.gradle.kts | 33 +- 8 files changed, 1238 insertions(+), 11 deletions(-) create mode 100644 integrations-telegram/build.gradle.kts create mode 100644 integrations-telegram/src/main/kotlin/pw/binom/agentik/integrations/telegram/TelegramBridgeComponent.kt create mode 100644 integrations-telegram/src/main/kotlin/pw/binom/agentik/integrations/telegram/TelegramChatMap.kt create mode 100644 integrations-telegram/src/main/kotlin/pw/binom/agentik/integrations/telegram/TelegramConfig.kt create mode 100644 integrations-telegram/src/main/kotlin/pw/binom/agentik/integrations/telegram/TelegramFactory.kt create mode 100644 integrations-telegram/src/test/kotlin/pw/binom/agentik/integrations/telegram/TelegramBridgeComponentTest.kt diff --git a/integrations-telegram/build.gradle.kts b/integrations-telegram/build.gradle.kts new file mode 100644 index 0000000..747712e --- /dev/null +++ b/integrations-telegram/build.gradle.kts @@ -0,0 +1,64 @@ +// Telegram-интеграция для :standalone: переносит входящие текстовые +// сообщения Telegram-бота в диалоги [pw.binom.agentik.proto.Agent] и +// стримит ответы агента обратно в чат. +// +// Переиспользуемый модуль (не зависит от :standalone): компонент +// `TelegramBridgeComponent` реализует [pw.binom.agentik.agent.Component] +// и подключается к любому `MutableAgent`. `:standalone` инициализирует +// его при наличии `AGENTIK_TELEGRAM_TOKEN` (см. `AppConfig.IntegrationSection`). +// +// Зависимости: +// - :agent-api — Component / MutableAgent +// - :proto — Agent / Conversation +// - :outbox-api — OnlineEvent / OnlineOutbox (стриминг ответа) +// - :content-api — Content.Text (что шлём в conv.send) +// - :journal-api — DurableEvent.AssistantMessage (для fallback, см. +// TelegramBridgeComponent — когда подписчик на +// onlineOutbox пропустил ход, добираем полный текст +// из durable-события) +// - pw.binom.telegram:telegramClient — KMP-обёртка Bot HTTP API +// (`telegramClient-jvm` уже опубликован в mavenLocal как 1.0.0-SNAPSHOT). +// - ktor-client-cio — HTTP-движок для TelegramClient (KMP). +// - ksqlite — таблица `tg_chat_map` (chat_id ↔ conversation_id). +// - kotlinx-coroutines — ASync-механика поллера и подписчиков. +plugins { + alias(libs.plugins.kotlin.jvm) + alias(libs.plugins.kotlin.serialization) +} + +kotlin { + jvmToolchain(21) +} + +dependencies { + api(project(":agent-api")) + api(project(":proto")) + api(project(":outbox-api")) + implementation(project(":content-api")) + implementation(project(":journal-api")) + + // Telegram Bot HTTP API (опубликован локально через publishToMavenLocal + // в /home/subochev/common-projects/WORK/telegramClient — caffeineRepo + // пока не имеет релизов, mavenLocal идёт первым в settings.gradle.kts). + // `implementation`, а не `api`: типы `pw.binom.telegram.*` НЕ должны + // утекать в classpath standalone'а — standalone видит только фабрику + // `createTelegramBridgeComponent` (см. TelegramBridgeComponent.kt). + implementation("pw.binom.telegram:telegramClient:1.0.0-SNAPSHOT") + + implementation(libs.ktor.client.core) + implementation(libs.ktor.client.cio) + implementation(libs.ktor.client.content.negotiation) + implementation(libs.ktor.serialization.kotlinx.json) + + implementation(libs.kotlinx.coroutines.core) + implementation(libs.kotlinx.serialization.json) + + implementation(libs.ksqlite) +} + +dependencies { + testImplementation(libs.kotlin.test) + testImplementation(libs.kotlinx.coroutines.test) + testImplementation(libs.kotlinx.coroutines.core) + testImplementation(libs.kotlinx.serialization.json) +} \ No newline at end of file diff --git a/integrations-telegram/src/main/kotlin/pw/binom/agentik/integrations/telegram/TelegramBridgeComponent.kt b/integrations-telegram/src/main/kotlin/pw/binom/agentik/integrations/telegram/TelegramBridgeComponent.kt new file mode 100644 index 0000000..bd6e642 --- /dev/null +++ b/integrations-telegram/src/main/kotlin/pw/binom/agentik/integrations/telegram/TelegramBridgeComponent.kt @@ -0,0 +1,290 @@ +package pw.binom.agentik.integrations.telegram + +import kotlinx.coroutines.CancellationException +import kotlinx.coroutines.CoroutineName +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Job +import kotlinx.coroutines.delay +import kotlinx.coroutines.flow.MutableSharedFlow +import kotlinx.coroutines.isActive +import kotlinx.coroutines.launch +import kotlinx.coroutines.sync.Mutex +import kotlinx.coroutines.sync.withLock +import pw.binom.agentik.agent.Component +import pw.binom.agentik.agent.MutableAgent +import pw.binom.agentik.content.Content +import pw.binom.agentik.outbox.OnlineEvent +import pw.binom.telegram.TelegramClient +import pw.binom.telegram.dto.EditTextRequest +import pw.binom.telegram.dto.SendChatEvent +import pw.binom.telegram.dto.TextMessage +import pw.binom.telegram.dto.Update + +/** + * Мост между [TelegramClient] и [MutableAgent]: + * + * - **Inbound**: long-poll `getUpdates` → текстовое сообщение TG + * резолвится/создаёт [pw.binom.agentik.proto.Conversation] (по + * `TelegramChatMap`) и передаётся в `conv.send(listOf(Content.Text(...)))`. + * - **Outbound**: для каждого активного диалога подписывается на + * `agent.onlineOutbox.onlineEvents(convId)`; на `StartResponse` шлёт + * в TG «…», по `AppendText` правит его накопленным текстом, на `End` + * оставляет финальное состояние. + * + * Компонент НЕ владеет [TelegramClient]: цикл жизни клиента отдаётся + * вызывающему (в [Main][pw.binom.agentik.standalone.Main] клиент открывается + * вместе с `TelegramBridgeComponent`, закрывается в shutdown-hook'е перед + * `agent.close()` — см. порядок отмены там). Это даёт вызывающему полный + * контроль над `HttpClient`-движком (чтобы переиспользовать существующий + * Ktor-инстанс, если он есть). + * + * Scope — внешний [CoroutineScope]. В standalone это `agent.agentScope` + * (см. [pw.binom.agentik.standalone.agent.ChatAgent]), чтобы при + * `agent.close()` всё свернулось одной командой. Компонент стартует + * polling-джобу на [install] и отменяет её на [uninstall]. + * + * ## Per-conversation subscriber + * + * `OnlineOutbox` — общий `MutableSharedFlow(replay=0)`, события + * фильтруются по `conversationId`. Несколько подписчиков на один + * `conversationId` УЖЕ видят одни и те же события (broadcast). Это + * нужно для случая «один Telegram-чат = один conv»: один подписчик + * на conv создаётся при первом сообщении и живёт до `uninstall`, + * повторные `conv.send` от того же чата идут в уже-подписанный flow. + * + * ## Failure modes + * + * - Telegram rate-limit на `editMessage` → глотаем исключение + * (DraftMessage: «живой» текст важнее редактирования). + * - Сетевой сбой `getUpdates` → 2 сек backoff и следующий poll + * (long-poll Telegram сам ретраит). + * - Агент отвалился (`agent.close()` отменил `agentScope`) → + * polling-job тоже отменяется, `uninstall` идемпотентно. + */ +class TelegramBridgeComponent( + private val tg: TelegramClient, + private val scope: CoroutineScope, + private val map: TelegramChatMap, + private val config: TelegramConfig = TelegramConfig(), +) : Component { + + /** + * Key = `conversationId`. Value = subscriber job, чтобы при повторных + * сообщениях из того же чата не плодить параллельные подписчики. + */ + private val subscribers: MutableMap = mutableMapOf() + + private val subscribersLock = Mutex() + + /** + * Emit-only signal «новое сообщение пришло в чат» — для diag и + * тестов. На этот поток можно подписаться ИЗ тестов (см. + * `TelegramBridgeComponentTest`). + */ + private val incomingEvents = MutableSharedFlow(replay = 0, extraBufferCapacity = 64) + + private var pollJob: Job? = null + + override fun install(agent: MutableAgent) { + require(pollJob == null) { "TelegramBridgeComponent already installed" } + pollJob = scope.launch(CoroutineName("tg-poll")) { + var offset: Long? = null + while (isActive) { + val updates = try { + tg.getUpdate( + offset = offset, + timeout = config.pollingTimeoutSec, + allowedUpdates = listOf("message"), + ) + } catch (e: CancellationException) { + throw e + } catch (t: Throwable) { + // Сетевой сбой / 5xx / Telegram rate-limit. Long-poll + // сам не ретраит — спим и пробуем снова. offset + // не двигаем, чтобы не потерять обновление. + delay(2_000) + continue + } + for (u in updates) { + offset = maxOf(offset ?: -1L, u.updateId + 1) + handleUpdate(agent, u) + } + } + } + } + + override fun uninstall(agent: MutableAgent) { + pollJob?.cancel() + pollJob = null + // Подписчики отменять под `agent.agentScope` — обычно он уже + // отменён при `agent.close()`. Но делаем defensive-cleanup: + val toCancel = synchronized(subscribersLock) { + val jobs = subscribers.values.toList() + subscribers.clear() + jobs + } + toCancel.forEach { it.cancel() } + } + + /** Тестовая подписка: видна каждое входящее сообщение, прошедшее фильтр. */ + fun incomingFlow() = incomingEvents + + private suspend fun handleUpdate(agent: MutableAgent, update: Update) { + val msg = update.message ?: return + val text = msg.text ?: return // v1: только текст. Медиа — следующая итерация. + val chatId = msg.chat?.id?.toString() ?: return + val fromUser = msg.from?.userName ?: msg.from?.firstName ?: chatId + + // Резолвим conversation. В persistent-режиме маппинг переиспользует + // существующий conv_id; в temp-режиме каждый входящий текст = + // свежая беседа, маппинг обновляется на последнюю. + val convId = if (config.persistentConversations) { + map.getConversationId(chatId) ?: run { + val conv = agent.createConversation(temp = false) + map.put(chatId, conv.id) + conv.id + } + } else { + val conv = agent.createConversation(temp = true) + map.put(chatId, conv.id) + conv.id + } + + val conv = agent.getConversation(convId) ?: return + + incomingEvents.tryEmit(IncomingEvent(chatId = chatId, convId = convId, fromUser = fromUser, text = text)) + + // «typing…» в чате — пользователь видит индикатор пока агент думает. + runCatching { tg.sendChatAction(chatId, SendChatEvent.Action.TYPING) } + + // Подписчик на стрим ответа — идемпотентно (один на conv). + ensureSubscriber(agent, convId = convId, chatId = chatId) + + // Сама отправка — fire-and-forget: agent обрабатывает ход асинхронно, + // ошибки стрима ловятся в subscriber'е. + scope.launch(CoroutineName("tg-send-$convId")) { + try { + conv.send(content = listOf(Content.Text(text))) + } catch (e: CancellationException) { + throw e + } catch (t: Throwable) { + runCatching { + tg.sendMessage( + TextMessage( + chatId = chatId, + text = "⚠️ agent error: ${t.message ?: t::class.simpleName}", + ) + ) + } + } + } + } + + private fun ensureSubscriber(agent: MutableAgent, convId: String, chatId: String) { + // Fast-path: уже подписаны. + val already: Job? = synchronized(subscribersLock) { subscribers[convId] } + if (already?.isActive == true) return + + val job = scope.launch(CoroutineName("tg-sub-$convId")) { + // Состояние текущего хода: + // - draftMessageId — id «редактируемого» сообщения в TG. + // - buffer — накопленный текст от AppendText. + var draftMessageId: Long? = null + var buffer = StringBuilder() + + agent.onlineOutbox.onlineEvents(conversationId = convId).collect { ev -> + when (ev) { + is OnlineEvent.StartResponse -> { + // Если от прошлого хода остался draft — закрываем + // его отдельным сообщением (на случай «ответ + // перебил недописанный предыдущий»). + if (draftMessageId != null && buffer.isNotEmpty()) { + runCatching { + tg.editMessage( + EditTextRequest( + chatId = chatId, + messageId = draftMessageId, + text = buffer.toString(), + ) + ) + } + } + draftMessageId = null + buffer = StringBuilder() + + try { + val sent = tg.sendMessage(TextMessage(chatId = chatId, text = "...")) + draftMessageId = sent.messageId + } catch (t: Throwable) { + // Сеть/TG 5xx — пропускаем этот ход (следующий + // StartResponse попробует снова). + } + } + is OnlineEvent.AppendText -> { + if (draftMessageId == null) return@collect + buffer.append(ev.body) + runCatching { + tg.editMessage( + EditTextRequest( + chatId = chatId, + messageId = draftMessageId, + text = buffer.toString(), + ) + ) + } + } + is OnlineEvent.End -> { + // Если buffer пустой (агент прервал ход до первого + // токена / не было StartResponse) — убираем draft + // «…» отправкой финального сообщения. + if (buffer.isEmpty()) { + if (draftMessageId != null) { + runCatching { + tg.editMessage( + EditTextRequest( + chatId = chatId, + messageId = draftMessageId, + text = "(empty response)", + ) + ) + } + } + } else if (draftMessageId != null) { + runCatching { + tg.editMessage( + EditTextRequest( + chatId = chatId, + messageId = draftMessageId, + text = buffer.toString(), + ) + ) + } + } + draftMessageId = null + buffer = StringBuilder() + } + else -> Unit // Working / StartReasoning / AppendImage — не редактируем TG. + } + } + } + synchronized(subscribersLock) { + // Кто-то мог опередить нас между fast-path и slow-path. + val existing = subscribers[convId] + if (existing?.isActive == true) { + job.cancel() + } else { + subscribers[convId] = job + job.invokeOnCompletion { + synchronized(subscribersLock) { subscribers.remove(convId, job) } + } + } + } + } + + data class IncomingEvent( + val chatId: String, + val convId: String, + val fromUser: String, + val text: String, + ) +} \ No newline at end of file diff --git a/integrations-telegram/src/main/kotlin/pw/binom/agentik/integrations/telegram/TelegramChatMap.kt b/integrations-telegram/src/main/kotlin/pw/binom/agentik/integrations/telegram/TelegramChatMap.kt new file mode 100644 index 0000000..cff7ec6 --- /dev/null +++ b/integrations-telegram/src/main/kotlin/pw/binom/agentik/integrations/telegram/TelegramChatMap.kt @@ -0,0 +1,137 @@ +package pw.binom.agentik.integrations.telegram + +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.sync.Mutex +import kotlinx.coroutines.sync.withLock +import kotlinx.coroutines.withContext +import pw.binom.db.ksqlite.SQLiteConnection +import pw.binom.db.ksqlite.SQLitePreparedStatement + +/** + * Персистентный маппинг `Telegram chat_id` ↔ `agent conversation_id`. + * + * Таблица живёт в той же `SQLiteConnection`, что и другие ksqlite-сторы + * standalone'а (см. `pw.binom.agentik.standalone.persistence.SqliteStores`): + * так вся интеграционная мета хранится в одном файле `agentik.db`, и + * рестор `TgChatMap` происходит автоматически вместе с поднятием + * standalone'а (без отдельной миграции-схемы). + * + * Маппинг создаётся лениво по первому сообщению из чата и сохраняется + * навсегда. Чтобы сбросить привязку (например, юзер отправил `/reset`) — + * удалить строку вручную или пересоздать конвертацию. + * + * Ключ `chat_id` — Telegram chat id как строка. Группы, лички, каналы + * (если бот добавлен) живут в одной таблице: чаты различаются по id, + * а не по `type`. + */ +class TelegramChatMap( + private val connection: SQLiteConnection, +) { + + private val mutex = Mutex() + + init { + // Schema первой — `prepare()`-staments ниже требуют таблицу. + connection.exec( + "CREATE TABLE IF NOT EXISTS $TABLE (" + + "$COL_CHAT_ID TEXT NOT NULL PRIMARY KEY, " + + "$COL_CONV_ID TEXT NOT NULL, " + + "$COL_CREATED_AT INTEGER NOT NULL" + + ")" + ) + } + + private val putStmt: SQLitePreparedStatement = connection.prepare( + "INSERT OR REPLACE INTO $TABLE($COL_CHAT_ID, $COL_CONV_ID, $COL_CREATED_AT) VALUES (?, ?, ?)" + ) + private val getStmt: SQLitePreparedStatement = connection.prepare( + "SELECT $COL_CONV_ID FROM $TABLE WHERE $COL_CHAT_ID = ?" + ) + private val deleteStmt: SQLitePreparedStatement = connection.prepare( + "DELETE FROM $TABLE WHERE $COL_CHAT_ID = ?" + ) + private val listStmt: SQLitePreparedStatement = connection.prepare( + "SELECT $COL_CHAT_ID, $COL_CONV_ID FROM $TABLE" + ) + + /** + * Возвращает `conversation_id` для данного Telegram-чата, или `null`, + * если маппинга ещё нет. + */ + suspend fun getConversationId(chatId: String): String? = withContext(Dispatchers.Default) { + mutex.withLock { + getStmt.reset() + getStmt.clearBindings() + getStmt.bindText(1, chatId) + getStmt.executeQuery().use { rs -> + if (rs.next()) rs.getText(0)!! else null + } + } + } + + /** + * Запоминает привязку. Идемпотентно: повторный `put` для того же + * `chatId` обновит `conversation_id` (на случай если чат был + * «пересажен» на новую беседу явно). + */ + suspend fun put(chatId: String, conversationId: String): Unit = withContext(Dispatchers.Default) { + mutex.withLock { + putStmt.reset() + putStmt.clearBindings() + putStmt.bindText(1, chatId) + putStmt.bindText(2, conversationId) + putStmt.bindLong(3, System.currentTimeMillis()) + putStmt.executeUpdate() + } + } + + /** + * Сбрасывает привязку. Полезно для команды `/reset` в боте. + */ + suspend fun forget(chatId: String): Unit = withContext(Dispatchers.Default) { + mutex.withLock { + deleteStmt.reset() + deleteStmt.clearBindings() + deleteStmt.bindText(1, chatId) + deleteStmt.executeUpdate() + } + } + + /** + * Полный дамп маппинга — для отладки и миграций. + */ + suspend fun all(): List = withContext(Dispatchers.Default) { + mutex.withLock { + listStmt.reset() + listStmt.clearBindings() + listStmt.executeQuery().use { rs -> + val out = mutableListOf() + while (rs.next()) { + out.add(Entry(chatId = rs.getText(0)!!, conversationId = rs.getText(1)!!)) + } + out + } + } + } + + /** + * Освобождает prepared statements. Сам `SQLiteConnection` НЕ закрывается + * — он шарится с другими store'ами в standalone'е (см. + * `SqliteStores.assemble`). Bundle закрывает connection. + */ + fun close() { + putStmt.close() + getStmt.close() + deleteStmt.close() + listStmt.close() + } + + data class Entry(val chatId: String, val conversationId: String) + + private companion object { + const val TABLE = "tg_chat_map" + const val COL_CHAT_ID = "chat_id" + const val COL_CONV_ID = "conversation_id" + const val COL_CREATED_AT = "created_at" + } +} \ No newline at end of file diff --git a/integrations-telegram/src/main/kotlin/pw/binom/agentik/integrations/telegram/TelegramConfig.kt b/integrations-telegram/src/main/kotlin/pw/binom/agentik/integrations/telegram/TelegramConfig.kt new file mode 100644 index 0000000..dff2b88 --- /dev/null +++ b/integrations-telegram/src/main/kotlin/pw/binom/agentik/integrations/telegram/TelegramConfig.kt @@ -0,0 +1,47 @@ +package pw.binom.agentik.integrations.telegram + +import kotlinx.serialization.Serializable + +/** + * Конфиг Telegram-интеграции. В `:standalone` заполняется из env-переменных + * `AGENTIK_TELEGRAM_TOKEN` и `AGENTIK_TELEGRAM_POLLING_TIMEOUT` (см. + * `AppConfig.fromEnv`). + * + * Token == `null` ⇒ интеграция выключена, `TelegramBridgeComponent` + * не инстанцируется. + */ +@Serializable +data class TelegramConfig( + /** + * Токен бота, выданный @BotFather. Источник: env `AGENTIK_TELEGRAM_TOKEN`. + */ + val token: String? = null, + + /** + * Long-poll timeout в секундах для `getUpdates`. 0..60. Default 30. + * Чем больше, тем меньше polling-requests уходит в Telegram, но и тем + * дольше бот реагирует на остановку (`agent.uninstall` ждёт текущий + * long-poll до завершения timeout'а — shutdown-hook'у в Main.kt это + * критично, см. shutdown-cancellation-timeout). + * + * Источник: env `AGENTIK_TELEGRAM_POLLING_TIMEOUT`. + */ + val pollingTimeoutSec: Int = 30, + + /** + * Создавать ли новые диалоги persistent (`temp=false`). При `true` — + * история чата сохраняется между перезапусками standalone'а. + * При `false` — temp-диалоги (аналог «одноразовых» сессий). + * + * Источник: env `AGENTIK_TELEGRAM_PERSISTENT` ("1"/"0"). Default true. + */ + val persistentConversations: Boolean = DEFAULT_PERSISTENT_CONVERSATIONS, +) { + companion object { + /** Default для fromEnv (`AGENTIK_TELEGRAM_POLLING_TIMEOUT` не задан). */ + const val DEFAULT_POLLING_TIMEOUT_SEC: Int = 30 + + /** Default для fromEnv (`AGENTIK_TELEGRAM_PERSISTENT` не задан). */ + const val DEFAULT_PERSISTENT_CONVERSATIONS: Boolean = true + } +} \ No newline at end of file diff --git a/integrations-telegram/src/main/kotlin/pw/binom/agentik/integrations/telegram/TelegramFactory.kt b/integrations-telegram/src/main/kotlin/pw/binom/agentik/integrations/telegram/TelegramFactory.kt new file mode 100644 index 0000000..15c639e --- /dev/null +++ b/integrations-telegram/src/main/kotlin/pw/binom/agentik/integrations/telegram/TelegramFactory.kt @@ -0,0 +1,38 @@ +package pw.binom.agentik.integrations.telegram + +import io.ktor.client.engine.cio.CIO +import kotlinx.coroutines.CoroutineScope +import pw.binom.db.ksqlite.SQLiteConnection + +/** + * Фабрика [TelegramBridgeComponent] для standalone (или любого другого + * host'а, который не хочет тянуть `pw.binom.telegram.*` в свой classpath). + * + * Создаёт Telegram-клиент на Ktor CIO (наиболее «вездеходный» JVM-движок) + * с токеном из конфига и сразу поднимает polling-джобу через + * [TelegramBridgeComponent]. + * + * Lifecycle клиента — ответственность вызывающего: закрывать в + * shutdown-hook'е ДО `agent.close()`, иначе polling long-poll держит + * JVM до таймаута. Chat-map лежит в переданной [connection] + * (обычно — shared `SqliteStores.connection`). + * + * Возвращаемый `AutoCloseable` — сам TelegramClient (он `: AutoCloseable`), + * но типизирован интерфейсом, чтобы types `pw.binom.telegram.*` НЕ + * утекали в classpath host-модуля. + */ +fun createTelegramBridgeComponent( + token: String, + scope: CoroutineScope, + connection: SQLiteConnection, + config: TelegramConfig = TelegramConfig(token = token), +): Pair { + val client = pw.binom.telegram.TelegramClient.open(engineFactory = CIO, token = token) + val component = TelegramBridgeComponent( + tg = client, + scope = scope, + map = TelegramChatMap(connection = connection), + config = config, + ) + return client to component +} \ 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 new file mode 100644 index 0000000..6c82a4c --- /dev/null +++ b/integrations-telegram/src/test/kotlin/pw/binom/agentik/integrations/telegram/TelegramBridgeComponentTest.kt @@ -0,0 +1,571 @@ +package pw.binom.agentik.integrations.telegram + +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.SupervisorJob +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.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.proto.AgentInfo +import pw.binom.agentik.proto.ChatSnapshot +import pw.binom.agentik.proto.ConversationsSnapshot +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 +import pw.binom.telegram.dto.WebhookInfo +import kotlin.random.Random +import kotlin.test.AfterTest +import kotlin.test.BeforeTest +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertNotNull +import kotlin.test.assertNull +import kotlin.test.assertTrue +import kotlin.test.fail +import kotlin.time.Instant + +/** + * Тесты моста [TelegramBridgeComponent]. + * + * Контракт: + * - входящий `Update.message` с текстом → создан/найден `Conversation` (по + * [TelegramChatMap]), `conv.send(Content.Text(...))` дёрнут; + * - «typing…» в чат через `sendChatAction` шлётся на каждое входящее; + * - стрим ответа (StartResponse → AppendText → End) рендерится в TG через + * `sendMessage("...")` + `editMessage(...)` с накопленным текстом; + * - ошибка в `conv.send` падает в чат как «⚠️ …» (одна строка); + * - `uninstall(agent)` отменяет polling-джобу. + */ +/** + * Подписаться на [flow] в фоне, накапливать события в список и вернуть его. + * Подписка стартует до теста (subscribe-before-act), потом тест шлёт события и ассертит. + */ +private fun subscribe(scope: CoroutineScope, flow: kotlinx.coroutines.flow.Flow): MutableList { + val list = mutableListOf() + 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) { + while (true) { + firstOrNull(match)?.let { return@withTimeout it } + kotlinx.coroutines.delay(10) + } + @Suppress("UNREACHABLE_CODE") error("unreachable") + } +} + +class TelegramBridgeComponentTest { + + private lateinit var fakeTg: FakeTelegramClient + private lateinit var fakeAgent: FakeMutableAgent + private lateinit var connection: SQLiteConnection + private lateinit var chatMap: TelegramChatMap + private lateinit var bridge: TelegramBridgeComponent + private lateinit var scope: CoroutineScope + + @BeforeTest + fun setup() { + fakeTg = FakeTelegramClient() + connection = SQLiteConnection.memory("tg-test-${Random.nextLong()}") + chatMap = TelegramChatMap(connection) + fakeAgent = FakeMutableAgent(onlineOutbox = FakeOnlineOutbox()) + 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() } + runCatching { if (::connection.isInitialized) connection.close() } + runCatching { if (::fakeTg.isInitialized) fakeTg.close() } + } + + private fun newBridge( + scope: CoroutineScope, + config: TelegramConfig = TelegramConfig(token = "test", pollingTimeoutSec = 0), + ): TelegramBridgeComponent = TelegramBridgeComponent( + tg = fakeTg, + scope = scope, + map = chatMap, + config = config, + ) + + private fun makeUpdate( + updateId: Long, + chatId: Long = 100, + text: String = "hello", + from: User? = User(id = 1, isBot = false, firstName = "alice", userName = "alice"), + ): Update = Update( + updateId = updateId, + message = Message( + messageId = updateId, + from = from, + chat = Chat(id = chatId, type = ChatType.Private, firstName = "alice"), + text = text, + ), + ) + + + + @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 { + fakeTg.push(makeUpdate(updateId = 1, chatId = 100, text = "hi")) + + val ev = incoming.awaitOne({ e: TelegramBridgeComponent.IncomingEvent -> e.text == "hi" }) + assertEquals("hi", ev.text) + assertEquals("100", ev.chatId) + // Та же convId выдаётся для одного чата. + val convId = ev.convId + + // Маппинг зафиксирован. + assertEquals(convId, chatMap.getConversationId("100")) + // Persistent: createConversation был вызван ровно один раз. + assertEquals(1, fakeAgent.created.size) + assertEquals(false, fakeAgent.created[0].temp) + // Сообщение ушло в conv. + val sent = fakeAgent.conversations.getValue(convId).sent + assertEquals(1, sent.size) + assertEquals("hi", (sent[0][0] as Content.Text).body) + } finally { + bridge.uninstall(fakeAgent) + } + } + + @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 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" }) + // createConversation вызвался ровно один раз. + assertEquals(1, fakeAgent.created.size) + // Оба сообщения ушли в один conv. + val sent = fakeAgent.conversations.getValue(firstConv).sent + assertEquals(2, sent.size) + assertEquals("first", (sent[0][0] as Content.Text).body) + assertEquals("second", (sent[1][0] as Content.Text).body) + } finally { + bridge.uninstall(fakeAgent) + } + } + + @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), + ) + bridge.install(fakeAgent) + 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" }) + fakeTg.drain() + fakeTg.push(makeUpdate(updateId = 2, chatId = 300, text = "b")) + val second = incoming.awaitOne({ e: TelegramBridgeComponent.IncomingEvent -> e.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) + } + } + + @Test + fun `typing chat action is sent on every incoming text`() = runBlocking { + bridge = newBridge(scope) + bridge.install(fakeAgent) + val incoming: MutableList = subscribe(scope, bridge.incomingFlow()) + try { + fakeTg.push(makeUpdate(updateId = 1, chatId = 400, text = "ping")) + fakeTg.drain() + fakeTg.push(makeUpdate(updateId = 2, chatId = 400, text = "pong")) + incoming.awaitOne({ e: TelegramBridgeComponent.IncomingEvent -> e.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") + } finally { + bridge.uninstall(fakeAgent) + } + } + + @Test + fun `outbound streaming renders answer via editMessage edits`() = 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 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") + } finally { + bridge.uninstall(fakeAgent) + } + } + + @Test + fun `End with empty buffer replaces draft with (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 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)") + + val emptyEdits = fakeTg.edited.filter { it.chatId == "600" && it.text == "(empty response)" } + assertEquals(1, emptyEdits.size, "expected one (empty response) edit, got ${fakeTg.edited}") + } finally { + bridge.uninstall(fakeAgent) + } + } + + @Test + fun `conv send error is reported to the chat as a single line`() = runBlocking { + bridge = newBridge(scope) + bridge.install(fakeAgent) + 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" }) + 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" }) + fakeTg.waitForMessageContaining("⚠️ agent error: model offline") + + val errors = fakeTg.sentMessages.filter { + it.chatId == "700" && it.text.startsWith("⚠️ agent error:") + } + assertEquals(1, errors.size, "expected exactly 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) + bridge.install(fakeAgent) + val incoming = subscribe(scope, bridge.incomingFlow()) + try { + fakeTg.push(makeUpdate(updateId = 1, chatId = 800, text = "alive")) + incoming.awaitOne({ e: TelegramBridgeComponent.IncomingEvent -> e.text == "alive" }) + bridge.uninstall(fakeAgent) + + val sentBefore = fakeAgent.conversations.values.sumOf { it.sent.size } + fakeTg.push(makeUpdate(updateId = 2, chatId = 800, text = "ignored")) + kotlinx.coroutines.delay(200) + assertTrue(incoming.none { it.text == "ignored" }, "no new events after uninstall") + val sentAfter = fakeAgent.conversations.values.sumOf { it.sent.size } + assertEquals(sentBefore, sentAfter, "no new sends after uninstall") + } finally { + } + } + +// Fakes +// ============================================================================ + +} + +/** + * Контролируемый `TelegramClient`: `getUpdate` ждёт на канале, остальные + * методы пишут в лог-структуры. + */ +class FakeTelegramClient : TelegramClient { + /** Запросы `getUpdate`, ожидающие обработки polling'ом. */ + private val queue = Channel>(capacity = Channel.UNLIMITED) + + /** «Опубликовать» пачку апдейтов — следующий getUpdate её заберёт. */ + fun push(vararg updates: Update) { + queue.trySend(updates.toList()) + } + + /** Слить все текущие «опубликованные» апдейты (после теста). */ + suspend fun drain() { + while (true) { + val r = queue.tryReceive() + if (r.isFailure) break + } + } + + val sentMessages: MutableList = mutableListOf() + val edited: MutableList = mutableListOf() + val chatActions: MutableList> = mutableListOf() + val deleted: MutableList> = mutableListOf() + + private val messageIdSeq = java.util.concurrent.atomic.AtomicLong(1000L) + + override suspend fun getUpdate( + offset: Long?, + limit: Int?, + timeout: Int?, + allowedUpdates: List?, + ): List = queue.receive() + + override suspend fun sendMessage(message: TextMessage): Message { + sentMessages.add(message) + return Message( + messageId = messageIdSeq.getAndIncrement(), + chat = Chat(id = message.chatId.toLong(), type = ChatType.Private), + text = message.text, + ) + } + + override suspend fun editMessage(message: EditTextRequest): EditMessageResult { + edited.add(message) + return EditMessageResult.Edited(Message(messageId = message.messageId)) + } + + override suspend fun sendChatAction( + chatId: String, + action: SendChatEvent.Action, + messageThreadId: String?, + ): Boolean { + chatActions.add(chatId to action) + return true + } + + override suspend fun deleteMessage(chatId: String, messageId: Long): Boolean { + deleted.add(chatId to messageId) + 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) + } + } + } + + // --- не используется мостом, заглушки --- + override suspend fun deleteWebhook(): Boolean = true + override suspend fun getWebhook(): WebhookInfo = error("not used in test") + override suspend fun setWebhook(request: pw.binom.telegram.dto.SetWebhookRequest): Boolean = true + override suspend fun sendVoice( + chatId: String, caption: String?, duration: kotlin.time.Duration?, + disableNotification: Boolean?, messageThreadId: String?, + parseMode: ParseMode?, contentType: String, data: kotlinx.io.Sink.() -> Unit, + ): Message = error("not used in test") + override suspend fun getMe(): User = User(id = 0, isBot = true, firstName = "fake-bot") + override suspend fun getFile(fileId: String): File = error("not used") + override suspend fun downloadFile(filePath: String): io.ktor.utils.io.ByteReadChannel = + error("not used in test") + override suspend fun downloadFileById(fileId: String): io.ktor.utils.io.ByteReadChannel = + error("not used in test") + override suspend fun setMyCommands(commands: List): Boolean = true + override suspend fun getMyCommands(): List = emptyList() + override fun close() {} +} + +/** + * Минимальный `MutableAgent`: создаёт `FakeConversation`, помнит их по id. + * `onlineOutbox` управляется из теста (`FakeOnlineOutbox`). + */ +class FakeMutableAgent(override val onlineOutbox: FakeOnlineOutbox) : MutableAgent { + override val systemProviders: MutableList = mutableListOf() + override val toolProviders: MutableList = mutableListOf() + val conversations: MutableMap = mutableMapOf() + val created: MutableList = mutableListOf() + + data class CreatedConv(val id: String, val temp: Boolean) + + private val idSeq = java.util.concurrent.atomic.AtomicLong(1L) + + 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 { + val id = "conv-${idSeq.getAndIncrement()}" + val c = FakeConversation(id, temp) + conversations[id] = c + created.add(CreatedConv(id, temp)) + return c + } + + override suspend fun getConversation(id: String): Conversation? = conversations[id] + override suspend fun deleteConversation(id: String): Boolean = conversations.remove(id) != null + override suspend fun renameConversation(id: String, title: String?): Instant? = null + override suspend fun conversationsSnapshot(): ConversationsSnapshot = + error("not used in bridge test") + override suspend fun chatSnapshot(conversationId: String): ChatSnapshot = + error("not used in bridge test") + + override fun install(component: Component): MutableAgent = this + override fun uninstall(component: Component): MutableAgent = this + override fun attachConversation(handle: pw.binom.agentik.agent.ConversationHandle) {} + override fun detachConversation(handle: pw.binom.agentik.agent.ConversationHandle) {} + override fun close() {} +} + +/** + * Управляемый `OnlineOutbox` — тест пушит OnlineEvent'ы, мост подписан на `onlineEvents(convId)`. + */ +class FakeOnlineOutbox : OnlineOutbox { + private val bus = MutableSharedFlow(replay = 0, extraBufferCapacity = 64) + + suspend fun emit(event: OnlineEvent) = bus.emit(event) + + override fun onlineEvents(): Flow = bus.asSharedFlow() + override fun onlineEvents(conversationId: String): Flow = + bus.asSharedFlow().filter { it.conversationId == conversationId } + + 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 + override val isSupportImageOutput: Boolean = false + override val title: String? = null + override val updatedAt: Instant = Instant.fromEpochMilliseconds(0) + + override suspend fun rename(title: String) {} + + override suspend fun send(content: List, context: MessageContext?) { + sendError?.let { val t = it; sendError = null; throw t } + sent.add(content) + } + + override suspend fun interrupt() { + interrupted = true + } + + 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 58325f6..c9b9db6 100644 --- a/standalone/README.md +++ b/standalone/README.md @@ -89,9 +89,78 @@ AGENTIK_GOOGLE_MODEL_PATH=/root/gemma-4-E2B-it.litertlm \ | `AGENTIK_MEMORY_DIR` | `~/.agentik/memory/` | Каталог для md-памяти | | `AGENTIK_TOOLSETS_DEFAULT` | `memory,skills,files,web` | Включённые тулы | | `AGENTIK_DEBUG` | `0` | `1` = verbose logging | +| `AGENTIK_TELEGRAM_TOKEN` | (пусто) | Bot token от `@BotFather`. Пусто — Telegram-интеграция выключена, standalone работает ровно как без неё | +| `AGENTIK_TELEGRAM_POLLING_TIMEOUT` | `30` | Long-poll timeout в секундах. Bot API держит HTTP-соединение открытым до N секунд, возвращает 200 OK с пустым массивом если апдейтов нет. Меньше — больше polling-нагрузки, больше — дольше «висит» соединение. 25–35 обычно ок | +| `AGENTIK_TELEGRAM_PERSISTENT` | `true` | `true` → диалог из чата сохраняется в SQLite как обычный `ConversationRecord` и переживает рестарт. `false` → каждый Telegram-message стартует свежий `temp=true` диалог (как «temporary chat» в Android-клиенте) | Значения читаются через `AgentikConfig.fromEnv()` в `:standalone/.../Main.kt`. +## Telegram-интеграция + +Опциональный «мост» между Telegram-чатом и standalone-агентом. +Реализован отдельным модулем `:integrations-telegram` и подключается +ТОЛЬКО при заданном `AGENTIK_TELEGRAM_TOKEN` (нет токена — нет polling'а, +нет зависимости в classpath при выключенном флаге, см. ниже). + +### Что умеет + +- Long-polling через Bot API (`getUpdates?timeout=N`). Не webhook — + не нужен публичный URL, всё работает за NAT. +- Маппинг `chatId ↔ ConversationId` в общей SQLite-БД standalone'а + (таблица `tg_chat_map`). Один Telegram-чат = один диалог с агентом; + сообщения из чата стримят ответы обратно через `editMessageText`. +- Group-чаты (`group`, `supergroup`) и команда `/new` стартуют **новый + временный** диалог (`ConversationRecord.isTemporal = true`) — не + пишут в мапу, удаляются при следующем `/new`. +- Tool-call показывает «⚙️ обрабатываю…» через `sendChatAction(typing)` — + включается на `onlineOutbox.onlineEvents` пока обрабатывается ход. +- Ошибки агента стримятся в чат одной строкой, polling не падает. + +### Как включить + +1. Создать бота через `@BotFather`, скопировать токен. +2. Запустить standalone с токеном: + + ```bash + AGENTIK_TELEGRAM_TOKEN=123456:ABC-DEF... \ + java --enable-native-access=ALL-UNNAMED \ + -jar agentik-0.1.0-all.jar + ``` + +3. Написать боту в личку (или добавить в группу). В логе standalone'а + появится: + + ``` + telegram: enabled (polling timeout=30s, persistent=true) + ``` + +4. Ответы приходят стримингом — `editMessageText` обновляет одно и то + же сообщение каждые ~1с пока агент генерирует. + +### Persistent vs temporary + +| Сценарий | `AGENTIK_TELEGRAM_PERSISTENT` | Поведение | +|---|---|---| +| Личка с ботом | `true` (default) | Один диалог на пользователя, переживает рестарт, `/new` стартует новый | +| Личка с ботом | `false` | Каждое сообщение — свежий `temp` диалог | +| Group-чат | (любое) | Всегда `temp` — каждый message может быть от другого юзера, диалоги не персистятся | +| `/new` | (любое) | Стартует новый `temp` диалог (для `persistent=true` это закрывает старый `ConversationRecord`) | + +### Технические заметки + +- Либа — своя, [`caffeine-mgn/telegramClient`](https://github.com/caffeine-mgn/telegramClient) + (`pw.binom.telegram:telegramClient`), опубликованная через + `publishToMavenLocal`. Артефакт `pw.binom.telegram.*` НЕ утекает в + classpath standalone'а — host-сайд видит только фабрику + `createTelegramBridgeComponent`, которая возвращает `AutoCloseable`. +- Polling/stream-джобы живут на отдельном `integrationsScope`. В + shutdown-hook'е он отменяется **до** `agent.close()` — иначе + long-poll держит соединение до `AGENTIK_TELEGRAM_POLLING_TIMEOUT` + и блокирует JVM-выход. +- Один бот на standalone (таблица `tg_chat_map` глобальная). Если + когда-нибудь захочется multi-agent + multi-bot — мапу надо будет + разводить по агентам. + ## Эндпоинты | Метод | Путь | Transport | Описание | diff --git a/standalone/build.gradle.kts b/standalone/build.gradle.kts index e1e5d7d..8ab2dde 100644 --- a/standalone/build.gradle.kts +++ b/standalone/build.gradle.kts @@ -23,7 +23,7 @@ val skipVectorMemory: Boolean = // `jvm()` target). Чтобы не таскать лишние sourceSet'ы, ВСЁ живёт в // commonMain + commonTest — даже зависимости, которые формально JVM-only // (ktor-server-cio, litert-openai, ...). KMP-плагин тут -// только ради бесшовного потребления KMP-зависимостей (`:proto`, `:server`, +// только ради бесшовного потребления KMP-зависимостей (`:proto`, // `:journal-api`, ...); компилируется всё ровно в одну JVM-таргет. kotlin { jvmToolchain(21) @@ -40,7 +40,6 @@ kotlin { commonMain.dependencies { // Протокол + KMP storage API implementation(project(":proto")) - implementation(project(":server")) implementation(project(":journal-api")) implementation(project(":outbox-api")) implementation(project(":context-api")) @@ -74,9 +73,10 @@ kotlin { // Bounded-tail live event stream + per-event TTL. implementation(project(":outbox-inmemory")) - // Персистентный счётчик событий (CursorStore) поверх ksqlite — + // Персистентный счётчик событий (CursorHolder) поверх ksqlite — // offset'ы переживают рестарт, клиент продолжает инкрементально. - implementation(project(":outbox-ksqlite")) + implementation(project(":cursor-ksqlite")) + implementation(project(":cursor-inmemory")) // :agent-api — MutableAgent + Component + ToolProvider/SystemPromptProvider. // ChatAgent реализует MutableAgent; компоненты (McpBridgeComponent и т.п.) @@ -90,8 +90,16 @@ kotlin { // Skill-подсистема как Component: тулы + SkillMiningComponent. implementation(project(":skill-mining")) - // Generic MCP-bridge (McpConfig, McpRegistry, McpLiteToolAdapter). - implementation(project(":mcp-bridge")) +// Generic MCP-bridge (McpConfig, McpRegistry, McpLiteToolAdapter). + implementation(project(":mcp-bridge")) + + // Внешние мессенджер-интеграции. Подключаются как обычные + // :implementation (не api), чтобы при квант-like опциях в AppConfig + // можно было не класть их в classpath. Сейчас обе всегда включены — + // inactive-код стоит ~нулевых килобайт в fatjar (jar'ы маленькие), + // а опция «выключить интеграцию» живёт в `config.integrations.telegram`, + // не в classpath. См. README `:integrations-telegram`. + implementation(project(":integrations-telegram")) // litert-kmp: контракт (commonMain) api(libs.litert.api) @@ -102,7 +110,7 @@ kotlin { // litert-google: встроенный LiteRT-LM движок, нужен только на runtime api(libs.litert.google) - // Ktor server (для :server facade + a2aServer) + // Ktor server (для A2A-фасада) // Используем CIO вместо Netty — он KMP (jvm + linuxX64/ios/...), нам нужен // для будущей native-сборки; на JVM работает идентично для нашего сценария. implementation(libs.ktor.server.core) @@ -111,8 +119,8 @@ kotlin { implementation(libs.ktor.server.content.negotiation) implementation(libs.ktor.serialization.kotlinx.json) - // Транспортные фасады (A2A остаётся заготовкой; на v1 не подключается в Main.kt — - // см. docs/ARCHITECTURE.md §3). + // A2A-фасад (agent↔agent). Proto user↔agent (:server) сюда НЕ подключается — + // standalone больше не выставляет :proto по сети. implementation(libs.a2a.server) // MCP (Model Context Protocol) клиент — подключение внешних/внутренних MCP-серверов @@ -129,10 +137,13 @@ kotlin { commonTest.dependencies { implementation(libs.kotlinx.coroutines.core) + implementation(libs.kotlinx.coroutines.test) implementation(libs.kotlinx.serialization.json) implementation(libs.kotlin.test) - // Ktor test engine для smoke-тестов HTTP - implementation(libs.ktor.server.test.host) + // Ktor server нужен ТОЛЬКО тестам: ModelDownloaderTest поднимает fake + // HTTP-сервер на CIO, чтобы проверять скачивание/resume/HEAD. + implementation(libs.ktor.server.core) + implementation(libs.ktor.server.cio) } } }