Compare commits
2 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| c140d0b758 | |||
| c486c7f9ab |
@@ -41,6 +41,7 @@ kotlin {
|
|||||||
implementation(libs.kotlinx.cli)
|
implementation(libs.kotlinx.cli)
|
||||||
|
|
||||||
implementation(libs.kotlinx.coroutines.core)
|
implementation(libs.kotlinx.coroutines.core)
|
||||||
|
implementation(libs.ktor.client.cio)
|
||||||
}
|
}
|
||||||
// :agentik-cli — commonMain-only (нет jvmMain/nativeMain разделения):
|
// :agentik-cli — commonMain-only (нет jvmMain/nativeMain разделения):
|
||||||
// весь код, включая platformEnv, лежит в commonMain.
|
// весь код, включая platformEnv, лежит в commonMain.
|
||||||
|
|||||||
@@ -0,0 +1,23 @@
|
|||||||
|
package pw.binom.agentik.cli
|
||||||
|
|
||||||
|
import io.ktor.client.HttpClient
|
||||||
|
import io.ktor.client.engine.cio.CIO
|
||||||
|
import pw.binom.agentik.client.agentikHttpClient
|
||||||
|
|
||||||
|
/**
|
||||||
|
* HTTP-клиент CLI: движок CIO + конфигурация agentik.
|
||||||
|
*
|
||||||
|
* Движок живёт здесь, а не в `:client`: библиотека не выбирает транспорт за
|
||||||
|
* потребителя. Таргеты `:agentik-cli` (jvm + linuxX64/macosX64/macosArm64/mingwX64)
|
||||||
|
* покрываются CIO.
|
||||||
|
*
|
||||||
|
* `requestTimeout = 0` — отключение встроенного request-таймаута CIO;
|
||||||
|
* defense-in-depth против обрыва долгих SSE-idle (основная защита —
|
||||||
|
* `noSseReadTimeout` в `:client`).
|
||||||
|
*
|
||||||
|
* [token] = `null` — авторизация выключена.
|
||||||
|
*/
|
||||||
|
internal fun defaultCliHttpClient(token: String? = null): HttpClient =
|
||||||
|
agentikHttpClient(engineFactory = CIO, token = token) {
|
||||||
|
engine { requestTimeout = 0 }
|
||||||
|
}
|
||||||
+2
-1
@@ -2,13 +2,14 @@ package pw.binom.agentik.cli.commands
|
|||||||
|
|
||||||
import kotlinx.cli.ArgType
|
import kotlinx.cli.ArgType
|
||||||
import pw.binom.agentik.cli.AgentikSubcommand
|
import pw.binom.agentik.cli.AgentikSubcommand
|
||||||
|
import pw.binom.agentik.cli.defaultCliHttpClient
|
||||||
import pw.binom.agentik.client.AgentikAgent
|
import pw.binom.agentik.client.AgentikAgent
|
||||||
|
|
||||||
class ConvDeleteSubcommand : ConvSubcommand("delete", "Удалить диалог") {
|
class ConvDeleteSubcommand : ConvSubcommand("delete", "Удалить диалог") {
|
||||||
val id by argument(ArgType.String, description = "ID диалога")
|
val id by argument(ArgType.String, description = "ID диалога")
|
||||||
|
|
||||||
override fun execute() = kotlinx.coroutines.runBlocking {
|
override fun execute() = kotlinx.coroutines.runBlocking {
|
||||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl)
|
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
|
||||||
val ok = agent.deleteConversation(id)
|
val ok = agent.deleteConversation(id)
|
||||||
if (ok) println("deleted: $id") else println("conversation not found: $id")
|
if (ok) println("deleted: $id") else println("conversation not found: $id")
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -3,6 +3,7 @@ package pw.binom.agentik.cli.commands
|
|||||||
import kotlinx.cli.ArgType
|
import kotlinx.cli.ArgType
|
||||||
import kotlinx.cli.default
|
import kotlinx.cli.default
|
||||||
import pw.binom.agentik.cli.AgentikSubcommand
|
import pw.binom.agentik.cli.AgentikSubcommand
|
||||||
|
import pw.binom.agentik.cli.defaultCliHttpClient
|
||||||
import pw.binom.agentik.client.AgentikAgent
|
import pw.binom.agentik.client.AgentikAgent
|
||||||
import pw.binom.agentik.proto.Agent
|
import pw.binom.agentik.proto.Agent
|
||||||
|
|
||||||
@@ -10,7 +11,7 @@ class ConvLsSubcommand : ConvSubcommand("ls", "Список диалогов а
|
|||||||
val limit by option(ArgType.Int, fullName = "limit", description = "Максимум диалогов").default(Agent.PAGE_SIZE)
|
val limit by option(ArgType.Int, fullName = "limit", description = "Максимум диалогов").default(Agent.PAGE_SIZE)
|
||||||
|
|
||||||
override fun execute() = kotlinx.coroutines.runBlocking {
|
override fun execute() = kotlinx.coroutines.runBlocking {
|
||||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl)
|
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
|
||||||
val convs = agent.getConversations(offset = 0, limit = limit.coerceAtMost(Agent.PAGE_SIZE))
|
val convs = agent.getConversations(offset = 0, limit = limit.coerceAtMost(Agent.PAGE_SIZE))
|
||||||
if (convs.isEmpty()) {
|
if (convs.isEmpty()) {
|
||||||
println("(no conversations)")
|
println("(no conversations)")
|
||||||
|
|||||||
+2
-1
@@ -3,13 +3,14 @@ package pw.binom.agentik.cli.commands
|
|||||||
import kotlinx.cli.ArgType
|
import kotlinx.cli.ArgType
|
||||||
import kotlinx.cli.default
|
import kotlinx.cli.default
|
||||||
import pw.binom.agentik.cli.AgentikSubcommand
|
import pw.binom.agentik.cli.AgentikSubcommand
|
||||||
|
import pw.binom.agentik.cli.defaultCliHttpClient
|
||||||
import pw.binom.agentik.client.AgentikAgent
|
import pw.binom.agentik.client.AgentikAgent
|
||||||
|
|
||||||
class ConvNewSubcommand : ConvSubcommand("new", "Создать диалог; печатает id") {
|
class ConvNewSubcommand : ConvSubcommand("new", "Создать диалог; печатает id") {
|
||||||
val temp by option(ArgType.Boolean, fullName = "temp", description = "Временный диалог").default(false)
|
val temp by option(ArgType.Boolean, fullName = "temp", description = "Временный диалог").default(false)
|
||||||
|
|
||||||
override fun execute() = kotlinx.coroutines.runBlocking {
|
override fun execute() = kotlinx.coroutines.runBlocking {
|
||||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl)
|
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
|
||||||
val conv = agent.createConversation(temp = temp)
|
val conv = agent.createConversation(temp = temp)
|
||||||
println(conv.id)
|
println(conv.id)
|
||||||
}
|
}
|
||||||
|
|||||||
+2
-1
@@ -2,6 +2,7 @@ package pw.binom.agentik.cli.commands
|
|||||||
|
|
||||||
import kotlinx.cli.ArgType
|
import kotlinx.cli.ArgType
|
||||||
import pw.binom.agentik.cli.AgentikSubcommand
|
import pw.binom.agentik.cli.AgentikSubcommand
|
||||||
|
import pw.binom.agentik.cli.defaultCliHttpClient
|
||||||
import pw.binom.agentik.client.AgentikAgent
|
import pw.binom.agentik.client.AgentikAgent
|
||||||
|
|
||||||
class ConvRenameSubcommand : ConvSubcommand("rename", "Переименовать диалог") {
|
class ConvRenameSubcommand : ConvSubcommand("rename", "Переименовать диалог") {
|
||||||
@@ -9,7 +10,7 @@ class ConvRenameSubcommand : ConvSubcommand("rename", "Переименоват
|
|||||||
val title by argument(ArgType.String, description = "Новое название")
|
val title by argument(ArgType.String, description = "Новое название")
|
||||||
|
|
||||||
override fun execute() = kotlinx.coroutines.runBlocking {
|
override fun execute() = kotlinx.coroutines.runBlocking {
|
||||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl)
|
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
|
||||||
val conv = agent.getConversation(id) ?: run {
|
val conv = agent.getConversation(id) ?: run {
|
||||||
println("conversation not found: $id")
|
println("conversation not found: $id")
|
||||||
return@runBlocking
|
return@runBlocking
|
||||||
|
|||||||
+2
-1
@@ -2,13 +2,14 @@ package pw.binom.agentik.cli.commands
|
|||||||
|
|
||||||
import kotlinx.cli.ArgType
|
import kotlinx.cli.ArgType
|
||||||
import pw.binom.agentik.cli.AgentikSubcommand
|
import pw.binom.agentik.cli.AgentikSubcommand
|
||||||
|
import pw.binom.agentik.cli.defaultCliHttpClient
|
||||||
import pw.binom.agentik.client.AgentikAgent
|
import pw.binom.agentik.client.AgentikAgent
|
||||||
|
|
||||||
class ConvShowSubcommand : ConvSubcommand("show", "Метаданные диалога") {
|
class ConvShowSubcommand : ConvSubcommand("show", "Метаданные диалога") {
|
||||||
val id by argument(ArgType.String, description = "ID диалога")
|
val id by argument(ArgType.String, description = "ID диалога")
|
||||||
|
|
||||||
override fun execute() = kotlinx.coroutines.runBlocking {
|
override fun execute() = kotlinx.coroutines.runBlocking {
|
||||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl)
|
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
|
||||||
val conv = agent.getConversation(id) ?: run {
|
val conv = agent.getConversation(id) ?: run {
|
||||||
println("conversation not found: $id")
|
println("conversation not found: $id")
|
||||||
return@runBlocking
|
return@runBlocking
|
||||||
|
|||||||
+2
-1
@@ -2,13 +2,14 @@ package pw.binom.agentik.cli.commands
|
|||||||
|
|
||||||
import kotlinx.cli.ArgType
|
import kotlinx.cli.ArgType
|
||||||
import pw.binom.agentik.cli.AgentikSubcommand
|
import pw.binom.agentik.cli.AgentikSubcommand
|
||||||
|
import pw.binom.agentik.cli.defaultCliHttpClient
|
||||||
import pw.binom.agentik.client.AgentikAgent
|
import pw.binom.agentik.client.AgentikAgent
|
||||||
|
|
||||||
class InterruptSubcommand : AgentikSubcommand("interrupt", "Прервать текущий ход диалога") {
|
class InterruptSubcommand : AgentikSubcommand("interrupt", "Прервать текущий ход диалога") {
|
||||||
val id by argument(ArgType.String, description = "ID диалога")
|
val id by argument(ArgType.String, description = "ID диалога")
|
||||||
|
|
||||||
override fun execute() = kotlinx.coroutines.runBlocking {
|
override fun execute() = kotlinx.coroutines.runBlocking {
|
||||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl)
|
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
|
||||||
val conv = agent.getConversation(id) ?: run {
|
val conv = agent.getConversation(id) ?: run {
|
||||||
println("conversation not found: $id")
|
println("conversation not found: $id")
|
||||||
return@runBlocking
|
return@runBlocking
|
||||||
|
|||||||
@@ -3,6 +3,7 @@ package pw.binom.agentik.cli.commands
|
|||||||
import kotlinx.cli.ArgType
|
import kotlinx.cli.ArgType
|
||||||
import kotlinx.cli.default
|
import kotlinx.cli.default
|
||||||
import pw.binom.agentik.cli.AgentikSubcommand
|
import pw.binom.agentik.cli.AgentikSubcommand
|
||||||
|
import pw.binom.agentik.cli.defaultCliHttpClient
|
||||||
import pw.binom.agentik.client.AgentikAgent
|
import pw.binom.agentik.client.AgentikAgent
|
||||||
import pw.binom.agentik.proto.Content
|
import pw.binom.agentik.proto.Content
|
||||||
import pw.binom.agentik.proto.Message
|
import pw.binom.agentik.proto.Message
|
||||||
@@ -13,7 +14,7 @@ class MsgsSubcommand : AgentikSubcommand("msgs", "Показать сообще
|
|||||||
val limit by option(ArgType.Int, fullName = "limit", description = "Максимум сообщений").default(100)
|
val limit by option(ArgType.Int, fullName = "limit", description = "Максимум сообщений").default(100)
|
||||||
|
|
||||||
override fun execute() = kotlinx.coroutines.runBlocking {
|
override fun execute() = kotlinx.coroutines.runBlocking {
|
||||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl)
|
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
|
||||||
val conv = agent.getConversation(id) ?: run {
|
val conv = agent.getConversation(id) ?: run {
|
||||||
println("conversation not found: $id")
|
println("conversation not found: $id")
|
||||||
return@runBlocking
|
return@runBlocking
|
||||||
|
|||||||
@@ -7,6 +7,7 @@ import kotlinx.coroutines.flow.onEach
|
|||||||
import kotlinx.coroutines.flow.takeWhile
|
import kotlinx.coroutines.flow.takeWhile
|
||||||
import kotlinx.coroutines.launch
|
import kotlinx.coroutines.launch
|
||||||
import pw.binom.agentik.cli.AgentikSubcommand
|
import pw.binom.agentik.cli.AgentikSubcommand
|
||||||
|
import pw.binom.agentik.cli.defaultCliHttpClient
|
||||||
import pw.binom.agentik.client.AgentikAgent
|
import pw.binom.agentik.client.AgentikAgent
|
||||||
import pw.binom.agentik.proto.Content
|
import pw.binom.agentik.proto.Content
|
||||||
import pw.binom.agentik.proto.Event
|
import pw.binom.agentik.proto.Event
|
||||||
@@ -17,7 +18,7 @@ class SendSubcommand : AgentikSubcommand("send", "Отправить user-ход
|
|||||||
val text by argument(ArgType.String, description = "Текст хода (все позиционные после <id> склеиваются пробелом)").vararg()
|
val text by argument(ArgType.String, description = "Текст хода (все позиционные после <id> склеиваются пробелом)").vararg()
|
||||||
|
|
||||||
override fun execute() = kotlinx.coroutines.runBlocking {
|
override fun execute() = kotlinx.coroutines.runBlocking {
|
||||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl)
|
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
|
||||||
val conv = agent.getConversation(id) ?: run {
|
val conv = agent.getConversation(id) ?: run {
|
||||||
println("conversation not found: $id")
|
println("conversation not found: $id")
|
||||||
return@runBlocking
|
return@runBlocking
|
||||||
|
|||||||
@@ -22,8 +22,7 @@ kotlin {
|
|||||||
commonMain.dependencies {
|
commonMain.dependencies {
|
||||||
api(project(":proto"))
|
api(project(":proto"))
|
||||||
|
|
||||||
implementation(libs.ktor.client.core)
|
api(libs.ktor.client.core)
|
||||||
implementation(libs.ktor.client.cio)
|
|
||||||
implementation(libs.ktor.client.content.negotiation)
|
implementation(libs.ktor.client.content.negotiation)
|
||||||
implementation(libs.ktor.serialization.kotlinx.json)
|
implementation(libs.ktor.serialization.kotlinx.json)
|
||||||
|
|
||||||
@@ -37,6 +36,7 @@ kotlin {
|
|||||||
implementation(libs.ktor.server.core)
|
implementation(libs.ktor.server.core)
|
||||||
implementation(libs.ktor.server.test.host)
|
implementation(libs.ktor.server.test.host)
|
||||||
implementation(libs.ktor.client.content.negotiation)
|
implementation(libs.ktor.client.content.negotiation)
|
||||||
|
implementation(libs.ktor.client.cio)
|
||||||
implementation(libs.ktor.server.cio)
|
implementation(libs.ktor.server.cio)
|
||||||
implementation(libs.ktor.server.sse)
|
implementation(libs.ktor.server.sse)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -8,9 +8,11 @@ import pw.binom.agentik.proto.Agent
|
|||||||
* (модуль `:server`).
|
* (модуль `:server`).
|
||||||
*
|
*
|
||||||
* ```
|
* ```
|
||||||
|
* val http = HttpClient(CIO) { applyAgentikDefaults(token = "s3cret") }
|
||||||
* val client = AgentikAgent(
|
* val client = AgentikAgent(
|
||||||
* id = "my-agent",
|
* id = "my-agent",
|
||||||
* baseUrl = "http://localhost:8080/agentik",
|
* baseUrl = "http://localhost:8080/agentik",
|
||||||
|
* httpClient = http,
|
||||||
* )
|
* )
|
||||||
* val conv = client.createConversation(temp = false)
|
* val conv = client.createConversation(temp = false)
|
||||||
* conv.send(listOf(Content.Text("hi")))
|
* conv.send(listOf(Content.Text("hi")))
|
||||||
@@ -21,26 +23,13 @@ import pw.binom.agentik.proto.Agent
|
|||||||
* агента не знает, поэтому клиент должен её знать сам (или взять из
|
* агента не знает, поэтому клиент должен её знать сам (или взять из
|
||||||
* конфига).
|
* конфига).
|
||||||
*
|
*
|
||||||
* [httpClient] по умолчанию — [defaultAgentikHttpClient] (платформо-зависимый
|
* Клиент приходит снаружи: `:client` не выбирает движок. Собрать [HttpClient]
|
||||||
* движок: CIO на JVM, libcurl на desktop-native). Можно передать свой.
|
* можно через [agentikHttpClient] (фабрика движка + опциональные движковые
|
||||||
|
* настройки) или вручную, применив к блоку конфигурации [applyAgentikDefaults]
|
||||||
|
* (JSON + опциональный Bearer-токен).
|
||||||
*/
|
*/
|
||||||
fun AgentikAgent(
|
fun AgentikAgent(
|
||||||
id: String,
|
id: String,
|
||||||
baseUrl: String,
|
baseUrl: String,
|
||||||
httpClient: HttpClient = defaultAgentikHttpClient(),
|
httpClient: HttpClient,
|
||||||
): Agent = AgentClient(httpClient = httpClient, baseUrl = baseUrl, id = id)
|
): Agent = AgentClient(httpClient = httpClient, baseUrl = baseUrl, id = id)
|
||||||
|
|
||||||
/**
|
|
||||||
* Дефолтный [HttpClient] для общения с `agentikAgent`. SSE-парсер ([readSse])
|
|
||||||
* живёт в общем коде и плагина `SSEClientContent` не требует.
|
|
||||||
*
|
|
||||||
* **Платформы:**
|
|
||||||
* - JVM: движок CIO. `engine { requestTimeout = 0 }` отключает встроенный
|
|
||||||
* 15-секундный request-таймаут движка (наш кастомный SSE-ридер не маркирует
|
|
||||||
* для долгих idle-стримов). Defense-in-depth: SSE-запросы в
|
|
||||||
* `ConversationClient.events`/`AgentClient.events` уже ставят
|
|
||||||
* `HttpTimeoutCapability` = INFINITE (см. [noSseReadTimeout]).
|
|
||||||
*
|
|
||||||
* Один движок CIO работает и на JVM, и на всех desktop-native (linux/macos/mingw).
|
|
||||||
* Реализация — в [HttpClientFactory.kt].
|
|
||||||
*/
|
|
||||||
@@ -1,18 +1,60 @@
|
|||||||
package pw.binom.agentik.client
|
package pw.binom.agentik.client
|
||||||
|
|
||||||
import io.ktor.client.HttpClient
|
import io.ktor.client.HttpClient
|
||||||
import io.ktor.client.engine.cio.CIO
|
import io.ktor.client.HttpClientConfig
|
||||||
|
import io.ktor.client.engine.HttpClientEngineConfig
|
||||||
|
import io.ktor.client.engine.HttpClientEngineFactory
|
||||||
|
import io.ktor.client.plugins.DefaultRequest
|
||||||
import io.ktor.client.plugins.contentnegotiation.ContentNegotiation
|
import io.ktor.client.plugins.contentnegotiation.ContentNegotiation
|
||||||
|
import io.ktor.client.request.header
|
||||||
|
import io.ktor.http.HttpHeaders
|
||||||
import io.ktor.serialization.kotlinx.json.json
|
import io.ktor.serialization.kotlinx.json.json
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Единый HTTP-клиент для JVM и всех 5 native-таргетов (:agentik-cli).
|
* Общая конфигурация HTTP-клиента agentik — платформо-независимая часть.
|
||||||
* CIO в ktor 3.x — KMP, поддерживает linuxX64/Arm64, macosX64/Arm64, mingwX64.
|
|
||||||
*
|
*
|
||||||
* `requestTimeout = 0` — defense-in-depth против read-таймаута на SSE:
|
* `:client` НЕ выбирает движок: его приносит потребитель. Здесь живёт только то,
|
||||||
* основная защита в `HttpRequestBuilder.noSseReadTimeout()` ([SseTimeout]).
|
* без чего клиент несовместим с `/agentik`:
|
||||||
|
* - JSON-конфиг [agentikJson] (обязан совпадать с серверным);
|
||||||
|
* - при заданном [token] — `Authorization: Bearer <token>` на ВСЕ запросы
|
||||||
|
* через [DefaultRequest] (накрывает 10 REST-вызовов и оба SSE-потока;
|
||||||
|
* заголовок живёт на клиенте, а не в отдельных запросах).
|
||||||
|
*
|
||||||
|
* `null` — авторизация выключена, заголовок не отправляется.
|
||||||
|
*
|
||||||
|
* Потребитель, знающий свой движок, добавляет к этому движковые настройки, напр.:
|
||||||
|
* ```
|
||||||
|
* val http = HttpClient(CIO) {
|
||||||
|
* engine { requestTimeout = 0 } // CIO-специфика, живёт у потребителя
|
||||||
|
* applyAgentikDefaults(token)
|
||||||
|
* }
|
||||||
|
* ```
|
||||||
*/
|
*/
|
||||||
fun defaultAgentikHttpClient(): HttpClient = HttpClient(CIO) {
|
fun HttpClientConfig<*>.applyAgentikDefaults(token: String? = null) {
|
||||||
engine { requestTimeout = 0 }
|
|
||||||
install(ContentNegotiation) { json(agentikJson) }
|
install(ContentNegotiation) { json(agentikJson) }
|
||||||
|
if (token != null) {
|
||||||
|
install(DefaultRequest) {
|
||||||
|
header(HttpHeaders.Authorization, "Bearer $token")
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Создаёт [HttpClient] из фабрики движка потребителя и сразу применяет к нему
|
||||||
|
* конфигурацию agentik ([applyAgentikDefaults]).
|
||||||
|
*
|
||||||
|
* Это точка, где `:client` НЕ привязан к реализации транспорта: [engineFactory]
|
||||||
|
* выбирает потребитель (CIO, OkHttp, Darwin, …), а `:client` только конфигурирует
|
||||||
|
* созданный клиент.
|
||||||
|
*
|
||||||
|
* [configure] — опциональный последний штрих потребителя (движковые настройки:
|
||||||
|
* таймауты, прокси, логирование). Вызывается ПОСЛЕ [applyAgentikDefaults].
|
||||||
|
*/
|
||||||
|
fun <T : HttpClientEngineConfig> agentikHttpClient(
|
||||||
|
engineFactory: HttpClientEngineFactory<T>,
|
||||||
|
token: String? = null,
|
||||||
|
configure: (HttpClientConfig<T>.() -> Unit)? = null,
|
||||||
|
): HttpClient = HttpClient(engineFactory) {
|
||||||
|
applyAgentikDefaults(token)
|
||||||
|
configure?.invoke(this)
|
||||||
|
}
|
||||||
@@ -0,0 +1,106 @@
|
|||||||
|
package pw.binom.agentik.client
|
||||||
|
|
||||||
|
import io.ktor.client.HttpClient
|
||||||
|
import io.ktor.client.engine.cio.CIO
|
||||||
|
import io.ktor.client.request.get
|
||||||
|
import io.ktor.client.statement.bodyAsText
|
||||||
|
import io.ktor.http.ContentType
|
||||||
|
import io.ktor.http.HttpHeaders
|
||||||
|
import io.ktor.http.HttpStatusCode
|
||||||
|
import io.ktor.server.application.call
|
||||||
|
import io.ktor.server.application.createRouteScopedPlugin
|
||||||
|
import io.ktor.server.cio.CIO as ServerCIO
|
||||||
|
import io.ktor.server.engine.EmbeddedServer
|
||||||
|
import io.ktor.server.engine.embeddedServer
|
||||||
|
import io.ktor.server.response.respondText
|
||||||
|
import io.ktor.server.routing.get
|
||||||
|
import io.ktor.server.routing.route
|
||||||
|
import io.ktor.server.routing.routing
|
||||||
|
import kotlinx.coroutines.runBlocking
|
||||||
|
import kotlin.test.Test
|
||||||
|
import kotlin.test.assertEquals
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Тесты клиентской части: [applyAgentikDefaults] с заданным `token` прикладывает
|
||||||
|
* `Authorization: Bearer <token>` ко всем запросам через плагин `DefaultRequest`,
|
||||||
|
* без токена — заголовок не отправляется.
|
||||||
|
*
|
||||||
|
* Сервер в тесте — локальный ktor-CIO с inline route-scoped Bearer-плагином (один в один
|
||||||
|
* как боевой [pw.binom.agentik.server.BearerTokenPlugin]). Тестовый `:server` не зависит
|
||||||
|
* от `:client`, поэтому боевой плагин тут переиспользовать нельзя — пересоздаём его
|
||||||
|
* минимально, контракт тот же.
|
||||||
|
*/
|
||||||
|
class BearerHeaderTest {
|
||||||
|
|
||||||
|
private val TestBearer = createRouteScopedPlugin(
|
||||||
|
name = "TestBearer",
|
||||||
|
createConfiguration = ::BearerCfg,
|
||||||
|
) {
|
||||||
|
val expected = pluginConfig.token
|
||||||
|
onCall { call ->
|
||||||
|
if (expected == null) return@onCall
|
||||||
|
if (call.request.headers[HttpHeaders.Authorization] != "Bearer $expected") {
|
||||||
|
call.respondText("Unauthorized", ContentType.Text.Plain, HttpStatusCode.Unauthorized)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private class BearerCfg {
|
||||||
|
var token: String? = null
|
||||||
|
}
|
||||||
|
|
||||||
|
private fun clientWith(token: String?): HttpClient =
|
||||||
|
HttpClient(CIO) { applyAgentikDefaults(token) }
|
||||||
|
|
||||||
|
private suspend fun startServer(): Pair<EmbeddedServer<*, *>, Int> {
|
||||||
|
val server = embeddedServer(ServerCIO, port = 0) {
|
||||||
|
routing {
|
||||||
|
route("/agentik") {
|
||||||
|
install(TestBearer) { token = "secret" }
|
||||||
|
get("/conversations") {
|
||||||
|
call.respondText("[]")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}.start(wait = false)
|
||||||
|
val port = server.engine.resolvedConnectors().first().port
|
||||||
|
return server to port
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun clientWithTokenAttachesBearerHeader() = runBlocking {
|
||||||
|
val (server, port) = startServer()
|
||||||
|
try {
|
||||||
|
val client = clientWith("secret")
|
||||||
|
val resp = client.get("http://127.0.0.1:$port/agentik/conversations")
|
||||||
|
assertEquals(HttpStatusCode.OK, resp.status)
|
||||||
|
assertEquals("[]", resp.bodyAsText())
|
||||||
|
} finally {
|
||||||
|
server.stop(100, 200)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun clientWithoutTokenGets401(): Unit = runBlocking {
|
||||||
|
val (server, port) = startServer()
|
||||||
|
try {
|
||||||
|
val client = clientWith(null)
|
||||||
|
val resp = client.get("http://127.0.0.1:$port/agentik/conversations")
|
||||||
|
assertEquals(HttpStatusCode.Unauthorized, resp.status)
|
||||||
|
} finally {
|
||||||
|
server.stop(100, 200)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun clientWithWrongTokenGets401(): Unit = runBlocking {
|
||||||
|
val (server, port) = startServer()
|
||||||
|
try {
|
||||||
|
val client = clientWith("wrong")
|
||||||
|
val resp = client.get("http://127.0.0.1:$port/agentik/conversations")
|
||||||
|
assertEquals(HttpStatusCode.Unauthorized, resp.status)
|
||||||
|
} finally {
|
||||||
|
server.stop(100, 200)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -37,6 +37,8 @@ kotlin {
|
|||||||
commonTest.dependencies {
|
commonTest.dependencies {
|
||||||
implementation(libs.kotlinx.coroutines.core)
|
implementation(libs.kotlinx.coroutines.core)
|
||||||
implementation(kotlin("test"))
|
implementation(kotlin("test"))
|
||||||
|
implementation(libs.ktor.server.test.host)
|
||||||
|
implementation(libs.ktor.server.cio)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,42 @@
|
|||||||
|
package pw.binom.agentik.server
|
||||||
|
|
||||||
|
import io.ktor.http.ContentType
|
||||||
|
import io.ktor.http.HttpHeaders
|
||||||
|
import io.ktor.http.HttpStatusCode
|
||||||
|
import io.ktor.server.application.createRouteScopedPlugin
|
||||||
|
import io.ktor.server.request.path
|
||||||
|
import io.ktor.server.response.respondText
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Конфиг плагина проверки `Authorization: Bearer <token>` для роута `agentikAgent`.
|
||||||
|
*
|
||||||
|
* По умолчанию [token] == null → плагин пропускает все запросы (см. [Module.kt]).
|
||||||
|
*/
|
||||||
|
internal class BearerTokenConfig {
|
||||||
|
var token: String? = null
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Route-scoped плагин: если в конфиге задан [BearerTokenConfig.token], на каждый
|
||||||
|
* запрос внутри ветки роута проверяет заголовок `Authorization: Bearer <token>`.
|
||||||
|
* При несовпадении отвечает `401 Unauthorized` (тело `Unauthorized`); дальнейшие
|
||||||
|
* обработчики не вызываются — ktor трактует отправленный ответ как завершение
|
||||||
|
* call-pipeline.
|
||||||
|
*
|
||||||
|
* `/health` всегда пропускается без проверки: это ручка liveness для
|
||||||
|
* балансировщика/мониторинга, закрывать её — сломать health-check.
|
||||||
|
*/
|
||||||
|
internal val BearerTokenPlugin = createRouteScopedPlugin(
|
||||||
|
name = "AgentikBearerToken",
|
||||||
|
createConfiguration = ::BearerTokenConfig,
|
||||||
|
) {
|
||||||
|
val expected = pluginConfig.token
|
||||||
|
onCall { call ->
|
||||||
|
if (expected == null) return@onCall
|
||||||
|
val path = call.request.path()
|
||||||
|
if (path.endsWith("/health")) return@onCall
|
||||||
|
if (call.request.headers[HttpHeaders.Authorization] != "Bearer $expected") {
|
||||||
|
call.respondText("Unauthorized", ContentType.Text.Plain, HttpStatusCode.Unauthorized)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -33,11 +33,16 @@ import pw.binom.agentik.proto.Agent
|
|||||||
* - `GET /events` — SSE: события агента
|
* - `GET /events` — SSE: события агента
|
||||||
* - `GET /health` — `"ok"`
|
* - `GET /health` — `"ok"`
|
||||||
*/
|
*/
|
||||||
fun Route.agentikAgent(agent: Agent, path: String = "/agentik") {
|
fun Route.agentikAgent(agent: Agent, path: String = "/agentik", token: String? = null) {
|
||||||
route(path) {
|
route(path) {
|
||||||
install(ContentNegotiation) {
|
install(ContentNegotiation) {
|
||||||
json(agentikJson)
|
json(agentikJson)
|
||||||
}
|
}
|
||||||
|
if (token != null) {
|
||||||
|
install(BearerTokenPlugin) {
|
||||||
|
this.token = token
|
||||||
|
}
|
||||||
|
}
|
||||||
agentikRoutes(agent)
|
agentikRoutes(agent)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,115 @@
|
|||||||
|
package pw.binom.agentik.server
|
||||||
|
|
||||||
|
import io.ktor.client.HttpClient
|
||||||
|
import io.ktor.client.engine.cio.CIO
|
||||||
|
import io.ktor.client.request.get
|
||||||
|
import io.ktor.client.request.header
|
||||||
|
import io.ktor.client.statement.bodyAsText
|
||||||
|
import io.ktor.http.HttpHeaders
|
||||||
|
import io.ktor.http.HttpStatusCode
|
||||||
|
import io.ktor.server.cio.CIO as ServerCIO
|
||||||
|
import io.ktor.server.engine.EmbeddedServer
|
||||||
|
import io.ktor.server.engine.embeddedServer
|
||||||
|
import io.ktor.server.routing.routing
|
||||||
|
import kotlinx.coroutines.flow.Flow
|
||||||
|
import kotlinx.coroutines.flow.emptyFlow
|
||||||
|
import kotlinx.coroutines.runBlocking
|
||||||
|
import pw.binom.agentik.proto.Agent
|
||||||
|
import pw.binom.agentik.proto.AgentEvent
|
||||||
|
import pw.binom.agentik.proto.Conversation
|
||||||
|
import kotlin.test.Test
|
||||||
|
import kotlin.test.assertEquals
|
||||||
|
import kotlin.time.Instant
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Тесты route-scoped плагина [BearerTokenPlugin]:
|
||||||
|
* - при `token != null` все роуты кроме `/health` требуют `Authorization: Bearer <token>`;
|
||||||
|
* - при `token == null` плагин не устанавливается, всё открыто;
|
||||||
|
* - `/health` всегда открыт, даже при заданном токене (liveness-ручка для балансировщика).
|
||||||
|
*/
|
||||||
|
class BearerTokenTest {
|
||||||
|
|
||||||
|
private class FakeAgent(override val id: String = "test") : Agent {
|
||||||
|
override fun createConversation(temp: Boolean): Conversation = TODO("not needed by tests")
|
||||||
|
override suspend fun getConversation(id: String): Conversation? = null
|
||||||
|
override suspend fun deleteConversation(id: String): Boolean = false
|
||||||
|
override suspend fun getConversations(offset: Int, limit: Int): List<Conversation> = emptyList()
|
||||||
|
override fun events(after: Instant): Flow<AgentEvent> = emptyFlow()
|
||||||
|
}
|
||||||
|
|
||||||
|
private suspend fun startServer(token: String?): Pair<EmbeddedServer<*, *>, Int> {
|
||||||
|
val server = embeddedServer(ServerCIO, port = 0) {
|
||||||
|
routing {
|
||||||
|
agentikAgent(FakeAgent(), path = "/agentik", token = token)
|
||||||
|
}
|
||||||
|
}.start(wait = false)
|
||||||
|
val port = server.engine.resolvedConnectors().first().port
|
||||||
|
return server to port
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun tokenRejectsRequestWithoutHeader() = runBlocking {
|
||||||
|
val (server, port) = startServer("secret")
|
||||||
|
try {
|
||||||
|
val client = HttpClient(CIO)
|
||||||
|
val resp = client.get("http://127.0.0.1:$port/agentik/conversations")
|
||||||
|
assertEquals(HttpStatusCode.Unauthorized, resp.status)
|
||||||
|
assertEquals("Unauthorized", resp.bodyAsText())
|
||||||
|
} finally {
|
||||||
|
server.stop(100, 200)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun tokenRejectsWrongHeader() = runBlocking {
|
||||||
|
val (server, port) = startServer("secret")
|
||||||
|
try {
|
||||||
|
val client = HttpClient(CIO)
|
||||||
|
val resp = client.get("http://127.0.0.1:$port/agentik/conversations") {
|
||||||
|
header(HttpHeaders.Authorization, "Bearer wrong")
|
||||||
|
}
|
||||||
|
assertEquals(HttpStatusCode.Unauthorized, resp.status)
|
||||||
|
} finally {
|
||||||
|
server.stop(100, 200)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun tokenAcceptsCorrectHeader() = runBlocking {
|
||||||
|
val (server, port) = startServer("secret")
|
||||||
|
try {
|
||||||
|
val client = HttpClient(CIO)
|
||||||
|
val resp = client.get("http://127.0.0.1:$port/agentik/conversations") {
|
||||||
|
header(HttpHeaders.Authorization, "Bearer secret")
|
||||||
|
}
|
||||||
|
assertEquals(HttpStatusCode.OK, resp.status)
|
||||||
|
} finally {
|
||||||
|
server.stop(100, 200)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun healthStaysOpenWithToken() = runBlocking {
|
||||||
|
val (server, port) = startServer("secret")
|
||||||
|
try {
|
||||||
|
val client = HttpClient(CIO)
|
||||||
|
val resp = client.get("http://127.0.0.1:$port/agentik/health")
|
||||||
|
assertEquals(HttpStatusCode.OK, resp.status)
|
||||||
|
assertEquals("ok", resp.bodyAsText())
|
||||||
|
} finally {
|
||||||
|
server.stop(100, 200)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun nullTokenMeansOpen() = runBlocking {
|
||||||
|
val (server, port) = startServer(null)
|
||||||
|
try {
|
||||||
|
val client = HttpClient(CIO)
|
||||||
|
val resp = client.get("http://127.0.0.1:$port/agentik/conversations")
|
||||||
|
assertEquals(HttpStatusCode.OK, resp.status)
|
||||||
|
} finally {
|
||||||
|
server.stop(100, 200)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -69,6 +69,8 @@ AGENTIK_GOOGLE_MODEL_PATH=/root/gemma-4-E2B-it.litertlm \
|
|||||||
|---|---|---|
|
|---|---|---|
|
||||||
| `AGENTIK_PORT` | `8080` | Порт HTTP-сервера |
|
| `AGENTIK_PORT` | `8080` | Порт HTTP-сервера |
|
||||||
| `AGENTIK_DB_PATH` | `./agentik.db` | Путь к SQLite |
|
| `AGENTIK_DB_PATH` | `./agentik.db` | Путь к SQLite |
|
||||||
|
| `AGENTIK_TOKEN` | (пусто) | Bearer-токен для HTTP-фасада `/agentik`. Пусто — авторизация выключена |
|
||||||
|
| `AGENTIK_A2A_TOKEN` | (пусто) | Bearer-токен для A2A-фасада `/a2a`. Пусто — авторизация выключена (независим от `AGENTIK_TOKEN`) |
|
||||||
| `AGENTIK_AGENT_ID` | `agentik` | ID агента (для multi-instance) |
|
| `AGENTIK_AGENT_ID` | `agentik` | ID агента (для multi-instance) |
|
||||||
| `AGENTIK_LLM_BACKEND` | `openai` | `openai` или `google` |
|
| `AGENTIK_LLM_BACKEND` | `openai` | `openai` или `google` |
|
||||||
| `AGENTIK_LLM_MODEL` | (выбирается по backend) | Имя модели |
|
| `AGENTIK_LLM_MODEL` | (выбирается по backend) | Имя модели |
|
||||||
|
|||||||
@@ -51,6 +51,8 @@ import pw.binom.agentik.llm.tools.SkillMiner
|
|||||||
*
|
*
|
||||||
* Вся конфигурация — [AppConfig.fromEnv] (см. [AppConfig]). Источники:
|
* Вся конфигурация — [AppConfig.fromEnv] (см. [AppConfig]). Источники:
|
||||||
* - AGENTIK_PORT / AGENTIK_DB_PATH
|
* - AGENTIK_PORT / AGENTIK_DB_PATH
|
||||||
|
* - AGENTIK_TOKEN — Bearer-токен для HTTP-фасада /agentik (пусто — авторизация выключена)
|
||||||
|
* - AGENTIK_A2A_TOKEN — Bearer-токен для A2A-фасада /a2a (пусто — авторизация выключена)
|
||||||
* - LLM: AGENTIK_LLM_BACKEND, OPENAI_* либо AGENTIK_GOOGLE_*
|
* - LLM: AGENTIK_LLM_BACKEND, OPENAI_* либо AGENTIK_GOOGLE_*
|
||||||
* - MCP: AGENTIK_MCP_CONFIG=<path>.json (формат Claude Desktop)
|
* - MCP: AGENTIK_MCP_CONFIG=<path>.json (формат Claude Desktop)
|
||||||
* - Skills: AGENTIK_SKILLS_DIR=<path> (папка с SKILL.md / *.yaml)
|
* - Skills: AGENTIK_SKILLS_DIR=<path> (папка с SKILL.md / *.yaml)
|
||||||
@@ -353,8 +355,13 @@ private fun runServer() {
|
|||||||
val server = embeddedServer(CIO, port = config.agent.port) {
|
val server = embeddedServer(CIO, port = config.agent.port) {
|
||||||
routing {
|
routing {
|
||||||
get("/health") { call.respondText("ok") }
|
get("/health") { call.respondText("ok") }
|
||||||
agentikAgent(agent, path = "/agentik")
|
agentikAgent(agent, path = "/agentik", token = config.agent.authToken)
|
||||||
a2aAgent(agentName = "agentik", handler = A2aBridge(agent), path = "/a2a")
|
a2aAgent(
|
||||||
|
agentName = "agentik",
|
||||||
|
handler = A2aBridge(agent),
|
||||||
|
path = "/a2a",
|
||||||
|
token = config.agent.a2aToken,
|
||||||
|
)
|
||||||
if (config.debug.endpoints) {
|
if (config.debug.endpoints) {
|
||||||
debugRoutes(
|
debugRoutes(
|
||||||
agent = agent,
|
agent = agent,
|
||||||
|
|||||||
@@ -50,6 +50,18 @@ data class AppConfig(
|
|||||||
* `null` — файл не читается, секция не добавляется.
|
* `null` — файл не читается, секция не добавляется.
|
||||||
*/
|
*/
|
||||||
val soulPath: String? = null,
|
val soulPath: String? = null,
|
||||||
|
/**
|
||||||
|
* Токен для HTTP-фасада agentik (`/agentik`). `null` — авторизация выключена,
|
||||||
|
* сервер открыт (поведение по умолчанию, обратная совместимость).
|
||||||
|
* Источник: env `AGENTIK_TOKEN`.
|
||||||
|
*/
|
||||||
|
val authToken: String? = null,
|
||||||
|
/**
|
||||||
|
* Токен для A2A-фасада (`/a2a`). НЕ связан с [authToken] — это разные
|
||||||
|
* интерфейсы и разные токены (см. библиотеку pw.binom.a2a).
|
||||||
|
* `null` — авторизация выключена. Источник: env `AGENTIK_A2A_TOKEN`.
|
||||||
|
*/
|
||||||
|
val a2aToken: String? = null,
|
||||||
)
|
)
|
||||||
|
|
||||||
/** Долговременная память. */
|
/** Долговременная память. */
|
||||||
@@ -163,6 +175,8 @@ data class AppConfig(
|
|||||||
dbPath = env("AGENTIK_DB_PATH")?.takeIf { it.isNotBlank() } ?: DEFAULT_DB_PATH,
|
dbPath = env("AGENTIK_DB_PATH")?.takeIf { it.isNotBlank() } ?: DEFAULT_DB_PATH,
|
||||||
skillsDir = env("AGENTIK_SKILLS_DIR")?.takeIf { it.isNotBlank() },
|
skillsDir = env("AGENTIK_SKILLS_DIR")?.takeIf { it.isNotBlank() },
|
||||||
soulPath = env("AGENTIK_SOUL")?.takeIf { it.isNotBlank() },
|
soulPath = env("AGENTIK_SOUL")?.takeIf { it.isNotBlank() },
|
||||||
|
authToken = env("AGENTIK_TOKEN")?.takeIf { it.isNotBlank() },
|
||||||
|
a2aToken = env("AGENTIK_A2A_TOKEN")?.takeIf { it.isNotBlank() },
|
||||||
),
|
),
|
||||||
llm = llm,
|
llm = llm,
|
||||||
mcp = McpConfig.fromEnv(env),
|
mcp = McpConfig.fromEnv(env),
|
||||||
|
|||||||
Reference in New Issue
Block a user