feat(standalone): :integrations-telegram — мост Telegram-чат ↔ agent
Новый модуль :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-таблицей,
описанием поведения и инструкцией по включению.
This commit is contained in:
@@ -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)
|
||||
}
|
||||
+290
@@ -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<String, Job> = mutableMapOf()
|
||||
|
||||
private val subscribersLock = Mutex()
|
||||
|
||||
/**
|
||||
* Emit-only signal «новое сообщение пришло в чат» — для diag и
|
||||
* тестов. На этот поток можно подписаться ИЗ тестов (см.
|
||||
* `TelegramBridgeComponentTest`).
|
||||
*/
|
||||
private val incomingEvents = MutableSharedFlow<IncomingEvent>(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,
|
||||
)
|
||||
}
|
||||
+137
@@ -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<Entry> = withContext(Dispatchers.Default) {
|
||||
mutex.withLock {
|
||||
listStmt.reset()
|
||||
listStmt.clearBindings()
|
||||
listStmt.executeQuery().use { rs ->
|
||||
val out = mutableListOf<Entry>()
|
||||
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"
|
||||
}
|
||||
}
|
||||
+47
@@ -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
|
||||
}
|
||||
}
|
||||
+38
@@ -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<AutoCloseable, TelegramBridgeComponent> {
|
||||
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
|
||||
}
|
||||
+571
@@ -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 <T> subscribe(scope: CoroutineScope, flow: kotlinx.coroutines.flow.Flow<T>): MutableList<T> {
|
||||
val list = mutableListOf<T>()
|
||||
scope.launch {
|
||||
flow.collect { list.add(it) }
|
||||
}
|
||||
return list
|
||||
}
|
||||
|
||||
/** Ждём, пока в [list] появится элемент, удовлетворяющий [match], или таймаут. */
|
||||
private suspend fun <T> List<T>.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<TelegramBridgeComponent.IncomingEvent> = 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<List<Update>>(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<TextMessage> = mutableListOf()
|
||||
val edited: MutableList<EditTextRequest> = mutableListOf()
|
||||
val chatActions: MutableList<Pair<String, SendChatEvent.Action>> = mutableListOf()
|
||||
val deleted: MutableList<Pair<String, Long>> = mutableListOf()
|
||||
|
||||
private val messageIdSeq = java.util.concurrent.atomic.AtomicLong(1000L)
|
||||
|
||||
override suspend fun getUpdate(
|
||||
offset: Long?,
|
||||
limit: Int?,
|
||||
timeout: Int?,
|
||||
allowedUpdates: List<String>?,
|
||||
): List<Update> = 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<BotCommand>): Boolean = true
|
||||
override suspend fun getMyCommands(): List<BotCommand> = emptyList()
|
||||
override fun close() {}
|
||||
}
|
||||
|
||||
/**
|
||||
* Минимальный `MutableAgent`: создаёт `FakeConversation`, помнит их по id.
|
||||
* `onlineOutbox` управляется из теста (`FakeOnlineOutbox`).
|
||||
*/
|
||||
class FakeMutableAgent(override val onlineOutbox: FakeOnlineOutbox) : MutableAgent {
|
||||
override val systemProviders: MutableList<pw.binom.agentik.agent.SystemPromptProvider> = mutableListOf()
|
||||
override val toolProviders: MutableList<pw.binom.agentik.agent.ToolProvider> = mutableListOf()
|
||||
val conversations: MutableMap<String, FakeConversation> = mutableMapOf()
|
||||
val created: MutableList<CreatedConv> = 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<OnlineEvent>(replay = 0, extraBufferCapacity = 64)
|
||||
|
||||
suspend fun emit(event: OnlineEvent) = bus.emit(event)
|
||||
|
||||
override fun onlineEvents(): Flow<OnlineEvent> = bus.asSharedFlow()
|
||||
override fun onlineEvents(conversationId: String): Flow<OnlineEvent> =
|
||||
bus.asSharedFlow().filter { it.conversationId == conversationId }
|
||||
|
||||
override fun close() {}
|
||||
}
|
||||
|
||||
class FakeConversation(
|
||||
override val id: String,
|
||||
override val isTemporal: Boolean = false,
|
||||
) : Conversation {
|
||||
val sent: MutableList<List<Content>> = 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<Content>, 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<pw.binom.agentik.proto.Message>()
|
||||
override fun close() {}
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// Flow helpers
|
||||
// ============================================================================
|
||||
|
||||
private suspend fun <T> Flow<T>.firstAsync(predicate: (T) -> Boolean = { true }): T =
|
||||
kotlinx.coroutines.withTimeout(2_000) {
|
||||
this@firstAsync.first { predicate(it) }
|
||||
}
|
||||
|
||||
private suspend fun <T> Flow<T>.firstAsyncOrNull(predicate: (T) -> Boolean): T? = try {
|
||||
kotlinx.coroutines.withTimeout(1_000) {
|
||||
this@firstAsyncOrNull.first { predicate(it) }
|
||||
}
|
||||
} catch (_: kotlinx.coroutines.TimeoutCancellationException) {
|
||||
null
|
||||
}
|
||||
@@ -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 | Описание |
|
||||
|
||||
+21
-10
@@ -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,9 +90,17 @@ kotlin {
|
||||
// Skill-подсистема как Component: тулы + SkillMiningComponent.
|
||||
implementation(project(":skill-mining"))
|
||||
|
||||
// Generic MCP-bridge (McpConfig, McpRegistry, McpLiteToolAdapter).
|
||||
// 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)
|
||||
// liteTool DSL (типизированные LiteTool через @Serializable args)
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user