diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..1e1b6b8 --- /dev/null +++ b/.gitignore @@ -0,0 +1,21 @@ +# Gradle +.gradle/ +build/ +**/build/ + +# Kotlin +*.iml +.kotlin/ + +# IDE +.idea/ +*.ipr +*.iws +out/ + +# OS +.DS_Store + +# Local tooling (Magic Context, IDE plugins, MCP configs) +.cortexkit/ +.veai/ \ No newline at end of file diff --git a/IRC-QUESTIONS.md b/IRC-QUESTIONS.md new file mode 100644 index 0000000..fc5d0c7 --- /dev/null +++ b/IRC-QUESTIONS.md @@ -0,0 +1,40 @@ +# IRC-транспорт — открытые вопросы + +По мере закрытия отмечаем `- N. [x]`. Закрытый вопрос остаётся в файле с принятым решением. + +- 1. [x] **История.** Принято: новый абстрактный метод `suspend fun getLatestMessages(offset: Int, limit: Int): List` в `:proto.Conversation` (offset = пропустить С КОНЦА, 0 = самые свежие). IRC-сервер не держит своего буфера, на `CHATHISTORY` дёргает агента. CAP `draft/chathistory` объявляем. +- 2. [x] **Tool/Error/Image события — раскладка по IRC.** Принято. Каждый `Event` мапится: + - `StartReasoning` → дропаем с провода + - `StartResponse(TEXT|IMAGE)` → CTCP `AGENTIK response-start {"type":"text"|"image"}` + - `AppendText(body)` → `PRIVMSG #chan :body` + - `AppendImage(body, mime)` → через `ImageStore` → CTCP `AGENTIK image {"url":..,"mime":..,"ttl":..}` + - `ToolCall` → CTCP `AGENTIK tool-call {json}` + - `ToolResult` → CTCP `AGENTIK tool-result {json}` + - `Error` → CTCP `AGENTIK error {json}` + - `End` → CTCP `AGENTIK end` + - `Interrupted` → CTCP `AGENTIK interrupted` +- 3. [x] **Interrupt.** **Упрощение:** команды `/stop` и `/interrupt` в `PRIVMSG` (т.е. `PRIVMSG #chan :/stop`) вызывают `Conversation.interrupt()`. Если `PRIVMSG` приходит во время активного размышления — сервер сначала зовёт `interrupt()`, затем `send(content)`. CTCP-вариант дропаем. +- 4. [x] **`AgentEvent.Created/Deleted/Renamed` маппинг.** Принято: Created = IRC `JOIN`-бродкаст; Deleted = `KICK` самого себя; Renamed = `TOPIC #foo :new title`. +- 5. [x] **NICK агента.** Принято: параметр в DSL, дефолт `"Agent"`. +- 6. [x] **Multi-user в канале.** Принято: **в канале всегда только наш агент и наш пользователь. Других не будет никогда.** +- 7. [x] **Маппинг канал ↔ Conversation.** Принято: имя IRC-канала = `Conversation.title`; `Conversation.id` = UUID, выдаётся через `CTCP AGENTIK id #foo`; при переименовании канала id стабилен. +- 8. [x] **Создание канала.** Принято: `JOIN #foo` → создаём `Conversation(title="foo", id=)`. Если уже есть — заходим. +- 9. [x] **Удаление канала.** Принято: `PART` закрывает сторону клиента; `CTCP AGENTIK delete #foo` — удаление Conversation-а. +- 10. [x] **Модуль.** Принято: `:irc-server`, KMP через kotlinx-io. Также модуль содержит HTTP staging-эндпоинт для картинок (см. п.13). +- 11. [x] **Аутентификация клиента.** Принято: без auth, любой может подключиться. +- 12. [x] **Capabilities (ircv3).** Принято: в первом проходе объявляем `server-time`, `message-tags`, `batch`, `draft/chathistory`. SASL не объявляем (п.11). Остальные CAPs (echo-message, labeled-response, standard-replies, multi-user stuff) добавляем инкрементально. +- 13. [x] **ImageStore.** Принято: `ImageStore` живёт в `:irc-server`, дефолтная in-memory реализация с **TTL 600 сек**, staging-порт **авто-pick** (0 → свободный). Клиент через IRC картинки **не шлёт** (для этого HTTP `:server`). Конкретную реализацию `ImageStore` пользователь сделает позже сам, в первом проходе — наша in-memory. + +## Все вопросы закрыты + +Итого решений по `:irc-server`: +- Модуль `:irc-server`, KMP через kotlinx-io. +- Канал IRC = `Conversation` (имя = title, UUID через CTCP `AGENTIK id`). +- В канале всегда только 1 пользователь + агент (ник `Agent` по умолчанию). +- `Event` → IRC: `PRIVMSG` для текста, CTCP `AGENTIK <имя> {json}` для всего остального. `StartReasoning` дропается. +- Interrupt через `/stop` / `/interrupt` в PRIVMSG; входящий PRIVMSG во время размышления = `interrupt()` + `send()`. +- История через `CHATHISTORY` (LATEST/BEFORE/BETWEEN/AFTER), сервер не буферизует, дёргает новый `:proto` метод `getLatestMessages(offset, limit)`. +- Картинки только agent → client через `ImageStore` + HTTP staging в том же модуле. +- CAPs: `server-time`, `message-tags`, `batch`, `draft/chathistory`. Без auth. + +Можно кодить. diff --git a/build.gradle.kts b/build.gradle.kts new file mode 100644 index 0000000..d508bc0 --- /dev/null +++ b/build.gradle.kts @@ -0,0 +1,8 @@ +plugins { + alias(libs.plugins.kotlin.multiplatform) apply false + alias(libs.plugins.kotlin.jvm) apply false + alias(libs.plugins.kotlin.serialization) apply false +} + +group = "pw.binom.agentik" +version = "0.1.0" diff --git a/client/build.gradle.kts b/client/build.gradle.kts new file mode 100644 index 0000000..f689274 --- /dev/null +++ b/client/build.gradle.kts @@ -0,0 +1,24 @@ +import org.jetbrains.kotlin.gradle.dsl.JvmTarget + +plugins { + alias(libs.plugins.kotlin.jvm) + alias(libs.plugins.kotlin.serialization) +} + +kotlin { + compilerOptions { + jvmTarget.set(JvmTarget.JVM_21) + } +} + +dependencies { + implementation(project(":proto")) + + 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) +} diff --git a/client/src/main/kotlin/pw/binom/agentik/client/AgentClient.kt b/client/src/main/kotlin/pw/binom/agentik/client/AgentClient.kt new file mode 100644 index 0000000..15423a7 --- /dev/null +++ b/client/src/main/kotlin/pw/binom/agentik/client/AgentClient.kt @@ -0,0 +1,78 @@ +package pw.binom.agentik.client + +import io.ktor.client.HttpClient +import io.ktor.client.call.body +import io.ktor.client.request.delete +import io.ktor.client.request.get +import io.ktor.client.request.parameter +import io.ktor.client.request.post +import io.ktor.client.request.setBody +import io.ktor.client.statement.bodyAsChannel +import io.ktor.http.ContentType +import io.ktor.http.HttpStatusCode +import io.ktor.http.contentType +import kotlinx.coroutines.flow.Flow +import kotlinx.coroutines.flow.flow +import kotlinx.coroutines.runBlocking +import pw.binom.agentik.proto.Agent +import pw.binom.agentik.proto.AgentEvent +import pw.binom.agentik.proto.Conversation +import kotlin.time.Instant + +/** + * HTTP-реализация [Agent]. Ходит в `:server`-фасад, см. `agentikAgent(...)`. + * + * Замечание по [createConversation]: интерфейс [Agent] объявлен не-suspend + * (in-process кейс этого не требует), но HTTP-вариант обязан ждать ответа + * POST `/conversations`. Используем `runBlocking` — это одноразовая + * операция (открытие чата), не горячий путь. В UI-контексте вызывающий сам + * решает, что делать. + */ +internal class AgentClient( + private val httpClient: HttpClient, + private val baseUrl: String, + override val id: String, +) : Agent { + + private val agentUrl: String = baseUrl.trimEnd('/') + + override fun createConversation(temp: Boolean): Conversation = + runBlocking { + val snapshot: ConversationSnapshot = httpClient.post("$agentUrl/conversations") { + contentType(ContentType.Application.Json) + setBody(RequestCreateConversation(temp)) + }.body() + ConversationClient(httpClient = httpClient, baseUrl = agentUrl, snapshot = snapshot) + } + + override suspend fun getConversation(id: String): Conversation? { + val response = httpClient.get("$agentUrl/conversations/$id") + if (response.status == HttpStatusCode.NotFound) return null + val snapshot = response.body() + return ConversationClient(httpClient, agentUrl, snapshot) + } + + override suspend fun deleteConversation(id: String): Boolean { + val response = httpClient.delete("$agentUrl/conversations/$id") + return response.status == HttpStatusCode.NoContent + } + + override suspend fun getConversations(offset: Int, limit: Int): List { + val snapshots = httpClient.get("$agentUrl/conversations") { + parameter("offset", offset) + parameter("limit", limit) + }.body>() + return snapshots.map { ConversationClient(httpClient, agentUrl, it) } + } + + override fun events(after: Instant): Flow = flow { + val response = httpClient.get("$agentUrl/events?after=$after") + check(response.status == HttpStatusCode.OK) { + "events: server returned ${response.status}" + } + readSse(response.bodyAsChannel()) + .collect { payload -> + emit(agentikJson.decodeFromString(AgentEvent.serializer(), payload)) + } + } +} diff --git a/client/src/main/kotlin/pw/binom/agentik/client/AgentikAgent.kt b/client/src/main/kotlin/pw/binom/agentik/client/AgentikAgent.kt new file mode 100644 index 0000000..6737836 --- /dev/null +++ b/client/src/main/kotlin/pw/binom/agentik/client/AgentikAgent.kt @@ -0,0 +1,43 @@ +package pw.binom.agentik.client + +import io.ktor.client.HttpClient +import io.ktor.client.engine.cio.CIO +import io.ktor.client.plugins.contentnegotiation.ContentNegotiation +import io.ktor.serialization.kotlinx.json.json +import pw.binom.agentik.proto.Agent + +/** + * Создаёт [Agent], который под капотом ходит в HTTP-фасад `agentikAgent` + * (модуль `:server`). + * + * ``` + * val client = AgentikAgent( + * id = "my-agent", + * baseUrl = "http://localhost:8080/agentik", + * ) + * val conv = client.createConversation(temp = false) + * conv.send(listOf(Content.Text("hi"))) + * conv.events(Instant.DISTANT_PAST).collect { ev -> ... } + * ``` + * + * [id] пробрасывается в реализацию [Agent.id] — сервер про идентичность + * агента не знает, поэтому клиент должен её знать сам (или взять из + * конфига). + * + * [httpClient] по умолчанию — [defaultAgentikHttpClient] (CIO + JSON + + * SSE). Можно передать свой, если нужен свой engine/логирование/аутентификация. + */ +fun AgentikAgent( + id: String, + baseUrl: String, + httpClient: HttpClient = defaultAgentikHttpClient(), +): Agent = AgentClient(httpClient = httpClient, baseUrl = baseUrl, id = id) + +/** + * Дефолтный [HttpClient] для общения с `agentikAgent`: CIO-движок и + * kotlinx-serialization с тем же wire-форматом, что на сервере. SSE-парсер + * (см. [readSse]) живёт в общем коде и плагина не требует. + */ +fun defaultAgentikHttpClient(): HttpClient = HttpClient(CIO) { + install(ContentNegotiation) { json(agentikJson) } +} diff --git a/client/src/main/kotlin/pw/binom/agentik/client/ConversationClient.kt b/client/src/main/kotlin/pw/binom/agentik/client/ConversationClient.kt new file mode 100644 index 0000000..931a82e --- /dev/null +++ b/client/src/main/kotlin/pw/binom/agentik/client/ConversationClient.kt @@ -0,0 +1,94 @@ +package pw.binom.agentik.client + +import io.ktor.client.HttpClient +import io.ktor.client.call.body +import io.ktor.client.request.get +import io.ktor.client.request.parameter +import io.ktor.client.request.patch +import io.ktor.client.request.post +import io.ktor.client.request.setBody +import io.ktor.client.statement.bodyAsChannel +import io.ktor.http.ContentType +import io.ktor.http.HttpStatusCode +import io.ktor.http.contentType +import kotlinx.coroutines.flow.Flow +import kotlinx.coroutines.flow.flow +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 kotlin.time.Instant + +/** + * HTTP-реализация [Conversation]. Ходит в `:server`-фасад под + * `/conversations/{id}/...`. + * + * `id` отдаётся синхронно (он в [snapshot], доступном сразу). Остальные + * поля (`title`, `isSupportImageInput`, ...) — тоже из snapshot. [snapshot] + * обновляется после [rename] (сервер возвращает свежий). + * + * **Caveat — `updatedAt`:** сервер бампит `updatedAt` на каждый + * `send`/`rename`, но клиент узнает об этом только при следующем + * `rename` или `getConversation`. Если нужна свежая свежесть после + * `send` — перезапроси через `Agent.getConversation(id)`. + * + * [close] — локальный no-op: сервер держит диалог живым. Удалить — + * через `Agent.deleteConversation(id)`. + */ +internal class ConversationClient( + private val httpClient: HttpClient, + private val baseUrl: String, + private var snapshot: ConversationSnapshot, +) : Conversation { + + override val id: String get() = snapshot.id + override val title: String? get() = snapshot.title + override val isSupportImageInput: Boolean get() = snapshot.isSupportImageInput + override val isSupportImageOutput: Boolean get() = snapshot.isSupportImageOutput + override val isTemporal: Boolean get() = snapshot.isTemporal + override val updatedAt: Instant get() = snapshot.updatedAt + + private val convUrl: String get() = "$baseUrl/conversations/$id" + + override suspend fun rename(title: String) { + val updated = httpClient.patch(convUrl) { + contentType(ContentType.Application.Json) + setBody(RequestRename(title)) + }.body() + snapshot = updated + } + + override suspend fun send(content: List) { + httpClient.post("$convUrl/messages") { + contentType(ContentType.Application.Json) + setBody(content) + } + } + + override suspend fun interrupt() { + httpClient.post("$convUrl/interrupt") + } + + override fun events(after: Instant): Flow = flow { + val response = httpClient.get("$convUrl/events?after=$after") + check(response.status == HttpStatusCode.OK) { + "events: server returned ${response.status}" + } + readSse(response.bodyAsChannel()) + .collect { payload -> + emit(agentikJson.decodeFromString(Event.serializer(), payload)) + } + } + + override suspend fun getMessages(after: Instant, offset: Int, limit: Int): List = + httpClient.get("$convUrl/messages") { + parameter("after", after.toString()) + parameter("offset", offset) + parameter("limit", limit) + }.body() + + override fun close() { + // Локальный no-op: диалог на сервере живёт, пока не вызван + // Agent.deleteConversation(id). См. [Conversation.close] KDoc. + } +} diff --git a/client/src/main/kotlin/pw/binom/agentik/client/Dto.kt b/client/src/main/kotlin/pw/binom/agentik/client/Dto.kt new file mode 100644 index 0000000..db9eabf --- /dev/null +++ b/client/src/main/kotlin/pw/binom/agentik/client/Dto.kt @@ -0,0 +1,26 @@ +package pw.binom.agentik.client + +import kotlinx.serialization.Serializable +import kotlin.time.Instant + +/** + * HTTP-снимок [pw.binom.agentik.proto.Conversation] — те же поля, что у + * интерфейса, но без методов. Зеркалит + * [pw.binom.agentik.server.ConversationSnapshot]. Дубликат сознательно: + * переедем в общий `:wire`, когда появится больше типов. + */ +@Serializable +data class ConversationSnapshot( + val id: String, + val isSupportImageInput: Boolean, + val isSupportImageOutput: Boolean, + val isTemporal: Boolean, + val title: String? = null, + val updatedAt: Instant, +) + +@Serializable +internal data class RequestCreateConversation(val temp: Boolean) + +@Serializable +internal data class RequestRename(val title: String) diff --git a/client/src/main/kotlin/pw/binom/agentik/client/Serialization.kt b/client/src/main/kotlin/pw/binom/agentik/client/Serialization.kt new file mode 100644 index 0000000..dc31722 --- /dev/null +++ b/client/src/main/kotlin/pw/binom/agentik/client/Serialization.kt @@ -0,0 +1,38 @@ +package pw.binom.agentik.client + +import kotlinx.serialization.KSerializer +import kotlinx.serialization.descriptors.PrimitiveKind +import kotlinx.serialization.descriptors.PrimitiveSerialDescriptor +import kotlinx.serialization.descriptors.SerialDescriptor +import kotlinx.serialization.encoding.Decoder +import kotlinx.serialization.encoding.Encoder +import kotlinx.serialization.json.Json +import kotlinx.serialization.modules.SerializersModule +import kotlin.time.Instant + +/** + * Зеркалит [pw.binom.agentik.server.InstantSerializer]. Дублируем сознательно: + * wire-формат компактный, альтернатива — отдельный `:wire`-модуль ради 10 строк. + */ +internal object InstantSerializer : KSerializer { + override val descriptor: SerialDescriptor = + PrimitiveSerialDescriptor("kotlin.time.Instant", PrimitiveKind.STRING) + + override fun serialize(encoder: Encoder, value: Instant) = + encoder.encodeString(value.toString()) + + override fun deserialize(decoder: Decoder): Instant = + Instant.parse(decoder.decodeString()) +} + +/** + * JSON-конфиг клиента. Должен **точно** совпадать с серверным `agentikJson` — + * один и тот же wire-формат с обеих сторон. + */ +internal val agentikJson: Json = Json { + ignoreUnknownKeys = true + explicitNulls = false + serializersModule = SerializersModule { + contextual(Instant::class, InstantSerializer) + } +} diff --git a/client/src/main/kotlin/pw/binom/agentik/client/Sse.kt b/client/src/main/kotlin/pw/binom/agentik/client/Sse.kt new file mode 100644 index 0000000..e787274 --- /dev/null +++ b/client/src/main/kotlin/pw/binom/agentik/client/Sse.kt @@ -0,0 +1,45 @@ +package pw.binom.agentik.client + +import io.ktor.utils.io.ByteReadChannel +import io.ktor.utils.io.readUTF8Line +import kotlinx.coroutines.flow.Flow +import kotlinx.coroutines.flow.flow + +/** + * Минимальный парсер Server-Sent Events, читающий канал до EOF и эмиттящий + * собранный `data:`-пейлоад каждого события. Достаточно для нашего wire-формата: + * сервер шлёт `data: \n\n`, имя события и прочие поля не используются. + * + * Формат (см. WHATWG): + * event: foo — игнор (у нас нет имён событий) + * data: {"k":1} — накапливается, многострочный `data:` склеивается через '\n' + * :comment — игнор + * id:/retry:/ — пустая строка = граница события; всё остальное игнор + * + * Поток закрывается, когда канал доходит до EOF; накопленный `data` (если есть) + * эмитится как финальный ивент. + */ +internal fun readSse(channel: ByteReadChannel): Flow = flow { + val data = StringBuilder() + while (!channel.isClosedForRead) { + val line = channel.readUTF8Line() ?: break + when { + line.isEmpty() -> { + if (data.isNotEmpty()) { + emit(data.toString()) + data.clear() + } + } + line.startsWith("data: ") -> { + if (data.isNotEmpty()) data.append('\n') + data.append(line.removePrefix("data: ")) + } + line.startsWith("data:") -> { + if (data.isNotEmpty()) data.append('\n') + data.append(line.removePrefix("data:")) + } + // event:, id:, retry:, ":" (comment) — игнорируем + } + } + if (data.isNotEmpty()) emit(data.toString()) +} diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md new file mode 100644 index 0000000..c776f9f --- /dev/null +++ b/docs/ARCHITECTURE.md @@ -0,0 +1,113 @@ +# Agentik — архитектура (на основе ограничений AGUI) + +## 1. Контекст + +agentik — runtime агента. Внешне он exposes два фасада: +- **AG-UI** — SSE-поток событий для UI-клиентов (`POST /agui`), транспорт Netty. +- **A2A** — JSON-RPC для agent↔agent (`POST /`, `message/send`), транспорт CIO. + +Ядро (сессии, чат-цикл, тулы, LLM, MCP) **не зависит от транспорта**: фасады лишь сериализуют/десериализуют его. Библиотеки из каталога `caffeine`: `agui` 0.1.0, `a2a` 1.0.0-SNAPSHOT; LLM — `pw.binom.openai` (ktor-impl). + +## 2. Ограничения AGUI (что диктует протокол) + +AGUI — push-протокол. Агент = `Agent { RunAgentInput -> Flow }`. + +Ключевые типы (из `pw.binom.agui.api`): + +| Тип | Смысл | +|---|---| +| `RunAgentInput { threadId, runId, state, messages, tools, context, forwardedProps }` | весь ввод хода. Клиент шлёт всё, сервер по умолчанию stateless. | +| `threadId` | **сессия/конверсация**. | +| `runId` | **один ход** в рамках сессии. | +| `Message { id, role: developer\|system\|assistant\|user\|tool, content, toolCalls?, toolCallId? }` | сообщение. У assistant — `toolCalls`; у tool — `toolCallId`. | +| `Tool { name, description, parameters: JSONSchema }` | тул (JSON Schema в `parameters`). | +| `State = Map` | нестрогий state хода (snapshot/delta). | +| `Context { description, value }` | кусок контекста. | + +События (`BaseEvent`, дискриминатор `type`): +- lifecycle: `RUN_STARTED{threadId,runId}` → … → `RUN_FINISHED` / `RUN_ERROR{message,code}` +- steps: `STEP_STARTED/FINISHED{stepName}` +- text: `TEXT_MESSAGE_START{messageId,role}` → `TEXT_MESSAGE_CONTENT{messageId,delta}` → `TEXT_MESSAGE_END{messageId}` +- tool: `TOOL_CALL_START{toolCallId,toolCallName,parentMessageId?}` → `TOOL_CALL_ARGS{toolCallId,delta}` → `TOOL_CALL_END{toolCallId}` → `TOOL_CALL_RESULT{toolCallId,content}` +- state: `STATE_SNAPSHOT{snapshot}` / `STATE_DELTA{delta}` / `MESSAGES_SNAPSHOT{messages}` +- прочее: `RAW{event}` / `CUSTOM{name,value}` + +Свободные варианты на стороне либы: +- stateless: `Agent { input -> flow }` (использует `input.messages/state/tools`). +- stateful: `AbstractAgent(agentId, threadId)` — держит `state` + `messages`, переопределяешь `runAgent(RunAgentParameters)`. + +Сервер (`aguiAgent(agent, path="/agui")`) не хранит ничего между запросами: декодирует `RunAgentInput`, вызывает `agent.run`, стримит события, оборачивает сбой в `RUN_ERROR`. + +## 3. Как вписываемся (принципы) + +1. **`threadId` = наша сессия.** Держим `SessionStore[threadId]` на сервере даже при stateless-сервере AGUI: на каждом ходе реконсилируем `input.messages` с сохранёнными и дописываем новые. Это даёт «настоящие сессии» с историей. +2. **`runId` = ход.** Один `RunAgentInput` → один цикл ядра → один поток событий. +3. **Ядро говорит на AGUI-событиях** — это и есть наш протокол вывода; фасады не переопределяют его. +4. **LLM — за интерфейсом.** Ядро не знает, какая LLM. Точка входа `LlmClient.complete(...)`. +5. **Тулы: внутренний реестр + MCP.** `ToolRegistry` объединяет «родные» тулы и тулы MCP-серверов (через bridge). LLM вызывает их; ядро исполняет и возвращает `TOOL_CALL_RESULT`. +6. **А2A = request/response.** Переиспользуем то же ядро, но собираем финальный текст и отдаём одним A2A `Message`. + +## 4. Пакеты / модули + +``` +pw.binom.agentik +├─ core +│ ├─ session Session(threadId, messages, state) + SessionStore (in-memory, by threadId) +│ ├─ tool ToolSpec, Tool, ToolRegistry (InMemoryToolRegistry, builtin tools) +│ ├─ llm LlmClient (interface) + LlmMessage/LlmToolCall/LlmResponse + EchoLlmClient (stub) +│ ├─ mcp McpClient (interface: listTools/callTool) + McpTool (bridge) + ToolRegistry.registerMcpClient +│ └─ engine AgentEngine: prompt -> LLM -> tool-loop -> AGUI events (RunStarted..RunFinished) +└─ fronts + ├─ agui AgentikAgent : Agent { input -> engine.run(input) } [вместо EchoAgent] + └─ a2a EngineA2aHandler : A2A handle(msg, contextId) -> Message [собирает финальный текст] +``` +Wiring (в `standalone`): `Main` собирает `AgentEngine(llm, tools+MCP, sessions)`, отдаёт его `AgentikAgent` (AGUI) и `EngineA2aHandler` (A2A). Конфигурация (LLM, MCP-серверы, тулы) — в одном месте (см. §8). + +## 5. Поток данных (AGUI) + +``` +client ──POST /agui {RunAgentInput}──> AgentikAgent + │ engine.run(input) + ▼ + SessionStore.getOrCreate(threadId) + merge input.messages, input.tools + │ + ┌─> llm.complete(system, history, tools, state) + │ └─ toolCalls? -> TOOL_CALL_* + ToolRegistry.call + TOOL_CALL_RESULT -> llm.complete (loop) + └─> text? -> TEXT_MESSAGE_* (последнее) + │ + persist final assistant msg in session + │ + client <──SSE [BaseEvent]──────┘ (RunStarted .. RunFinished / RunError) +``` + +## 6. A2A-фасад (request/response) + +`A2A handle(request: Message, contextId): Message`: +1. `text = request.parts.filterIsInstance().join("")`. +2. `RunAgentInput(threadId = contextId ?: new, runId = randomRunId(), messages = [Message(USER, text)])`. +3. `events = engine.run(input)` (синхронно, `runBlocking { ...toList() }`). +4. `answer = events.filterIsInstance().join { it.delta }`. +5. вернуть `Message(role=AGENT, parts=[TextPart(answer)])`. + +`contextId` ↔ `threadId` — общая сессия. A2A не стримит, поэтому события ядра сворачиваются в финальный ответ. + +## 7. Зоны ответственности + +**Ядро (скелет уже закладываем):** +`Session/SessionStore`, `Tool/ToolRegistry`, `LlmClient`(iface)+stub, `McpClient`(iface)+bridge, `AgentEngine`, `AgentikAgent`, `EngineA2aHandler`, `Main`. + +**Ты (поверх ядра):** +- реализация `LlmClient` (через `pw.binom.openai:ktor-impl`), +- конкретные «родные» тулы, +- транспорт `McpClient` (подключение к MCP-серверам). + +## 8. Открытые решения (нужны решения) + +1. **Персист сессий** — in-memory или на диск/БД (сейчас in-memory). +2. **Конфигурация** — MCP-серверы, LLM endpoint/ключ, системный промпт: env / yaml / файл (сейчас env-заготовки: `AGENTIK_A2A_AGENTS`). +3. **Ограничения** — `maxToolRounds`, таймауты LLM/тулов, лимиты размера. +4. **Аутентификация** — Bearer на обоих фасадах (а2a/agui принимают `token`), пока не включена. +5. **Структура модулей** — держать всё в `standalone` (по пакетам) или вынести `core`/`mcp` в отдельные подмодули (для переиспользования). +6. **Поведение сессии** — «клиент источник правды» (replace) vs «сервер источник правды» (append) при реконсилии `messages`. +``` diff --git a/gradle/libs.versions.toml b/gradle/libs.versions.toml new file mode 100644 index 0000000..e525557 --- /dev/null +++ b/gradle/libs.versions.toml @@ -0,0 +1,40 @@ +[versions] +kotlin = "2.4.20" +kotlinx-serialization = "1.11.0" +kotlinx-coroutines = "1.11.0" +kotlinx-datetime = "0.8.0" +ktor = "3.1.3" +agui = "0.1.0" +a2a = "1.0.0-SNAPSHOT" + +[plugins] +kotlin-multiplatform = { id = "org.jetbrains.kotlin.multiplatform", version.ref = "kotlin" } +kotlin-jvm = { id = "org.jetbrains.kotlin.jvm", version.ref = "kotlin" } +kotlin-serialization = { id = "org.jetbrains.kotlin.plugin.serialization", version.ref = "kotlin" } + +[libraries] +# --- AG-UI (pw.binom.agui) — KMP: jvm + native --- +agui-api = { module = "pw.binom.agui:api", version.ref = "agui" } +agui-client = { module = "pw.binom.agui:client", version.ref = "agui" } +agui-server = { module = "pw.binom.agui:server", version.ref = "agui" } + +# --- A2A (pw.binom.a2a) — shared: KMP (jvm + linuxX64); client/server: JVM-only --- +a2a-shared = { module = "pw.binom.a2a:shared", version.ref = "a2a" } +a2a-client = { module = "pw.binom.a2a:client", version.ref = "a2a" } +a2a-server = { module = "pw.binom.a2a:server", version.ref = "a2a" } + +# --- Ktor (сервер, JVM) --- +ktor-server-core = { module = "io.ktor:ktor-server-core", version.ref = "ktor" } +ktor-server-sse = { module = "io.ktor:ktor-server-sse", version.ref = "ktor" } +ktor-server-netty = { module = "io.ktor:ktor-server-netty", version.ref = "ktor" } +ktor-server-content-negotiation = { module = "io.ktor:ktor-server-content-negotiation", version.ref = "ktor" } +ktor-serialization-kotlinx-json = { module = "io.ktor:ktor-serialization-kotlinx-json", version.ref = "ktor" } +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" } + +# --- commons --- +kotlinx-coroutines-core = { module = "org.jetbrains.kotlinx:kotlinx-coroutines-core", version.ref = "kotlinx-coroutines" } +kotlinx-serialization-core = { module = "org.jetbrains.kotlinx:kotlinx-serialization-core", version.ref = "kotlinx-serialization" } +kotlinx-serialization-json = { module = "org.jetbrains.kotlinx:kotlinx-serialization-json", version.ref = "kotlinx-serialization" } +kotlinx-datetime = { module = "org.jetbrains.kotlinx:kotlinx-datetime", version.ref = "kotlinx-datetime" } diff --git a/gradle/wrapper/gradle-wrapper.jar b/gradle/wrapper/gradle-wrapper.jar new file mode 100644 index 0000000..7f93135 Binary files /dev/null and b/gradle/wrapper/gradle-wrapper.jar differ diff --git a/gradle/wrapper/gradle-wrapper.properties b/gradle/wrapper/gradle-wrapper.properties new file mode 100644 index 0000000..c61a118 --- /dev/null +++ b/gradle/wrapper/gradle-wrapper.properties @@ -0,0 +1,7 @@ +distributionBase=GRADLE_USER_HOME +distributionPath=wrapper/dists +distributionUrl=https\://services.gradle.org/distributions/gradle-9.4.1-bin.zip +networkTimeout=10000 +validateDistributionUrl=true +zipStoreBase=GRADLE_USER_HOME +zipStorePath=wrapper/dists diff --git a/gradlew b/gradlew new file mode 100755 index 0000000..0adc8e1 --- /dev/null +++ b/gradlew @@ -0,0 +1,249 @@ +#!/bin/sh + +# +# Copyright © 2015-2021 the original authors. +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# https://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. +# + +############################################################################## +# +# Gradle start up script for POSIX generated by Gradle. +# +# Important for running: +# +# (1) You need a POSIX-compliant shell to run this script. If your /bin/sh is +# noncompliant, but you have some other compliant shell such as ksh or +# bash, then to run this script, type that shell name before the whole +# command line, like: +# +# ksh Gradle +# +# Busybox and similar reduced shells will NOT work, because this script +# requires all of these POSIX shell features: +# * functions; +# * expansions «$var», «${var}», «${var:-default}», «${var+SET}», +# «${var#prefix}», «${var%suffix}», and «$( cmd )»; +# * compound commands having a testable exit status, especially «case»; +# * various built-in commands including «command», «set», and «ulimit». +# +# Important for patching: +# +# (2) This script targets any POSIX shell, so it avoids extensions provided +# by Bash, Ksh, etc; in particular arrays are avoided. +# +# The "traditional" practice of packing multiple parameters into a +# space-separated string is a well documented source of bugs and security +# problems, so this is (mostly) avoided, by progressively accumulating +# options in "$@", and eventually passing that to Java. +# +# Where the inherited environment variables (DEFAULT_JVM_OPTS, JAVA_OPTS, +# and GRADLE_OPTS) rely on word-splitting, this is performed explicitly; +# see the in-line comments for details. +# +# There are tweaks for specific operating systems such as AIX, CygWin, +# Darwin, MinGW, and NonStop. +# +# (3) This script is generated from the Groovy template +# https://github.com/gradle/gradle/blob/HEAD/subprojects/plugins/src/main/resources/org/gradle/api/internal/plugins/unixStartScript.txt +# within the Gradle project. +# +# You can find Gradle at https://github.com/gradle/gradle/. +# +############################################################################## + +# Attempt to set APP_HOME + +# Resolve links: $0 may be a link +app_path=$0 + +# Need this for daisy-chained symlinks. +while + APP_HOME=${app_path%"${app_path##*/}"} # leaves a trailing /; empty if no leading path + [ -h "$app_path" ] +do + ls=$( ls -ld "$app_path" ) + link=${ls#*' -> '} + case $link in #( + /*) app_path=$link ;; #( + *) app_path=$APP_HOME$link ;; + esac +done + +# This is normally unused +# shellcheck disable=SC2034 +APP_BASE_NAME=${0##*/} +# Discard cd standard output in case $CDPATH is set (https://github.com/gradle/gradle/issues/25036) +APP_HOME=$( cd "${APP_HOME:-./}" > /dev/null && pwd -P ) || exit + +# Use the maximum available, or set MAX_FD != -1 to use that value. +MAX_FD=maximum + +warn () { + echo "$*" +} >&2 + +die () { + echo + echo "$*" + echo + exit 1 +} >&2 + +# OS specific support (must be 'true' or 'false'). +cygwin=false +msys=false +darwin=false +nonstop=false +case "$( uname )" in #( + CYGWIN* ) cygwin=true ;; #( + Darwin* ) darwin=true ;; #( + MSYS* | MINGW* ) msys=true ;; #( + NONSTOP* ) nonstop=true ;; +esac + +CLASSPATH=$APP_HOME/gradle/wrapper/gradle-wrapper.jar + + +# Determine the Java command to use to start the JVM. +if [ -n "$JAVA_HOME" ] ; then + if [ -x "$JAVA_HOME/jre/sh/java" ] ; then + # IBM's JDK on AIX uses strange locations for the executables + JAVACMD=$JAVA_HOME/jre/sh/java + else + JAVACMD=$JAVA_HOME/bin/java + fi + if [ ! -x "$JAVACMD" ] ; then + die "ERROR: JAVA_HOME is set to an invalid directory: $JAVA_HOME + +Please set the JAVA_HOME variable in your environment to match the +location of your Java installation." + fi +else + JAVACMD=java + if ! command -v java >/dev/null 2>&1 + then + die "ERROR: JAVA_HOME is not set and no 'java' command could be found in your PATH. + +Please set the JAVA_HOME variable in your environment to match the +location of your Java installation." + fi +fi + +# Increase the maximum file descriptors if we can. +if ! "$cygwin" && ! "$darwin" && ! "$nonstop" ; then + case $MAX_FD in #( + max*) + # In POSIX sh, ulimit -H is undefined. That's why the result is checked to see if it worked. + # shellcheck disable=SC3045 + MAX_FD=$( ulimit -H -n ) || + warn "Could not query maximum file descriptor limit" + esac + case $MAX_FD in #( + '' | soft) :;; #( + *) + # In POSIX sh, ulimit -n is undefined. That's why the result is checked to see if it worked. + # shellcheck disable=SC3045 + ulimit -n "$MAX_FD" || + warn "Could not set maximum file descriptor limit to $MAX_FD" + esac +fi + +# Collect all arguments for the java command, stacking in reverse order: +# * args from the command line +# * the main class name +# * -classpath +# * -D...appname settings +# * --module-path (only if needed) +# * DEFAULT_JVM_OPTS, JAVA_OPTS, and GRADLE_OPTS environment variables. + +# For Cygwin or MSYS, switch paths to Windows format before running java +if "$cygwin" || "$msys" ; then + APP_HOME=$( cygpath --path --mixed "$APP_HOME" ) + CLASSPATH=$( cygpath --path --mixed "$CLASSPATH" ) + + JAVACMD=$( cygpath --unix "$JAVACMD" ) + + # Now convert the arguments - kludge to limit ourselves to /bin/sh + for arg do + if + case $arg in #( + -*) false ;; # don't mess with options #( + /?*) t=${arg#/} t=/${t%%/*} # looks like a POSIX filepath + [ -e "$t" ] ;; #( + *) false ;; + esac + then + arg=$( cygpath --path --ignore --mixed "$arg" ) + fi + # Roll the args list around exactly as many times as the number of + # args, so each arg winds up back in the position where it started, but + # possibly modified. + # + # NB: a `for` loop captures its iteration list before it begins, so + # changing the positional parameters here affects neither the number of + # iterations, nor the values presented in `arg`. + shift # remove old arg + set -- "$@" "$arg" # push replacement arg + done +fi + + +# Add default JVM options here. You can also use JAVA_OPTS and GRADLE_OPTS to pass JVM options to this script. +DEFAULT_JVM_OPTS='"-Xmx64m" "-Xms64m"' + +# Collect all arguments for the java command; +# * $DEFAULT_JVM_OPTS, $JAVA_OPTS, and $GRADLE_OPTS can contain fragments of +# shell script including quotes and variable substitutions, so put them in +# double quotes to make sure that they get re-expanded; and +# * put everything else in single quotes, so that it's not re-expanded. + +set -- \ + "-Dorg.gradle.appname=$APP_BASE_NAME" \ + -classpath "$CLASSPATH" \ + org.gradle.wrapper.GradleWrapperMain \ + "$@" + +# Stop when "xargs" is not available. +if ! command -v xargs >/dev/null 2>&1 +then + die "xargs is not available" +fi + +# Use "xargs" to parse quoted args. +# +# With -n1 it outputs one arg per line, with the quotes and backslashes removed. +# +# In Bash we could simply go: +# +# readarray ARGS < <( xargs -n1 <<<"$var" ) && +# set -- "${ARGS[@]}" "$@" +# +# but POSIX shell has neither arrays nor command substitution, so instead we +# post-process each arg (as a line of input to sed) to backslash-escape any +# character that might be a shell metacharacter, then use eval to reverse +# that process (while maintaining the separation between arguments), and wrap +# the whole thing up as a single "set" statement. +# +# This will of course break if any of these variables contains a newline or +# an unmatched quote. +# + +eval "set -- $( + printf '%s\n' "$DEFAULT_JVM_OPTS $JAVA_OPTS $GRADLE_OPTS" | + xargs -n1 | + sed ' s~[^-[:alnum:]+,./:=@_]~\\&~g; ' | + tr '\n' ' ' + )" '"$@"' + +exec "$JAVACMD" "$@" diff --git a/gradlew.bat b/gradlew.bat new file mode 100644 index 0000000..93e3f59 --- /dev/null +++ b/gradlew.bat @@ -0,0 +1,92 @@ +@rem +@rem Copyright 2015 the original author or authors. +@rem +@rem Licensed under the Apache License, Version 2.0 (the "License"); +@rem you may not use this file except in compliance with the License. +@rem You may obtain a copy of the License at +@rem +@rem https://www.apache.org/licenses/LICENSE-2.0 +@rem +@rem Unless required by applicable law or agreed to in writing, software +@rem distributed under the License is distributed on an "AS IS" BASIS, +@rem WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +@rem See the License for the specific language governing permissions and +@rem limitations under the License. +@rem + +@if "%DEBUG%"=="" @echo off +@rem ########################################################################## +@rem +@rem Gradle startup script for Windows +@rem +@rem ########################################################################## + +@rem Set local scope for the variables with windows NT shell +if "%OS%"=="Windows_NT" setlocal + +set DIRNAME=%~dp0 +if "%DIRNAME%"=="" set DIRNAME=. +@rem This is normally unused +set APP_BASE_NAME=%~n0 +set APP_HOME=%DIRNAME% + +@rem Resolve any "." and ".." in APP_HOME to make it shorter. +for %%i in ("%APP_HOME%") do set APP_HOME=%%~fi + +@rem Add default JVM options here. You can also use JAVA_OPTS and GRADLE_OPTS to pass JVM options to this script. +set DEFAULT_JVM_OPTS="-Xmx64m" "-Xms64m" + +@rem Find java.exe +if defined JAVA_HOME goto findJavaFromJavaHome + +set JAVA_EXE=java.exe +%JAVA_EXE% -version >NUL 2>&1 +if %ERRORLEVEL% equ 0 goto execute + +echo. +echo ERROR: JAVA_HOME is not set and no 'java' command could be found in your PATH. +echo. +echo Please set the JAVA_HOME variable in your environment to match the +echo location of your Java installation. + +goto fail + +:findJavaFromJavaHome +set JAVA_HOME=%JAVA_HOME:"=% +set JAVA_EXE=%JAVA_HOME%/bin/java.exe + +if exist "%JAVA_EXE%" goto execute + +echo. +echo ERROR: JAVA_HOME is set to an invalid directory: %JAVA_HOME% +echo. +echo Please set the JAVA_HOME variable in your environment to match the +echo location of your Java installation. + +goto fail + +:execute +@rem Setup the command line + +set CLASSPATH=%APP_HOME%\gradle\wrapper\gradle-wrapper.jar + + +@rem Execute Gradle +"%JAVA_EXE%" %DEFAULT_JVM_OPTS% %JAVA_OPTS% %GRADLE_OPTS% "-Dorg.gradle.appname=%APP_BASE_NAME%" -classpath "%CLASSPATH%" org.gradle.wrapper.GradleWrapperMain %* + +:end +@rem End local scope for the variables with windows NT shell +if %ERRORLEVEL% equ 0 goto mainEnd + +:fail +rem Set variable GRADLE_EXIT_CONSOLE if you need the _script_ return code instead of +rem the _cmd.exe /c_ return code! +set EXIT_CODE=%ERRORLEVEL% +if %EXIT_CODE% equ 0 set EXIT_CODE=1 +if not ""=="%GRADLE_EXIT_CONSOLE%" exit %EXIT_CODE% +exit /b %EXIT_CODE% + +:mainEnd +if "%OS%"=="Windows_NT" endlocal + +:omega diff --git a/proto/build.gradle.kts b/proto/build.gradle.kts new file mode 100644 index 0000000..e6c2df9 --- /dev/null +++ b/proto/build.gradle.kts @@ -0,0 +1,30 @@ +plugins { + alias(libs.plugins.kotlin.multiplatform) + alias(libs.plugins.kotlin.serialization) +} + +kotlin { + jvmToolchain(21) + + // "Все возможные цели сборки": jvm + весь натив. Зеркалит набор AG-UI api. + jvm() + macosX64() + macosArm64() + iosX64() + iosArm64() + iosSimulatorArm64() + linuxX64() + linuxArm64() + mingwX64() + + sourceSets { + commonMain.dependencies { + api(libs.kotlinx.coroutines.core) + api(libs.kotlinx.datetime) + api(libs.kotlinx.serialization.core) + } + commonTest.dependencies { + implementation(kotlin("test")) + } + } +} diff --git a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Agent.kt b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Agent.kt new file mode 100644 index 0000000..3112a08 --- /dev/null +++ b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Agent.kt @@ -0,0 +1,56 @@ +package pw.binom.agentik.proto + +import kotlinx.coroutines.flow.Flow +import kotlinx.coroutines.flow.flow +import kotlin.time.Instant + +/** + * Ядро собственного протокола agentik (замена AG-UI). Явно stateful. + * + * Транспортно-агностично. [Agent] — фабрика stateful-диалогов: + * [createConversation] возвращает [Conversation], который сам хранит историю + * и которому отправляют ходы через [Conversation.send]. + */ +public interface Agent { + + /** Идентификатор агента. */ + val id: String + + /** Создаёт новый stateful-диалог с агентом. */ + fun createConversation(temp: Boolean): Conversation + + /** Диалог по идентификатору; `null`, если не найден. */ + suspend fun getConversation(id: String): Conversation? + + /** Удаляет диалог. Возвращает `true`, если диалог существовал и удалён. */ + suspend fun deleteConversation(id: String): Boolean + + /** Страница диалогов: не более [limit] штук, начиная с [offset]-го. */ + suspend fun getConversations(offset: Int, limit: Int): List + + /** Все диалоги, начиная с [offset], как поток: подгружает по [PAGE_SIZE] за раз. */ + fun getConversations(offset: Int = 0): Flow = flow { + var skip = offset + while (true) { + val page = getConversations(skip, PAGE_SIZE) + if (page.isEmpty()) break + page.forEach { emit(it) } + skip += page.size + } + } + + /** + * Live-подписка на изменения в множестве диалогов агента: создание, + * удаление, переименование (см. [AgentEvent]). События внутри конкретного + * диалога приходят через [Conversation.events]. + * + * **Не реплеит** прошлое — для снимка множества используй [getConversations] + * или [getConversation]. + */ + fun events(after: Instant): Flow + + companion object { + + const val PAGE_SIZE: Int = 100 + } +} diff --git a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/AgentEvent.kt b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/AgentEvent.kt new file mode 100644 index 0000000..0f07d0f --- /dev/null +++ b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/AgentEvent.kt @@ -0,0 +1,43 @@ +package pw.binom.agentik.proto + +import kotlinx.serialization.SerialName +import kotlinx.serialization.Serializable +import kotlin.time.Instant + +/** + * Live-события уровня [Agent]: изменения в множестве диалогов + * (создание, удаление, переименование). События, происходящие **внутри** + * конкретного диалога, приходят через [Conversation.events], а не сюда. + * + * Каждое событие несёт [date] — момент эмиссии в UTC. Семантика подписки + * идентична [Conversation.events]: поток **не реплеит** прошлое, для бэкфилла + * используются `getConversations`/`getConversation`. + */ +@Serializable +sealed interface AgentEvent { + /** Момент эмиссии события в UTC. */ + val date: Instant + + /** + * Создан новый диалог. Передаётся его id — handle можно получить через + * [Agent.getConversation]. Подписчик после [Created] может сразу открыть + * live-подписку на этот диалог через [Conversation.events]. + */ + @Serializable + @SerialName("created") + data class Created(override val date: Instant, val conversationId: String) : AgentEvent + + /** + * Диалог удалён. Переданный [Conversation]-handle реализация обязана + * закрыть (`close()`) до эмиссии этого события — после [Deleted] + * пользоваться handle нельзя. + */ + @Serializable + @SerialName("deleted") + data class Deleted(override val date: Instant, val id: String) : AgentEvent + + /** У диалога сменился заголовок. */ + @Serializable + @SerialName("renamed") + data class Renamed(override val date: Instant, val id: String, val title: String?) : AgentEvent +} diff --git a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Content.kt b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Content.kt new file mode 100644 index 0000000..68ac71c --- /dev/null +++ b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Content.kt @@ -0,0 +1,18 @@ +package pw.binom.agentik.proto + +import kotlinx.serialization.SerialName +import kotlinx.serialization.Serializable + +@Serializable +sealed interface Content { + @Serializable + @SerialName("text") + class Text(val body: String) : Content + + /** + * Картинка. [mime] — MIME-тип, например `"image/png"`, `"image/jpeg"`. + */ + @Serializable + @SerialName("image") + class Image(val data: ByteArray, val mime: String) : Content +} diff --git a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Conversation.kt b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Conversation.kt new file mode 100644 index 0000000..1f45030 --- /dev/null +++ b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Conversation.kt @@ -0,0 +1,88 @@ +package pw.binom.agentik.proto + +import kotlinx.coroutines.flow.Flow +import kotlinx.coroutines.flow.flow +import kotlin.time.Instant + +/** + * Stateful-диалог клиента и [Agent]. Хранит собственную историю: на каждый + * [send] агенту не нужно пересылать транскрипт — он уже живёт внутри + * [Conversation]. + */ +interface Conversation : AutoCloseable { + val id: String + + val isSupportImageInput: Boolean + val isSupportImageOutput: Boolean + + /** + * Признак временного диалога: не персистится между перезапусками агента, + * живёт только в памяти текущего процесса. + */ + val isTemporal: Boolean + + val title: String? + + /** + * Момент последнего изменения диалога (любой [send], [rename] и т.п.) в UTC. + * Используется для сортировки списка диалогов по свежести. + */ + val updatedAt: Instant + + suspend fun rename(title: String) + + /** + * Ставит новый user-ход в очередь. Возвращает управление сразу — поток + * событий ответа приходит через [events]. + * + * Если в момент вызова выполняется другой ход, новый встаёт в очередь + * за ним. Чтобы отменить текущий — вызови [interrupt] перед [send]. + */ + suspend fun send(content: List) + + /** + * Прерывает текущий исполняемый ход (best-effort: LLM-stream прибивается, + * in-flight tool может доехать или отвалиться). В [events] эмитится + * [Event.Interrupted], затем может начаться следующий ход из очереди. + * + * Если хода нет — no-op. + */ + suspend fun interrupt() + + /** + * Live-подписка на всё, что происходит в диалоге, начиная с [after]. + * + * **Не реплеит** события, произошедшие до [after] — для бэкфилла + * используй [getMessages]. Если [after] — момент последнего виденного + * клиентом события, поток продолжается «с того места». + * + * Подписки независимы: каждый вызов возвращает свой [Flow], отмена одного + * не влияет на других подписчиков и на сам диалог. + */ + fun events(after: Instant): Flow + + /** Страница истории: не более [limit] сообщений после [after], начиная с [offset]-го. */ + suspend fun getMessages(after: Instant, offset: Int, limit: Int): List + + /** Все сообщения после [after] начиная с [offset], как поток: подгружает по [PAGE_SIZE] за раз. */ + fun getMessages(after: Instant, offset: Int = 0): Flow = flow { + var skip = offset + while (true) { + val page = getMessages(after, skip, PAGE_SIZE) + if (page.isEmpty()) break + page.forEach { emit(it) } + skip += page.size + } + } + + /** + * Освобождает ресурсы диалога (подписки, сетевые хэндлы). Идемпотентно. + * После [close] дальнейшие вызовы [send]/[interrupt]/[events]/[getMessages]/[rename] не определены. + */ + override fun close() + + companion object { + + const val PAGE_SIZE: Int = 100 + } +} diff --git a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Event.kt b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Event.kt new file mode 100644 index 0000000..af3d131 --- /dev/null +++ b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Event.kt @@ -0,0 +1,88 @@ +package pw.binom.agentik.proto + +import kotlinx.serialization.SerialName +import kotlinx.serialization.Serializable +import kotlin.time.Instant + +/** + * Элемент live-потока [Conversation.events]. + * + * Каждое событие несёт [date] — момент эмиссии в UTC. Используется клиентом + * для трекинга «где остановился» при обрыве/переподключении и для разрешения + * порядка при равных timestamps. + * + * Базовая структура хода: + * `StartReasoning?` → `StartResponse(TEXT|IMAGE)` → ...контент... → `End` | `Interrupted` | `Error`. + * `StartReasoning` может отсутствовать, если агент не показывал рассуждения. + */ +@Serializable +sealed interface Event { + /** Момент эмиссии события в UTC. */ + val date: Instant + + @Serializable + enum class ResponseType { + @SerialName("text") TEXT, + @SerialName("image") IMAGE + } + + /** Ассистент начал рассуждение (опциональный маркер; контент рассуждения приходит через [AppendText]). */ + @Serializable + @SerialName("start_reasoning") + data class StartReasoning(override val date: Instant) : Event + + /** Начало ответа ассистента заданного типа. После него идут соответствующие `Append*`/`Tool*`-события, потом [End]/[Interrupted]/[Error]. */ + @Serializable + @SerialName("start_response") + data class StartResponse(override val date: Instant, val responseType: ResponseType) : Event + + /** Ход завершён нормально. Соответствующий [Message.AssistantMessage] появится в `getMessages`. */ + @Serializable + @SerialName("end") + data class End(override val date: Instant) : Event + + /** Ход прерван через [Conversation.interrupt]. Частичный ответ НЕ сохраняется в истории. */ + @Serializable + @SerialName("interrupted") + data class Interrupted(override val date: Instant) : Event + + @Serializable + @SerialName("append_text") + data class AppendText(override val date: Instant, val body: String) : Event + + @Serializable + @SerialName("append_image") + data class AppendImage(override val date: Instant, val body: ByteArray, val mime: String) : Event + + /** + * Агент начал вызов тула. Аргументы приходят целиком — стриминга нет. + * [id] совпадает с id соответствующего [Message.ToolCall] в истории + * после завершения хода. + */ + @Serializable + @SerialName("tool_call") + data class ToolCall( + override val date: Instant, + val id: String, + val title: String?, + val toolName: String, + val toolArgs: String, + ) : Event + + /** + * Результат вызова тула. Приходит целиком после завершения исполнения. + * [id] совпадает с [ToolCall.id], к которому относится результат, и + * с id [Message.ToolResult] в истории. + */ + @Serializable + @SerialName("tool_result") + data class ToolResult(override val date: Instant, val id: String, val result: String?) : Event + + /** + * Ошибка хода. После неё поток завершается; дальнейшие события могут + * прийти, но ход считается проваленным. + */ + @Serializable + @SerialName("error") + data class Error(override val date: Instant, val message: String, val code: String? = null) : Event +} diff --git a/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Message.kt b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Message.kt new file mode 100644 index 0000000..4a908c8 --- /dev/null +++ b/proto/src/commonMain/kotlin/pw/binom/agentik/proto/Message.kt @@ -0,0 +1,42 @@ +package pw.binom.agentik.proto + +import kotlinx.serialization.SerialName +import kotlinx.serialization.Serializable +import kotlin.time.Instant + +@Serializable +sealed interface Message { + /** + * Уникальный идентификатор сообщения в рамках диалога. Стабилен между + * стримом [Event] и историей: id, пришедший в [Event.ToolCall], равен + * id соответствующего [ToolCall] в истории после завершения хода. + */ + val id: String + + /** + * Дата сообщения в UTC + */ + val date: Instant + + @Serializable + @SerialName("user_message") + class UserMessage(override val id: String, val content: List, override val date: Instant) : Message + + @Serializable + @SerialName("assistant_message") + class AssistantMessage(override val id: String, val content: List, override val date: Instant) : Message + + @Serializable + @SerialName("tool_call") + class ToolCall( + override val id: String, + val title: String?, + val toolName: String, + val toolArgs: String, + override val date: Instant + ) : Message + + @Serializable + @SerialName("tool_result") + class ToolResult(override val id: String, val result: String?, override val date: Instant) : Message +} diff --git a/server/build.gradle.kts b/server/build.gradle.kts new file mode 100644 index 0000000..e99a707 --- /dev/null +++ b/server/build.gradle.kts @@ -0,0 +1,23 @@ +import org.jetbrains.kotlin.gradle.dsl.JvmTarget + +plugins { + alias(libs.plugins.kotlin.jvm) + alias(libs.plugins.kotlin.serialization) +} + +kotlin { + compilerOptions { + jvmTarget.set(JvmTarget.JVM_21) + } +} + +dependencies { + implementation(project(":proto")) + + implementation(libs.ktor.server.core) + implementation(libs.ktor.server.content.negotiation) + implementation(libs.ktor.serialization.kotlinx.json) + + implementation(libs.kotlinx.coroutines.core) + implementation(libs.kotlinx.serialization.json) +} diff --git a/server/src/main/kotlin/pw/binom/agentik/server/Dto.kt b/server/src/main/kotlin/pw/binom/agentik/server/Dto.kt new file mode 100644 index 0000000..934b9c6 --- /dev/null +++ b/server/src/main/kotlin/pw/binom/agentik/server/Dto.kt @@ -0,0 +1,34 @@ +package pw.binom.agentik.server + +import kotlinx.serialization.Serializable +import pw.binom.agentik.proto.Conversation +import kotlin.time.Instant + +/** + * HTTP-снимок [Conversation] — те же поля, что у интерфейса, но без методов. + * Сериализуется в JSON и обратно. + */ +@Serializable +data class ConversationSnapshot( + val id: String, + val isSupportImageInput: Boolean, + val isSupportImageOutput: Boolean, + val isTemporal: Boolean, + val title: String? = null, + val updatedAt: Instant, +) + +internal fun Conversation.snapshot(): ConversationSnapshot = ConversationSnapshot( + id = id, + isSupportImageInput = isSupportImageInput, + isSupportImageOutput = isSupportImageOutput, + isTemporal = isTemporal, + title = title, + updatedAt = updatedAt, +) + +@Serializable +internal data class RequestCreateConversation(val temp: Boolean) + +@Serializable +internal data class RequestRename(val title: String) diff --git a/server/src/main/kotlin/pw/binom/agentik/server/Module.kt b/server/src/main/kotlin/pw/binom/agentik/server/Module.kt new file mode 100644 index 0000000..e00d264 --- /dev/null +++ b/server/src/main/kotlin/pw/binom/agentik/server/Module.kt @@ -0,0 +1,43 @@ +package pw.binom.agentik.server + +import io.ktor.serialization.kotlinx.json.json +import io.ktor.server.application.install +import io.ktor.server.plugins.contentnegotiation.ContentNegotiation +import io.ktor.server.routing.Route +import io.ktor.server.routing.route +import pw.binom.agentik.proto.Agent + +/** + * Встраивает HTTP/SSE-фасад протокола agentik в твой Ktor-роутинг. + * + * Использование: + * ``` + * embeddedServer(Netty, port = 8080) { + * routing { + * agentikAgent(MyAgent()) // все роуты под /agentik + * agentikAgent(MyAgent(), "/api/chat") // или под произвольным префиксом + * } + * }.start(wait = true) + * ``` + * + * Под префиксом [path] монтируются: + * - `POST /conversations` — создать диалог + * - `GET /conversations` — список + * - `GET /conversations/{id}` — один диалог + * - `PATCH /conversations/{id}` — переименовать + * - `DELETE /conversations/{id}` — удалить + * - `POST /conversations/{id}/messages` — `send` (202 Accepted) + * - `POST /conversations/{id}/interrupt` — `interrupt` + * - `GET /conversations/{id}/messages` — история + * - `GET /conversations/{id}/events` — SSE: события хода + * - `GET /events` — SSE: события агента + * - `GET /health` — `"ok"` + */ +fun Route.agentikAgent(agent: Agent, path: String = "/agentik") { + route(path) { + install(ContentNegotiation) { + json(agentikJson) + } + agentikRoutes(agent) + } +} diff --git a/server/src/main/kotlin/pw/binom/agentik/server/Routes.kt b/server/src/main/kotlin/pw/binom/agentik/server/Routes.kt new file mode 100644 index 0000000..df3696d --- /dev/null +++ b/server/src/main/kotlin/pw/binom/agentik/server/Routes.kt @@ -0,0 +1,160 @@ +package pw.binom.agentik.server + +import io.ktor.http.ContentType +import io.ktor.http.HttpStatusCode +import io.ktor.server.application.ApplicationCall +import io.ktor.server.application.call +import io.ktor.server.request.receive +import io.ktor.server.response.respond +import io.ktor.server.response.respondText +import io.ktor.server.response.respondTextWriter +import io.ktor.server.routing.Route +import io.ktor.server.routing.delete +import io.ktor.server.routing.get +import io.ktor.server.routing.patch +import io.ktor.server.routing.post +import kotlinx.coroutines.flow.Flow +import kotlinx.coroutines.flow.catch +import kotlinx.serialization.KSerializer +import kotlinx.serialization.json.Json +import pw.binom.agentik.proto.Agent +import pw.binom.agentik.proto.AgentEvent +import pw.binom.agentik.proto.Conversation +import pw.binom.agentik.proto.Content +import pw.binom.agentik.proto.Event +import kotlin.time.Instant + +internal fun Route.agentikRoutes(agent: Agent) { + + get("/health") { + call.respondText("ok") + } + + // ---- Agent: множество диалогов ---- + + get("/conversations") { + val offset = call.request.queryParameters["offset"]?.toIntOrNull() ?: 0 + val limit = call.request.queryParameters["limit"]?.toIntOrNull() ?: Agent.PAGE_SIZE + call.respond(agent.getConversations(offset, limit).map { it.snapshot() }) + } + + post("/conversations") { + val req = call.receive() + val c = agent.createConversation(req.temp) + call.respond(HttpStatusCode.Created, c.snapshot()) + } + + get("/conversations/{id}") { + val id = call.parameters["id"]!! + val c = agent.getConversation(id) + if (c == null) call.respond(HttpStatusCode.NotFound) + else call.respond(c.snapshot()) + } + + delete("/conversations/{id}") { + val id = call.parameters["id"]!! + call.respond(if (agent.deleteConversation(id)) HttpStatusCode.NoContent else HttpStatusCode.NotFound) + } + + // ---- Conversation ---- + + patch("/conversations/{id}") { + val id = call.parameters["id"]!! + val c = agent.getConversation(id) + if (c == null) { + call.respond(HttpStatusCode.NotFound) + return@patch + } + val req = call.receive() + c.rename(req.title) + call.respond(c.snapshot()) + } + + post("/conversations/{id}/messages") { + val id = call.parameters["id"]!! + val c = agent.getConversation(id) + if (c == null) { + call.respond(HttpStatusCode.NotFound) + return@post + } + val content = call.receive>() + c.send(content) + call.respond(HttpStatusCode.Accepted) + } + + post("/conversations/{id}/interrupt") { + val id = call.parameters["id"]!! + val c = agent.getConversation(id) + if (c == null) { + call.respond(HttpStatusCode.NotFound) + return@post + } + c.interrupt() + call.respond(HttpStatusCode.Accepted) + } + + get("/conversations/{id}/messages") { + val id = call.parameters["id"]!! + val c = agent.getConversation(id) + if (c == null) { + call.respond(HttpStatusCode.NotFound) + return@get + } + val after = call.parseAfter() ?: return@get + val offset = call.request.queryParameters["offset"]?.toIntOrNull() ?: 0 + val limit = call.request.queryParameters["limit"]?.toIntOrNull() ?: Conversation.PAGE_SIZE + call.respond(c.getMessages(after, offset, limit)) + } + + // ---- SSE ---- + + get("/conversations/{id}/events") { + val id = call.parameters["id"]!! + val c = agent.getConversation(id) + if (c == null) { + call.respond(HttpStatusCode.NotFound) + return@get + } + val after = call.parseAfter() ?: return@get + call.streamJsonSse(c.events(after), Event.serializer()) + } + + get("/events") { + val after = call.parseAfter() ?: return@get + call.streamJsonSse(agent.events(after), AgentEvent.serializer()) + } +} + +// ---------- helpers ---------- + +/** + * Парсит query-параметр `after` как ISO-8601 [Instant]. Отсутствие = [Instant.DISTANT_PAST]. + * При невалидном значении отвечает 400 и возвращает `null`. + */ +private suspend fun ApplicationCall.parseAfter(): Instant? { + val raw = request.queryParameters["after"] + if (raw == null) return Instant.DISTANT_PAST + return try { + Instant.parse(raw) + } catch (_: IllegalArgumentException) { + respond(HttpStatusCode.BadRequest, "Invalid 'after' (expected ISO-8601): $raw") + null + } +} + +private suspend fun ApplicationCall.streamJsonSse( + flow: Flow, + serializer: KSerializer, + json: Json = agentikJson, +) { + respondTextWriter(contentType = ContentType.Text.EventStream) { + flow.catch { /* клиент отвалился — глушим */ } + .collect { value -> + val s = json.encodeToString(serializer, value) + write("data: ") + write(s) + write("\n\n") + flush() + } + } +} diff --git a/server/src/main/kotlin/pw/binom/agentik/server/Serialization.kt b/server/src/main/kotlin/pw/binom/agentik/server/Serialization.kt new file mode 100644 index 0000000..d05fc40 --- /dev/null +++ b/server/src/main/kotlin/pw/binom/agentik/server/Serialization.kt @@ -0,0 +1,30 @@ +package pw.binom.agentik.server + +import kotlinx.serialization.KSerializer +import kotlinx.serialization.descriptors.PrimitiveKind +import kotlinx.serialization.descriptors.PrimitiveSerialDescriptor +import kotlinx.serialization.descriptors.SerialDescriptor +import kotlinx.serialization.encoding.Decoder +import kotlinx.serialization.encoding.Encoder +import kotlinx.serialization.json.Json +import kotlinx.serialization.modules.SerializersModule +import kotlin.time.Instant + +internal object InstantSerializer : KSerializer { + override val descriptor: SerialDescriptor = + PrimitiveSerialDescriptor("pw.binom.agentik.Instant", PrimitiveKind.STRING) + + override fun serialize(encoder: Encoder, value: Instant) = + encoder.encodeString(value.toString()) + + override fun deserialize(decoder: Decoder): Instant = + Instant.parse(decoder.decodeString()) +} + +internal val agentikJson: Json = Json { + ignoreUnknownKeys = true + explicitNulls = false + serializersModule = SerializersModule { + contextual(Instant::class, InstantSerializer) + } +} diff --git a/settings.gradle.kts b/settings.gradle.kts new file mode 100644 index 0000000..19a26af --- /dev/null +++ b/settings.gradle.kts @@ -0,0 +1,32 @@ +val caffeineRepo = providers.gradleProperty("caffeineRepo").getOrElse("http://192.168.76.117/repository/caffeine") + +pluginManagement { + repositories { + gradlePluginPortal() + mavenCentral() + google() + } +} + +dependencyResolutionManagement { + repositories { + mavenCentral() + google() + // Home Nexus, репо "caffeine": pw.binom.* (AG-UI, A2A, ...) + maven { + name = "caffeine" + url = uri(caffeineRepo) + setAllowInsecureProtocol(true) + } + } +} + +rootProject.name = "agentik" + +include(":standalone") +// Собственный протокол agentik (замена AG-UI). Пока в нём пилим, потом вынесем. +include(":proto") +// Ktor-сервер, экспонирующий Agent по HTTP (SSE + JSON). +include(":server") +// Ktor-клиент, превращающий HTTP-фасад в `Agent`/`Conversation`. +include(":client") diff --git a/standalone/build.gradle.kts b/standalone/build.gradle.kts new file mode 100644 index 0000000..72c6c5e --- /dev/null +++ b/standalone/build.gradle.kts @@ -0,0 +1,47 @@ +import org.jetbrains.kotlin.gradle.ExperimentalKotlinGradlePluginApi + +plugins { + alias(libs.plugins.kotlin.multiplatform) +} + +kotlin { + jvmToolchain(21) + + jvm { + @OptIn(ExperimentalKotlinGradlePluginApi::class) + binaries { + executable { + mainClass.set("pw.binom.agentik.standalone.MainKt") + } + } + } + + sourceSets { + commonMain.dependencies { + // AG-UI: протокол (события, RunAgentInput, Agent) — KMP, jvm + native + implementation(libs.agui.api) + // proto: наш in-house протокол (KMP) + implementation(project(":proto")) + } + jvmMain.dependencies { + // AG-UI: Ktor-хелперы маршрута + движок Netty (JVM) + implementation(libs.agui.server) + implementation(libs.ktor.server.core) + implementation(libs.ktor.server.sse) + implementation(libs.ktor.server.netty) + implementation(libs.kotlinx.coroutines.core) + implementation(libs.kotlinx.serialization.json) + // A2A (pw.binom.a2a): shared — протокол, server + client — JVM + implementation(libs.a2a.shared) + implementation(libs.a2a.server) + implementation(libs.a2a.client) + implementation(libs.ktor.client.core) + implementation(libs.ktor.client.cio) + // server: Ktor-фасад нашего :proto (JVM) + implementation(project(":server")) + } + commonTest.dependencies { + implementation(kotlin("test")) + } + } +} \ No newline at end of file diff --git a/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/PlaceholderTest.kt b/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/PlaceholderTest.kt new file mode 100644 index 0000000..4bf886b --- /dev/null +++ b/standalone/src/commonTest/kotlin/pw/binom/agentik/standalone/PlaceholderTest.kt @@ -0,0 +1,11 @@ +package pw.binom.agentik.standalone + +import kotlin.test.Test +import kotlin.test.assertTrue + +class PlaceholderTest { + @Test + fun placeholder() { + assertTrue(true) + } +} diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/A2aOutbound.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/A2aOutbound.kt new file mode 100644 index 0000000..d5b6f65 --- /dev/null +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/A2aOutbound.kt @@ -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() + + fun remoteAgents(): List = + 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() + } +} diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/EchoA2aHandler.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/EchoA2aHandler.kt new file mode 100644 index 0000000..2766f3f --- /dev/null +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/EchoA2aHandler.kt @@ -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().joinToString("") { it.text } + return Message( + role = Role.AGENT, + parts = listOf(TextPart("echo: $text")), + contextId = contextId, + ) + } +} diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/EchoAgent.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/EchoAgent.kt new file mode 100644 index 0000000..1c059a7 --- /dev/null +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/EchoAgent.kt @@ -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 = 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)) + } +} diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/EchoProtoAgent.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/EchoProtoAgent.kt new file mode 100644 index 0000000..71a6306 --- /dev/null +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/EchoProtoAgent.kt @@ -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() + private val agentEvents = MutableSharedFlow(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 = + conversations.values.sortedByDescending { it.updatedAt }.drop(offset).take(limit) + + override fun events(after: Instant): Flow = 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() + private val events = MutableSharedFlow(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) { + 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 = flow { + events.asSharedFlow().collect { event -> + if (event.date > after) emit(event) + } + } + + override suspend fun getMessages(after: Instant, offset: Int, limit: Int): List = + mutex.withLock { messages.asSequence().filter { it.date > after }.drop(offset).take(limit).toList() } + + override fun close() { + activeJob?.cancel() + scope.coroutineContext[Job]?.cancel() + } +} diff --git a/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/Main.kt b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/Main.kt new file mode 100644 index 0000000..08db0fc --- /dev/null +++ b/standalone/src/jvmMain/kotlin/pw/binom/agentik/standalone/Main.kt @@ -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> = arrayOf( + Conversation::class.java, + Message::class.java, +)