Bring up :proto protocol + :server (Ktor) + :client (HTTP) modules; wire :server into standalone with EchoProtoAgent
Major additions: * :proto (KMP submodule) — in-house stateful protocol replacing AG-UI. Agent owns conversation transcript; Conversation.events(after) is a live, replay-free stream; backfill via Conversation.getMessages(after, offset, limit). Each Event carries an Instant date for client-side resume tracking. Sealed hierarchies (Content/Message/Event/AgentEvent) annotated @Serializable with snake_case @SerialName JSON discriminators so the wire format is decoupled from Kotlin class names. * :server (JVM, Ktor 3.1.3) — REST+SSE facade for Agent. Public entry: Route.agentikAgent(agent, path = "/agentik"). Endpoints: create/list/get/patch/delete conversations, POST messages (202), POST interrupt, GET messages, GET conversation events (SSE), GET agent events (SSE), GET /health. Custom Instant serializer for kotlin.time.Instant registered contextually on agentikJson (ISO-8601, ignoreUnknownKeys=true, explicitNulls=false). * :client (JVM, Ktor HTTP Client + CIO) — mirror of :server returning a pw.binom.agentik.proto.Agent backed by HTTP calls. Custom SSE parser since ktor-client-sse is not on the 3.1.3 client classpath. * standalone — EchoProtoAgent (in-memory Agent for :proto), EchoAgent (existing AG-UI echo), both mounted on the same Netty embedded server on port 8080 (/agui and /agentik); A2A stays on its own CIO engine on 8081. EchoProtoAgent smoke-tested end-to-end against :server: all 11 endpoints, including live SSE delivery of StartResponse/AppendText/End event triplets and Agent-level Created/Deleted events. Design notes pinned in: * agentik/IRC-QUESTIONS.md — closed 13-item checklist for the upcoming :irc-server transport (channel = conversation, CTCP for structural events, draft/chathistory for backfill, ImageStore side-channel, etc). * docs/ARCHITECTURE.md — overall layout snapshot.
This commit is contained in:
@@ -0,0 +1,11 @@
|
||||
package pw.binom.agentik.standalone
|
||||
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertTrue
|
||||
|
||||
class PlaceholderTest {
|
||||
@Test
|
||||
fun placeholder() {
|
||||
assertTrue(true)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,48 @@
|
||||
package pw.binom.agentik.standalone
|
||||
|
||||
import pw.binom.a2a.client.A2AClient
|
||||
import pw.binom.a2a.model.AgentCard
|
||||
import pw.binom.a2a.model.Message
|
||||
import pw.binom.a2a.model.Role
|
||||
import pw.binom.a2a.model.Task
|
||||
import pw.binom.a2a.model.TextPart
|
||||
|
||||
data class RemoteAgent(val name: String, val baseUrl: String, val token: String? = null)
|
||||
|
||||
object A2aOutbound {
|
||||
private val clients = mutableMapOf<String, A2AClient>()
|
||||
|
||||
fun remoteAgents(): List<RemoteAgent> =
|
||||
System.getenv("AGENTIK_A2A_AGENTS")
|
||||
?.split(";")
|
||||
?.mapNotNull { entry ->
|
||||
if (!entry.contains("=")) return@mapNotNull null
|
||||
val name = entry.substringBefore("=").trim()
|
||||
val rest = entry.substringAfter("=").split(",", limit = 2)
|
||||
val baseUrl = rest.getOrNull(0)?.trim().orEmpty()
|
||||
if (name.isEmpty() || baseUrl.isEmpty()) return@mapNotNull null
|
||||
val token = rest.getOrNull(1)?.trim()?.takeIf { it.isNotEmpty() }
|
||||
RemoteAgent(name = name, baseUrl = baseUrl, token = token)
|
||||
}
|
||||
?: emptyList()
|
||||
|
||||
private fun clientFor(name: String): A2AClient =
|
||||
clients.getOrPut(name) {
|
||||
val remote =
|
||||
remoteAgents().firstOrNull { it.name == name }
|
||||
?: error("Remote agent $name is not configured (AGENTIK_A2A_AGENTS)")
|
||||
A2AClient.create(baseUrl = remote.baseUrl, bearerToken = remote.token)
|
||||
}
|
||||
|
||||
suspend fun send(name: String, text: String, contextId: String? = null): Task {
|
||||
val message = Message(role = Role.USER, parts = listOf(TextPart(text = text)), contextId = contextId)
|
||||
return clientFor(name).sendMessage(message, contextId)
|
||||
}
|
||||
|
||||
suspend fun agentCard(name: String): AgentCard = clientFor(name).agentCard()
|
||||
|
||||
fun closeAll() {
|
||||
clients.values.forEach { it.close() }
|
||||
clients.clear()
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,21 @@
|
||||
package pw.binom.agentik.standalone
|
||||
|
||||
import pw.binom.a2a.model.Message
|
||||
import pw.binom.a2a.model.Role
|
||||
import pw.binom.a2a.model.TextPart
|
||||
import pw.binom.a2a.server.AgentHandler
|
||||
|
||||
/**
|
||||
* A2A-обработчик-заглушка: эхоит входящее сообщение.
|
||||
* Здесь позже будет реальный агент (LLM / инструменты).
|
||||
*/
|
||||
object EchoA2aHandler : AgentHandler {
|
||||
override suspend fun handle(request: Message, contextId: String?): Message {
|
||||
val text = request.parts.filterIsInstance<TextPart>().joinToString("") { it.text }
|
||||
return Message(
|
||||
role = Role.AGENT,
|
||||
parts = listOf(TextPart("echo: $text")),
|
||||
contextId = contextId,
|
||||
)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,33 @@
|
||||
package pw.binom.agentik.standalone
|
||||
|
||||
import pw.binom.agui.api.agent.Agent
|
||||
import pw.binom.agui.api.event.BaseEvent
|
||||
import pw.binom.agui.api.event.RunFinishedEvent
|
||||
import pw.binom.agui.api.event.RunStartedEvent
|
||||
import pw.binom.agui.api.event.TextMessageContentEvent
|
||||
import pw.binom.agui.api.event.TextMessageEndEvent
|
||||
import pw.binom.agui.api.event.TextMessageStartEvent
|
||||
import pw.binom.agui.api.message.MessageRole
|
||||
import pw.binom.agui.api.run.RunAgentInput
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.flow.flow
|
||||
|
||||
/**
|
||||
* Заглушка AG-UI агента: отвечает последовательностью событий
|
||||
* RUN_STARTED -> TEXT_MESSAGE_* -> RUN_FINISHED. Пока — эхо последнего
|
||||
* пользовательского сообщения. Дальше сюда зайдёт реальный LLM-агент.
|
||||
*/
|
||||
object EchoAgent : Agent {
|
||||
override fun run(input: RunAgentInput): Flow<BaseEvent> = flow {
|
||||
emit(RunStartedEvent(threadId = input.threadId, runId = input.runId))
|
||||
|
||||
val userText = input.messages.lastOrNull { it.role == MessageRole.USER }?.content.orEmpty()
|
||||
val messageId = "msg-${input.runId}"
|
||||
|
||||
emit(TextMessageStartEvent(messageId = messageId, role = MessageRole.ASSISTANT))
|
||||
emit(TextMessageContentEvent(messageId = messageId, delta = "echo: $userText"))
|
||||
emit(TextMessageEndEvent(messageId = messageId))
|
||||
|
||||
emit(RunFinishedEvent(threadId = input.threadId, runId = input.runId))
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,147 @@
|
||||
package pw.binom.agentik.standalone
|
||||
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.Job
|
||||
import kotlinx.coroutines.SupervisorJob
|
||||
import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.flow.MutableSharedFlow
|
||||
import kotlinx.coroutines.flow.asSharedFlow
|
||||
import kotlinx.coroutines.flow.flow
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.sync.Mutex
|
||||
import kotlinx.coroutines.sync.withLock
|
||||
import pw.binom.agentik.proto.Agent
|
||||
import pw.binom.agentik.proto.AgentEvent
|
||||
import pw.binom.agentik.proto.Content
|
||||
import pw.binom.agentik.proto.Conversation
|
||||
import pw.binom.agentik.proto.Event
|
||||
import pw.binom.agentik.proto.Message
|
||||
import java.util.UUID
|
||||
import java.util.concurrent.ConcurrentHashMap
|
||||
import kotlin.time.Instant
|
||||
|
||||
/**
|
||||
* Минимальная in-memory реализация [Agent] из нашего :proto: эхо-агент,
|
||||
* повторяет текст пользователя в [AppendText].
|
||||
*
|
||||
* Не persistent, без очередей ходов и без реальной отмены — stub для проверки
|
||||
* что HTTP-фасад `:server` корректно подключён к агенту и прокидывает все
|
||||
* события/историю согласно контракту.
|
||||
*/
|
||||
class EchoProtoAgent(override val id: String = "echo-proto") : Agent {
|
||||
|
||||
private val conversations = ConcurrentHashMap<String, EchoProtoConversation>()
|
||||
private val agentEvents = MutableSharedFlow<AgentEvent>(replay = 0, extraBufferCapacity = 64)
|
||||
|
||||
override fun createConversation(temp: Boolean): Conversation {
|
||||
val c = EchoProtoConversation(isTemporal = temp)
|
||||
conversations[c.id] = c
|
||||
agentEvents.tryEmit(
|
||||
AgentEvent.Created(
|
||||
date = c.updatedAt,
|
||||
conversationId = c.id,
|
||||
)
|
||||
)
|
||||
return c
|
||||
}
|
||||
|
||||
override suspend fun getConversation(id: String): Conversation? = conversations[id]
|
||||
|
||||
override suspend fun deleteConversation(id: String): Boolean {
|
||||
val removed = conversations.remove(id) ?: return false
|
||||
removed.close()
|
||||
agentEvents.tryEmit(AgentEvent.Deleted(date = Instant.fromEpochMilliseconds(System.currentTimeMillis()), id = id))
|
||||
return true
|
||||
}
|
||||
|
||||
override suspend fun getConversations(offset: Int, limit: Int): List<Conversation> =
|
||||
conversations.values.sortedByDescending { it.updatedAt }.drop(offset).take(limit)
|
||||
|
||||
override fun events(after: Instant): Flow<AgentEvent> = flow {
|
||||
agentEvents.asSharedFlow().collect { event ->
|
||||
if (event.date > after) emit(event)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
internal class EchoProtoConversation(
|
||||
override val id: String = UUID.randomUUID().toString(),
|
||||
override val isTemporal: Boolean,
|
||||
) : Conversation {
|
||||
|
||||
override val isSupportImageInput: Boolean = false
|
||||
override val isSupportImageOutput: Boolean = false
|
||||
override var title: String? = null
|
||||
private set
|
||||
override var updatedAt: Instant = Instant.fromEpochMilliseconds(System.currentTimeMillis())
|
||||
private set
|
||||
|
||||
private val messages = mutableListOf<Message>()
|
||||
private val events = MutableSharedFlow<Event>(replay = 0, extraBufferCapacity = 64)
|
||||
private val mutex = Mutex()
|
||||
private var activeJob: Job? = null
|
||||
private val scope = CoroutineScope(SupervisorJob() + Dispatchers.Default)
|
||||
|
||||
override suspend fun rename(title: String) = mutex.withLock {
|
||||
this.title = title
|
||||
updatedAt = Instant.fromEpochMilliseconds(System.currentTimeMillis())
|
||||
}
|
||||
|
||||
override suspend fun send(content: List<Content>) {
|
||||
val userMessage = Message.UserMessage(
|
||||
id = "um-${UUID.randomUUID()}",
|
||||
content = content,
|
||||
date = Instant.fromEpochMilliseconds(System.currentTimeMillis()),
|
||||
)
|
||||
mutex.withLock {
|
||||
messages.add(userMessage)
|
||||
updatedAt = Instant.fromEpochMilliseconds(System.currentTimeMillis())
|
||||
}
|
||||
|
||||
// stub-политика: прерываем предыдущий ход, запускаем новый.
|
||||
activeJob?.cancel()
|
||||
|
||||
activeJob = scope.launch {
|
||||
delay(20)
|
||||
val now = Instant.fromEpochMilliseconds(System.currentTimeMillis())
|
||||
val echo = content.joinToString(separator = " ") { c ->
|
||||
when (c) {
|
||||
is Content.Text -> c.body
|
||||
is Content.Image -> "[image ${c.mime}, ${c.data.size}B]"
|
||||
}
|
||||
}
|
||||
events.emit(Event.StartResponse(date = now, responseType = Event.ResponseType.TEXT))
|
||||
val body = "echo: $echo"
|
||||
events.emit(Event.AppendText(date = now, body = body))
|
||||
val assistant = Message.AssistantMessage(
|
||||
id = "am-${UUID.randomUUID()}",
|
||||
content = listOf(Content.Text(body)),
|
||||
date = Instant.fromEpochMilliseconds(System.currentTimeMillis()),
|
||||
)
|
||||
mutex.withLock { messages.add(assistant) }
|
||||
events.emit(Event.End(date = Instant.fromEpochMilliseconds(System.currentTimeMillis())))
|
||||
}
|
||||
}
|
||||
|
||||
override suspend fun interrupt() {
|
||||
activeJob?.cancel()
|
||||
activeJob = null
|
||||
events.emit(Event.Interrupted(date = Instant.fromEpochMilliseconds(System.currentTimeMillis())))
|
||||
}
|
||||
|
||||
override fun events(after: Instant): Flow<Event> = flow {
|
||||
events.asSharedFlow().collect { event ->
|
||||
if (event.date > after) emit(event)
|
||||
}
|
||||
}
|
||||
|
||||
override suspend fun getMessages(after: Instant, offset: Int, limit: Int): List<Message> =
|
||||
mutex.withLock { messages.asSequence().filter { it.date > after }.drop(offset).take(limit).toList() }
|
||||
|
||||
override fun close() {
|
||||
activeJob?.cancel()
|
||||
scope.coroutineContext[Job]?.cancel()
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,66 @@
|
||||
package pw.binom.agentik.standalone
|
||||
|
||||
import io.ktor.server.engine.embeddedServer
|
||||
import io.ktor.server.netty.Netty
|
||||
import io.ktor.server.response.respondText
|
||||
import io.ktor.server.routing.get
|
||||
import io.ktor.server.routing.routing
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import pw.binom.a2a.server.A2AServer
|
||||
import pw.binom.agentik.proto.Conversation
|
||||
import pw.binom.agentik.proto.Message
|
||||
import pw.binom.agentik.server.agentikAgent
|
||||
import pw.binom.agui.server.aguiAgent
|
||||
|
||||
/**
|
||||
* standalone-контейнер agentik:
|
||||
* - AG-UI: встраиваемый Ktor (Netty), порт AGENTIK_PORT (default 8080)
|
||||
* POST /agui -> text/event-stream (протокол AG-UI, агент [EchoAgent])
|
||||
* GET /health -> "ok"
|
||||
* - A2A: встраиваемый Ktor (CIO), порт AGENTIK_A2A_PORT (default 8081)
|
||||
* POST / -> JSON-RPC (message/send, tasks/get, tasks/cancel)
|
||||
* GET /.well-known/agent-card.json
|
||||
* - :server (proto): встраиваемый Ktor (Netty), порт AGENTIK_PORT (default 8080) — общий с AG-UI
|
||||
* POST /agentik/conversations -> ConversationSnapshot (201)
|
||||
* GET /agentik/conversations -> [ConversationSnapshot]
|
||||
* GET /agentik/conversations/{id} -> ConversationSnapshot
|
||||
* PATCH /agentik/conversations/{id} -> ConversationSnapshot
|
||||
* DELETE /agentik/conversations/{id} -> 204
|
||||
* POST /agentik/conversations/{id}/messages -> 202
|
||||
* POST /agentik/conversations/{id}/interrupt -> 202
|
||||
* GET /agentik/conversations/{id}/messages -> [Message]
|
||||
* GET /agentik/conversations/{id}/events -> text/event-stream (SSE)
|
||||
* GET /agentik/events -> text/event-stream (SSE, Agent-level)
|
||||
*
|
||||
* Для обращения к другим агентам: pw.binom.a2a.client.A2AClient.create(baseUrl, token).
|
||||
* Для in-process вызова :server: pw.binom.agentik.client.AgentikAgent(id, baseUrl, httpClient).
|
||||
*/
|
||||
fun main() {
|
||||
val aguiPort = System.getenv("AGENTIK_PORT")?.toIntOrNull() ?: 8080
|
||||
val a2aPort = System.getenv("AGENTIK_A2A_PORT")?.toIntOrNull() ?: 8081
|
||||
|
||||
// A2A: agent <-> agent (CIO). Нестреляющий старт — свой event-loop.
|
||||
val a2a = A2AServer.create(port = a2aPort, agentName = "agentik", handler = EchoA2aHandler)
|
||||
runBlocking { a2a.start() }
|
||||
println("A2A server -> http://localhost:$a2aPort/ (JSON-RPC: message/send, tasks/get, tasks/cancel)")
|
||||
|
||||
// shared port: AG-UI (Netty) + :server (Netty) живут на 8080, A2A — на 8081.
|
||||
val agui =
|
||||
embeddedServer(Netty, port = aguiPort) {
|
||||
routing {
|
||||
get("/health") { call.respondText("ok") }
|
||||
aguiAgent(EchoAgent, path = "/agui")
|
||||
agentikAgent(EchoProtoAgent(), path = "/agentik")
|
||||
}
|
||||
}
|
||||
println("AG-UI -> http://localhost:$aguiPort/agui (SSE), /health")
|
||||
println(":server proto -> http://localhost:$aguiPort/agentik/... (REST+SSE, агент [EchoProtoAgent])")
|
||||
agui.start(wait = true)
|
||||
}
|
||||
|
||||
// References for IDE noise suppression when source not auto-imported.
|
||||
@Suppress("unused")
|
||||
private val keepReferences: Array<Class<*>> = arrayOf(
|
||||
Conversation::class.java,
|
||||
Message::class.java,
|
||||
)
|
||||
Reference in New Issue
Block a user