standalone: add MCP client + tool-call loop

- 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<NamedTool> 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
This commit is contained in:
2026-09-13 15:59:58 +03:00
parent afdfb37e35
commit 4c66947c25
12 changed files with 752 additions and 51 deletions
@@ -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=<path>.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()
})
@@ -31,6 +31,7 @@ class ChatAgent(
private val stores: SqliteStores,
private val llm: LiteLlm,
private val llmConfig: LlmConfig,
private val tools: List<NamedTool> = emptyList(),
) : ProtoAgent, AutoCloseable {
private val agentEvents = MutableSharedFlow<AgentEvent>(
@@ -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 }
}
}
@@ -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<NamedTool> = 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<String, NamedTool> = tools.associateBy { it.name }
private val events = MutableSharedFlow<ProtoEvent>(
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<LiteContentPart> = parts
var loopGuard = 0
while (loopGuard++ < MAX_TOOL_LOOPS) {
val collectedCalls = mutableListOf<LiteToolCall>()
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 { "<empty result>" }
} 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, Any?>): String = encodeToolArgs(arguments)
companion object {
private const val MAX_TOOL_LOOPS = 16
}
}
internal fun List<Content>.toLiteContents(): List<LiteContentPart> = 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,
@@ -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)
@@ -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<String> = emptyList(),
val env: Map<String, String> = emptyMap(),
) : McpServerSpec
/** http (Streamable HTTP transport): подключаемся к существующему MCP-серверу. */
data class Http(
override val name: String,
val url: String,
val headers: Map<String, String> = 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<McpServerSpec>,
) {
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
}
}
}
@@ -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<McpServerSpec>,
private val httpClient: HttpClient = defaultHttpClient(),
private val clientName: String = "agentik",
private val clientVersion: String = "1.0.0",
) : AutoCloseable {
private val connected: MutableMap<String, ConnectedServer> = ConcurrentHashMap()
@Volatile private var closed = false
/** Все [LiteTool] со всех подключённых серверов. */
val allTools: List<LiteTool> by lazy {
connected.values.flatMap { it.tools }
}
/** Все [LiteTool] с именами (server__tool), которые видит LLM. */
val namedTools: List<NamedTool> by lazy {
allTools.filterIsInstance<McpLiteToolAdapter>().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<LiteTool>,
)
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": "<server>__<tool>",
* "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<String, Any?> {
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<String, Any?> =
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()
}
}
}
@@ -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<NamedTool> = 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<LiteMessage>()
override val history: List<LiteMessage> get() = hist.toList()
override fun sendStream(prompt: String) = sendStreamContents(listOf(LiteContentPart.Text(prompt)))
override fun sendStreamContents(contents: List<LiteContentPart>): Flow<LiteDelta> {
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<LiteContentPart>): 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<LiteDelta> = error("not used")
override fun close() {}
}
@@ -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)
}
}
@@ -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)
}
}