a2a: подключить A2A-транспорт (POST /a2a + agent-card) через A2aBridge
:standalone декларировал зависимость a2a-server, но mount не было (e2e нашёл 404 на /a2a). A2aBridge (AgentHandler) гоняет A2A-context на :proto-диалог: contextId -> Conversation (пустой/неизвестный -> новый), ответ = склеенные AppendText хода (подписка на events() до send, стоп по End/Interrupted/Error), id диалога в metadata.agentikConversationId.
This commit is contained in:
@@ -0,0 +1,98 @@
|
||||
package pw.binom.agentik.standalone
|
||||
|
||||
import kotlinx.coroutines.CompletableDeferred
|
||||
import kotlinx.coroutines.async
|
||||
import kotlinx.coroutines.coroutineScope
|
||||
import kotlinx.serialization.json.JsonPrimitive
|
||||
import kotlinx.serialization.json.buildJsonObject
|
||||
import mu.KotlinLogging
|
||||
import pw.binom.a2a.model.Message
|
||||
import pw.binom.a2a.model.Role
|
||||
import pw.binom.a2a.model.TextPart
|
||||
import pw.binom.a2a.server.AgentHandler
|
||||
import pw.binom.agentik.proto.Agent
|
||||
import pw.binom.agentik.proto.Content
|
||||
import pw.binom.agentik.proto.Conversation
|
||||
import pw.binom.agentik.proto.Event
|
||||
import java.util.concurrent.ConcurrentHashMap
|
||||
|
||||
private val log = KotlinLogging.logger {}
|
||||
|
||||
/**
|
||||
* Адаптер :proto [Agent] к A2A [AgentHandler] (pw.binom.a2a:server).
|
||||
*
|
||||
* Контекст A2A (contextId) мапится на диалог :proto:
|
||||
* - неизвестный/пустой contextId -> createConversation(temp=false);
|
||||
* - известный contextId -> getConversation(id); если диалог пропал (удаляли/рестартились)
|
||||
* -> новый диалог, контекст пересоздаётся.
|
||||
* - id внутреннего диалога отдаётся клиенту в `metadata.agentikConversationId` ответа.
|
||||
*
|
||||
* Ответ A2A = склеенные [Event.AppendText] нашего хода. Подписку на [Conversation.events]
|
||||
* открываем ДО [Conversation.send] (иначе события начала хода могут быть упущены),
|
||||
* завершение хода ждём по [Event.End] / [Event.Interrupted] / [Event.Error].
|
||||
*
|
||||
* Ограничение v1: tool-события и картинки в A2A-ответ не транслируются;
|
||||
* при нескольких ходов в очереди за контекстом текст предыдущего хода
|
||||
* может попасть в ответ.
|
||||
*/
|
||||
class A2aBridge(private val agent: Agent) : AgentHandler {
|
||||
|
||||
private val contextToConversation = ConcurrentHashMap<String, String>()
|
||||
|
||||
override suspend fun handle(request: Message, contextId: String?): Message = coroutineScope {
|
||||
val text = request.parts
|
||||
.filterIsInstance<TextPart>()
|
||||
.joinToString("\n") { it.text }
|
||||
val conv = resolveConversation(contextId)
|
||||
|
||||
val since = conv.updatedAt
|
||||
val reply = StringBuilder()
|
||||
val turnDone = CompletableDeferred<Unit>()
|
||||
val subscription = async {
|
||||
conv.events(since).collect { e ->
|
||||
when (e) {
|
||||
is Event.AppendText -> reply.append(e.body)
|
||||
is Event.End, is Event.Interrupted -> turnDone.complete(Unit)
|
||||
is Event.Error ->
|
||||
if (!turnDone.completeExceptionally(
|
||||
IllegalStateException("agent turn failed: ${e.message}")
|
||||
)
|
||||
) {}
|
||||
else -> {}
|
||||
}
|
||||
}
|
||||
}
|
||||
conv.send(listOf(Content.Text(text)))
|
||||
try {
|
||||
turnDone.await()
|
||||
} finally {
|
||||
subscription.cancel()
|
||||
}
|
||||
log.info { "a2a context=$contextId conv=${conv.id} reply=${reply.length} chars" }
|
||||
Message(
|
||||
role = Role.AGENT,
|
||||
parts = listOf(TextPart(reply.toString().ifEmpty { "(empty response)" })),
|
||||
metadata = buildJsonObject {
|
||||
put("agentikConversationId", JsonPrimitive(conv.id))
|
||||
},
|
||||
)
|
||||
}
|
||||
|
||||
private suspend fun resolveConversation(contextId: String?): Conversation {
|
||||
if (contextId.isNullOrBlank()) {
|
||||
val conv = agent.createConversation(temp = false)
|
||||
contextToConversation[conv.id] = conv.id
|
||||
log.info { "a2a: new context -> conv=${conv.id}" }
|
||||
return conv
|
||||
}
|
||||
val existing = contextToConversation[contextId]
|
||||
if (existing != null) {
|
||||
agent.getConversation(existing)?.let { return it }
|
||||
log.info { "a2a: context=$contextId conv=$existing lost, creating new" }
|
||||
}
|
||||
val conv = agent.createConversation(temp = false)
|
||||
contextToConversation[contextId] = conv.id
|
||||
log.info { "a2a: context=$contextId -> conv=${conv.id}" }
|
||||
return conv
|
||||
}
|
||||
}
|
||||
@@ -10,6 +10,7 @@ import io.ktor.server.routing.get
|
||||
import io.ktor.server.routing.routing
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.io.files.Path
|
||||
import pw.binom.a2a.server.a2aAgent
|
||||
import pw.binom.agentik.memory.MemoryReviewer
|
||||
import pw.binom.agentik.memory.MemorySystem
|
||||
import pw.binom.agentik.memory.md.openMdMemorySystem
|
||||
@@ -41,6 +42,8 @@ import java.io.File
|
||||
* GET /agentik/conversations/{id}/messages -> [Message]
|
||||
* GET /agentik/conversations/{id}/events -> text/event-stream (SSE)
|
||||
* GET /agentik/events -> text/event-stream (SSE, Agent-level)
|
||||
* POST /a2a/ -> A2A JSON-RPC (message/send, tasks/get, tasks/cancel)
|
||||
* GET /a2a/.well-known/agent-card.json -> AgentCard
|
||||
* GET /health -> "ok"
|
||||
*
|
||||
* Вся конфигурация — [AgentikConfig.fromEnv] (см. [AgentikConfig]). Источники:
|
||||
@@ -215,6 +218,7 @@ fun main() {
|
||||
routing {
|
||||
get("/health") { call.respondText("ok") }
|
||||
agentikAgent(agent, path = "/agentik")
|
||||
a2aAgent(agentName = "agentik", handler = A2aBridge(agent), path = "/a2a")
|
||||
if (config.debugEndpoints) {
|
||||
debugRoutes(
|
||||
agent = agent,
|
||||
@@ -231,6 +235,8 @@ fun main() {
|
||||
println(" GET /health")
|
||||
println(" POST /agentik/conversations -> 201")
|
||||
println(" GET /agentik/conversations/{id}/events -> SSE")
|
||||
println(" POST /a2a/ -> A2A JSON-RPC (message/send, tasks/get, tasks/cancel)")
|
||||
println(" GET /a2a/.well-known/agent-card.json -> AgentCard")
|
||||
println(" storage: ${config.dbPath}")
|
||||
println(" llm: ${config.llm.backend} ${config.llm.modelInfo()}")
|
||||
println(" mcp: ${mcpRegistry.allTools.size} tools from ${mcpRegistry.connectedServerCount} servers")
|
||||
|
||||
Reference in New Issue
Block a user