refactor(standalone): extract modules, event-driven background, AppConfig
ci / JVM build + tests (push) Failing after 2m5s
ci / JVM build + tests (push) Failing after 2m5s
Standalone refactor — modularity + correctness improvements after
STANDALONE-REVIEW findings. Touches ~30 files. Build green, 178 tests pass.
(1) Module extractions — generic components out of :standalone:
• :llm-tools (new KMP module, package pw.binom.agentik.llm.tools)
- LlmReflector, SkillMiner, LlmMemoryReviewer, LiteLlmContextCompactor
- Parsers: ReflectionParser, SkillMiningParser, ReviewDecisionParser
- Prompts: ReflectionPrompts, SkillMiningPrompts, ReviewPrompts
• :mcp-bridge (new JVM module, package pw.binom.agentik.mcp.bridge)
- McpConfig, McpRegistry, McpLiteToolAdapter
• NamedTool moved from :standalone to :agent-toolsets/commonMain
- Generic (name + LiteTool) wrapper, used by both :mcp-bridge
and :standalone's tool dispatcher
:standalone loses ~1400 lines, depends on the two new modules.
(2) Background work → event-driven (no more interval-polling):
• New :standalone/agent/BackgroundEvents.kt — internal event bus:
- ToolCallEvent.Succeeded/Failed (emitted by ToolDispatcher after invoke)
- CompactionEvent.Triggered (emitted by CompactionCoordinator pre-delete)
- ConversationLifecycleEvent.Closing (emitted by ConversationLoop.close)
• BackgroundScheduler rewritten as event subscriber:
- On Closing: final reflection + skill mining (last-chance extraction)
- On Compaction (turnsToDelete > 10): skill mining (debounced 60s)
- On ToolFailure x2 in 60s window: reflection (debounced 5min)
- Dropped: maybeScheduleReview/Reflection/SkillMining (interval-based)
- Dropped config: memoryReviewInterval, reflectionInterval, skillMiningInterval
• ToolDispatcher emits ToolCallEvent after each invoke.
• CompactionCoordinator emits CompactionEvent before workingMemory.compact().
• ConversationLoop.close() emits Closing BEFORE agentScope.cancel() so the
subscription gets to run final reflection/mining.
Net effect: typical 30-turn conversation runs ~38 LLM calls (was: 30 main +
3 review + 3 reflection + 2 mining). With event-driven, review/mining only fire
when their triggers actually make sense (compaction about to delete, or
conversation closing).
(3) AppConfig single source of truth:
• Replaces AgentikConfig + LlmConfig.fromEnv + McpConfig.fromEnv with one
AppConfig.fromEnv() that reads all ~25 env vars in a single pass.
• Sections: AgentSection, LlmSection, McpSection, MemorySection,
EmbeddingSection, ReflectionSection, SkillMiningSection, DebugSection.
• OPENAI_CONTEXT_WINDOW / AGENTIK_GOOGLE_CONTEXT_WINDOW no longer
read twice (was a bug per STANDALONE-REVIEW E3).
(4) Other fixes inherited from earlier waves:
• Hardening — size caps on user-input boundaries:
MAX_MEMORY_CONTENT_LEN=32KB, MAX_SKILL_BODY_LEN=64KB,
MAX_MCP_CONFIG_BYTES=1MB, MAX_A2A_REPLY_LEN=10MB, MAX_PORT=65535,
blank-rejection in LlmConfig.requireEnv, URL/command validation.
• Single scope — :standalone/agent/ConversationLoop has one
agentScope (was: scope + backgroundScope).
• liteConvRef race fix — capture-then-use pattern replaces !!-after-read;
close() + runTurn.finally race on LiteConv JNI handled via
AtomicReference.getAndSet.
• SkillMiner.maxTurns / LlmReflector.maxTurns exposed as public (needed
by BackgroundScheduler for prompt sizing).
• Tests: MemoryWiringTest updated for new compaction-triggered review
behavior; all parser/test imports updated for new packages.
Test results: 178/178 in :standalone, 36/36 in :agent-toolsets — all green.
This commit is contained in:
@@ -0,0 +1,39 @@
|
||||
// Generic MCP (Model Context Protocol) bridge — переиспользуемый модуль,
|
||||
// который превращает любой MCP-сервер (stdio subprocess или HTTP endpoint)
|
||||
// в набор [LiteTool]-адаптеров.
|
||||
//
|
||||
// Вынесен из :standalone — MCP не специфичен для standalone'а, это generic
|
||||
// мост между MCP-SDK и litert-kmp. Может переиспользоваться в :agentik-cli
|
||||
// или :agentik-tui когда те снова включатся.
|
||||
//
|
||||
// Зависимости:
|
||||
// - :agent-toolsets для NamedTool (обёртка для LiteTool + имя-как-видит-модель)
|
||||
// - litert.api для LiteTool контракта
|
||||
// - MCP SDK (JVM-only)
|
||||
// - Ktor client (для StreamableHttpClientTransport)
|
||||
// - kotlinx-serialization для парсинга конфига
|
||||
plugins {
|
||||
alias(libs.plugins.kotlin.jvm)
|
||||
alias(libs.plugins.kotlin.serialization)
|
||||
}
|
||||
|
||||
kotlin {
|
||||
jvmToolchain(21)
|
||||
}
|
||||
|
||||
dependencies {
|
||||
implementation(project(":agent-toolsets"))
|
||||
|
||||
api(libs.litert.api)
|
||||
|
||||
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)
|
||||
|
||||
implementation(libs.kotlinx.coroutines.core)
|
||||
implementation(libs.kotlinx.serialization.json)
|
||||
|
||||
implementation(libs.kotlin.logging)
|
||||
}
|
||||
@@ -0,0 +1,105 @@
|
||||
package pw.binom.agentik.mcp.bridge
|
||||
|
||||
|
||||
import kotlinx.serialization.SerialName
|
||||
import kotlinx.serialization.Serializable
|
||||
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)
|
||||
*/
|
||||
@Serializable
|
||||
sealed interface McpServerSpec {
|
||||
/** Уникальное имя сервера, как его видно в логах и в tool-prefix. */
|
||||
val name: String
|
||||
|
||||
/** stdio: спавним процесс, читаем его stdout, пишем в stdin. */
|
||||
@Serializable
|
||||
@SerialName("stdio")
|
||||
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-серверу. */
|
||||
@Serializable
|
||||
@SerialName("http")
|
||||
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`. Если не задан или файл не существует — пустой список.
|
||||
*/
|
||||
@Serializable
|
||||
data class McpConfig(
|
||||
val servers: List<McpServerSpec> = emptyList(),
|
||||
) {
|
||||
val isEmpty: Boolean get() = servers.isEmpty()
|
||||
|
||||
companion object {
|
||||
private val log = mu.KotlinLogging.logger {}
|
||||
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()) {
|
||||
log.warn { "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)
|
||||
}
|
||||
log.warn { "MCP server '$name' has neither 'url' nor 'command' — skipped" }
|
||||
return null
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,259 @@
|
||||
package pw.binom.agentik.mcp.bridge
|
||||
|
||||
import mu.KotlinLogging
|
||||
|
||||
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.booleanOrNull
|
||||
import kotlinx.serialization.json.buildJsonObject
|
||||
import kotlinx.serialization.json.doubleOrNull
|
||||
import kotlinx.serialization.json.intOrNull
|
||||
import kotlinx.serialization.json.longOrNull
|
||||
import kotlinx.serialization.json.put
|
||||
import pw.binom.litert.LiteTool
|
||||
import java.util.concurrent.ConcurrentHashMap
|
||||
import pw.binom.agentik.toolsets.NamedTool
|
||||
|
||||
/**
|
||||
* Реестр подключённых MCP-серверов.
|
||||
*
|
||||
* На старте подключается ко всем [McpServerSpec] из [McpConfig], у каждого запрашивает
|
||||
* список tools и оборачивает их в [LiteTool]-адаптеры ([McpLiteToolAdapter]). Все
|
||||
* адаптеры собираются в [allTools], который `ChatConversation` подмешивает в
|
||||
* [pw.binom.litert.LiteConversationConfig.tools].
|
||||
*
|
||||
* [close] убивает stdio-процессы и закрывает HTTP-клиент.
|
||||
*/
|
||||
private val log = KotlinLogging.logger {}
|
||||
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 ->
|
||||
log.warn { "MCP server '${spec.name}' failed to connect: ${e.message}" }
|
||||
}
|
||||
}
|
||||
log.warn { "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(" ")
|
||||
log.warn { "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 -> {
|
||||
log.warn { "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-дескриптор — flat OpenAPI-спецификация
|
||||
* (формат, который напрямую принимает LiteRT-LM; litert-openai сам оборачивает
|
||||
* её в OpenAI-формат):
|
||||
*
|
||||
* ```json
|
||||
* {
|
||||
* "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 {
|
||||
// Flat OpenAPI-спецификация (name/description/parameters) — формат LiteRT-LM.
|
||||
// litert-openai оборачивает её в OpenAI-формат сам (normalizeToolDescriptor).
|
||||
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) {
|
||||
log.warn { "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) }
|
||||
|
||||
/**
|
||||
* Преобразует [JsonPrimitive] в типизированное значение (Boolean/Int/Long/Float/Double),
|
||||
* и только в крайнем случае — в String. Без этого MCP-сервер получает все аргументы
|
||||
* как строки и отвергает их по JSON-schema (например, `max_length:500` → `"500"` →
|
||||
* «'500' is not of type 'integer'»).
|
||||
*/
|
||||
private fun jsonElementToAny(el: JsonElement): Any? = when (el) {
|
||||
is JsonPrimitive ->
|
||||
el.booleanOrNull
|
||||
?: el.intOrNull
|
||||
?: el.longOrNull
|
||||
?: el.doubleOrNull
|
||||
?: if (el.isString) el.content else el.content
|
||||
is JsonArray -> el.map { jsonElementToAny(it) }
|
||||
is JsonObject -> el.toAnyMap()
|
||||
else -> null
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user