From 4c66947c253a043efa87aa8874a4413e48ac2a1a Mon Sep 17 00:00:00 2001 From: subochev Date: Sun, 13 Sep 2026 15:59:58 +0300 Subject: [PATCH] standalone: add MCP client + tool-call loop MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Add mcp/McpConfig + McpRegistry wrapping io.modelcontextprotocol:kotlin-sdk-client 0.15.0 - Supports stdio (uvx/npx/python) and streamable HTTP transports - Claude Desktop-compatible JSON config (AGENTIK_MCP_CONFIG) - server__tool name prefix to avoid collisions between servers - Add agent/NamedTool (name + LiteTool pair) and tools: List on ChatAgent/ChatConversation - Implement tool-call loop in ChatConversation.runTurn: - delta.toolCalls -> emit ToolCall event -> persist audit -> execute tool -> emit ToolResult -> persist -> liteConv.addToolResult(callId, name, result) - separate tc-/tr- prefixes keep SQL PRIMARY KEY unique while toolCallId FK is preserved - 8 McpConfig + 4 McpRegistry unit tests; +1 ChatAgentTest tool-loop test (44/44 total) - e2e verified: real MCP fetch server (mcp-server-fetch) + litellm local/codding -> LLM calls fetch__fetch, MCP exec, result fed back, conversation continues docs/STANDALONE.md: drop 'no tools / no MCP' from §8; replace 'Подключить тул (v2)' stub with full in-agent + MCP recipe and tool-loop algorithm in §7 --- docs/STANDALONE.md | 34 ++- gradle/libs.versions.toml | 4 + standalone/build.gradle.kts | 7 + .../pw/binom/agentik/standalone/Main.kt | 11 +- .../agentik/standalone/agent/ChatAgent.kt | 5 +- .../standalone/agent/ChatConversation.kt | 159 +++++++++--- .../agentik/standalone/agent/NamedTool.kt | 12 + .../binom/agentik/standalone/mcp/McpConfig.kt | 95 +++++++ .../agentik/standalone/mcp/McpRegistry.kt | 243 ++++++++++++++++++ .../agentik/standalone/agent/ChatAgentTest.kt | 97 ++++++- .../agentik/standalone/mcp/McpConfigTest.kt | 88 +++++++ .../agentik/standalone/mcp/McpRegistryTest.kt | 48 ++++ 12 files changed, 752 insertions(+), 51 deletions(-) create mode 100644 standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/NamedTool.kt create mode 100644 standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/mcp/McpConfig.kt create mode 100644 standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/mcp/McpRegistry.kt create mode 100644 standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/mcp/McpConfigTest.kt create mode 100644 standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/mcp/McpRegistryTest.kt diff --git a/docs/STANDALONE.md b/docs/STANDALONE.md index caaa137..e0b0da0 100644 --- a/docs/STANDALONE.md +++ b/docs/STANDALONE.md @@ -219,6 +219,7 @@ suspend fun touch(id: String, now: Instant) | `AGENTIK_DB_PATH` | путь к SQLite-файлу | `./agentik.db` | | `AGENTIK_LLM_BACKEND` | `openai` или `google` | `openai` | | `AGENTIK_SYSTEM_PROMPT` | текст системного промпта | «Ты полезный ассистент. Отвечай кратко и по делу.» | +| `AGENTIK_MCP_CONFIG` | путь к `mcp.json` в формате Claude Desktop (`{"mcpServers":{"name":{"command":"...","args":[...]}` или `"url":"..."}`) | не задан (MCP выключен) | ### Backend `openai` @@ -248,14 +249,37 @@ suspend fun touch(id: String, now: Instant) ## 7. Расширение -### Подключить тул (v2) +### Подключить тулы (MCP / in-agent) ```kotlin -val tools: List = listOf(ReadFileTool(Path.of("/work"))) -val cfg = LiteConversationConfig(systemInstruction = "…", tools = tools) +// 1) In-agent tool (нативный LiteTool): +val echoTool = object : LiteTool { + override fun describe(): String = + """{"type":"object","properties":{"x":{"type":"string"}},"required":["x"]}""" + override fun invoke(args: String): String = "echoed: $args" +} +val tools = listOf(NamedTool("echo", echoTool)) + +// 2) MCP (stdio / streamable HTTP) — конфиг в Claude Desktop-формате: +val mcp = McpConfig.fromEnv() // читает AGENTIK_MCP_CONFIG=path/to/mcp.json +val registry = McpRegistry.fromConfig(mcp) // стартует все серверы, лист LiteTool'ов +val tools = registry.namedTools // server__tool префикс автоматически + +// 3) В обоих случаях: +val agent = ChatAgent(id, stores, llm, llmConfig, tools = tools) ``` -LiteConversation сообщит модели о доступных тулах. После появления `LiteContentPart.ToolResult` (v2) появится и `ChatConversation.addToolResult(...)` — ручной tool-loop на on-device движке без пересоздания беседы. +`LiteConversation` принимает `tools = ...` в `LiteConversationConfig`. На каждый `delta.toolCalls` из модели `ChatConversation.runTurn`: + +1. Эмитит `Event.ToolCall(callId, toolName, argsJson)` клиенту (по SSE) +2. Записывает `MessageRecord.ToolCall` в audit + working memory (если не temp) +3. Вызывает `tool.invoke(argsJson)` +4. Эмитит `Event.ToolResult(resultId, resultText)` +5. Записывает `MessageRecord.ToolResult` (с `toolCallId = callId`) +6. Кормит `liteConv.addToolResult(callId, name, result)` в LiteConversation (KV-cache выживает между итерациями) +7. Цикл повторяется до `delta.toolCalls.isEmpty()` + +ID у `ToolCall` и `ToolResult` разные (`tc-…` / `tr-…`), но `MessageRecord.ToolResult.toolCallId` указывает на `MessageRecord.ToolCall.id` той же логической пары. Этим достигается уникальность PK в таблице `message`. ### Добавить ещё один транспорт @@ -283,9 +307,7 @@ LiteConversation сообщит модели о доступных тулах. ## 8. Что НЕ делает `standalone` сегодня -* **Нет инструментов (тулов).** Поддержка ждёт `LiteContentPart.ToolResult` в litert-api (v2). До тех пор агент — text-only чат. * **Нет суммаризации.** `WorkingMemoryStore.compact` уже есть, но без LLM-вызова для генерации текста суммаризации. -* **Нет MCP-клиента.** * **Нет авторизации.** Все эндпоинты открыты. * **Нет инкрементальной догрузки старых сообщений.** `getMessages(after)` работает с offset/limit, но без «схлопывания» (compaction в визуальной истории — задача клиента). diff --git a/gradle/libs.versions.toml b/gradle/libs.versions.toml index a0ead7d..2754d5c 100644 --- a/gradle/libs.versions.toml +++ b/gradle/libs.versions.toml @@ -46,6 +46,10 @@ ktor-client-core = { module = "io.ktor:ktor-client-core", version.ref = "ktor" } ktor-client-cio = { module = "io.ktor:ktor-client-cio", version.ref = "ktor" } ktor-client-content-negotiation = { module = "io.ktor:ktor-client-content-negotiation", version.ref = "ktor" } ktor-server-test-host = { module = "io.ktor:ktor-server-test-host", version.ref = "ktor" } +ktor-client-sse = { module = "io.ktor:ktor-client-sse", version.ref = "ktor" } + +# --- Model Context Protocol (MCP) --- +mcp-sdk-client = { module = "io.modelcontextprotocol:kotlin-sdk-client", version = "0.15.0" } kotlin-test = { module = "org.jetbrains.kotlin:kotlin-test", version.ref = "kotlin" } diff --git a/standalone/build.gradle.kts b/standalone/build.gradle.kts index f8650fe..0d22f7e 100644 --- a/standalone/build.gradle.kts +++ b/standalone/build.gradle.kts @@ -57,6 +57,13 @@ kotlin { // Транспортные фасады implementation(libs.agui.server) implementation(libs.a2a.server) + + // MCP (Model Context Protocol) клиент — подключение внешних/внутренних MCP-серверов + implementation(libs.mcp.sdk.client) + implementation(libs.ktor.client.core) + implementation(libs.ktor.client.cio) + implementation(libs.ktor.client.content.negotiation) + implementation(libs.ktor.serialization.kotlinx.json) } commonTest.dependencies { diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/Main.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/Main.kt index 659084e..fcbc25a 100644 --- a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/Main.kt +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/Main.kt @@ -8,8 +8,9 @@ import io.ktor.server.routing.routing import pw.binom.agentik.server.agentikAgent import pw.binom.agentik.standalone.agent.ChatAgent import pw.binom.agentik.standalone.llm.LlmConfig +import pw.binom.agentik.standalone.mcp.McpConfig +import pw.binom.agentik.standalone.mcp.McpRegistry import pw.binom.agentik.standalone.persistence.sqlite.SqliteStores - /** * standalone-контейнер agentik: * - :server (proto): встраиваемый Ktor (Netty), порт AGENTIK_PORT (default 8080) @@ -26,7 +27,9 @@ import pw.binom.agentik.standalone.persistence.sqlite.SqliteStores * GET /health -> "ok" * * Хранилище — SQLite (env: AGENTIK_DB_PATH, default `./agentik.db`, `:memory:` для тестов). - * LLM — litert-openai (env: OPENAI_BASE_URL, OPENAI_API_KEY, OPENAI_MODEL). + * LLM — litert-openai (env: OPENAI_BASE_URL, OPENAI_API_KEY, OPENAI_MODEL) или litert-google + * (env: AGENTIK_LLM_BACKEND=google, AGENTIK_GOOGLE_MODEL_PATH). + * MCP — AGENTIK_MCP_CONFIG=.json (формат Claude Desktop). * * System prompt — AGENTIK_SYSTEM_PROMPT (default: встроенный `Ты полезный ассистент...`). */ @@ -37,11 +40,13 @@ fun main() { val llmConfig = LlmConfig.fromEnv() val llm = llmConfig.createLlm() val stores = SqliteStores.open(dbPath = dbPath) + val mcpRegistry = McpRegistry.fromConfig(McpConfig.fromEnv()) val agent = ChatAgent( id = "agentik", stores = stores, llm = llm, llmConfig = llmConfig, + tools = mcpRegistry.namedTools, ) val server = embeddedServer(Netty, port = port) { @@ -56,8 +61,10 @@ fun main() { println(" GET /agentik/conversations/{id}/events -> SSE") println(" storage: $dbPath") println(" llm: ${llmConfig.backend} ${llmConfig.modelInfo()}") + println(" mcp: ${mcpRegistry.allTools.size} tools from ${mcpRegistry.connectedServerCount} servers") Runtime.getRuntime().addShutdownHook(Thread { agent.close() + mcpRegistry.close() stores.close() llm.close() }) diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt index 79168d4..fb4940a 100644 --- a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ChatAgent.kt @@ -31,6 +31,7 @@ class ChatAgent( private val stores: SqliteStores, private val llm: LiteLlm, private val llmConfig: LlmConfig, + private val tools: List = emptyList(), ) : ProtoAgent, AutoCloseable { private val agentEvents = MutableSharedFlow( @@ -71,7 +72,7 @@ class ChatAgent( ) } } - val conv = ChatConversation(record = rec, stores = stores, llm = llm, systemPrompt = llmConfig.systemPrompt) + val conv = ChatConversation(record = rec, stores = stores, llm = llm, systemPrompt = llmConfig.systemPrompt, tools = tools) runBlocking { liveLock.withLock { live[conv.id] = conv } } @@ -82,7 +83,7 @@ class ChatAgent( override suspend fun getConversation(id: String): ProtoConversation? { liveLock.withLock { live[id] }?.let { if (!it.isClosed) return it } val rec = stores.conversations.get(id) ?: return null - return ChatConversation(record = rec, stores = stores, llm = llm, systemPrompt = llmConfig.systemPrompt).also { + return ChatConversation(record = rec, stores = stores, llm = llm, systemPrompt = llmConfig.systemPrompt, tools = tools).also { liveLock.withLock { live[id] = it } } } diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ChatConversation.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ChatConversation.kt index bd13b3c..c05746d 100644 --- a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ChatConversation.kt +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/ChatConversation.kt @@ -30,34 +30,35 @@ import pw.binom.litert.LiteConversationConfig import pw.binom.litert.LiteLlm import pw.binom.litert.LiteMessage import pw.binom.litert.LiteRole +import pw.binom.litert.LiteToolCall import kotlin.time.Instant /** * Stateful [pw.binom.agentik.proto.Conversation] поверх [LiteLlm]. * - * Контракт: один [LiteConversation] живёт столько же, сколько [ChatConversation]. - * Это требование абстракции LiteLlm — у реализаций внутри `LiteConversation` хранится - * KV-cache движка (LiteRT-LM) либо инкрементальная история (litert-openai). Пересоздание - * на каждый send сломало бы и то, и другое. + * Один [LiteConversation] живёт всю беседу (требование абстракции LiteLlm — у реализаций + * внутри `LiteConversation` хранится KV-cache движка / инкрементальная история). * * На каждый [send]: - * 1. Записывает user-сообщение в audit log (`message`) + working memory. - * 2. Создаёт [LiteConversation] **один раз** при первом send (initial messages = - * системный промпт + пары User/Assistant из working memory, **без** только что - * добавленного user-сообщения — мы передадим его через sendStreamContents). - * 3. Стримит ответ [LiteDelta] в [sendStreamContents] — эмитит [ProtoEvent.AppendText]. - * Движок сам добавляет user/assistant к своей внутренней истории. - * 4. По завершении записывает AssistantMessage в audit + working memory. + * 1. Записывает user-сообщение в audit log + working memory (если не temp). + * 2. Создаёт [LiteConversation] **один раз** при первом send с initial messages из + * working memory (без только что записанного user-сообщения — мы его отдадим через + * [LiteConversation.sendStreamContents]). + * 3. Стримит ответ [LiteDelta]: + * - text → [ProtoEvent.AppendText] + * - toolCalls → выполняем через [tools] (MCP), эмитим [ProtoEvent.ToolCall]/[ProtoEvent.ToolResult], + * подаём результат через [LiteConversation.addToolResult], продолжаем стрим + * до тех пор, пока модель не перестанет вызывать тулы. + * 4. По завершении записывает AssistantMessage в audit + working memory. * - * Если LiteConversation упал при инициализации или во время send — закрываем его, - * следующий send попробует создать заново. После [interrupt] LiteConversation жив; - * отменяется только текущий send. + * При ошибке LiteConversation выбрасывается и пересоздаётся на следующий [send]. */ class ChatConversation( record: ConversationRecord, private val stores: SqliteStores, private val llm: LiteLlm, private val systemPrompt: String, + private val tools: List = emptyList(), ) : ProtoConversation, AutoCloseable { private var record: ConversationRecord = record @@ -73,9 +74,11 @@ class ChatConversation( private val messageStore: MessageStore get() = stores.messages private val workingMemory: WorkingMemoryStore get() = stores.workingMemory + private val toolsByName: Map = tools.associateBy { it.name } + private val events = MutableSharedFlow( replay = 0, - extraBufferCapacity = 128, + extraBufferCapacity = 256, onBufferOverflow = BufferOverflow.DROP_OLDEST, ) @@ -151,6 +154,14 @@ class ChatConversation( scope.cancel() } + /** + * Один ход: user → (assistant → tool → ... → assistant)*. + * + * Tool-loop: после каждого `sendStreamContents` смотрим `delta.toolCalls`. Если есть — + * исполняем, подаём результат через `addToolResult`, делаем ещё один send (с пустым + * user-сообщением как триггером продолжения — модель уже знает, что делать дальше + * по tool-results в истории), повторяем. Защита от зацикливания — [MAX_TOOL_LOOPS]. + */ private suspend fun runTurn(userRecord: MessageRecord.UserMessage, turnStarted: Instant) { emitEvent(ProtoEvent.StartReasoning(date = turnStarted)) emitEvent(ProtoEvent.StartResponse(date = now(), responseType = ProtoEvent.ResponseType.TEXT)) @@ -181,19 +192,44 @@ class ChatConversation( } val reply = StringBuilder() - try { - liteConv.sendStreamContents(parts).collect { delta -> - if (delta.text.isNotEmpty()) { - reply.append(delta.text) - emitEvent(ProtoEvent.AppendText(date = now(), body = delta.text)) + var currentParts: List = parts + var loopGuard = 0 + + while (loopGuard++ < MAX_TOOL_LOOPS) { + val collectedCalls = mutableListOf() + try { + liteConv.sendStreamContents(currentParts).collect { delta -> + if (delta.text.isNotEmpty()) { + reply.append(delta.text) + emitEvent(ProtoEvent.AppendText(date = now(), body = delta.text)) + } + if (delta.toolCalls.isNotEmpty()) { + collectedCalls.addAll(delta.toolCalls) + } } + } catch (e: kotlinx.coroutines.CancellationException) { + throw e + } catch (e: Throwable) { + emitEvent(ProtoEvent.Error(date = now(), message = e.message ?: e.javaClass.simpleName)) + emitEvent(ProtoEvent.End(date = now())) + return } - } catch (e: kotlinx.coroutines.CancellationException) { - throw e - } catch (e: Throwable) { - emitEvent(ProtoEvent.Error(date = now(), message = e.message ?: e.javaClass.simpleName)) - emitEvent(ProtoEvent.End(date = now())) - return + + if (collectedCalls.isEmpty()) break + + for (call in collectedCalls) { + executeToolCall(liteConv, call) + } + + // Continuation: send a no-op user message so the engine produces the next + // assistant response (which will see the tool results we just fed via + // addToolResult in its history). The leading newline + space is a benign + // trigger — every LLM treats it as "please continue". + currentParts = listOf(LiteContentPart.Text(" ")) + } + + if (loopGuard >= MAX_TOOL_LOOPS) { + System.err.println("[agentik] tool loop hit MAX_TOOL_LOOPS=$MAX_TOOL_LOOPS for $id — bailing") } val assistantId = newId("msg") @@ -223,14 +259,64 @@ class ChatConversation( emitEvent(ProtoEvent.End(date = assistantAt)) } + /** + * Один tool-call: эмитим Event.ToolCall, выполняем tool (MCP), эмитим Event.ToolResult, + * пишем в audit + working memory, подаём результат в LiteConversation. + */ + private suspend fun executeToolCall(liteConv: LiteConversation, call: LiteToolCall) { + val callId = newId("tc") + val resultId = newId("tr") + val argsJson = encodeArgsJson(call.arguments) + val nowTs = now() + + emitEvent(ProtoEvent.ToolCall(date = nowTs, id = callId, title = null, toolName = call.name, toolArgs = argsJson)) + + if (!record.isTemporal) { + messageStore.append( + MessageRecord.ToolCall( + id = callId, + conversationId = id, + toolName = call.name, + toolTitle = null, + toolArgsJson = argsJson, + createdAt = nowTs, + ), + ) + } + + val tool = toolsByName[call.name] + val resultText: String = if (tool == null) { + System.err.println("[agentik] tool '${call.name}' requested but not registered") + "[tool not found: ${call.name}]" + } else { + try { + tool.tool.invoke(argsJson).ifBlank { "" } + } catch (e: Throwable) { + System.err.println("[agentik] tool '${call.name}' threw: ${e.message}") + "[tool error: ${e.message ?: e.javaClass.simpleName}]" + } + } + + emitEvent(ProtoEvent.ToolResult(date = now(), id = resultId, result = resultText)) + + if (!record.isTemporal) { + messageStore.append( + MessageRecord.ToolResult( + id = resultId, + conversationId = id, + toolCallId = callId, + result = resultText, + createdAt = now(), + ), + ) + } + + liteConv.addToolResult(callId = callId, name = call.name, result = resultText) + } + /** * Возвращает существующий [LiteConversation] или создаёт новый, инициализированный * системным промптом и прошлыми User/Assistant из working memory. - * - * [excludeUserSourceId] — если задан, исключает одну строку (свежее user-сообщение, - * уже записанное в audit + working memory, но ещё не отправленное в LLM — мы отдадим - * его через [LiteConversation.sendStreamContents]). Это предотвращает дублирование - * "user → user" в LiteConversation history. */ private suspend fun getOrCreateLiteConversation(excludeUserSourceId: String? = null): LiteConversation { liteConv?.let { return it } @@ -241,8 +327,6 @@ class ChatConversation( ?.let { (it.entry as WorkingMemoryEntry.System).text } ?: systemPrompt - // Берём пары (User, Assistant) из прошлой истории, исключая свежее user-сообщение, - // которое отправим через sendStreamContents. val pastTurns = if (record.isTemporal) emptyList() else wm .filter { row -> val isUserOrAssistant = row.entry is WorkingMemoryEntry.User || row.entry is WorkingMemoryEntry.Assistant @@ -260,6 +344,7 @@ class ChatConversation( val config = LiteConversationConfig( systemInstruction = resolvedSystemPrompt.takeIf { it.isNotBlank() }, initialMessages = pastTurns, + tools = tools.map { it.tool }, ) return llm.createConversation(config).also { liteConv = it } @@ -273,6 +358,12 @@ class ChatConversation( Instant.fromEpochMilliseconds(System.currentTimeMillis()) private fun newId(prefix: String): String = "$prefix-${java.util.UUID.randomUUID()}" + + private fun encodeArgsJson(arguments: Map): String = encodeToolArgs(arguments) + + companion object { + private const val MAX_TOOL_LOOPS = 16 + } } internal fun List.toLiteContents(): List = map { it.toLite() } @@ -315,8 +406,6 @@ internal fun MessageRecord.toProto(): ProtoMessage = when (this) { date = createdAt, result = result, ) - // Синтетические строки рабочей памяти: отдаём клиенту как текст ассистента - // (Summary — это результат суммаризации, по форме — ответ модели) или юзера (System). is MessageRecord.Summary -> ProtoMessage.AssistantMessage( id = id, date = createdAt, diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/NamedTool.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/NamedTool.kt new file mode 100644 index 0000000..c0f8814 --- /dev/null +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/agent/NamedTool.kt @@ -0,0 +1,12 @@ +package pw.binom.agentik.standalone.agent + +import pw.binom.litert.LiteTool + +/** + * (имя-как-видит-модель) → [LiteTool]. + * + * Имя используется как ключ для матчинга `LiteToolCall.name` (приходящего от LLM) + * с конкретной реализацией тула. Для MCP-адаптеров имя имеет формат `server__tool`, + * чтобы избежать коллизий между разными MCP-серверами. + */ +data class NamedTool(val name: String, val tool: LiteTool) diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/mcp/McpConfig.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/mcp/McpConfig.kt new file mode 100644 index 0000000..c0fd15e --- /dev/null +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/mcp/McpConfig.kt @@ -0,0 +1,95 @@ +package pw.binom.agentik.standalone.mcp + +import kotlinx.serialization.json.Json +import kotlinx.serialization.json.JsonObject +import kotlinx.serialization.json.JsonPrimitive +import kotlinx.serialization.json.jsonArray +import kotlinx.serialization.json.jsonObject +import kotlinx.serialization.json.jsonPrimitive +import java.io.File + +/** + * Описание одного MCP-сервера (см. [McpRegistry]). + * + * Два варианта транспорта: + * - [Stdio]: запустить процесс (`command` + `args`), общаться через stdin/stdout + * - [Http]: подключиться к удалённому MCP-серверу по URL (Streamable HTTP) + */ +sealed interface McpServerSpec { + /** Уникальное имя сервера, как его видно в логах и в tool-prefix. */ + val name: String + + /** stdio: спавним процесс, читаем его stdout, пишем в stdin. */ + data class Stdio( + override val name: String, + val command: String, + val args: List = emptyList(), + val env: Map = emptyMap(), + ) : McpServerSpec + + /** http (Streamable HTTP transport): подключаемся к существующему MCP-серверу. */ + data class Http( + override val name: String, + val url: String, + val headers: Map = emptyMap(), + ) : McpServerSpec +} + +/** + * Конфигурация MCP-слоя standalone'а. + * + * Парсится из JSON-файла в формате, совместимом с Claude Desktop + * (`mcpServers.{name}.command/args` для stdio, `mcpServers.{name}.url` для http). + * + * Путь к файлу — env `AGENTIK_MCP_CONFIG`. Если не задан или файл не существует — пустой список. + */ +data class McpConfig( + val servers: List, +) { + val isEmpty: Boolean get() = servers.isEmpty() + + companion object { + private val json = Json { ignoreUnknownKeys = true } + + fun fromEnv(env: (String) -> String? = System::getenv): McpConfig { + val path = env("AGENTIK_MCP_CONFIG")?.takeIf { it.isNotBlank() } ?: return empty() + val file = File(path) + if (!file.exists()) { + System.err.println("[agentik] AGENTIK_MCP_CONFIG points to missing file: $path") + return empty() + } + return fromJson(file.readText()) + } + + fun empty(): McpConfig = McpConfig(servers = emptyList()) + + fun fromJson(raw: String): McpConfig { + val root = json.parseToJsonElement(raw).jsonObject + val mcpServers = root["mcpServers"]?.jsonObject ?: return empty() + val servers = mcpServers.entries.mapNotNull { (name, spec) -> parseServer(name, spec.jsonObject) } + return McpConfig(servers) + } + + private fun parseServer(name: String, spec: JsonObject): McpServerSpec? { + val url = (spec["url"] as? JsonPrimitive)?.jsonPrimitive?.content + if (url != null) { + val headers = (spec["headers"] as? JsonObject)?.entries + ?.associate { (k, v) -> k to (v as JsonPrimitive).jsonPrimitive.content } + ?: emptyMap() + return McpServerSpec.Http(name = name, url = url, headers = headers) + } + val command = (spec["command"] as? JsonPrimitive)?.jsonPrimitive?.content + if (command != null) { + val args = (spec["args"] as? kotlinx.serialization.json.JsonArray) + ?.map { (it as JsonPrimitive).jsonPrimitive.content } + ?: emptyList() + val env = (spec["env"] as? JsonObject)?.entries + ?.associate { (k, v) -> k to (v as JsonPrimitive).jsonPrimitive.content } + ?: emptyMap() + return McpServerSpec.Stdio(name = name, command = command, args = args, env = env) + } + System.err.println("[agentik] MCP server '$name' has neither 'url' nor 'command' — skipped") + return null + } + } +} diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/mcp/McpRegistry.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/mcp/McpRegistry.kt new file mode 100644 index 0000000..020800f --- /dev/null +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/mcp/McpRegistry.kt @@ -0,0 +1,243 @@ +package pw.binom.agentik.standalone.mcp + +import io.ktor.client.HttpClient +import io.ktor.client.engine.cio.CIO +import io.ktor.client.plugins.defaultRequest +import io.modelcontextprotocol.kotlin.sdk.client.Client +import io.modelcontextprotocol.kotlin.sdk.client.ClientOptions +import io.modelcontextprotocol.kotlin.sdk.client.StdioClientTransport +import io.modelcontextprotocol.kotlin.sdk.client.StreamableHttpClientTransport +import io.modelcontextprotocol.kotlin.sdk.shared.Transport +import io.modelcontextprotocol.kotlin.sdk.types.CallToolResult +import io.modelcontextprotocol.kotlin.sdk.types.Implementation +import io.modelcontextprotocol.kotlin.sdk.types.TextContent +import io.modelcontextprotocol.kotlin.sdk.types.Tool +import io.modelcontextprotocol.kotlin.sdk.types.ToolSchema +import kotlinx.coroutines.runBlocking +import kotlinx.io.asSink +import kotlinx.io.asSource +import kotlinx.io.buffered +import kotlinx.serialization.json.Json +import kotlinx.serialization.json.JsonArray +import kotlinx.serialization.json.JsonElement +import kotlinx.serialization.json.JsonObject +import kotlinx.serialization.json.JsonPrimitive +import kotlinx.serialization.json.buildJsonObject +import kotlinx.serialization.json.put +import pw.binom.agentik.standalone.agent.NamedTool +import pw.binom.litert.LiteTool +import java.util.concurrent.ConcurrentHashMap + +/** + * Реестр подключённых MCP-серверов. + * + * На старте подключается ко всем [McpServerSpec] из [McpConfig], у каждого запрашивает + * список tools и оборачивает их в [LiteTool]-адаптеры ([McpLiteToolAdapter]). Все + * адаптеры собираются в [allTools], который `ChatConversation` подмешивает в + * [pw.binom.litert.LiteConversationConfig.tools]. + * + * [close] убивает stdio-процессы и закрывает HTTP-клиент. + */ +class McpRegistry( + private val servers: List, + private val httpClient: HttpClient = defaultHttpClient(), + private val clientName: String = "agentik", + private val clientVersion: String = "1.0.0", +) : AutoCloseable { + + private val connected: MutableMap = ConcurrentHashMap() + @Volatile private var closed = false + + /** Все [LiteTool] со всех подключённых серверов. */ + val allTools: List by lazy { + connected.values.flatMap { it.tools } + } + + /** Все [LiteTool] с именами (server__tool), которые видит LLM. */ + val namedTools: List by lazy { + allTools.filterIsInstance().map { NamedTool(it.fullName, it) } + } + + /** Количество успешно подключённых серверов. */ + val connectedServerCount: Int get() = connected.size + + init { + if (servers.isNotEmpty()) { + for (spec in servers) { + runCatching { + val server = connectOne(spec) + connected[spec.name] = server + }.onFailure { e -> + System.err.println("[agentik] MCP server '${spec.name}' failed to connect: ${e.message}") + } + } + System.err.println("[agentik] MCP registry: ${connected.size}/${servers.size} servers connected, ${allTools.size} tools total") + } + } + + private fun connectOne(spec: McpServerSpec): ConnectedServer { + var ownedProcess: Process? = null + val transport: Transport = when (spec) { + is McpServerSpec.Stdio -> { + val cmd = (listOf(spec.command) + spec.args).joinToString(" ") + System.err.println("[agentik] MCP stdio '$spec.name': $cmd") + val pb = ProcessBuilder(buildList { add(spec.command); addAll(spec.args) }) + .redirectErrorStream(false) + spec.env.forEach { (k, v) -> pb.environment()[k] = v } + val proc = pb.start() + ownedProcess = proc + StdioClientTransport( + input = proc.inputStream.asSource().buffered(), + output = proc.outputStream.asSink().buffered(), + error = proc.errorStream.asSource().buffered(), + ) + } + is McpServerSpec.Http -> { + System.err.println("[agentik] MCP http '$spec.name': ${spec.url}") + StreamableHttpClientTransport( + client = httpClient.config { + if (spec.headers.isNotEmpty()) { + defaultRequest { + spec.headers.forEach { (k, v) -> + headers { append(k, v) } + } + } + } + }, + url = spec.url, + ) + } + } + + val client = Client( + clientInfo = Implementation(name = clientName, version = clientVersion, title = null, websiteUrl = null, icons = emptyList()), + ) + runBlocking { client.connect(transport) } + val mcpTools = runBlocking { client.listTools().tools } + val liteTools = mcpTools.map { tool -> McpLiteToolAdapter(spec.name, tool, client) } + return ConnectedServer(spec, client, transport, ownedProcess, liteTools) + } + + override fun close() { + if (closed) return + closed = true + connected.values.forEach { server -> + runCatching { runBlocking { server.transport.close() } } + server.ownedProcess?.let { runCatching { it.destroyForcibly() } } + } + connected.clear() + runCatching { httpClient.close() } + } + + private data class ConnectedServer( + val spec: McpServerSpec, + val client: Client, + val transport: Transport, + val ownedProcess: Process?, + val tools: List, + ) + + companion object { + private fun defaultHttpClient(): HttpClient = HttpClient(CIO) + + fun fromConfig(config: McpConfig): McpRegistry = + McpRegistry(servers = config.servers) + } +} + +/** + * Адаптер MCP [Tool] → [pw.binom.litert.LiteTool]. + * + * Имя тула префиксуется именем сервера через `__`, чтобы избежать коллизий + * между MCP-серверами (например, оба могут иметь tool `search`). + * + * [describe] сериализует tool в JSON-дескриптор в формате, который litert-openai + * и litert-google принимают как function-calling definition: + * + * ```json + * { + * "type": "function", + * "function": { + * "name": "__", + * "description": "...", + * "parameters": { "type": "object", "properties": {...}, "required": [...] } + * } + * } + * ``` + * + * [invoke] вызывает [Client.callTool] по оригинальному (непрефиксованному) имени тула + * на нужном MCP-сервере и схлопывает [CallToolResult] в плоский текст: + * для каждого блока контента возвращается либо `text`, либо JSON-представление. + * Если `isError == true` — текст префиксуется `[tool error]`. + */ +internal class McpLiteToolAdapter( + private val serverName: String, + private val tool: Tool, + private val client: Client, +) : LiteTool { + + internal val fullName: String = "${serverName}__${tool.name}" + + override fun describe(): String = + buildJsonObject { + put("type", "function") + put("function", buildJsonObject { + put("name", fullName) + put("description", tool.description ?: "") + put("parameters", tool.inputSchema.toJsonSchema()) + }) + }.toString() + + override fun invoke(arguments: String): String { + val argsMap = parseArgsJson(arguments, tool.name) + val result: CallToolResult = runBlocking { client.callTool(tool.name, argsMap) } + return renderResult(result) + } + + private fun renderResult(result: CallToolResult): String { + val isError = result.isError == true + val parts = result.content.map { block -> + when (block) { + is TextContent -> block.text + else -> Json.encodeToString(JsonElement.serializer(), JsonPrimitive(block.toString())) + } + } + val text = parts.joinToString("\n").ifEmpty { "[]" } + return if (isError) "[tool error] $text" else text + } + + private fun ToolSchema?.toJsonSchema(): JsonElement { + if (this == null) return buildJsonObject { put("type", "object") } + val props = this.properties ?: buildJsonObject { } + val reqs = this.required ?: emptyList() + return buildJsonObject { + put("type", this@toJsonSchema.type.ifEmpty { "object" }) + put("properties", props) + put("required", JsonArray(reqs.map { JsonPrimitive(it) })) + } + } + + companion object { + private val json = Json { ignoreUnknownKeys = true; isLenient = true } + + private fun parseArgsJson(raw: String, toolName: String): Map { + if (raw.isBlank()) return emptyMap() + return try { + val parsed = json.parseToJsonElement(raw) + if (parsed !is JsonObject) emptyMap() else parsed.toAnyMap() + } catch (e: Throwable) { + System.err.println("[agentik] MCP tool '$toolName' got invalid args JSON: ${e.message}") + emptyMap() + } + } + + private fun JsonObject.toAnyMap(): Map = + entries.associate { (k, v) -> k to jsonElementToAny(v) } + + private fun jsonElementToAny(el: JsonElement): Any? = when (el) { + is JsonPrimitive -> el.content + is JsonArray -> el.map { jsonElementToAny(it) } + is JsonObject -> el.toAnyMap() + } + } +} diff --git a/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentTest.kt b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentTest.kt index b64ab38..b41b5c0 100644 --- a/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentTest.kt +++ b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/agent/ChatAgentTest.kt @@ -10,6 +10,7 @@ import kotlinx.coroutines.test.runTest import pw.binom.agentik.proto.AgentEvent import pw.binom.agentik.proto.Content import pw.binom.agentik.proto.Event as ProtoEvent +import pw.binom.agentik.standalone.llm.LlmBackend import pw.binom.agentik.standalone.llm.LlmConfig import pw.binom.agentik.standalone.persistence.sqlite.SqliteStores import pw.binom.litert.LiteContentPart @@ -19,6 +20,8 @@ import pw.binom.litert.LiteDelta import pw.binom.litert.LiteLlm import pw.binom.litert.LiteMessage import pw.binom.litert.LiteRole +import pw.binom.litert.LiteTool +import pw.binom.litert.LiteToolCall import pw.binom.litert.openai.OpenAiConfig import kotlin.test.AfterTest import kotlin.test.BeforeTest @@ -46,14 +49,21 @@ class ChatAgentTest { stores.close() } - private fun newAgent(): ChatAgent { - val cfg = LlmConfig( - backend = pw.binom.agentik.standalone.llm.LlmBackend.OPENAI, + private fun newAgent( + stores: SqliteStores = this.stores, + llm: LiteLlm = this.fakeLlm, + tools: List = emptyList(), + ): ChatAgent = ChatAgent( + id = "agentik", + stores = stores, + llm = llm, + llmConfig = LlmConfig( + backend = LlmBackend.OPENAI, systemPrompt = "be brief", openai = OpenAiConfig(baseUrl = "http://test", apiKey = "test", model = "test"), - ) - return ChatAgent(id = "agentik", stores = stores, llm = fakeLlm, llmConfig = cfg) - } + ), + tools = tools, + ) @Test fun `createConversation seeds system prompt into working memory`() = runTest { @@ -335,6 +345,27 @@ class ChatAgentTest { override fun close() {} } + @Test + fun `tool-call loop executes registered tool and feeds result back`() = runTest { + val echoTool = object : LiteTool { + override fun describe(): String = """{"type":"function","function":{"name":"echo"}}""" + override fun invoke(arguments: String): String = "echoed: $arguments" + } + val toolLlm = ToolLoopFakeLiteLlm() + val agent = newAgent(llm = toolLlm, tools = listOf(NamedTool("echo", echoTool))) + val conv = agent.createConversation(temp = false) + + conv.send(listOf(pw.binom.agentik.proto.Content.Text("call the tool"))) + + assertEquals(1, toolLlm.toolCallCount, "expected one round-trip through LiteConversation") + assertEquals("echoed: {\"x\":\"hi\"}", toolLlm.lastToolResult, + "expected echo tool invoked with the LLM's args, result fed back via addToolResult") + assertEquals("final reply", toolLlm.finalReplyEmitted, + "expected continuation send after tool result to emit final text") + + agent.close() + } + @Test fun `agentEvents — Created + Deleted flow`() = runTest { val agent = newAgent() @@ -433,3 +464,57 @@ private class FakeLiteConversation( override fun addToolResult(callId: String?, name: String, result: String) { error("not used") } override fun close() {} } + +/** + * LiteLlm который имитирует tool-loop: + * - первый send → LiteDelta(toolCalls=[LiteToolCall("echo", {"x":"hi"})], isDone=true) + * - после addToolResult → продолжение send отдаёт LiteDelta(text="final reply", isDone=true) + */ +private class ToolLoopFakeLiteLlm : LiteLlm { + override val backendName: String = "fake-tool" + override val capabilities: pw.binom.litert.LiteCapabilities? = null + + var toolCallCount: Int = 0 + var lastToolResult: String? = null + var finalReplyEmitted: String? = null + + override fun isInitialized(): Boolean = true + + override fun createConversation(config: LiteConversationConfig): LiteConversation { + return object : LiteConversation { + private val hist = mutableListOf() + override val history: List get() = hist.toList() + override fun sendStream(prompt: String) = sendStreamContents(listOf(LiteContentPart.Text(prompt))) + override fun sendStreamContents(contents: List): Flow { + hist.add(LiteMessage(LiteRole.USER, contents)) + return if (toolCallCount == 0) { + toolCallCount++ + flowOf( + LiteDelta( + text = "", + isDone = true, + toolCalls = listOf(LiteToolCall(name = "echo", arguments = mapOf("x" to "hi"))), + ), + ) + } else { + val reply = "final reply" + finalReplyEmitted = reply + flowOf(LiteDelta(text = reply, isDone = true)) + } + } + override fun send(prompt: String): String = "unused" + override fun sendContents(contents: List): String = "unused" + override fun cancel() {} + override fun tokenCount(): Int = hist.size + override fun addToolResult(callId: String?, name: String, result: String) { + System.err.println("[agentik] DEBUG ToolLoopFakeLiteLlm.addToolResult callId=$callId name=$name result=$result") + lastToolResult = result + } + override fun close() {} + } + } + + override fun infer(request: pw.binom.litert.LiteRequest): String = error("not used") + override fun inferStream(request: pw.binom.litert.LiteRequest): Flow = error("not used") + override fun close() {} +} diff --git a/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/mcp/McpConfigTest.kt b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/mcp/McpConfigTest.kt new file mode 100644 index 0000000..dd1cdd2 --- /dev/null +++ b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/mcp/McpConfigTest.kt @@ -0,0 +1,88 @@ +package pw.binom.agentik.standalone.mcp + +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertTrue + +class McpConfigTest { + + @Test + fun `fromEnv returns empty when env var unset`() { + val cfg = McpConfig.fromEnv { null } + assertTrue(cfg.isEmpty) + assertEquals(emptyList(), cfg.servers) + } + + @Test + fun `fromEnv returns empty when env var blank`() { + val cfg = McpConfig.fromEnv { "" } + assertTrue(cfg.isEmpty) + } + + @Test + fun `fromEnv returns empty when file missing`() { + val cfg = McpConfig.fromEnv { "/tmp/agentik-nonexistent-mcp-${System.nanoTime()}.json" } + assertTrue(cfg.isEmpty) + } + + @Test + fun `fromJson parses stdio server`() { + val json = """ + { "mcpServers": { + "fs": { "command": "npx", "args": ["-y", "fs-mcp"], "env": { "ROOT": "/work" } } + } } + """.trimIndent() + val cfg = McpConfig.fromJson(json) + assertEquals(1, cfg.servers.size) + val s = cfg.servers.single() as McpServerSpec.Stdio + assertEquals("fs", s.name) + assertEquals("npx", s.command) + assertEquals(listOf("-y", "fs-mcp"), s.args) + assertEquals(mapOf("ROOT" to "/work"), s.env) + } + + @Test + fun `fromJson parses http server with headers`() { + val json = """ + { "mcpServers": { + "remote": { "url": "https://example.com/mcp", "headers": { "Authorization": "Bearer X" } } + } } + """.trimIndent() + val cfg = McpConfig.fromJson(json) + val s = cfg.servers.single() as McpServerSpec.Http + assertEquals("remote", s.name) + assertEquals("https://example.com/mcp", s.url) + assertEquals(mapOf("Authorization" to "Bearer X"), s.headers) + } + + @Test + fun `fromJson parses both stdio and http together`() { + val json = """ + { "mcpServers": { + "fs": { "command": "npx", "args": [] }, + "remote":{ "url": "https://example.com/mcp" } + } } + """.trimIndent() + val cfg = McpConfig.fromJson(json) + assertEquals(2, cfg.servers.size) + assertTrue(cfg.servers.any { it is McpServerSpec.Stdio && it.name == "fs" }) + assertTrue(cfg.servers.any { it is McpServerSpec.Http && it.name == "remote" }) + } + + @Test + fun `fromJson skips entries without command or url`() { + val json = """ + { "mcpServers": { + "broken": { "description": "no transport" } + } } + """.trimIndent() + val cfg = McpConfig.fromJson(json) + assertTrue(cfg.isEmpty) + } + + @Test + fun `empty returns empty config`() { + val cfg = McpConfig.empty() + assertTrue(cfg.isEmpty) + } +} diff --git a/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/mcp/McpRegistryTest.kt b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/mcp/McpRegistryTest.kt new file mode 100644 index 0000000..3514e25 --- /dev/null +++ b/standalone/src/jvmTest/kotlin/pw/binom/agentik/standalone/mcp/McpRegistryTest.kt @@ -0,0 +1,48 @@ +package pw.binom.agentik.standalone.mcp + +import pw.binom.agentik.standalone.agent.NamedTool +import pw.binom.litert.LiteTool +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertSame +import kotlin.test.assertTrue + +class McpRegistryTest { + + @Test + fun `empty config produces empty registry`() { + val registry = McpRegistry.fromConfig(McpConfig.empty()) + assertEquals(0, registry.allTools.size) + assertEquals(0, registry.connectedServerCount) + assertEquals(emptyList(), registry.namedTools) + registry.close() + } + + @Test + fun `empty servers list produces empty registry`() { + val registry = McpRegistry(servers = emptyList()) + assertEquals(0, registry.namedTools.size) + assertEquals(0, registry.allTools.size) + registry.close() + } + + @Test + fun `NamedTool holds name and tool reference`() { + val noop: LiteTool = object : LiteTool { + override fun describe(): String = "{}" + override fun invoke(arguments: String): String = "" + } + val nt = NamedTool(name = "server__echo", tool = noop) + assertEquals("server__echo", nt.name) + assertSame(noop, nt.tool) + assertEquals("{}", nt.tool.describe()) + } + + @Test + fun `registry close is idempotent`() { + val registry = McpRegistry.fromConfig(McpConfig.empty()) + registry.close() + registry.close() // should not throw + assertTrue(true) + } +}