Compare commits
24 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| c140d0b758 | |||
| c486c7f9ab | |||
| 9d310c5fd0 | |||
| 741ad8963d | |||
| 2634e0e204 | |||
| ac5d209fce | |||
| 25771a0c33 | |||
| 78cbe9b463 | |||
| 65e05612a1 | |||
| 850ee99cb6 | |||
| b5b21d146a | |||
| ee0b9d8341 | |||
| 9d826a4e81 | |||
| db3c49099c | |||
| ddd9d076c1 | |||
| 7a47131f6f | |||
| 9196102f68 | |||
| 6b12dd2c5b | |||
| eed1ab9a17 | |||
| b27ac622b4 | |||
| 14b46087dd | |||
| 098c97c7bd | |||
| 4b8e5bb0bd | |||
| d75289ac56 |
+16
-21
@@ -17,6 +17,14 @@ concurrency:
|
||||
group: ${{ github.workflow }}-${{ github.ref }}
|
||||
cancel-in-progress: true
|
||||
|
||||
# UTF-8 обязателен: в именах тестов есть типографские символы (—), а Kotlin-компилятор
|
||||
# создаёт .class-файлы с именем теста. При LANG=C sun.jnu.encoding = ASCII, и компилятор
|
||||
# падает с "InvalidPathException: Malformed input or input contains unmappable characters"
|
||||
# (проверено локально: LANG=C → BUILD FAILED, LANG=C.UTF-8 → BUILD SUCCESSFUL).
|
||||
env:
|
||||
LANG: C.UTF-8
|
||||
LC_ALL: C.UTF-8
|
||||
|
||||
jobs:
|
||||
build-jvm:
|
||||
name: JVM build + tests
|
||||
@@ -65,25 +73,12 @@ jobs:
|
||||
./gradlew :agentik-cli:shadowJar \
|
||||
-Dorg.gradle.jvmargs=-Xmx4096M \
|
||||
--no-daemon --no-watch-fs --stacktrace
|
||||
test -f agentik-cli/build/libs/agentik-cli-all.jar \
|
||||
&& echo "shadowJar OK: $(du -h agentik-cli/build/libs/agentik-cli-all.jar)"
|
||||
test -f agentik-cli/build/libs/agentik-cli-*-all.jar \
|
||||
&& echo "shadowJar OK: $(du -h agentik-cli/build/libs/agentik-cli-*-all.jar)"
|
||||
|
||||
- name: Build :agentik-tui shadowJar
|
||||
shell: bash
|
||||
run: |
|
||||
./gradlew :agentik-tui:shadowJar \
|
||||
-Dorg.gradle.jvmargs=-Xmx4096M \
|
||||
--no-daemon --no-watch-fs --stacktrace
|
||||
test -f agentik-tui/build/libs/agentik-tui-all.jar \
|
||||
&& echo "shadowJar OK: $(du -h agentik-tui/build/libs/agentik-tui-all.jar)"
|
||||
|
||||
- name: Upload shadowJars
|
||||
uses: actions/upload-artifact@v4
|
||||
with:
|
||||
name: agentik-jars
|
||||
path: |
|
||||
standalone/build/libs/standalone-all.jar
|
||||
agentik-cli/build/libs/agentik-cli-all.jar
|
||||
agentik-tui/build/libs/agentik-tui-all.jar
|
||||
if-no-files-found: error
|
||||
retention-days: 7
|
||||
# Шага "Upload shadowJars" здесь нет сознательно: upload-artifact@v4 требует
|
||||
# @actions/artifact v2, который на GHES/Gitea-раннере падает с
|
||||
# "GHESNotSupportedError: @actions/artifact v2.0.0+ ... not supported on GHES"
|
||||
# и валит весь джоб уже ПОСЛЕ успешной сборки и зелёных тестов.
|
||||
# У соседних репо (asr-kmp, litert-kmp) артефакты наружу тоже не выгружаются —
|
||||
# проверка сборки ограничивается test -f на jar (шаги выше).
|
||||
|
||||
+27
-113
@@ -1,14 +1,20 @@
|
||||
# Триггерится при публикации релиза в Gitea. Публикует все KMP-библиотеки
|
||||
# (jvm + native таргеты) в домашний Nexus-репозиторий "caffeine", а также
|
||||
# собирает fatjar'ы запускаемых модулей и прикрепляет их к релизу как
|
||||
# бинарные ассеты.
|
||||
# (jvm + native таргеты) в домашний Nexus-репозиторий "caffeine".
|
||||
#
|
||||
# Требуемые Gitea Action Variables:
|
||||
# BINOM_REPO_URL — например http://192.168.76.117/repository/caffeine/
|
||||
# Требуемые Gitea Action Secrets:
|
||||
# BINOM_REPO_USER, BINOM_REPO_PASSWORD — креды Nexus с правами на публикацию.
|
||||
# RELEASE_TOKEN — Gitea API-токен с правами на запись в релиз (используется
|
||||
# forgejo-release action для прикрепления fatjar'ов к release).
|
||||
# Fatjar-ы запускаемых модулей (:standalone, :agentik-cli).
|
||||
# :agentik-tui был исключён из сборки 2026-09-17 (см. settings.gradle.kts).
|
||||
# НЕ собираются и НЕ крепятся к релизу здесь. Сборка артефактов
|
||||
# выполняется локально из исходников (или руками через `./gradlew
|
||||
# :<module>:shadowJar`) и загружается в релиз через Gitea UI / API
|
||||
# отдельно от этого workflow.
|
||||
#
|
||||
# Версия публикации = имя тега релиза (без префикса 'v'). Релиз с именем "3"
|
||||
# публикует pw.binom.agentik:*:3 в Nexus. Ничего хардкодить не нужно —
|
||||
# версия берётся из тега каждый раз.
|
||||
#
|
||||
# Публикация выполняется общим composite-action'ом subochev/devops/publish@main
|
||||
# (тот же, что у asr-kmp / litert-kmp / embedder-kmp / a2a-protocol) — credentials
|
||||
# BINOM_REPO_* берутся им из Gitea Action Variables (owner_id=0, глобальные).
|
||||
name: release
|
||||
|
||||
on:
|
||||
@@ -19,6 +25,15 @@ concurrency:
|
||||
group: release-${{ github.ref }}
|
||||
cancel-in-progress: false
|
||||
|
||||
# UTF-8 обязателен: генерация POM/Kotlin-метаданных и имена тестовых классов
|
||||
# содержат не-ASCII символы; при LANG=C sun.jnu.encoding = ASCII и сборка
|
||||
# падает с "InvalidPathException: Malformed input or input contains unmappable
|
||||
# characters" (проверено локально 19.09.2026: LANG=C → BUILD FAILED,
|
||||
# LANG=C.UTF-8 → BUILD SUCCESSFUL).
|
||||
env:
|
||||
LANG: C.UTF-8
|
||||
LC_ALL: C.UTF-8
|
||||
|
||||
jobs:
|
||||
publish-libraries:
|
||||
name: Publish KMP libraries → caffeine Nexus
|
||||
@@ -28,108 +43,7 @@ jobs:
|
||||
- name: Checkout
|
||||
uses: actions/checkout@v4
|
||||
|
||||
- name: Setup JDK 21
|
||||
uses: actions/setup-java@v4
|
||||
- name: Publish libraries (all KMP targets, all modules) to Nexus
|
||||
uses: https://git.binom.pw/subochev/devops/publish@main
|
||||
with:
|
||||
java-version: '21'
|
||||
distribution: 'adopt'
|
||||
|
||||
- name: Gradle cache
|
||||
uses: actions/cache@v4
|
||||
with:
|
||||
path: |
|
||||
~/.gradle/caches
|
||||
~/.gradle/wrapper
|
||||
.gradle
|
||||
key: ${{ runner.os }}-gradle-agentik-${{ hashFiles('**/*.gradle.kts', '**/gradle/libs.versions.toml', 'gradle/wrapper/gradle-wrapper.properties') }}
|
||||
restore-keys: |
|
||||
${{ runner.os }}-gradle-agentik-
|
||||
|
||||
- name: Publish libraries (all KMP targets, all modules)
|
||||
shell: bash
|
||||
env:
|
||||
BINOM_REPO_USER: ${{ secrets.BINOM_REPO_USER }}
|
||||
BINOM_REPO_PASSWORD: ${{ secrets.BINOM_REPO_PASSWORD }}
|
||||
BINOM_REPO_URL: ${{ vars.BINOM_REPO_URL }}
|
||||
run: |
|
||||
./gradlew \
|
||||
"-Pversion=${GITEA_REF_NAME}" \
|
||||
"-Pbinom.repo.url=${BINOM_REPO_URL}" \
|
||||
"-Pbinom.repo.user=${BINOM_REPO_USER}" \
|
||||
"-Pbinom.repo.password=${BINOM_REPO_PASSWORD}" \
|
||||
publish \
|
||||
-Dorg.gradle.jvmargs=-Xmx4096M \
|
||||
--no-daemon --no-watch-fs --stacktrace
|
||||
|
||||
build-fatjars:
|
||||
name: Build runnable fatjars
|
||||
runs-on: ubuntu-latest
|
||||
timeout-minutes: 30
|
||||
steps:
|
||||
- name: Checkout
|
||||
uses: actions/checkout@v4
|
||||
|
||||
- name: Setup JDK 21
|
||||
uses: actions/setup-java@v4
|
||||
with:
|
||||
java-version: '21'
|
||||
distribution: 'adopt'
|
||||
|
||||
- name: Gradle cache
|
||||
uses: actions/cache@v4
|
||||
with:
|
||||
path: |
|
||||
~/.gradle/caches
|
||||
~/.gradle/wrapper
|
||||
.gradle
|
||||
key: ${{ runner.os }}-gradle-agentik-${{ hashFiles('**/*.gradle.kts', '**/gradle/libs.versions.toml', 'gradle/wrapper/gradle-wrapper.properties') }}
|
||||
restore-keys: |
|
||||
${{ runner.os }}-gradle-agentik-
|
||||
|
||||
- name: Build :standalone shadowJar
|
||||
shell: bash
|
||||
run: |
|
||||
./gradlew :standalone:shadowJar \
|
||||
-Pdisable-javadoc=true \
|
||||
-Dorg.gradle.jvmargs=-Xmx4096M \
|
||||
--no-daemon --no-watch-fs --stacktrace
|
||||
|
||||
- name: Build :agentik-cli shadowJar
|
||||
shell: bash
|
||||
run: |
|
||||
./gradlew :agentik-cli:shadowJar \
|
||||
-Pdisable-javadoc=true \
|
||||
-Dorg.gradle.jvmargs=-Xmx4096M \
|
||||
--no-daemon --no-watch-fs --stacktrace
|
||||
|
||||
- name: Build :agentik-tui shadowJar
|
||||
shell: bash
|
||||
run: |
|
||||
./gradlew :agentik-tui:shadowJar \
|
||||
-Pdisable-javadoc=true \
|
||||
-Dorg.gradle.jvmargs=-Xmx4096M \
|
||||
--no-daemon --no-watch-fs --stacktrace
|
||||
|
||||
- name: Upload fatjars as release assets
|
||||
uses: actions/upload-artifact@v4
|
||||
with:
|
||||
name: agentik-fatjars
|
||||
path: |
|
||||
standalone/build/libs/standalone-*-all.jar
|
||||
agentik-cli/build/libs/agentik-cli-*-all.jar
|
||||
agentik-tui/build/libs/agentik-tui-*-all.jar
|
||||
if-no-files-found: error
|
||||
retention-days: 90
|
||||
|
||||
- name: Attach to release
|
||||
uses: https://git.binom.pw/actions/forgejo-release@v1
|
||||
if: startsWith(github.ref, 'refs/tags/')
|
||||
with:
|
||||
url: ${{ github.server_url }}
|
||||
repo: ${{ github.repository }}
|
||||
token: ${{ secrets.RELEASE_TOKEN }}
|
||||
tag: ${{ github.ref_name }}
|
||||
files: |
|
||||
standalone/build/libs/standalone-*-all.jar
|
||||
agentik-cli/build/libs/agentik-cli-*-all.jar
|
||||
agentik-tui/build/libs/agentik-tui-*-all.jar
|
||||
version: ${{ gitea.ref_name }}
|
||||
|
||||
@@ -20,6 +20,9 @@ out/
|
||||
.cortexkit/
|
||||
.veai/
|
||||
|
||||
# Internal review scratch dir (review/validation .md файлы, .tasks структура)
|
||||
.tasks/
|
||||
|
||||
# Runtime / test artifacts
|
||||
agentik.db
|
||||
agentik.db-shm
|
||||
|
||||
@@ -21,8 +21,8 @@ agentik/
|
||||
├── storage-inmemory/ in-memory реализация для тестов и Android
|
||||
├── storage-sqlite/ SQLite реализация для production
|
||||
├── agent-toolsets/ ядро tool-calls с cooperative cancel + concurrency budget
|
||||
├── agentik-cli/ JVM REPL-клиент (JLine) к /agentik
|
||||
├── agentik-tui/ Compose-for-Mosaic TUI-клиент (desktop) к /agentik
|
||||
├── agentik-cli/ JVM one-shot CLI-клиент (kotlinx.cli) к /agentik
|
||||
├── ~~agentik-tui/~~ ~~Compose-for-Mosaic TUI-клиент (desktop)~~ — исключён 2026-09-17
|
||||
└── standalone/ single-jar HTTP-сервер со всеми transport'ами и движками
|
||||
```
|
||||
|
||||
@@ -70,10 +70,7 @@ java --enable-native-access=ALL-UNNAMED -jar agentik-0.1.0-all.jar
|
||||
|
||||
```bash
|
||||
# CLI
|
||||
java --enable-native-access=ALL-UNNAMED -jar agentik-cli-0.1.0-all.jar
|
||||
|
||||
# TUI
|
||||
java --enable-native-access=ALL-UNNAMED -jar agentik-tui-0.1.0-all.jar
|
||||
java --enable-native-access=ALL-UNNAMED -jar agentik-cli-0.1.0-SNAPSHOT-all.jar --help
|
||||
|
||||
# curl
|
||||
curl http://localhost:8080/health
|
||||
@@ -83,8 +80,7 @@ curl http://localhost:8080/health
|
||||
|
||||
- Запускаемые:
|
||||
- [`:standalone`](standalone/README.md) — single-jar HTTP-сервер.
|
||||
- [`:agentik-cli`](agentik-cli/README.md) — REPL-клиент (JLine).
|
||||
- [`:agentik-tui`](agentik-tui/README.md) — Compose-for-Mosaic TUI.
|
||||
- [`:agentik-cli`](agentik-cli/README.md) — one-shot CLI-клиент (kotlinx.cli), JVM + 4 native.
|
||||
- Библиотеки (контракты и реализации):
|
||||
- [`:proto`](proto/README.md) — stateful KMP-протокол.
|
||||
- [`:server`](server/README.md) — HTTP/SSE фасад `:proto`.
|
||||
|
||||
+4
-1
@@ -1,4 +1,4 @@
|
||||
package pw.binom.agentik.standalone.agent
|
||||
package pw.binom.agentik.toolsets
|
||||
|
||||
import pw.binom.litert.LiteTool
|
||||
|
||||
@@ -8,5 +8,8 @@ import pw.binom.litert.LiteTool
|
||||
* Имя используется как ключ для матчинга `LiteToolCall.name` (приходящего от LLM)
|
||||
* с конкретной реализацией тула. Для MCP-адаптеров имя имеет формат `server__tool`,
|
||||
* чтобы избежать коллизий между разными MCP-серверами.
|
||||
*
|
||||
* Перенесён из `:standalone/agent/NamedTool.kt` — это generic data-класс,
|
||||
* должен жить рядом с другими тулами в `:agent-toolsets`.
|
||||
*/
|
||||
data class NamedTool(val name: String, val tool: LiteTool)
|
||||
+146
-53
@@ -1,83 +1,176 @@
|
||||
# `:agentik-cli` — JVM CLI клиент к `/agentik`
|
||||
# `:agentik-cli` — one-shot CLI клиент к `/agentik`
|
||||
|
||||
## Что это
|
||||
|
||||
JVM-only REPL-клиент к серверу `:standalone` через `:client`
|
||||
над HTTP+SSE:
|
||||
**One-shot subcommand CLI** (Kotlin Multiplatform) к серверу
|
||||
`:standalone` через `:client` по HTTP+SSE. Один вызов — одна команда:
|
||||
стрим ответа `send` идёт в stdout построчно, никакого embedded-REPL.
|
||||
|
||||
- Нативный REPL с JLine (стрелки влево/вправо/вверх, история,
|
||||
Ctrl-D/E).
|
||||
- Подписка на live-стрим событий агента.
|
||||
- Slash-команды: `/new /list /switch /rename /rm /interrupt /history
|
||||
/pwd /help /exit /quit`.
|
||||
- Persistent session id в `~/.agentik/cli-state.json`.
|
||||
Решает: быстрый способ дёрнуть агента из shell-скрипта или руками,
|
||||
не поднимая отдельную TUI-сессии.
|
||||
|
||||
Решает: быстрый способ проверить агента руками из терминала.
|
||||
Используется в CI-смоук-тестах и для daily-driver.
|
||||
## Платформы
|
||||
|
||||
| Платформа | Артефакт | Размер | Статус |
|
||||
|---|---|---|---|
|
||||
| `jvm` (JRE 21) | `*-all.jar` | ~7 МБ | ✓ собирается и работает |
|
||||
| `linuxX64` | `.kexe` | ~5 МБ | ✓ собирается и работает на этом хосте |
|
||||
| `macosX64` | `.kexe` | — | собирается на macOS-раннере |
|
||||
| `macosArm64` | `.kexe` | — | собирается на macOS-arm64-раннере |
|
||||
| `mingwX64` | `.exe` | ~6 МБ | ✓ собирается (cross-compile с Linux) |
|
||||
| `linuxArm64` | — | — | **нет** — kotlinx-cli 0.3.6 не публикует klib для linuxArm64 |
|
||||
| `iOS` | — | — | нет смысла на iOS |
|
||||
|
||||
## Подкоманды
|
||||
|
||||
```
|
||||
agentik-cli <command> [args...]
|
||||
|
||||
Команды верхнего уровня:
|
||||
conv <subcommand> операции над диалогами (см. ниже)
|
||||
msgs <id> [--limit N] показать сообщения
|
||||
send <id> <text...> отправить ход, стримит response-события в stdout
|
||||
interrupt <id> прервать текущий ход
|
||||
info показать конфиг (server URL + agent id)
|
||||
|
||||
Подкоманды `conv`:
|
||||
conv ls список диалогов
|
||||
conv new [--temp] создать диалог, печатает id
|
||||
conv show <id> метаданные диалога
|
||||
conv delete <id> удалить диалог
|
||||
conv rename <id> <title> переименовать
|
||||
```
|
||||
|
||||
`--server URL` и `--id ID` (env: `AGENTIK_SERVER`, `AGENTIK_AGENT_ID`)
|
||||
задаются **после** имени subcommand'а — kotlinx.cli не шарит опции
|
||||
родителя в subcommand. Примеры:
|
||||
|
||||
```bash
|
||||
agentik-cli conv ls --server http://192.168.76.166:8080/agentik
|
||||
agentik-cli conv new --server http://localhost:8080/agentik
|
||||
agentik-cli send --server http://localhost:8080/agentik conv-abc "привет"
|
||||
agentik-cli info # через AGENTIK_SERVER env-переменную
|
||||
```
|
||||
|
||||
## Как запустить
|
||||
|
||||
### Требования
|
||||
|
||||
- JVM 21+ (на машине должна быть JAVA_HOME или `java` в PATH).
|
||||
- Запущенный `:standalone` (по умолчанию `http://localhost:8080/agentik`).
|
||||
|
||||
### Запуск из готового fatjar
|
||||
### JVM (fatjar)
|
||||
|
||||
```bash
|
||||
java --enable-native-access=ALL-UNNAMED -jar agentik-cli-0.1.0-all.jar \
|
||||
--server http://192.168.76.166:8080/agentik
|
||||
./gradlew :agentik-cli:shadowJar
|
||||
java --enable-native-access=ALL-UNNAMED \
|
||||
-jar agentik-cli/build/libs/agentik-cli-0.1.0-SNAPSHOT-all.jar conv --help
|
||||
```
|
||||
|
||||
### Запуск через Gradle (dev)
|
||||
### Native linuxX64
|
||||
|
||||
```bash
|
||||
./gradlew :agentik-cli:run --args="--server http://localhost:8080/agentik"
|
||||
./gradlew :agentik-cli:linkReleaseExecutableLinuxX64
|
||||
./agentik-cli/build/bin/linuxX64/releaseExecutable/agentik-cli.kexe conv --help
|
||||
```
|
||||
|
||||
## Параметры CLI
|
||||
### Native macOS / Windows
|
||||
|
||||
| Флаг | ENV | Что делает |
|
||||
|---|---|---|
|
||||
| `--server URL` | `AGENTIK_SERVER` | URL `/agentik` (default `http://localhost:8080/agentik`) |
|
||||
| `--id ID` | `USER`/`USERNAME` | Имя агента (default — текущий пользователь) |
|
||||
| `--no-history` | — | Не восстанавливать последнюю диалог после запуска |
|
||||
| `--help` | — | Показывает help и выходит |
|
||||
На Linux-хосте `macosX64`/`macosArm64` линкуются пустыми (нужен
|
||||
macOS-раннер, Apple Mach-O формат). `mingwX64` собирается через
|
||||
кросс-компиляцию.
|
||||
|
||||
## Slash-команды (внутри REPL)
|
||||
CI-ноут: запускать `./gradlew :agentik-cli:linkReleaseExecutableMacosX64
|
||||
:agentik-cli:linkReleaseExecutableMacosArm64` на `macos-latest`
|
||||
раннере Gitea Actions.
|
||||
|
||||
| Команда | Синонимы | Что делает |
|
||||
|---|---|---|
|
||||
| `/help` | | Показывает help |
|
||||
| `/new [title]` | | Создать диалог |
|
||||
| `/list` | `/ls` | Список диалогов |
|
||||
| `/switch <id>` | `/sw`, `/cd` | Переключиться на диалог |
|
||||
| `/rename <title>` | | Переименовать текущий диалог |
|
||||
| `/rm [id]` | `/delete` | Удалить (текущий или по id) |
|
||||
| `/interrupt` | `/stop`, `/cancel` | Прервать текущий ход |
|
||||
| `/history` | `/h`, `/hist` | Показывает историю текущего диалога |
|
||||
| `/pwd` | | Путь к state-file |
|
||||
| `/exit`, `/quit` | | Выйти |
|
||||
## Примеры
|
||||
|
||||
## Переменные среды (пробрасываются серверу через `--server`)
|
||||
```bash
|
||||
# Список диалогов (таблица)
|
||||
agentik-cli conv ls --server http://localhost:8080/agentik
|
||||
|
||||
См. [`../standalone/README.md`](../standalone/README.md). На стороне
|
||||
клиента они **не** интерпретируются — это лишь настройки запуска
|
||||
агента. CLI только знает, по какому URL стучаться.
|
||||
# Создать диалог
|
||||
ID=$(agentik-cli conv new --server http://localhost:8080/agentik)
|
||||
echo "new conv: $ID"
|
||||
|
||||
## Известное ограничение
|
||||
# Переименовать
|
||||
agentik-cli conv rename --server http://localhost:8080/agentik "$ID" "мой чат"
|
||||
|
||||
SSE event-stream в не-TTY ssh-сессии (без `-tt`) закрывается на
|
||||
default-таймауте Ktor. Используйте либо `ssh -tt`, либо нативный
|
||||
terminal. Это upstream-особенность Ktor SSE.
|
||||
# Отправить ход и стримить ответ
|
||||
agentik-cli send --server http://localhost:8080/agentik "$ID" "2+2"
|
||||
|
||||
# Показать последние N сообщений
|
||||
agentik-cli msgs --server http://localhost:8080/agentik "$ID" --limit 10
|
||||
|
||||
# Прервать активный ход
|
||||
agentik-cli interrupt --server http://localhost:8080/agentik "$ID"
|
||||
|
||||
# Удалить
|
||||
agentik-cli conv delete --server http://localhost:8080/agentik "$ID"
|
||||
|
||||
# Через env-переменную
|
||||
AGENTIK_SERVER=http://localhost:8080/agentik agentik-cli info
|
||||
```
|
||||
|
||||
## Формат вывода `send`
|
||||
|
||||
Каждое SSE-событие печатается отдельной строкой `event <Type> ...` —
|
||||
пригодно для парсинга через `awk`/`jq`-обёртки:
|
||||
|
||||
```
|
||||
event StartReasoning
|
||||
event StartResponse TEXT
|
||||
event AppendText \n\n
|
||||
event AppendText Привет!
|
||||
event End
|
||||
```
|
||||
|
||||
Терминальные события (`End`, `Interrupted`, `Error`) тоже
|
||||
печатаются; CLI выходит сразу после `End`.
|
||||
|
||||
## Почему kotlinx.cli (а не clikt)
|
||||
|
||||
- **kotlinx.cli 0.3.6** (JetBrains, KMP) — единственный зрелый
|
||||
arg-parser, который стабильно линкуется под `linux_x64` +
|
||||
`macos_x64`/`macos_arm64` + `mingw_x64`. Минус: нет `linux_arm64`.
|
||||
- **clikt-multiplatform 5.x** (ajalt) — имеет linuxArm64, но
|
||||
ломается на native linker: `duplicate symbol selfAndAncestors`
|
||||
между `clikt` и `clikt-mordant` commonMain (issue
|
||||
[ajalt/clikt#598](https://github.com/ajalt/clikt/issues/598)).
|
||||
Workaround `kotlin.native.cacheKind.linuxX64=none` замедляет
|
||||
сборку на порядки и не решает проблему до конца. Поэтому clikt
|
||||
отвергнут.
|
||||
|
||||
## Платформенные детали
|
||||
|
||||
- **entryPoint на K/N** — это FQN функции **без** суффикса `Kt`
|
||||
(т.е. `pw.binom.agentik.cli.main`, а не `MainKt.main`). JVM
|
||||
convention `MainKt.main` тут не работает — K/N линкер ищет
|
||||
функцию по `package.main`.
|
||||
- **`platformEnv(key)`** для чтения env-переменных:
|
||||
- JVM: `System.getenv(key)` через `jvmMain` actual.
|
||||
- Native: `getenv(key)` из `platform.posix` через
|
||||
`kotlinx.cinterop.toKString()` (`nativeMain` actual,
|
||||
требует `@OptIn(ExperimentalForeignApi::class)`).
|
||||
- **Stdout / exit code** — работают на K/N через корутины.
|
||||
|
||||
## Готчасы kotlinx.cli
|
||||
|
||||
- **Вложенные subcommands + parent.execute().** В kotlinx.cli 0.3.6
|
||||
`parent.execute()` вызывается ПОСЛЕ `leaf.execute()` всегда,
|
||||
когда leaf был достигнут через parent. Поэтому `ConvCommand.execute()`
|
||||
сделан no-op (`override fun execute() = Unit`), иначе вывод
|
||||
дочерней команды дублируется выводом родителя. Дочерние команды
|
||||
смотрятся через `agentik-cli conv --help`.
|
||||
- **strictSubcommandOptionsOrder.** Без этого флага `conv new --server ...`
|
||||
парсится как `conv [--server ...]` + позиционный аргумент `new`
|
||||
на уровне родителя — и дочерняя команда не запускается.
|
||||
В `ArgParser` сразу включается `strictSubcommandOptionsOrder = true`.
|
||||
|
||||
## Тесты
|
||||
|
||||
```
|
||||
./gradlew :agentik-cli:jvmTest
|
||||
```
|
||||
Тесты для подкоманд пока не написаны (TODO). Базовый smoke
|
||||
покрывается руками против живого сервера.
|
||||
|
||||
23 теста: парсер slash-команд, event-рендер, state-repository.
|
||||
```bash
|
||||
./gradlew :agentik-cli:jvmTest # 0/0 — пока пусто
|
||||
```
|
||||
|
||||
## Версии
|
||||
|
||||
|
||||
@@ -1,7 +1,6 @@
|
||||
import org.jetbrains.kotlin.gradle.ExperimentalKotlinGradlePluginApi
|
||||
|
||||
import com.github.jengelman.gradle.plugins.shadow.tasks.ShadowJar
|
||||
import org.gradle.api.artifacts.ConfigurationContainer
|
||||
|
||||
plugins {
|
||||
alias(libs.plugins.kotlin.multiplatform)
|
||||
@@ -12,67 +11,68 @@ plugins {
|
||||
kotlin {
|
||||
jvmToolchain(21)
|
||||
|
||||
// Suppress Beta-предупреждения от expect/actual объектов — фича стабильна с Kotlin 1.9,
|
||||
// но компилятор всё ещё требует -Xexpect-actual-classes, чтобы не ныть.
|
||||
compilerOptions {
|
||||
freeCompilerArgs.add("-Xexpect-actual-classes")
|
||||
}
|
||||
|
||||
// "Все возможные цели сборки": jvm + весь натив. Зеркалит набор :server/:proto.
|
||||
// commonMain зависит только от :proto (KMP). jvmMain подключает :client (JVM-only)
|
||||
// и JLine — там же и `:client`'s AgentClient. nativeMain пока получает stub actual,
|
||||
// расширять будем через ktor-client-* {curl,darwin,winhttp} когда дойдёт очередь.
|
||||
// Native-таргеты, которые покрывает kotlinx.cli 0.3.6 (см. его .module
|
||||
// в Maven Central): linux_x64, macos_x64, macos_arm64, mingw_x64.
|
||||
// linuxArm64 не входит — kotlinx.cli 0.3.6 для него не публикуется
|
||||
// (последний релиз 2023-09, KMP-targets зафиксированы). clikt-multiplatform
|
||||
// 5.x имеет linuxArm64, но ломается на duplicate symbol `selfAndAncestors`
|
||||
// между clikt и clikt-mordant при линковке native (issue ajalt/clikt#598),
|
||||
// поэтому clikt отвергнут.
|
||||
//
|
||||
// iOS не входит: :agentik-cli бессмыслен на iOS, а :client (единственный
|
||||
// его потребитель) тоже без iOS.
|
||||
jvm()
|
||||
macosX64()
|
||||
macosArm64()
|
||||
iosX64()
|
||||
iosArm64()
|
||||
iosSimulatorArm64()
|
||||
linuxX64()
|
||||
linuxArm64()
|
||||
mingwX64()
|
||||
listOf(
|
||||
linuxX64(),
|
||||
macosX64(),
|
||||
macosArm64(),
|
||||
mingwX64(),
|
||||
)
|
||||
|
||||
sourceSets {
|
||||
commonMain.dependencies {
|
||||
implementation(project(":proto"))
|
||||
implementation(project(":client"))
|
||||
|
||||
// kotlinx.cli 0.3.6 — KMP subcommand-парсер от JetBrains.
|
||||
// clikt 5.x имеет upstream-баг: `duplicate symbol selfAndAncestors`
|
||||
// между `clikt` и `clikt-mordant` при линковке native. kotlinx.cli
|
||||
// таких проблем нет.
|
||||
implementation(libs.kotlinx.cli)
|
||||
|
||||
implementation(libs.kotlinx.coroutines.core)
|
||||
implementation(libs.kotlinx.serialization.core)
|
||||
implementation(libs.kotlinx.serialization.json)
|
||||
}
|
||||
jvmMain.dependencies {
|
||||
// :client JVM-only (ktor-cio). Подключаем только в jvmMain.
|
||||
implementation(project(":client"))
|
||||
// JLine для readline с историей и completion.
|
||||
implementation(libs.jline)
|
||||
}
|
||||
commonTest.dependencies {
|
||||
implementation(kotlin("test"))
|
||||
implementation(libs.kotlinx.coroutines.core)
|
||||
// runTest { } — suspend test runner для commonTest.
|
||||
implementation("org.jetbrains.kotlinx:kotlinx-coroutines-test:1.11.0")
|
||||
}
|
||||
jvmTest.dependencies {
|
||||
// JUnit нужен в jvmTest — kotlin-test на JVM = JUnit4.
|
||||
implementation("junit:junit:4.13.2")
|
||||
implementation(libs.ktor.client.cio)
|
||||
}
|
||||
// :agentik-cli — commonMain-only (нет jvmMain/nativeMain разделения):
|
||||
// весь код, включая platformEnv, лежит в commonMain.
|
||||
}
|
||||
|
||||
@OptIn(ExperimentalKotlinGradlePluginApi::class)
|
||||
jvm {
|
||||
binaries {
|
||||
executable {
|
||||
mainClass.set("pw.binom.agentik.cli.MainKt")
|
||||
mainClass.set("pw.binom.agentik.cli.AgentikCliKt")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// entryPoint на K/N — это FQN функции БЕЗ суффикса `Kt`
|
||||
// (Java/Kotlin convention `MainKt.main` тут не работает, линкер K/N ищет
|
||||
// функцию как `package.main`). На JVM суффикс `Kt` сохраняется через
|
||||
// mainClass.set(...) выше.
|
||||
listOf(
|
||||
linuxX64(),
|
||||
macosX64(),
|
||||
macosArm64(),
|
||||
mingwX64(),
|
||||
).forEach {
|
||||
it.binaries.executable {
|
||||
entryPoint = "pw.binom.agentik.cli.main"
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// --- Fatjar (uberjar) ---
|
||||
//
|
||||
// По аналогии с :standalone: shadowJar берёт `jvmJar` + `jvmRuntimeClasspath`.
|
||||
// Shadow 8.x не авторегистрирует shadowJar в KMP-проектах — нужно явно register.
|
||||
|
||||
// Fatjar — аналог :standalone.
|
||||
val shadowJarTask = tasks.register<ShadowJar>("shadowJar") {
|
||||
archiveBaseName.set("agentik-cli")
|
||||
archiveClassifier.set("all")
|
||||
@@ -80,20 +80,13 @@ val shadowJarTask = tasks.register<ShadowJar>("shadowJar") {
|
||||
group = "build"
|
||||
|
||||
from(tasks.named("jvmJar"))
|
||||
val cc = try {
|
||||
@Suppress("UNCHECKED_CAST")
|
||||
configurations as org.gradle.api.artifacts.ConfigurationContainer
|
||||
} catch (_: ClassCastException) {
|
||||
@Suppress("UNCHECKED_CAST")
|
||||
(project as org.gradle.api.Project).configurations as org.gradle.api.artifacts.ConfigurationContainer
|
||||
}
|
||||
from(cc.getByName("jvmRuntimeClasspath"))
|
||||
from(project.configurations.getByName("jvmRuntimeClasspath"))
|
||||
|
||||
mergeServiceFiles()
|
||||
duplicatesStrategy = DuplicatesStrategy.EXCLUDE
|
||||
|
||||
manifest {
|
||||
attributes["Main-Class"] = "pw.binom.agentik.cli.MainKt"
|
||||
attributes["Main-Class"] = "pw.binom.agentik.cli.AgentikCliKt"
|
||||
attributes["Implementation-Title"] = "agentik-cli"
|
||||
attributes["Implementation-Version"] = project.version.toString()
|
||||
}
|
||||
|
||||
@@ -1,323 +1,87 @@
|
||||
package pw.binom.agentik.cli
|
||||
|
||||
import kotlinx.coroutines.CompletableDeferred
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.cancel
|
||||
import kotlinx.coroutines.flow.first
|
||||
import kotlinx.coroutines.isActive
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import pw.binom.agentik.proto.Agent
|
||||
import pw.binom.agentik.proto.Content
|
||||
import pw.binom.agentik.proto.Conversation
|
||||
import pw.binom.agentik.proto.Event
|
||||
import pw.binom.agentik.proto.Message
|
||||
import kotlin.time.Instant
|
||||
import kotlinx.cli.ArgParser
|
||||
import kotlinx.cli.ArgType
|
||||
import kotlinx.cli.ExperimentalCli
|
||||
import kotlinx.cli.Subcommand
|
||||
import kotlinx.cli.default
|
||||
import pw.binom.agentik.cli.commands.ConvCommand
|
||||
import pw.binom.agentik.cli.commands.InfoSubcommand
|
||||
import pw.binom.agentik.cli.commands.InterruptSubcommand
|
||||
import pw.binom.agentik.cli.commands.MsgsSubcommand
|
||||
import pw.binom.agentik.cli.commands.SendSubcommand
|
||||
|
||||
/**
|
||||
* Главный класс REPL.
|
||||
* Default `server` URL: env `AGENTIK_SERVER` или `http://localhost:8080/agentik`.
|
||||
* Default `agent id`: env `AGENTIK_AGENT_ID` или `cli`.
|
||||
*
|
||||
* Управляет:
|
||||
* - текущим диалогом ([currentConv]) + позицией в его event-stream ([lastEventAt]);
|
||||
* - фоновым job'ом, слушающим events и рендерящим их через [EventRenderer].
|
||||
* - персистентностью сессии (восстановление последнего диалога при перезапуске CLI).
|
||||
*
|
||||
* Один ход = один заход в REPL: пока идёт turn, REPL ждёт его завершения.
|
||||
* `/interrupt` стучится в [Conversation.interrupt] — фоновый подписчик событий
|
||||
* увидит [Event.Interrupted] и сам завершится.
|
||||
* Используется в `runAgentikCli` и в каждом subcommand'е для своего
|
||||
* `--server`/`--id` (иначе subcommand не видит значения родителя).
|
||||
*/
|
||||
class AgentikCli internal constructor(private val config: CliConfig) {
|
||||
internal fun defaultServerUrl(): String = platformEnv("AGENTIK_SERVER") ?: "http://localhost:8080/agentik"
|
||||
internal fun defaultAgentId(): String = platformEnv("AGENTIK_AGENT_ID") ?: "cli"
|
||||
|
||||
private val agent: Agent = CliPlatform.openAgent(baseUrl = config.server, id = config.id)
|
||||
private val terminal: CliTerminal = CliPlatform.openTerminal(
|
||||
historyFile = if (config.historyEnabled) stateFilePath() else null,
|
||||
prompt = "agentik> ",
|
||||
)
|
||||
private val sessionRepo = SessionRepository(
|
||||
filePath = if (config.historyEnabled) stateFilePath() else null,
|
||||
io = CliPlatform.sessionIo(),
|
||||
/**
|
||||
* Корневой [ArgParser] `agentik-cli`. Один вызов — одна команда.
|
||||
*
|
||||
* ```
|
||||
* agentik-cli <command> [args...]
|
||||
*
|
||||
* Commands:
|
||||
* conv ls|new|show|delete|rename операции над диалогами
|
||||
* msgs <id> [--limit N] показать сообщения
|
||||
* send <id> <text...> отправить ход, стримит response-события в stdout
|
||||
* interrupt <id> прервать текущий ход
|
||||
* info показать конфиг
|
||||
*
|
||||
* `--server` и `--id` задаются ПОСЛЕ имени subcommand'а (т.е.
|
||||
* `agentik-cli conv ls --server http://...`), не до — kotlinx.cli не
|
||||
* шарит опции родителя в subcommand.
|
||||
*
|
||||
* Вложенные subcommands (`conv ls`, `conv new`, ...) реализованы
|
||||
* через [Subcommand.subcommands]: `conv` сам — subcommand, и его
|
||||
* дочерние команды (`ls`, `new`, `show`, `delete`, `rename`)
|
||||
* регистрируются у него.
|
||||
*/
|
||||
@OptIn(ExperimentalCli::class)
|
||||
fun runAgentikCli(args: Array<String>) {
|
||||
val parser = ArgParser(
|
||||
programName = "agentik-cli",
|
||||
// Все аргументы после имени subcommand должны передаваться
|
||||
// В subcommand-парсер, а не парситься на уровне родителя.
|
||||
// Без этого `conv new --server ...` парсится как `conv [--server ...]`
|
||||
// + аргумент "new" → execute родителя, без вложенной команды.
|
||||
strictSubcommandOptionsOrder = true,
|
||||
)
|
||||
|
||||
private var currentConv: Conversation? = null
|
||||
private var currentTitle: String? = null
|
||||
private var lastEventAt: Instant = Instant.DISTANT_PAST
|
||||
val conv = ConvCommand()
|
||||
parser.subcommands(
|
||||
conv,
|
||||
MsgsSubcommand(),
|
||||
SendSubcommand(),
|
||||
InterruptSubcommand(),
|
||||
InfoSubcommand(),
|
||||
)
|
||||
|
||||
private val scope = CoroutineScope(Dispatchers.Default)
|
||||
|
||||
suspend fun run() {
|
||||
try {
|
||||
// Восстановление сессии.
|
||||
val saved = sessionRepo.load()
|
||||
if (saved != null) {
|
||||
val conv = runCatching { agent.getConversation(saved.conversationId) }
|
||||
.getOrNull()
|
||||
if (conv != null) {
|
||||
currentConv = conv
|
||||
currentTitle = conv.title
|
||||
lastEventAt = saved.lastEventAt
|
||||
terminal.printSystem(
|
||||
"восстановлен диалог ${shorten(conv.id)}" +
|
||||
" (${conv.title ?: "без названия"})",
|
||||
)
|
||||
} else {
|
||||
terminal.printSystem(
|
||||
"прошлый диалог ${shorten(saved.conversationId)} больше не существует",
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
printBanner()
|
||||
|
||||
// Главный цикл.
|
||||
while (scope.isActive) {
|
||||
terminal.print(prompt())
|
||||
val line = terminal.readLine() ?: break // EOF → выходим
|
||||
val trimmed = line.trim()
|
||||
if (trimmed.isEmpty()) continue
|
||||
|
||||
if (trimmed.startsWith("/")) {
|
||||
when (val r = parseSlash(trimmed.substring(1))) {
|
||||
is ParseResult.Success -> {
|
||||
if (handleCommand(r.command) == CommandResult.Exit) break
|
||||
}
|
||||
is ParseResult.Failure -> terminal.printSystem(r.message)
|
||||
}
|
||||
} else {
|
||||
handleUserMessage(trimmed)
|
||||
}
|
||||
}
|
||||
} finally {
|
||||
terminal.printSystem("до свидания.")
|
||||
currentConv?.close()
|
||||
terminal.close()
|
||||
sessionRepo.close()
|
||||
scope.cancel()
|
||||
}
|
||||
}
|
||||
|
||||
// ============================================================ banner / prompt
|
||||
|
||||
private suspend fun printBanner() {
|
||||
terminal.println()
|
||||
terminal.println("agentik-cli — id=${config.id} — type /help")
|
||||
terminal.println("server: ${config.server}")
|
||||
when (val c = currentConv) {
|
||||
null -> terminal.println("диалог: не выбран — начните с /new или /switch <id>")
|
||||
else -> terminal.println("диалог: ${shorten(c.id)} (${c.title ?: "без названия"})")
|
||||
}
|
||||
terminal.println()
|
||||
}
|
||||
|
||||
private fun prompt(): String = "agentik${if (currentConv != null) "" else " (-)"}> "
|
||||
|
||||
private suspend fun printHelp() {
|
||||
terminal.println(
|
||||
"""
|
||||
|Slash-команды:
|
||||
| /help эта справка
|
||||
| /new [title] создать новый диалог
|
||||
| /list, /ls список диалогов (новые сверху)
|
||||
| /switch <id>, /sw переключиться на диалог по id
|
||||
| /rename <title> переименовать текущий диалог
|
||||
| /delete [<id>], /rm удалить диалог (по id или текущий)
|
||||
| /history, /h последние сообщения текущего диалога
|
||||
| /interrupt, /stop прервать текущий ход
|
||||
| /pwd показать текущий диалог
|
||||
| /exit, /quit выйти (Ctrl-D тоже)
|
||||
|
|
||||
|Любой ввод без ведущего `/` отправляется агенту в текущий диалог.
|
||||
""".trimMargin(),
|
||||
)
|
||||
}
|
||||
|
||||
// ============================================================ command dispatch
|
||||
|
||||
private suspend fun handleCommand(cmd: SlashCommand): CommandResult = when (cmd) {
|
||||
SlashCommand.Help -> { printHelp(); CommandResult.Continue }
|
||||
SlashCommand.Exit, SlashCommand.Quit -> CommandResult.Exit
|
||||
is SlashCommand.New -> { handleNew(cmd.title); CommandResult.Continue }
|
||||
SlashCommand.List -> { handleList(); CommandResult.Continue }
|
||||
is SlashCommand.Switch -> { handleSwitch(cmd.id); CommandResult.Continue }
|
||||
is SlashCommand.Rename -> { handleRename(cmd.title); CommandResult.Continue }
|
||||
is SlashCommand.Delete -> { handleDelete(cmd.id); CommandResult.Continue }
|
||||
SlashCommand.Interrupt -> { handleInterrupt(); CommandResult.Continue }
|
||||
SlashCommand.History -> { handleHistory(); CommandResult.Continue }
|
||||
SlashCommand.Pwd -> { handlePwd(); CommandResult.Continue }
|
||||
}
|
||||
|
||||
private suspend fun handleNew(title: String?) {
|
||||
val conv = agent.createConversation(temp = false)
|
||||
if (title != null) conv.rename(title)
|
||||
currentConv = conv
|
||||
currentTitle = title ?: conv.title
|
||||
lastEventAt = Instant.DISTANT_PAST
|
||||
terminal.printSystem("создан диалог ${shorten(conv.id)}" + if (title != null) " — «$title»" else "")
|
||||
sessionRepo.save(conv.id, lastEventAt)
|
||||
}
|
||||
|
||||
private suspend fun handleList() {
|
||||
terminal.println("диалоги (новые сверху):")
|
||||
agent.getConversations(offset = 0).collect { conv ->
|
||||
val marker = if (conv.id == currentConv?.id) "*" else " "
|
||||
val title = conv.title ?: "(без названия)"
|
||||
terminal.println(" $marker ${shorten(conv.id)} $title [${conv.updatedAt}]")
|
||||
}
|
||||
}
|
||||
|
||||
private suspend fun handleSwitch(id: String) {
|
||||
val conv = agent.getConversation(id)
|
||||
if (conv == null) {
|
||||
terminal.printSystem("диалог $id не найден")
|
||||
return
|
||||
}
|
||||
currentConv?.close()
|
||||
currentConv = conv
|
||||
currentTitle = conv.title
|
||||
lastEventAt = Instant.DISTANT_PAST
|
||||
sessionRepo.save(conv.id, lastEventAt)
|
||||
terminal.printSystem("переключились на ${shorten(conv.id)} (${conv.title ?: "без названия"})")
|
||||
}
|
||||
|
||||
private suspend fun handleRename(title: String) {
|
||||
val c = currentConv ?: run {
|
||||
terminal.printSystem("нет активного диалога — /new")
|
||||
return
|
||||
}
|
||||
c.rename(title)
|
||||
currentTitle = title
|
||||
terminal.printSystem("заголовок: $title")
|
||||
}
|
||||
|
||||
private suspend fun handleDelete(id: String?) {
|
||||
val target = id ?: currentConv?.id
|
||||
if (target == null) {
|
||||
terminal.printSystem("нет диалога для удаления")
|
||||
return
|
||||
}
|
||||
val ok = agent.deleteConversation(target)
|
||||
if (ok) {
|
||||
terminal.printSystem("удалён ${shorten(target)}")
|
||||
if (target == currentConv?.id) {
|
||||
currentConv?.close()
|
||||
currentConv = null
|
||||
currentTitle = null
|
||||
sessionRepo.clear()
|
||||
}
|
||||
} else {
|
||||
terminal.printSystem("диалог ${shorten(target)} не найден")
|
||||
}
|
||||
}
|
||||
|
||||
private suspend fun handleInterrupt() {
|
||||
val c = currentConv ?: run {
|
||||
terminal.printSystem("нет активного диалога")
|
||||
return
|
||||
}
|
||||
c.interrupt()
|
||||
terminal.printSystem("прерывание отправлено")
|
||||
}
|
||||
|
||||
private suspend fun handlePwd() {
|
||||
val c = currentConv ?: run {
|
||||
terminal.printSystem("диалог: не выбран")
|
||||
return
|
||||
}
|
||||
terminal.printSystem("id: ${c.id}")
|
||||
terminal.printSystem("title: ${c.title ?: "—"}")
|
||||
terminal.printSystem("updatedAt: ${c.updatedAt}")
|
||||
terminal.printSystem("temporal: ${c.isTemporal}")
|
||||
}
|
||||
|
||||
private suspend fun handleHistory() {
|
||||
val c = currentConv ?: run {
|
||||
terminal.printSystem("нет активного диалога")
|
||||
return
|
||||
}
|
||||
terminal.println("история:")
|
||||
c.getMessages(after = Instant.DISTANT_PAST).collect { msg -> renderHistoryMessage(msg) }
|
||||
}
|
||||
|
||||
private suspend fun renderHistoryMessage(msg: Message) {
|
||||
val prefix = " [${msg.date}] "
|
||||
when (msg) {
|
||||
is Message.UserMessage ->
|
||||
terminal.println(prefix + "user | " + msg.content.text())
|
||||
is Message.AssistantMessage ->
|
||||
terminal.println(prefix + "agent | " + msg.content.text())
|
||||
is Message.ToolCall ->
|
||||
terminal.println(prefix + "tool>${msg.toolName} | ${msg.toolArgs.take(160)}")
|
||||
is Message.ToolResult ->
|
||||
terminal.println(prefix + "tool< | " + (msg.result?.take(160) ?: "null"))
|
||||
is Message.Error ->
|
||||
terminal.println(prefix + "<error${msg.code?.let { "/$it" } ?: ""}> ${msg.message}")
|
||||
}
|
||||
}
|
||||
|
||||
private fun List<Content>.text(): String =
|
||||
joinToString(separator = "") { c ->
|
||||
when (c) {
|
||||
is Content.Text -> c.body
|
||||
is Content.Image -> "[image:${c.mime}:${c.data.size}B]"
|
||||
}
|
||||
}
|
||||
|
||||
// ============================================================ user-message
|
||||
|
||||
private suspend fun handleUserMessage(text: String) {
|
||||
val conv = currentConv ?: run {
|
||||
terminal.printSystem("нет активного диалога — /new")
|
||||
return
|
||||
}
|
||||
|
||||
terminal.println() // пустая строка для визуального отделения блока
|
||||
|
||||
val renderer = EventRenderer(terminal)
|
||||
val turnFinished = CompletableDeferred<Unit>()
|
||||
|
||||
// Подписчик events: принимает события и обновляет lastEventAt,
|
||||
// по терминальному событию закрывает Deferred.
|
||||
val eventsJob = scope.launch {
|
||||
try {
|
||||
conv.events(after = lastEventAt).collect { ev ->
|
||||
renderer.render(ev)
|
||||
if (ev.date > lastEventAt) {
|
||||
lastEventAt = ev.date
|
||||
sessionRepo.save(conv.id, lastEventAt)
|
||||
}
|
||||
if (ev is Event.End || ev is Event.Interrupted || ev is Event.Error) {
|
||||
if (!turnFinished.isCompleted) turnFinished.complete(Unit)
|
||||
}
|
||||
}
|
||||
} catch (t: Throwable) {
|
||||
if (!turnFinished.isCompleted) turnFinished.complete(Unit)
|
||||
if (t !is kotlinx.coroutines.CancellationException) {
|
||||
terminal.printSystem("[events stream error] ${t.message}")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
try {
|
||||
conv.send(listOf(Content.Text(text)))
|
||||
turnFinished.await()
|
||||
} catch (t: Throwable) {
|
||||
terminal.printSystem("[send error] ${t.message}")
|
||||
} finally {
|
||||
eventsJob.cancel()
|
||||
renderer.close()
|
||||
terminal.println()
|
||||
}
|
||||
}
|
||||
|
||||
// ============================================================ utils
|
||||
|
||||
private fun shorten(id: String): String = id.take(8)
|
||||
|
||||
private fun stateFilePath(): String? {
|
||||
val home = CliPlatform.homeDir() ?: return null
|
||||
val dir = "$home/.agentik"
|
||||
return "$dir/cli-state.json"
|
||||
}
|
||||
parser.parse(args)
|
||||
}
|
||||
|
||||
private enum class CommandResult { Continue, Exit }
|
||||
/**
|
||||
* Базовый класс subcommand'а: каждый subcommand владеет своим `--server`/`--id`,
|
||||
* чтобы значения родительских флагов были ему доступны (kotlinx.cli не шарит
|
||||
* свойства родителя в subcommand).
|
||||
*/
|
||||
abstract class AgentikSubcommand(name: String, description: String) : Subcommand(name, description) {
|
||||
val serverUrl: String by option(
|
||||
ArgType.String, fullName = "server", shortName = "s",
|
||||
description = "Base URL агента (env AGENTIK_SERVER)",
|
||||
).default(defaultServerUrl())
|
||||
val agentId: String by option(
|
||||
ArgType.String, fullName = "id", shortName = "i",
|
||||
description = "Идентификатор агента (env AGENTIK_AGENT_ID)",
|
||||
).default(defaultAgentId())
|
||||
}
|
||||
|
||||
fun main(args: Array<String>) {
|
||||
runAgentikCli(args)
|
||||
}
|
||||
|
||||
@@ -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 }
|
||||
}
|
||||
@@ -1,52 +0,0 @@
|
||||
package pw.binom.agentik.cli
|
||||
|
||||
import pw.binom.agentik.proto.Agent
|
||||
|
||||
/**
|
||||
* Платформенные зависимости CLI. Все вещи, требующие JVM-stdlib или
|
||||
* нативных API (терминал, env, файловое IO для state-файла, HTTP-клиент),
|
||||
* предоставляются здесь как `expect/actual`.
|
||||
*
|
||||
* Текущий статус: jvmMain полностью реализован (JLine + `java.io` + `:client`),
|
||||
* nativeMain — заглушки (подключение native ktor-движков и termios — отдельная задача).
|
||||
*/
|
||||
expect object CliPlatform {
|
||||
fun openAgent(baseUrl: String, id: String): Agent
|
||||
|
||||
fun openTerminal(
|
||||
historyFile: String?,
|
||||
prompt: String,
|
||||
): CliTerminal
|
||||
|
||||
/** HOME/USERPROFILE для пути пути state-файла; null если недоступна. */
|
||||
fun homeDir(): String?
|
||||
|
||||
/** Переменная среды (native API). Для jvmMain — `System.getenv`. */
|
||||
fun env(key: String): String?
|
||||
|
||||
/** Файловое IO для session-state; nativeMain возвращает no-op. */
|
||||
fun sessionIo(): SessionIo
|
||||
}
|
||||
|
||||
/**
|
||||
* Абстракция терминала, нужная для REPL. suspend-методы, чтобы не блокировать
|
||||
* event-loop агентного цикла во время ожидания ввода.
|
||||
*/
|
||||
interface CliTerminal {
|
||||
val prompt: String
|
||||
|
||||
/** Следующая строка пользователя (без prompt). null = EOF (Ctrl-D/Ctrl-Z). */
|
||||
suspend fun readLine(): String?
|
||||
|
||||
/** Печатает строку + перевод строки. */
|
||||
suspend fun println(text: String = "")
|
||||
|
||||
/** Печатает строку без перевода (для streamed chunks). */
|
||||
suspend fun print(text: String)
|
||||
|
||||
/** Подсветить prompt (символы-разделители сообщений, системные баннеры и т.п.). */
|
||||
suspend fun printSystem(text: String)
|
||||
|
||||
/** Закрыть терминал: restore raw mode, flush history file, ... */
|
||||
fun close()
|
||||
}
|
||||
@@ -1,82 +0,0 @@
|
||||
package pw.binom.agentik.cli
|
||||
|
||||
import pw.binom.agentik.proto.Event
|
||||
|
||||
/**
|
||||
* Печатает [Event] в человеко-читаемом виде через [CliTerminal].
|
||||
*
|
||||
* Дизайн:
|
||||
* - [Event.StartReasoning] — просто системный маркер; текст мысли НЕ выводим
|
||||
* отдельным форматом (см. proto: reasonig текст идёт через [Event.AppendText]).
|
||||
* - [Event.StartResponse] с `responseType=TEXT` — начало печати ответа; закрытие
|
||||
* происходит при [Event.End] или [Event.Interrupted].
|
||||
* - [Event.AppendText] — кусок текста, печатается БЕЗ перевода строки (чанки).
|
||||
* - [Event.AppendImage] — выводим как `[image: <mime>, <bytes> bytes]` placeholder.
|
||||
* Реальный рендеринг сделаем позже через iTerm/Kitty протоколы.
|
||||
* - [Event.End] / [Event.Interrupted] — закрывают текущий блок.
|
||||
* - [Event.Error] — отдельный системный блок `[error: …]`.
|
||||
*/
|
||||
class EventRenderer(private val terminal: CliTerminal) {
|
||||
|
||||
/** Трекает открыт ли сейчас «блок ответа» (после [Event.StartResponse], до [Event.End]). */
|
||||
private var responseOpen = false
|
||||
|
||||
suspend fun render(event: Event) {
|
||||
when (event) {
|
||||
is Event.StartReasoning -> {
|
||||
terminal.printSystem("…thinking…")
|
||||
if (responseOpen) {
|
||||
terminal.println()
|
||||
responseOpen = false
|
||||
}
|
||||
}
|
||||
|
||||
is Event.StartResponse -> {
|
||||
if (responseOpen) terminal.println()
|
||||
responseOpen = true
|
||||
// Без префикса — текст будет стримиться дальше через AppendText.
|
||||
}
|
||||
|
||||
is Event.AppendText -> {
|
||||
terminal.print(event.body)
|
||||
}
|
||||
|
||||
is Event.AppendImage -> {
|
||||
terminal.print("[image:${event.mime}:${event.body.size} bytes]")
|
||||
}
|
||||
|
||||
is Event.Interrupted -> {
|
||||
if (responseOpen) {
|
||||
terminal.println()
|
||||
terminal.printSystem("[interrupted]")
|
||||
responseOpen = false
|
||||
} else {
|
||||
terminal.printSystem("[interrupted]")
|
||||
}
|
||||
}
|
||||
|
||||
is Event.End -> {
|
||||
if (responseOpen) {
|
||||
terminal.println()
|
||||
responseOpen = false
|
||||
}
|
||||
}
|
||||
|
||||
is Event.Error -> {
|
||||
terminal.println()
|
||||
terminal.printSystem("[error${event.code?.let { "/$it" } ?: ""}] ${event.message}")
|
||||
if (responseOpen) responseOpen = false
|
||||
}
|
||||
|
||||
else -> {
|
||||
// ToolCall/ToolResult — это «структура» диалога, в текстовом стриме
|
||||
// не показываем; в веб-UI будет по-другому.
|
||||
terminal.printSystem("[event:${event::class.simpleName}]")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fun close() {
|
||||
responseOpen = false
|
||||
}
|
||||
}
|
||||
@@ -1,105 +0,0 @@
|
||||
package pw.binom.agentik.cli
|
||||
|
||||
import kotlinx.coroutines.runBlocking
|
||||
|
||||
/**
|
||||
* Точка входа CLI. Поддерживает аргументы командной строки:
|
||||
*
|
||||
* ```
|
||||
* agentik-cli [--server URL] [--id ID] [--no-history] [--help]
|
||||
*
|
||||
* --server URL базовый URL сервера agentik (default $AGENTIK_SERVER или
|
||||
* http://localhost:8080/agentik)
|
||||
* --id ID идентификатор этого клиента (default "cli:$USER")
|
||||
* --no-history не сохранять состояние в ~/.agentik/cli-state.json
|
||||
* --help, -h распечатать usage и выйти
|
||||
* ```
|
||||
*
|
||||
* Без аргументов — стартует REPL.
|
||||
*/
|
||||
fun main(args: Array<String>) = runBlocking {
|
||||
val cfg = parseCliArgs(args)
|
||||
if (cfg == null) {
|
||||
printUsage()
|
||||
return@runBlocking
|
||||
}
|
||||
AgentikCli(cfg).run()
|
||||
}
|
||||
|
||||
/**
|
||||
* Конфигурация CLI, вычисленная из аргументов + переменных среды.
|
||||
* Доступна из других файлов commonMain (видна как `internal` внутри модуля).
|
||||
*/
|
||||
internal data class CliConfig(
|
||||
val server: String,
|
||||
val id: String,
|
||||
val historyEnabled: Boolean,
|
||||
)
|
||||
|
||||
private fun parseCliArgs(args: Array<String>): CliConfig? {
|
||||
var server: String? = null
|
||||
var id: String? = null
|
||||
var historyEnabled = true
|
||||
|
||||
var i = 0
|
||||
while (i < args.size) {
|
||||
when (val a = args[i]) {
|
||||
"--help", "-h", "help" -> return null
|
||||
"--server", "-s" -> {
|
||||
require(i + 1 < args.size) { "$a требует URL" }
|
||||
server = args[i + 1]; i += 2
|
||||
}
|
||||
"--id" -> {
|
||||
require(i + 1 < args.size) { "$a требует значение" }
|
||||
id = args[i + 1]; i += 2
|
||||
}
|
||||
"--no-history" -> { historyEnabled = false; i++ }
|
||||
"--" -> i++ // разделитель; остальное игнорируем
|
||||
else -> error("неизвестный аргумент: $a (введите --help)")
|
||||
}
|
||||
}
|
||||
|
||||
val resolvedServer = server
|
||||
?: CliPlatform.env("AGENTIK_SERVER")
|
||||
?: "http://localhost:8080/agentik"
|
||||
val resolvedId = id ?: "cli:${CliPlatform.env("USER") ?: CliPlatform.env("USERNAME") ?: "anon"}"
|
||||
|
||||
return CliConfig(
|
||||
server = resolvedServer,
|
||||
id = resolvedId,
|
||||
historyEnabled = historyEnabled,
|
||||
)
|
||||
}
|
||||
|
||||
private fun printUsage() {
|
||||
val defaultServer = CliPlatform.env("AGENTIK_SERVER") ?: "http://localhost:8080/agentik"
|
||||
val defaultUser = CliPlatform.env("USER") ?: CliPlatform.env("USERNAME") ?: "anon"
|
||||
println("""
|
||||
agentik-cli — REPL поверх протокола agentik
|
||||
|
||||
Использование:
|
||||
agentik-cli [--server URL] [--id ID] [--no-history]
|
||||
|
||||
Аргументы:
|
||||
--server, -s URL базовый URL (default: $defaultServer)
|
||||
--id ID идентификатор клиента (default: cli:${defaultUser})
|
||||
--no-history не сохранять состояние в ~/.agentik/cli-state.json
|
||||
--help, -h эта справка
|
||||
|
||||
Переменные среды:
|
||||
AGENTIK_SERVER базовый URL агента (используется если --server не задан)
|
||||
HOME для пути ~/.agentik/cli-state.json
|
||||
|
||||
В REPL:
|
||||
/help список slash-команд
|
||||
/new [title] создать диалог (title опционально)
|
||||
/list, /ls список диалогов
|
||||
/switch <id>, /sw <id> переключиться на диалог
|
||||
/rename <title> переименовать текущий диалог
|
||||
/delete [<id>], /rm удалить (по id или текущий)
|
||||
/history, /h последние сообщения текущего диалога
|
||||
/interrupt, /stop прервать текущий ход
|
||||
/pwd показать текущий диалог
|
||||
/exit, /quit выйти (Ctrl-D тоже работает)
|
||||
""".trimIndent())
|
||||
}
|
||||
@@ -0,0 +1,3 @@
|
||||
package pw.binom.agentik.cli
|
||||
|
||||
internal expect fun platformEnv(key: String): String?
|
||||
@@ -1,81 +0,0 @@
|
||||
package pw.binom.agentik.cli
|
||||
|
||||
import kotlinx.serialization.Serializable
|
||||
import kotlinx.serialization.json.Json
|
||||
import kotlin.time.Instant
|
||||
|
||||
/**
|
||||
* Состояние CLI между запусками: последний выбранный диалог и момент последнего
|
||||
* увиденного [Event.date] в его потоке (для корректного `events(after)` после рестарта).
|
||||
*
|
||||
* Доступ к диску инкапсулирован в платформенный [CliPlatform] — commonMain ничего
|
||||
* не знает про `java.io.File`/`NSFileManager`, чтобы KMP-сборка собиралась
|
||||
* под все цели. Файл: `$HOME/.agentik/cli-state.json`.
|
||||
*/
|
||||
internal class SessionRepository internal constructor(
|
||||
private val filePath: String?,
|
||||
private val io: SessionIo,
|
||||
) {
|
||||
|
||||
@Serializable
|
||||
private data class State(
|
||||
val conversationId: String,
|
||||
val lastEventAt: String,
|
||||
)
|
||||
|
||||
private val json = Json { prettyPrint = true; ignoreUnknownKeys = true }
|
||||
|
||||
/** Открывается ленивым чтением. [save] ещё не было — файл может отсутствовать. */
|
||||
private var cached: State? = null
|
||||
|
||||
fun load(): SavedSession? {
|
||||
val path = filePath ?: return null
|
||||
val raw = io.readAll(path) ?: return null
|
||||
return runCatching {
|
||||
val state = json.decodeFromString(State.serializer(), raw)
|
||||
cached = state
|
||||
SavedSession(
|
||||
conversationId = state.conversationId,
|
||||
lastEventAt = Instant.parse(state.lastEventAt),
|
||||
)
|
||||
}.getOrNull()
|
||||
}
|
||||
|
||||
fun save(conversationId: String, lastEventAt: Instant) {
|
||||
val path = filePath ?: return
|
||||
val state = State(
|
||||
conversationId = conversationId,
|
||||
lastEventAt = lastEventAt.toString(),
|
||||
)
|
||||
cached = state
|
||||
val body = json.encodeToString(State.serializer(), state)
|
||||
io.writeAtomic(path, body)
|
||||
}
|
||||
|
||||
fun clear() {
|
||||
val path = filePath ?: return
|
||||
io.delete(path)
|
||||
cached = null
|
||||
}
|
||||
|
||||
fun close() {
|
||||
// для совместимости с будущим in-memory state; пока no-op
|
||||
}
|
||||
}
|
||||
|
||||
internal data class SavedSession(
|
||||
val conversationId: String,
|
||||
val lastEventAt: Instant,
|
||||
)
|
||||
|
||||
/**
|
||||
* Минимальный платформо-зависимый IO-интерфейс для одного файла. Реализации
|
||||
* в jvmMain (`java.io.File` + atomic `tmp → rename`) и в nativeMain (пока no-op-stub).
|
||||
*
|
||||
* public, потому что его возвращает public [CliPlatform.sessionIo].
|
||||
*/
|
||||
interface SessionIo {
|
||||
fun readAll(path: String): String?
|
||||
fun writeAtomic(path: String, body: String)
|
||||
fun delete(path: String)
|
||||
}
|
||||
@@ -1,90 +0,0 @@
|
||||
package pw.binom.agentik.cli
|
||||
|
||||
/**
|
||||
* Slash-команды REPL'а. Первая буква `/` не хранится — парсер уже её отрезал.
|
||||
*
|
||||
* Свободный ввод (без `/` в начале) — это сообщение пользователя агенту в
|
||||
* текущий диалог и НЕ разбирается в [parse].
|
||||
*/
|
||||
sealed interface SlashCommand {
|
||||
data object Help : SlashCommand
|
||||
data object Exit : SlashCommand
|
||||
data object Quit : SlashCommand // синоним Exit
|
||||
|
||||
/** Создать новый диалог; опционально — заголовок. */
|
||||
data class New(val title: String?) : SlashCommand
|
||||
|
||||
/** Список диалогов (cold flow — печатаем по мере прихода страниц). */
|
||||
data object List : SlashCommand
|
||||
|
||||
/** Подключиться к существующему диалогу по id. */
|
||||
data class Switch(val id: String) : SlashCommand
|
||||
|
||||
/** Переименовать текущий диалог. */
|
||||
data class Rename(val title: String) : SlashCommand
|
||||
|
||||
/** Удалить диалог (по id или текущий). */
|
||||
data class Delete(val id: String?) : SlashCommand
|
||||
|
||||
/** Прервать текущий ход. no-op если хода нет. */
|
||||
data object Interrupt : SlashCommand
|
||||
|
||||
/** Показать последние сообщения текущего диалога (cold flow). */
|
||||
data object History : SlashCommand
|
||||
|
||||
/** Показать информацию о текущем диалоге. */
|
||||
data object Pwd : SlashCommand
|
||||
}
|
||||
|
||||
/**
|
||||
* Парсит строку (без ведущего `/`) в [SlashCommand] либо возвращает [Result.Failure]
|
||||
* с сообщением об ошибке.
|
||||
*
|
||||
* Команды нечувствительны к регистру (команда `/LIST` == `/list`).
|
||||
*/
|
||||
fun parseSlash(input: String): ParseResult {
|
||||
val s = input.trim()
|
||||
if (s.isEmpty()) return ParseResult.Failure("пустая команда (введите /help)")
|
||||
|
||||
// Разбиваем на команду и её аргументы. Поддерживаем склейку: /new foo bar → new "foo bar"
|
||||
val firstSpace = s.indexOfAny(charArrayOf(' ', '\t'))
|
||||
val cmd = if (firstSpace < 0) s else s.substring(0, firstSpace)
|
||||
val rest = if (firstSpace < 0) "" else s.substring(firstSpace + 1).trim()
|
||||
val args = if (rest.isEmpty()) emptyList() else rest.split(' ').filter { it.isNotEmpty() }
|
||||
|
||||
val command: SlashCommand? = when (cmd.lowercase()) {
|
||||
"help", "?" -> SlashCommand.Help
|
||||
"exit" -> SlashCommand.Exit
|
||||
"quit", "q" -> SlashCommand.Quit
|
||||
"new" -> SlashCommand.New(rest.takeIf { it.isNotEmpty() })
|
||||
"list", "ls" -> SlashCommand.List
|
||||
"switch", "sw", "cd" -> args.firstOrNull()?.let { SlashCommand.Switch(it) }
|
||||
"rename", "mv", "title" -> rest.takeIf { it.isNotEmpty() }?.let { SlashCommand.Rename(it) }
|
||||
"delete", "rm" -> SlashCommand.Delete(args.firstOrNull())
|
||||
"interrupt", "stop", "cancel" -> SlashCommand.Interrupt
|
||||
"history", "hist", "h" -> SlashCommand.History
|
||||
"pwd", "where" -> SlashCommand.Pwd
|
||||
else -> null
|
||||
}
|
||||
if (command != null) return ParseResult.Success(command)
|
||||
|
||||
// Не нашли команду: либо неизвестная, либо не хватает аргумента.
|
||||
val cmdLower = cmd.lowercase()
|
||||
return when (cmdLower) {
|
||||
"switch", "sw", "cd" -> ParseResult.Failure("укажите id диалога: /switch <id>")
|
||||
"rename", "mv", "title" -> ParseResult.Failure("укажите заголовок: /rename <title>")
|
||||
else -> ParseResult.Failure("неизвестная команда: /$cmd (введите /help)")
|
||||
}
|
||||
}
|
||||
|
||||
sealed interface ParseResult {
|
||||
data class Success(val command: SlashCommand) : ParseResult
|
||||
data class Failure(val message: String) : ParseResult
|
||||
}
|
||||
|
||||
/** Удобный helper для тестов и общего кода. */
|
||||
fun parseSlashOrNull(input: String): SlashCommand? =
|
||||
when (val r = parseSlash(input)) {
|
||||
is ParseResult.Success -> r.command
|
||||
is ParseResult.Failure -> null
|
||||
}
|
||||
@@ -0,0 +1,32 @@
|
||||
package pw.binom.agentik.cli.commands
|
||||
|
||||
import kotlinx.cli.ExperimentalCli
|
||||
import kotlinx.cli.Subcommand
|
||||
import pw.binom.agentik.cli.AgentikSubcommand
|
||||
|
||||
/**
|
||||
* Родительская группа `conv`: операции над диалогами.
|
||||
*
|
||||
* Сама команда `agentik-cli conv` (без подкоманды) — no-op:
|
||||
* в kotlinx.cli parent.execute() вызывается ПОСЛЕ leaf.execute(),
|
||||
* поэтому любая работа в execute() дублирует вывод дочерней команды.
|
||||
* Для просмотра дочерних команд есть `agentik-cli conv --help`.
|
||||
*
|
||||
* Дочерние команды регистрируются через [subcommands] в конструкторе.
|
||||
*/
|
||||
@OptIn(ExperimentalCli::class)
|
||||
class ConvCommand : Subcommand("conv", "Операции над диалогами") {
|
||||
init {
|
||||
subcommands(
|
||||
ConvLsSubcommand(),
|
||||
ConvNewSubcommand(),
|
||||
ConvShowSubcommand(),
|
||||
ConvDeleteSubcommand(),
|
||||
ConvRenameSubcommand(),
|
||||
)
|
||||
}
|
||||
|
||||
override fun execute() = Unit
|
||||
}
|
||||
|
||||
abstract class ConvSubcommand(name: String, description: String) : AgentikSubcommand(name, description)
|
||||
+16
@@ -0,0 +1,16 @@
|
||||
package pw.binom.agentik.cli.commands
|
||||
|
||||
import kotlinx.cli.ArgType
|
||||
import pw.binom.agentik.cli.AgentikSubcommand
|
||||
import pw.binom.agentik.cli.defaultCliHttpClient
|
||||
import pw.binom.agentik.client.AgentikAgent
|
||||
|
||||
class ConvDeleteSubcommand : ConvSubcommand("delete", "Удалить диалог") {
|
||||
val id by argument(ArgType.String, description = "ID диалога")
|
||||
|
||||
override fun execute() = kotlinx.coroutines.runBlocking {
|
||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
|
||||
val ok = agent.deleteConversation(id)
|
||||
if (ok) println("deleted: $id") else println("conversation not found: $id")
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,33 @@
|
||||
package pw.binom.agentik.cli.commands
|
||||
|
||||
import kotlinx.cli.ArgType
|
||||
import kotlinx.cli.default
|
||||
import pw.binom.agentik.cli.AgentikSubcommand
|
||||
import pw.binom.agentik.cli.defaultCliHttpClient
|
||||
import pw.binom.agentik.client.AgentikAgent
|
||||
import pw.binom.agentik.proto.Agent
|
||||
|
||||
class ConvLsSubcommand : ConvSubcommand("ls", "Список диалогов агента") {
|
||||
val limit by option(ArgType.Int, fullName = "limit", description = "Максимум диалогов").default(Agent.PAGE_SIZE)
|
||||
|
||||
override fun execute() = kotlinx.coroutines.runBlocking {
|
||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
|
||||
val convs = agent.getConversations(offset = 0, limit = limit.coerceAtMost(Agent.PAGE_SIZE))
|
||||
if (convs.isEmpty()) {
|
||||
println("(no conversations)")
|
||||
return@runBlocking
|
||||
}
|
||||
println("ID UPDATED-AT TITLE FLAGS")
|
||||
convs.forEach { c ->
|
||||
val flags = buildString {
|
||||
if (c.isTemporal) append('T')
|
||||
if (c.isSupportImageInput) append('I')
|
||||
if (c.isSupportImageOutput) append('O')
|
||||
if (isEmpty()) append('-')
|
||||
}
|
||||
val title = c.title ?: "(untitled)"
|
||||
println("${c.id.padEnd(38)} ${c.updatedAt.toString().padEnd(22)} ${title.take(30).padEnd(31)} $flags")
|
||||
}
|
||||
println("--- ${convs.size} conversation(s)")
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,17 @@
|
||||
package pw.binom.agentik.cli.commands
|
||||
|
||||
import kotlinx.cli.ArgType
|
||||
import kotlinx.cli.default
|
||||
import pw.binom.agentik.cli.AgentikSubcommand
|
||||
import pw.binom.agentik.cli.defaultCliHttpClient
|
||||
import pw.binom.agentik.client.AgentikAgent
|
||||
|
||||
class ConvNewSubcommand : ConvSubcommand("new", "Создать диалог; печатает id") {
|
||||
val temp by option(ArgType.Boolean, fullName = "temp", description = "Временный диалог").default(false)
|
||||
|
||||
override fun execute() = kotlinx.coroutines.runBlocking {
|
||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
|
||||
val conv = agent.createConversation(temp = temp)
|
||||
println(conv.id)
|
||||
}
|
||||
}
|
||||
+25
@@ -0,0 +1,25 @@
|
||||
package pw.binom.agentik.cli.commands
|
||||
|
||||
import kotlinx.cli.ArgType
|
||||
import pw.binom.agentik.cli.AgentikSubcommand
|
||||
import pw.binom.agentik.cli.defaultCliHttpClient
|
||||
import pw.binom.agentik.client.AgentikAgent
|
||||
|
||||
class ConvRenameSubcommand : ConvSubcommand("rename", "Переименовать диалог") {
|
||||
val id by argument(ArgType.String, description = "ID диалога")
|
||||
val title by argument(ArgType.String, description = "Новое название")
|
||||
|
||||
override fun execute() = kotlinx.coroutines.runBlocking {
|
||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
|
||||
val conv = agent.getConversation(id) ?: run {
|
||||
println("conversation not found: $id")
|
||||
return@runBlocking
|
||||
}
|
||||
try {
|
||||
conv.rename(title)
|
||||
} finally {
|
||||
conv.close()
|
||||
}
|
||||
println("renamed: $id -> $title")
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,28 @@
|
||||
package pw.binom.agentik.cli.commands
|
||||
|
||||
import kotlinx.cli.ArgType
|
||||
import pw.binom.agentik.cli.AgentikSubcommand
|
||||
import pw.binom.agentik.cli.defaultCliHttpClient
|
||||
import pw.binom.agentik.client.AgentikAgent
|
||||
|
||||
class ConvShowSubcommand : ConvSubcommand("show", "Метаданные диалога") {
|
||||
val id by argument(ArgType.String, description = "ID диалога")
|
||||
|
||||
override fun execute() = kotlinx.coroutines.runBlocking {
|
||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
|
||||
val conv = agent.getConversation(id) ?: run {
|
||||
println("conversation not found: $id")
|
||||
return@runBlocking
|
||||
}
|
||||
try {
|
||||
println("id: ${conv.id}")
|
||||
println("title: ${conv.title ?: "(untitled)"}")
|
||||
println("updatedAt: ${conv.updatedAt}")
|
||||
println("isTemporal: ${conv.isTemporal}")
|
||||
println("isSupportImageInput: ${conv.isSupportImageInput}")
|
||||
println("isSupportImageOutput: ${conv.isSupportImageOutput}")
|
||||
} finally {
|
||||
conv.close()
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,10 @@
|
||||
package pw.binom.agentik.cli.commands
|
||||
|
||||
import pw.binom.agentik.cli.AgentikSubcommand
|
||||
|
||||
class InfoSubcommand : AgentikSubcommand("info", "Показать server URL и agent id") {
|
||||
override fun execute() {
|
||||
println("server: $serverUrl")
|
||||
println("id: $agentId")
|
||||
}
|
||||
}
|
||||
+24
@@ -0,0 +1,24 @@
|
||||
package pw.binom.agentik.cli.commands
|
||||
|
||||
import kotlinx.cli.ArgType
|
||||
import pw.binom.agentik.cli.AgentikSubcommand
|
||||
import pw.binom.agentik.cli.defaultCliHttpClient
|
||||
import pw.binom.agentik.client.AgentikAgent
|
||||
|
||||
class InterruptSubcommand : AgentikSubcommand("interrupt", "Прервать текущий ход диалога") {
|
||||
val id by argument(ArgType.String, description = "ID диалога")
|
||||
|
||||
override fun execute() = kotlinx.coroutines.runBlocking {
|
||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
|
||||
val conv = agent.getConversation(id) ?: run {
|
||||
println("conversation not found: $id")
|
||||
return@runBlocking
|
||||
}
|
||||
try {
|
||||
conv.interrupt()
|
||||
println("interrupted: $id")
|
||||
} finally {
|
||||
conv.close()
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,55 @@
|
||||
package pw.binom.agentik.cli.commands
|
||||
|
||||
import kotlinx.cli.ArgType
|
||||
import kotlinx.cli.default
|
||||
import pw.binom.agentik.cli.AgentikSubcommand
|
||||
import pw.binom.agentik.cli.defaultCliHttpClient
|
||||
import pw.binom.agentik.client.AgentikAgent
|
||||
import pw.binom.agentik.proto.Content
|
||||
import pw.binom.agentik.proto.Message
|
||||
import kotlin.time.Instant
|
||||
|
||||
class MsgsSubcommand : AgentikSubcommand("msgs", "Показать сообщения диалога") {
|
||||
val id by argument(ArgType.String, description = "ID диалога")
|
||||
val limit by option(ArgType.Int, fullName = "limit", description = "Максимум сообщений").default(100)
|
||||
|
||||
override fun execute() = kotlinx.coroutines.runBlocking {
|
||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
|
||||
val conv = agent.getConversation(id) ?: run {
|
||||
println("conversation not found: $id")
|
||||
return@runBlocking
|
||||
}
|
||||
try {
|
||||
val msgs = conv.getMessages(Instant.DISTANT_PAST, offset = 0, limit = limit)
|
||||
.sortedBy { it.date }
|
||||
msgs.forEach { m -> println(formatMessage(m)) }
|
||||
println("--- ${msgs.size} message(s)")
|
||||
} finally {
|
||||
conv.close()
|
||||
}
|
||||
}
|
||||
|
||||
private fun formatMessage(m: Message): String =
|
||||
"[${m.date}] ${m.role().padEnd(11)} ${m.bodyOneLine()}"
|
||||
|
||||
private fun Message.role(): String = when (this) {
|
||||
is Message.UserMessage -> "[user]"
|
||||
is Message.AssistantMessage -> "[assistant]"
|
||||
is Message.ToolCall -> "[tool_call]"
|
||||
is Message.ToolResult -> "[tool_result]"
|
||||
is Message.Error -> "[error]"
|
||||
}
|
||||
|
||||
private fun Message.bodyOneLine(): String = when (this) {
|
||||
is Message.UserMessage -> content.joinToString(" ") { c -> c.toOneLine() }
|
||||
is Message.AssistantMessage -> content.joinToString(" ") { c -> c.toOneLine() }
|
||||
is Message.ToolCall -> "tool=$toolName args=$toolArgs"
|
||||
is Message.ToolResult -> "id=$id result=${result ?: "<null>"}"
|
||||
is Message.Error -> "code=${code ?: "?"} message=$message"
|
||||
}
|
||||
|
||||
private fun Content.toOneLine(): String = when (this) {
|
||||
is Content.Text -> body.replace('\n', ' ').take(200)
|
||||
is Content.Image -> "<image ${data.size}B $mime>"
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,64 @@
|
||||
package pw.binom.agentik.cli.commands
|
||||
|
||||
import kotlinx.cli.ArgType
|
||||
import kotlinx.cli.vararg
|
||||
import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.flow.onEach
|
||||
import kotlinx.coroutines.flow.takeWhile
|
||||
import kotlinx.coroutines.launch
|
||||
import pw.binom.agentik.cli.AgentikSubcommand
|
||||
import pw.binom.agentik.cli.defaultCliHttpClient
|
||||
import pw.binom.agentik.client.AgentikAgent
|
||||
import pw.binom.agentik.proto.Content
|
||||
import pw.binom.agentik.proto.Event
|
||||
import kotlin.time.Instant
|
||||
|
||||
class SendSubcommand : AgentikSubcommand("send", "Отправить user-ход и стримить ответ") {
|
||||
val id by argument(ArgType.String, description = "ID диалога")
|
||||
val text by argument(ArgType.String, description = "Текст хода (все позиционные после <id> склеиваются пробелом)").vararg()
|
||||
|
||||
override fun execute() = kotlinx.coroutines.runBlocking {
|
||||
val agent = AgentikAgent(id = agentId, baseUrl = serverUrl, httpClient = defaultCliHttpClient())
|
||||
val conv = agent.getConversation(id) ?: run {
|
||||
println("conversation not found: $id")
|
||||
return@runBlocking
|
||||
}
|
||||
try {
|
||||
// Подписываемся на поток событий ДО send: события, отправленные
|
||||
// до подписки, не реплеятся (shared-flow без replay).
|
||||
val eventsJob = launch {
|
||||
conv.events(Instant.DISTANT_PAST)
|
||||
// onEach печатает и терминальный event, takeWhile лишь
|
||||
// завершает сбор после него.
|
||||
.onEach { ev -> emit(ev) }
|
||||
.takeWhile { ev -> !isTerminal(ev) }
|
||||
.collect { }
|
||||
}
|
||||
// Даём SSE-подписке установиться, затем шлём ход.
|
||||
delay(200)
|
||||
conv.send(listOf(Content.Text(text.joinToString(" "))))
|
||||
eventsJob.join()
|
||||
} finally {
|
||||
conv.close()
|
||||
}
|
||||
}
|
||||
|
||||
private fun isTerminal(ev: Event): Boolean =
|
||||
ev is Event.End || ev is Event.Interrupted || ev is Event.Error
|
||||
|
||||
private fun emit(ev: Event) {
|
||||
when (ev) {
|
||||
is Event.StartReasoning -> println("event StartReasoning")
|
||||
is Event.StartResponse -> println("event StartResponse ${ev.responseType}")
|
||||
is Event.AppendText -> println("event AppendText ${escape(ev.body)}")
|
||||
is Event.AppendImage -> println("event AppendImage <${ev.body.size}B ${ev.mime}>")
|
||||
is Event.ToolCall -> println("event ToolCall ${ev.id} ${ev.toolName} ${escape(ev.toolArgs)}")
|
||||
is Event.ToolResult -> println("event ToolResult ${ev.id} ${escape(ev.result ?: "")}")
|
||||
is Event.End -> println("event End")
|
||||
is Event.Interrupted -> println("event Interrupted")
|
||||
is Event.Error -> println("event Error ${ev.code ?: ""} ${escape(ev.message)}")
|
||||
}
|
||||
}
|
||||
|
||||
private fun escape(s: String): String = s.replace("\n", "\\n").replace("\r", "\\r")
|
||||
}
|
||||
@@ -1,81 +0,0 @@
|
||||
package pw.binom.agentik.cli
|
||||
|
||||
import kotlinx.coroutines.test.runTest
|
||||
import pw.binom.agentik.proto.Event
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertTrue
|
||||
import kotlin.time.Instant
|
||||
|
||||
/**
|
||||
* Подменяем [CliTerminal] простой in-memory реализацией и проверяем,
|
||||
* что события рендерятся в правильном формате.
|
||||
*/
|
||||
class EventRendererTest {
|
||||
|
||||
private class FakeTerminal : CliTerminal {
|
||||
override val prompt: String = ">"
|
||||
val out = StringBuilder()
|
||||
override suspend fun readLine(): String? = null
|
||||
override suspend fun println(text: String) { out.appendLine(text) }
|
||||
override suspend fun print(text: String) { out.append(text) }
|
||||
override suspend fun printSystem(text: String) { out.appendLine("· $text") }
|
||||
override fun close() {}
|
||||
fun text() = out.toString()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `simple response stream`() = runTest {
|
||||
val t = FakeTerminal()
|
||||
val r = EventRenderer(t)
|
||||
r.render(Event.StartResponse(Instant.DISTANT_PAST, Event.ResponseType.TEXT))
|
||||
r.render(Event.AppendText(Instant.DISTANT_PAST, "Привет"))
|
||||
r.render(Event.AppendText(Instant.DISTANT_PAST, ", мир!"))
|
||||
r.render(Event.End(Instant.DISTANT_PAST))
|
||||
|
||||
// StartResponse открывает блок, AppendText без \n, End закрывает \n
|
||||
val text = t.text()
|
||||
assertTrue(text.contains("Привет, мир!"), "got: $text")
|
||||
// после End должен быть перевод строки
|
||||
assertTrue(text.endsWith("\n"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `interrupted closes block`() = runTest {
|
||||
val t = FakeTerminal()
|
||||
val r = EventRenderer(t)
|
||||
r.render(Event.StartResponse(Instant.DISTANT_PAST, Event.ResponseType.TEXT))
|
||||
r.render(Event.AppendText(Instant.DISTANT_PAST, "Частично"))
|
||||
r.render(Event.Interrupted(Instant.DISTANT_PAST))
|
||||
val text = t.text()
|
||||
assertTrue(text.contains("Частично"))
|
||||
assertTrue(text.contains("· [interrupted]"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `error before response`() = runTest {
|
||||
val t = FakeTerminal()
|
||||
val r = EventRenderer(t)
|
||||
r.render(Event.Error(Instant.DISTANT_PAST, message = "что-то сломалось", code = "500"))
|
||||
val text = t.text()
|
||||
assertTrue(text.contains("· [error/500] что-то сломалось"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `start_reasoning is printed as system line`() = runTest {
|
||||
val t = FakeTerminal()
|
||||
val r = EventRenderer(t)
|
||||
r.render(Event.StartReasoning(Instant.DISTANT_PAST))
|
||||
assertTrue(t.text().contains("· …thinking…"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `image append renders placeholder`() = runTest {
|
||||
val t = FakeTerminal()
|
||||
val r = EventRenderer(t)
|
||||
r.render(Event.StartResponse(Instant.DISTANT_PAST, Event.ResponseType.IMAGE))
|
||||
r.render(Event.AppendImage(Instant.DISTANT_PAST, body = ByteArray(64), mime = "image/png"))
|
||||
r.render(Event.End(Instant.DISTANT_PAST))
|
||||
assertTrue(t.text().contains("[image:image/png:64 bytes]"))
|
||||
}
|
||||
}
|
||||
@@ -1,101 +0,0 @@
|
||||
package pw.binom.agentik.cli
|
||||
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertIs
|
||||
import kotlin.test.assertTrue
|
||||
|
||||
class SlashCommandTest {
|
||||
|
||||
@Test
|
||||
fun `help is parsed`() {
|
||||
assertIs<SlashCommand.Help>(parseSlashOrNull("help"))
|
||||
assertIs<SlashCommand.Help>(parseSlashOrNull("?"))
|
||||
assertIs<SlashCommand.Help>(parseSlashOrNull("HELP"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `exit and quit alias`() {
|
||||
assertIs<SlashCommand.Exit>(parseSlashOrNull("exit"))
|
||||
assertIs<SlashCommand.Quit>(parseSlashOrNull("q"))
|
||||
assertIs<SlashCommand.Quit>(parseSlashOrNull("Quit"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `new without title`() {
|
||||
assertIs<SlashCommand.New>(parseSlashOrNull("new")).let {
|
||||
assertEquals(null, it.title)
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `new with multi-word title`() {
|
||||
val cmd = parseSlashOrNull("new my cool chat")
|
||||
assertIs<SlashCommand.New>(cmd)
|
||||
assertEquals("my cool chat", cmd.title)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `switch requires id`() {
|
||||
val r = parseSlash("sw")
|
||||
assertIs<ParseResult.Failure>(r)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `switch with id`() {
|
||||
val cmd = parseSlashOrNull("switch abc123")
|
||||
assertIs<SlashCommand.Switch>(cmd)
|
||||
assertEquals("abc123", cmd.id)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `rename requires title`() {
|
||||
val r = parseSlash("rename")
|
||||
assertIs<ParseResult.Failure>(r)
|
||||
// А "rename " (с пробелом, но без слов после) — это уже успех с пустым title?
|
||||
// У нас: rest = "" → takeIf { it.isNotEmpty() } → null → Failure. ОК.
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `rename with title`() {
|
||||
val cmd = parseSlashOrNull("rename my new title ")
|
||||
assertIs<SlashCommand.Rename>(cmd)
|
||||
assertEquals("my new title", cmd.title) // trim() делает своё
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `delete may have id or not`() {
|
||||
assertIs<SlashCommand.Delete>(parseSlashOrNull("rm")).let {
|
||||
assertEquals(null, it.id)
|
||||
}
|
||||
assertIs<SlashCommand.Delete>(parseSlashOrNull("delete abc")).let {
|
||||
assertEquals("abc", it.id)
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `unknown command fails`() {
|
||||
val r = parseSlash("foobar")
|
||||
assertIs<ParseResult.Failure>(r)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `empty command fails`() {
|
||||
val r = parseSlash("")
|
||||
assertIs<ParseResult.Failure>(r)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `command is case insensitive`() {
|
||||
assertIs<SlashCommand.List>(parseSlashOrNull("LIST"))
|
||||
assertIs<SlashCommand.Interrupt>(parseSlashOrNull("STOP"))
|
||||
assertIs<SlashCommand.Pwd>(parseSlashOrNull("PWD"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `interrupt synonyms`() {
|
||||
assertIs<SlashCommand.Interrupt>(parseSlashOrNull("interrupt"))
|
||||
assertIs<SlashCommand.Interrupt>(parseSlashOrNull("stop"))
|
||||
assertIs<SlashCommand.Interrupt>(parseSlashOrNull("cancel"))
|
||||
}
|
||||
}
|
||||
@@ -1,151 +0,0 @@
|
||||
package pw.binom.agentik.cli
|
||||
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.withContext
|
||||
import org.jline.reader.EndOfFileException
|
||||
import org.jline.reader.LineReader
|
||||
import org.jline.reader.LineReaderBuilder
|
||||
import org.jline.reader.UserInterruptException
|
||||
import org.jline.terminal.TerminalBuilder
|
||||
import pw.binom.agentik.client.AgentikAgent
|
||||
import pw.binom.agentik.proto.Agent
|
||||
import java.io.File
|
||||
import java.nio.file.Files
|
||||
import java.nio.file.StandardCopyOption
|
||||
|
||||
actual object CliPlatform {
|
||||
actual fun openAgent(baseUrl: String, id: String): Agent =
|
||||
AgentikAgent(id = id, baseUrl = baseUrl)
|
||||
|
||||
actual fun openTerminal(historyFile: String?, prompt: String): CliTerminal =
|
||||
JLineTerminal(historyFile = historyFile, prompt = prompt)
|
||||
|
||||
actual fun homeDir(): String? =
|
||||
System.getenv("HOME") ?: System.getenv("USERPROFILE")
|
||||
|
||||
actual fun env(key: String): String? = System.getenv(key)
|
||||
|
||||
actual fun sessionIo(): SessionIo = JvmSessionIo
|
||||
}
|
||||
|
||||
/**
|
||||
* Реализация [SessionIo] поверх `java.io.File` + atomic `tmp → rename`.
|
||||
* tmp-файл пишется в той же директории, что и целевой, чтобы rename
|
||||
* был атомарным в рамках одного раздела (POSIX rename(2) и Windows
|
||||
* MoveFileEx — атомарны внутри одного тома).
|
||||
*/
|
||||
private object JvmSessionIo : SessionIo {
|
||||
override fun readAll(path: String): String? {
|
||||
val f = File(path)
|
||||
if (!f.exists()) return null
|
||||
return runCatching { f.readText() }.getOrNull()
|
||||
}
|
||||
|
||||
override fun writeAtomic(path: String, body: String) {
|
||||
val target = File(path)
|
||||
target.parentFile?.mkdirs()
|
||||
val tmp = File(path + ".tmp")
|
||||
tmp.writeText(body)
|
||||
if (!tmp.renameTo(target)) {
|
||||
// fallback: Windows-специфика — renameTo может не перезаписать существующий.
|
||||
runCatching { Files.move(tmp.toPath(), target.toPath(), StandardCopyOption.REPLACE_EXISTING, StandardCopyOption.ATOMIC_MOVE) }
|
||||
.getOrElse { target.writeText(tmp.readText()); tmp.delete() }
|
||||
}
|
||||
} override fun delete(path: String) {
|
||||
runCatching { File(path).delete() }
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Реализация [CliTerminal] поверх JLine ([LineReader]).
|
||||
*
|
||||
* JLine-3 API:
|
||||
* - [TerminalBuilder.builder().system(true).build()] — открыть системный TTY.
|
||||
* - [LineReader] поверх Terminal — readline-редактор (стрелки, history, Ctrl-A/E).
|
||||
* - [LineReader.readLine(prompt)] — suspend-free, блокирующий IO; мы оборачиваем
|
||||
* в [withContext] [Dispatchers.IO], чтобы не держать event-loop.
|
||||
* - [DefaultHistory] (org.jline.reader.history.DefaultHistory) + история из файла.
|
||||
*/
|
||||
private class JLineTerminal(
|
||||
historyFile: String?,
|
||||
override val prompt: String,
|
||||
) : CliTerminal {
|
||||
|
||||
private val terminal = TerminalBuilder.builder()
|
||||
.system(true)
|
||||
.jna(true)
|
||||
.build()
|
||||
|
||||
private val historyImpl: org.jline.reader.History? = run {
|
||||
if (historyFile == null) null else try {
|
||||
val history = org.jline.reader.impl.history.DefaultHistory()
|
||||
val histFile = File(historyFile)
|
||||
histFile.parentFile?.mkdirs()
|
||||
history.load()
|
||||
if (histFile.exists()) {
|
||||
history.append(histFile.toPath(), true)
|
||||
}
|
||||
history
|
||||
} catch (t: Throwable) {
|
||||
null
|
||||
}
|
||||
}
|
||||
|
||||
private val reader: LineReader = LineReaderBuilder.builder()
|
||||
.terminal(terminal)
|
||||
.apply { if (historyImpl != null) history(historyImpl) }
|
||||
.build()
|
||||
|
||||
private val historyFilePath: java.nio.file.Path? =
|
||||
historyFile?.let { File(it).toPath() }
|
||||
|
||||
override suspend fun readLine(): String? = withContext(Dispatchers.IO) {
|
||||
try {
|
||||
val line = reader.readLine(prompt)
|
||||
// Сохраняем history при каждой строке — дешево, и при Ctrl-D / Ctrl-C
|
||||
// ничего не теряется.
|
||||
flushHistory()
|
||||
line
|
||||
} catch (_: UserInterruptException) {
|
||||
// Ctrl-C: трактуем как «всё, выходим», как и EOF.
|
||||
flushHistory()
|
||||
null
|
||||
} catch (_: EndOfFileException) {
|
||||
// Ctrl-D на пустой строке.
|
||||
flushHistory()
|
||||
null
|
||||
}
|
||||
}
|
||||
|
||||
override suspend fun println(text: String): Unit = withContext(Dispatchers.IO) {
|
||||
terminal.writer().println(text)
|
||||
terminal.writer().flush()
|
||||
}
|
||||
|
||||
override suspend fun print(text: String): Unit = withContext(Dispatchers.IO) {
|
||||
terminal.writer().print(text)
|
||||
terminal.writer().flush()
|
||||
}
|
||||
|
||||
override suspend fun printSystem(text: String): Unit = withContext(Dispatchers.IO) {
|
||||
terminal.writer().println("· $text")
|
||||
terminal.writer().flush()
|
||||
}
|
||||
|
||||
private fun flushHistory() {
|
||||
val hf = historyFilePath ?: return
|
||||
val h = historyImpl ?: return
|
||||
runCatching {
|
||||
h.save()
|
||||
if (!h.isEmpty) {
|
||||
// читаем из .tmp и дописываем
|
||||
h.append(hf, true)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
override fun close() {
|
||||
runCatching { flushHistory() }
|
||||
runCatching { terminal.close() }
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,3 @@
|
||||
package pw.binom.agentik.cli
|
||||
|
||||
internal actual fun platformEnv(key: String): String? = System.getenv(key)
|
||||
@@ -1,108 +0,0 @@
|
||||
package pw.binom.agentik.cli
|
||||
|
||||
import org.junit.After
|
||||
import org.junit.Before
|
||||
import java.io.File
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertNotNull
|
||||
import kotlin.test.assertNull
|
||||
import kotlin.test.assertTrue
|
||||
import kotlin.time.Instant
|
||||
|
||||
/**
|
||||
* Интеграционный тест на реальном временном файле. Только JVM: использует
|
||||
* [java.io.File] для IO-интерфейса. На native-таргетах тест не собирается —
|
||||
* TODO: переписать на kotlinx-io Files и перенести в commonTest.
|
||||
*/
|
||||
class SessionRepositoryTest {
|
||||
|
||||
private lateinit var tmp: File
|
||||
|
||||
@Before
|
||||
fun setUp() {
|
||||
tmp = File.createTempFile("agentik-cli-state", ".json")
|
||||
tmp.delete()
|
||||
}
|
||||
|
||||
@After
|
||||
fun tearDown() {
|
||||
if (tmp.exists()) tmp.delete()
|
||||
File(tmp.path + ".tmp").delete()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `load returns null when file missing`() {
|
||||
val repo = SessionRepository(tmp.path, JvmIo)
|
||||
assertNull(repo.load())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `save then load roundtrip`() {
|
||||
val repo = SessionRepository(tmp.path, JvmIo)
|
||||
val savedAt = Instant.parse("2026-09-16T10:00:00Z")
|
||||
repo.save(conversationId = "abcd-1234", lastEventAt = savedAt)
|
||||
repo.close()
|
||||
|
||||
val repo2 = SessionRepository(tmp.path, JvmIo)
|
||||
val restored = repo2.load()
|
||||
assertNotNull(restored)
|
||||
assertEquals("abcd-1234", restored.conversationId)
|
||||
assertEquals(savedAt, restored.lastEventAt)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `save overwrites previous state`() {
|
||||
val repo = SessionRepository(tmp.path, JvmIo)
|
||||
repo.save("conv-1", Instant.parse("2026-09-16T10:00:00Z"))
|
||||
repo.save("conv-2", Instant.parse("2026-09-16T11:00:00Z"))
|
||||
repo.close()
|
||||
|
||||
val restored = SessionRepository(tmp.path, JvmIo).load()
|
||||
assertNotNull(restored)
|
||||
assertEquals("conv-2", restored.conversationId)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `null filepath means no-op`() {
|
||||
val repo = SessionRepository(null, JvmIo)
|
||||
repo.save("conv-X", Instant.parse("2026-09-16T10:00:00Z"))
|
||||
// Не должно ни читать, ни писать.
|
||||
assertNull(repo.load())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `clear deletes file`() {
|
||||
val repo = SessionRepository(tmp.path, JvmIo)
|
||||
repo.save("conv-Z", Instant.parse("2026-09-16T10:00:00Z"))
|
||||
repo.close()
|
||||
assertTrue(tmp.exists())
|
||||
|
||||
val repo2 = SessionRepository(tmp.path, JvmIo)
|
||||
repo2.clear()
|
||||
assertTrue(!tmp.exists())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `corrupt json is ignored (does not throw)`() {
|
||||
File(tmp.path).writeText("this is not json")
|
||||
val repo = SessionRepository(tmp.path, JvmIo)
|
||||
assertNull(repo.load())
|
||||
}
|
||||
}
|
||||
|
||||
// JVM-only helper: реализация [SessionIo] поверх `java.io.File` для теста.
|
||||
// В продакшен-коде на jvmMain ровно такая же логика.
|
||||
private object JvmIo : SessionIo {
|
||||
override fun readAll(path: String): String? {
|
||||
val f = File(path); if (!f.exists()) return null
|
||||
return runCatching { f.readText() }.getOrNull()
|
||||
}
|
||||
override fun writeAtomic(path: String, body: String) {
|
||||
val target = File(path); target.parentFile?.mkdirs()
|
||||
val tmp = File(path + ".tmp")
|
||||
tmp.writeText(body)
|
||||
if (!tmp.renameTo(target)) target.writeText(tmp.readText()).also { tmp.delete() }
|
||||
}
|
||||
override fun delete(path: String) { File(path).delete() }
|
||||
}
|
||||
@@ -1,36 +0,0 @@
|
||||
package pw.binom.agentik.cli
|
||||
|
||||
import pw.binom.agentik.proto.Agent
|
||||
|
||||
/**
|
||||
* Платформо-зависимая реализация для native-целей.
|
||||
*
|
||||
* Текущий статус: stub. native HTTP требует подключения ktor-client-core +
|
||||
* платформенных engine'ов (ktor-client-darwin для Apple, ktor-client-curl для
|
||||
* linux/mingw, ktor-client-okhttp для Android в перспективе) и переиспользования
|
||||
* уже существующего `:client` SSE-парсера. Native readline требует termios
|
||||
* через `kotlinx.cinterop` — добавим, когда дойдут руки.
|
||||
*
|
||||
* Пока запустить агента из native-бинаря CLI нельзя, но проект компилируется
|
||||
* под все 8 KMP-целей — структурная готовность соблюдена.
|
||||
*/
|
||||
actual object CliPlatform {
|
||||
actual fun openAgent(baseUrl: String, id: String): Agent =
|
||||
error("agentik-cli native target is not implemented yet (baseUrl=$baseUrl)")
|
||||
|
||||
actual fun openTerminal(historyFile: String?, prompt: String): CliTerminal =
|
||||
error("agentik-cli native target is not implemented yet (prompt=$prompt)")
|
||||
|
||||
actual fun homeDir(): String? = null
|
||||
|
||||
actual fun env(key: String): String? = null
|
||||
|
||||
actual fun sessionIo(): SessionIo = NoopSessionIo
|
||||
}
|
||||
|
||||
/** Минимальный no-op-IO для native-целей пока не подключён реальный движок. */
|
||||
private object NoopSessionIo : SessionIo {
|
||||
override fun readAll(path: String): String? = null
|
||||
override fun writeAtomic(path: String, body: String) {}
|
||||
override fun delete(path: String) {}
|
||||
}
|
||||
@@ -0,0 +1,8 @@
|
||||
package pw.binom.agentik.cli
|
||||
|
||||
import kotlinx.cinterop.ExperimentalForeignApi
|
||||
import kotlinx.cinterop.toKString
|
||||
import platform.posix.getenv
|
||||
|
||||
@OptIn(ExperimentalForeignApi::class)
|
||||
internal actual fun platformEnv(key: String): String? = getenv(key)?.toKString()
|
||||
@@ -24,6 +24,23 @@ kotlin {
|
||||
linuxArm64()
|
||||
mingwX64()
|
||||
|
||||
// Native executables. По умолчанию Kotlin/Native для каждого target'а
|
||||
// собирает только .klib (библиотеку) — для запускаемого .kexe надо
|
||||
// явно попросить binaries.executable(). entryPoint нужно задать явно:
|
||||
// KMP-линкер ищет функцию по FQN (без `Kt`-суффикса), а Kotlin/Native
|
||||
// добавляет суффикс только для файлов с именем `Main.kt`, поэтому
|
||||
// указываем точку входа как `pw.binom.agentik.tui.main` (без суффикса).
|
||||
//
|
||||
// Применяем к каждому из linuxX64/macosX64/macosArm64/linuxArm64/mingwX64
|
||||
// явно (а не через targets.withType), потому что targets DSL в KMP не
|
||||
// поддерживает реифицированный withType<KotlinNativeTarget>().
|
||||
@OptIn(ExperimentalKotlinGradlePluginApi::class)
|
||||
listOf(linuxX64(), linuxArm64(), macosX64(), macosArm64(), mingwX64()).forEach {
|
||||
it.binaries.executable {
|
||||
entryPoint = "pw.binom.agentik.tui.main"
|
||||
}
|
||||
}
|
||||
|
||||
sourceSets {
|
||||
commonMain.dependencies {
|
||||
implementation(project(":proto"))
|
||||
@@ -34,6 +51,11 @@ kotlin {
|
||||
// JetBrains Compose runtime — тащит Mosaic как обёртку.
|
||||
implementation(libs.mosaic.runtime)
|
||||
implementation(libs.mosaic.tty.terminal)
|
||||
|
||||
// Health-check в Main.kt: Ktor CIO на JVM, на native не собирается —
|
||||
// там работает stub actual через expect/actual.
|
||||
implementation(libs.ktor.client.core)
|
||||
implementation(libs.ktor.client.cio)
|
||||
}
|
||||
jvmMain.dependencies {
|
||||
implementation(project(":client"))
|
||||
@@ -41,6 +63,7 @@ kotlin {
|
||||
commonTest.dependencies {
|
||||
implementation(kotlin("test"))
|
||||
implementation(libs.kotlinx.coroutines.core)
|
||||
implementation(libs.kotlinx.coroutines.test)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -3,18 +3,19 @@ package pw.binom.agentik.tui
|
||||
import androidx.compose.runtime.Composable
|
||||
import androidx.compose.runtime.collectAsState
|
||||
import androidx.compose.runtime.getValue
|
||||
import com.jakewharton.mosaic.layout.KeyEvent
|
||||
import com.jakewharton.mosaic.layout.drawBehind
|
||||
import com.jakewharton.mosaic.layout.onPreviewKeyEvent
|
||||
import com.jakewharton.mosaic.modifier.Modifier
|
||||
import com.jakewharton.mosaic.ui.Box
|
||||
import com.jakewharton.mosaic.ui.Column
|
||||
import com.jakewharton.mosaic.ui.Row
|
||||
import com.jakewharton.mosaic.ui.Text
|
||||
import com.jakewharton.mosaic.ui.TextStyle
|
||||
import pw.binom.agentik.tui.ui.Footer
|
||||
import pw.binom.agentik.tui.ui.Header
|
||||
import pw.binom.agentik.tui.ui.HelpOverlay
|
||||
import pw.binom.agentik.tui.ui.HistoryPanel
|
||||
import pw.binom.agentik.tui.ui.InputLine
|
||||
|
||||
/**
|
||||
* Корневая Compose-композиция TUI.
|
||||
* Корневая Compose-композиция TUI. Содержит только каркас + глобальный key-handler;
|
||||
* каждый регион (header/history/input/footer/help) — отдельный компонент в `ui/`.
|
||||
*
|
||||
* Layout (минимальный):
|
||||
* ```
|
||||
@@ -28,6 +29,9 @@ import com.jakewharton.mosaic.ui.TextStyle
|
||||
* │ FOOTER: ↑↓ scroll Tab focus Enter send F1 help … │
|
||||
* └─────────────────────────────────────────────────────────┘
|
||||
* ```
|
||||
*
|
||||
* Глобальные клавиши (Tab/Shift-Tab/F1/Esc) обрабатываются здесь.
|
||||
* Клавиши внутри строки ввода — в [InputLine] (через свой `onPreviewKeyEvent`).
|
||||
*/
|
||||
@Composable
|
||||
internal fun App(state: AppState) {
|
||||
@@ -50,105 +54,8 @@ internal fun App(state: AppState) {
|
||||
Header(state, focusIndex)
|
||||
HistoryPanel(state)
|
||||
InputLine(state)
|
||||
Footer(state, showHelp)
|
||||
Footer(showHelp)
|
||||
}
|
||||
}
|
||||
if (showHelp) HelpOverlay()
|
||||
}
|
||||
|
||||
@Composable
|
||||
private fun Header(state: AppState, focusIndex: Int) {
|
||||
val title by state.currentTitle.collectAsState()
|
||||
val convId by state.currentConversationId.collectAsState()
|
||||
val focusLabel = when (focusIndex) { 0 -> "input"; 1 -> "history"; 2 -> "sidebar"; else -> "?" }
|
||||
val convStr = convId?.let { " · ${it.take(8)}…" } ?: ""
|
||||
val titleStr = title ?: "(нет диалога)"
|
||||
Text(
|
||||
value = " agentik · ${state.config.id}$convStr · $titleStr · focus=$focusLabel ",
|
||||
textStyle = TextStyle.Bold + TextStyle.Invert,
|
||||
)
|
||||
}
|
||||
|
||||
@Composable
|
||||
private fun HistoryPanel(state: AppState) {
|
||||
val messages by state.messages.collectAsState()
|
||||
val scroll by state.historyScroll.collectAsState()
|
||||
val rendered = if (messages.isEmpty()) {
|
||||
" (пока пусто)\n Tab — переключить фокус, F1 — подсказки.\n"
|
||||
} else {
|
||||
messages.joinToString("") { renderMessage(it) }
|
||||
}
|
||||
Text(value = rendered)
|
||||
}
|
||||
|
||||
private fun renderMessage(m: TuiMessage): String = when (m) {
|
||||
is TuiMessage.System -> " ── ${m.text}\n"
|
||||
is TuiMessage.User -> " > ${m.text}\n"
|
||||
is TuiMessage.Assistant -> " ╰ ${m.text}\n"
|
||||
is TuiMessage.AssistantStreaming -> " ╰ ${m.text} ▍\n"
|
||||
is TuiMessage.ToolCall -> " ⚙ ${m.toolName}${if (!m.title.isNullOrEmpty()) ": ${m.title}" else ""}\n"
|
||||
is TuiMessage.ToolResult -> " ↳ ${m.result.take(200)}${if (m.result.length > 200) "…" else ""}\n"
|
||||
}
|
||||
|
||||
@Composable
|
||||
private fun InputLine(state: AppState) {
|
||||
val text by state.input.collectAsState()
|
||||
val cursor by state.cursor.collectAsState()
|
||||
val streaming by state.streaming.collectAsState()
|
||||
val cursorPos = cursor.coerceIn(0, text.length)
|
||||
val before = text.substring(0, cursorPos)
|
||||
val cursorChar = if (cursorPos < text.length) text[cursorPos].toString() else " "
|
||||
val afterStart = if (cursorPos < text.length) cursorPos + 1 else cursorPos
|
||||
val after = text.substring(afterStart.coerceAtMost(text.length))
|
||||
val prompt = if (streaming) " ⋯" else " >"
|
||||
|
||||
Text(
|
||||
value = "$prompt $before|$cursorChar|${after}",
|
||||
modifier = Modifier
|
||||
.onPreviewKeyEvent { ev -> handleInputKey(state, ev) }
|
||||
.drawBehind {
|
||||
// Snapshot-read state в drawBehind чтобы changes триггерили redraw.
|
||||
state.input.let { /* touch */ }
|
||||
},
|
||||
)
|
||||
}
|
||||
|
||||
private fun handleInputKey(state: AppState, ev: KeyEvent): Boolean {
|
||||
if (ev.alt || ev.ctrl) return false
|
||||
return when (ev.key) {
|
||||
"Enter" -> state.submitInput() != null
|
||||
"Backspace" -> { state.inputBackspace(); true }
|
||||
"Delete" -> { state.inputDelete(); true }
|
||||
"Left", "ArrowLeft" -> { state.inputMoveCursor(-1); true }
|
||||
"Right", "ArrowRight" -> { state.inputMoveCursor(+1); true }
|
||||
"Home" -> { state.inputCursorHome(); true }
|
||||
"End" -> { state.inputCursorEnd(); true }
|
||||
else -> {
|
||||
val s = ev.key
|
||||
if (s.length == 1) { state.inputInsert(s); true }
|
||||
else false
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Composable
|
||||
private fun Footer(state: AppState, showHelp: Boolean) {
|
||||
val hint = if (showHelp) " ↑ наверху help-оверлей ↑ "
|
||||
else " Tab focus ↑↓ scroll Enter send Esc clear F1 help Ctrl-D exit "
|
||||
Text(value = hint, textStyle = TextStyle.Italic)
|
||||
}
|
||||
|
||||
@Composable
|
||||
private fun HelpOverlay() {
|
||||
Column(modifier = Modifier) {
|
||||
Text(value = " --- HELP ---", textStyle = TextStyle.Bold + TextStyle.Invert)
|
||||
Text(value = " Tab / Shift-Tab переключить фокус (history / input / sidebar)")
|
||||
Text(value = " ↑ / ↓ скролл истории / курсор в input")
|
||||
Text(value = " ← / → курсор в input")
|
||||
Text(value = " Enter отправить сообщение")
|
||||
Text(value = " Backspace / Del удалить символ")
|
||||
Text(value = " Esc очистить input")
|
||||
Text(value = " Ctrl-D / Ctrl-C выход")
|
||||
Text(value = " F1 toggle help", textStyle = TextStyle.Italic)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -15,6 +15,10 @@ import kotlin.time.Instant
|
||||
* (см. samples/snake в репо Mosaic).
|
||||
*/
|
||||
internal class AppState(val config: TuiConfig) {
|
||||
/** Бэкенд, прикреплённый из TuiApp — маршрутизирует submitInput → send. */
|
||||
private var backend: TuiBackend? = null
|
||||
fun attachBackend(b: TuiBackend) { backend = b }
|
||||
|
||||
/** Зона фокуса: 0 = input, 1 = history, 2 = sidebar. */
|
||||
private val _focusIndex = MutableStateFlow(0)
|
||||
val focusIndex: StateFlow<Int> = _focusIndex.asStateFlow()
|
||||
@@ -101,9 +105,28 @@ internal class AppState(val config: TuiConfig) {
|
||||
_messages.value = _messages.value + TuiMessage.User(text = text, ts = nowInstant())
|
||||
inputClear()
|
||||
_streaming.value = true
|
||||
backend?.onUserMessage(text)
|
||||
return text
|
||||
}
|
||||
|
||||
fun setStreaming(v: Boolean) { _streaming.value = v }
|
||||
|
||||
fun setConversation(id: String, title: String?) {
|
||||
_currentConversationId.value = id
|
||||
_currentTitle.value = title
|
||||
_messages.value = emptyList()
|
||||
_historyScroll.value = 0
|
||||
_streaming.value = false
|
||||
}
|
||||
|
||||
fun postToolCall(toolName: String, title: String?, args: String) {
|
||||
_messages.value = _messages.value + TuiMessage.ToolCall(toolName = toolName, title = title, args = args, ts = nowInstant())
|
||||
}
|
||||
|
||||
fun postToolResult(toolName: String, result: String) {
|
||||
_messages.value = _messages.value + TuiMessage.ToolResult(toolName = toolName, result = result, ts = nowInstant())
|
||||
}
|
||||
|
||||
fun appendAssistant(chunk: String) {
|
||||
val list = _messages.value.toMutableList()
|
||||
val last = list.lastOrNull()
|
||||
|
||||
@@ -1,6 +1,12 @@
|
||||
package pw.binom.agentik.tui
|
||||
|
||||
import io.ktor.client.HttpClient
|
||||
import io.ktor.client.engine.cio.CIO
|
||||
import io.ktor.client.plugins.HttpTimeout
|
||||
import io.ktor.client.request.get
|
||||
import io.ktor.client.statement.bodyAsText
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import pw.binom.agentik.proto.Agent
|
||||
|
||||
/**
|
||||
* Точка входа TUI-клиента agentik.
|
||||
@@ -9,14 +15,62 @@ import kotlinx.coroutines.runBlocking
|
||||
* agentik-tui [--server URL] [--id ID] [--no-history] [--help]
|
||||
* ```
|
||||
*
|
||||
* Без аргументов — стартует Compose-Mosaic UI.
|
||||
* Перед запуском UI — обязательный health-check: `GET {server}/health`.
|
||||
* Если сервер недоступен — печатаем понятную ошибку и выходим с кодом 1.
|
||||
* Если OK — создаём [Agent] через платформенную actual и запускаем
|
||||
* [TuiApp].
|
||||
*/
|
||||
fun main(args: Array<String>) = runBlocking {
|
||||
val cfg = parseCliArgs(args) ?: run {
|
||||
printUsage()
|
||||
return@runBlocking
|
||||
}
|
||||
TuiApp(cfg).run()
|
||||
checkServer(cfg.server)
|
||||
val agent = platformCreateAgent(cfg.server, cfg.id)
|
||||
TuiApp(cfg, agent).run()
|
||||
}
|
||||
|
||||
/**
|
||||
* Делает синхронный GET `{baseUrl}/health`. Внутри [route(path)] на сервере
|
||||
* `/health` зарегистрирован под тем же path-prefix'ом, что и сам API
|
||||
* (например, baseUrl = `http://localhost:8080/agentik` → health = …/agentik/health).
|
||||
*
|
||||
* При любой ошибке (connect refused, timeout, не-200 ответ, не `"ok"`) —
|
||||
* бросает [IllegalStateException] с понятным сообщением. [runBlocking]-обёртка
|
||||
* в [main] разворачивает её в stack-trace и `exit 1`.
|
||||
*/
|
||||
private suspend fun checkServer(baseUrl: String) {
|
||||
val healthUrl = "${baseUrl.trimEnd('/')}/health"
|
||||
val client = HttpClient(CIO) {
|
||||
install(HttpTimeout) {
|
||||
requestTimeoutMillis = 5_000
|
||||
connectTimeoutMillis = 3_000
|
||||
}
|
||||
expectSuccess = false
|
||||
}
|
||||
try {
|
||||
val response = client.get(healthUrl)
|
||||
if (response.status.value !in 200..299) {
|
||||
throw IllegalStateException("сервер ответил HTTP ${response.status.value} на GET $healthUrl")
|
||||
}
|
||||
val body = response.bodyAsText().trim()
|
||||
if (body != "ok") {
|
||||
throw IllegalStateException("сервер ответил неожиданным телом на GET $healthUrl: '$body'")
|
||||
}
|
||||
} catch (e: IllegalStateException) {
|
||||
throw e
|
||||
} catch (e: Exception) {
|
||||
// На JVM сюда упадут java.net.ConnectException, UnknownHostException,
|
||||
// io.ktor.client.network.sockets.ConnectTimeoutException и т.п.
|
||||
// На нативе native stub падает раньше в platformCreateAgent, так что
|
||||
// сюда мы попадём только под JVM-actual.
|
||||
throw IllegalStateException(
|
||||
"ошибка health-check $healthUrl: ${e::class.simpleName} — ${e.message ?: "(нет сообщения)"}",
|
||||
e,
|
||||
)
|
||||
} finally {
|
||||
client.close()
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -70,6 +124,12 @@ private fun parseCliArgs(args: Array<String>): TuiConfig? {
|
||||
*/
|
||||
internal expect fun platformEnv(key: String): String?
|
||||
|
||||
/**
|
||||
* Создаёт платформенную реализацию [Agent]. JVM actual подключает `:client`
|
||||
* и ходит в HTTP-фасад; native actual пока возвращает stub (см. Platform.native.kt).
|
||||
*/
|
||||
internal expect fun platformCreateAgent(baseUrl: String, id: String): Agent
|
||||
|
||||
private fun printUsage() {
|
||||
val defaultServer = platformEnv("AGENTIK_SERVER") ?: "http://localhost:8080/agentik"
|
||||
val defaultUser = platformEnv("USER") ?: platformEnv("USERNAME") ?: "anon"
|
||||
@@ -86,11 +146,15 @@ private fun printUsage() {
|
||||
--no-history не сохранять состояние
|
||||
--help, -h эта справка
|
||||
|
||||
Переменные среды:
|
||||
AGENTIK_SERVER базовый URL (эквивалент --server)
|
||||
USER / USERNAME используется в id клиента по умолчанию
|
||||
|
||||
В UI:
|
||||
Tab / Shift-Tab переключить фокус между историей и вводом
|
||||
↑ / ↓ скроллить историю / двигать курсор в инпуте
|
||||
← / → двинуть курсор в инпуте
|
||||
Enter отправить сообщение
|
||||
Enter отправить сообщение (создаст новый диалог, если их нет)
|
||||
Ctrl-C / Ctrl-D выйти
|
||||
F1 показать подсказки по горячим клавишам
|
||||
""".trimIndent())
|
||||
|
||||
@@ -1,27 +1,31 @@
|
||||
package pw.binom.agentik.tui
|
||||
|
||||
import androidx.compose.runtime.Composable
|
||||
import androidx.compose.runtime.LaunchedEffect
|
||||
import androidx.compose.runtime.remember
|
||||
import com.jakewharton.mosaic.runMosaicBlocking
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Job
|
||||
import kotlinx.coroutines.SupervisorJob
|
||||
import kotlinx.coroutines.cancel
|
||||
import kotlinx.coroutines.launch
|
||||
import pw.binom.agentik.proto.Agent
|
||||
import kotlin.coroutines.CoroutineContext
|
||||
|
||||
/**
|
||||
* Корневая точка запуска UI. Стартует Mosaic-рантайм и ждёт завершения приложения.
|
||||
* Корневая точка запуска UI. Стартует Mosaic-рантайм, монтирует [TuiBackend] в
|
||||
* его coroutine-scope и ждёт завершения приложения.
|
||||
*
|
||||
* По дизайну — singleton: все остальные модули (UI, бэкенд-корутины) живут внутри
|
||||
* одной Compose-композиции и пользуются её [CoroutineScope].
|
||||
*
|
||||
* Реальный бэкенд (TuiBackend) подключается в следующем коммите: сейчас
|
||||
* стартует на пустом [Agent]-заглушке для smoke-теста.
|
||||
* Бэкенд — единый singleton на процесс: UI-композиция, сетевые подписки и
|
||||
* coroutine job'ы делят scope [runMosaicBlocking] (через [LaunchedEffect]).
|
||||
*/
|
||||
internal class TuiApp(private val config: TuiConfig) {
|
||||
internal class TuiApp(
|
||||
private val config: TuiConfig,
|
||||
private val agent: Agent,
|
||||
) {
|
||||
fun run() {
|
||||
runMosaicBlocking {
|
||||
val state = remember { AppState(config) }
|
||||
val backend = remember { TuiBackend(state = state, agent = agent) }
|
||||
LaunchedEffect(backend) {
|
||||
backend.start(this)
|
||||
}
|
||||
state.attachBackend(backend)
|
||||
App(state = state)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,138 @@
|
||||
package pw.binom.agentik.tui
|
||||
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Job
|
||||
import kotlinx.coroutines.flow.collect
|
||||
import kotlinx.coroutines.launch
|
||||
import pw.binom.agentik.proto.Agent
|
||||
import pw.binom.agentik.proto.Content
|
||||
import pw.binom.agentik.proto.Conversation
|
||||
import pw.binom.agentik.proto.Event
|
||||
import kotlin.coroutines.CoroutineContext
|
||||
import kotlin.time.Instant
|
||||
|
||||
/**
|
||||
* Backend-логика TUI: мост между [Agent] и [AppState].
|
||||
*
|
||||
* Жизненный цикл:
|
||||
* 1. На старте [start] — health-check сделан в [Main] ДО Mosaic; здесь только
|
||||
* пост-сообщение "connected to …".
|
||||
* 2. Подписка на [Agent.events] — обновление списка диалогов в sidebar.
|
||||
* 3. При [onUserMessage] — если текущего диалога нет, создаём
|
||||
* [createConversation] (temp=false, чтобы он персистился на сервере), затем
|
||||
* [send]. Подписка на [Conversation.events] идёт сразу при создании/открытии.
|
||||
*
|
||||
* Дизайн: один backend-объект на процесс, живёт в [runMosaicBlocking]-scope.
|
||||
*/
|
||||
internal class TuiBackend(
|
||||
private val state: AppState,
|
||||
private val agent: Agent,
|
||||
) {
|
||||
/** Текущий открытый диалог, либо `null`, если ещё не выбран. */
|
||||
private var current: Conversation? = null
|
||||
|
||||
/** Активная джоба подписки на [Conversation.events]. */
|
||||
private var eventsJob: Job? = null
|
||||
|
||||
/** Последний виденный момент событий — для переподписки при reconnect. */
|
||||
private var lastSeenAt: Instant = Instant.DISTANT_PAST
|
||||
|
||||
/**
|
||||
* Запускает фоновые подписки в scope [scope] (передаётся из Mosaic
|
||||
* LaunchedEffect'а — это scope recomposer'а, живёт до закрытия UI).
|
||||
*/
|
||||
fun start(scope: CoroutineScope) {
|
||||
this.scope = scope
|
||||
state.postSystem("подключено к ${state.config.server}")
|
||||
scope.launch {
|
||||
try {
|
||||
agent.events(Instant.DISTANT_PAST).collect { /* sidebar refresh */ }
|
||||
} catch (_: kotlinx.coroutines.CancellationException) {
|
||||
// штатная отмена при закрытии UI
|
||||
} catch (e: Exception) {
|
||||
state.postSystem("ошибка live-events: ${e.message ?: e::class.simpleName}")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private lateinit var scope: CoroutineScope
|
||||
|
||||
/**
|
||||
* Обработка пользовательского сообщения, отправленного из input.
|
||||
*
|
||||
* Если текущего диалога нет — создаём его; затем `send`. Подписка на
|
||||
* события конкретного диалога стартует в [ensureConversation].
|
||||
*/
|
||||
fun onUserMessage(text: String) {
|
||||
scope.launch {
|
||||
try {
|
||||
val conv = ensureConversation()
|
||||
conv.send(listOf(Content.Text(text)))
|
||||
} catch (e: Exception) {
|
||||
state.postSystem("ошибка отправки: ${e.message ?: e::class.simpleName}")
|
||||
state.setStreaming(false)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Создаёт [Conversation], если ещё не было; открывает подписку на её события.
|
||||
*/
|
||||
private suspend fun ensureConversation(): Conversation {
|
||||
current?.let { return it }
|
||||
val conv = agent.createConversation(temp = false)
|
||||
state.setConversation(id = conv.id, title = conv.title)
|
||||
subscribeEvents(conv, Instant.DISTANT_PAST)
|
||||
current = conv
|
||||
return conv
|
||||
}
|
||||
|
||||
/**
|
||||
* Подписывается на [Conversation.events] и перенаправляет их в [state].
|
||||
*/
|
||||
private fun subscribeEvents(conv: Conversation, from: Instant) {
|
||||
eventsJob?.cancel()
|
||||
eventsJob = scope.launch {
|
||||
conv.events(from).collect { ev -> dispatch(ev) }
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Маппинг [Event] → [AppState] (что показать в TUI).
|
||||
*
|
||||
* - AppendText → дописывает в последний ассистентский чанк
|
||||
* - StartReasoning / StartResponse → новый streaming-чанк
|
||||
* - End → закрывает streaming
|
||||
* - Interrupted → закрывает streaming + системное сообщение
|
||||
* - ToolCall / ToolResult → сообщения в историю
|
||||
* - Error → системное сообщение
|
||||
*/
|
||||
private fun dispatch(ev: Event) {
|
||||
lastSeenAt = ev.date
|
||||
when (ev) {
|
||||
is Event.AppendText -> state.appendAssistant(ev.body)
|
||||
is Event.StartReasoning -> {
|
||||
state.postSystem("… думаю")
|
||||
}
|
||||
is Event.StartResponse -> state.setStreaming(true)
|
||||
is Event.End -> state.finishAssistant()
|
||||
is Event.Interrupted -> {
|
||||
state.finishAssistant()
|
||||
state.postSystem("прервано")
|
||||
}
|
||||
is Event.AppendImage -> {
|
||||
state.postSystem("[картинка: ${ev.mime}, ${ev.body.size} байт]")
|
||||
}
|
||||
is Event.ToolCall -> {
|
||||
state.postToolCall(toolName = ev.toolName, title = null, args = ev.toolArgs)
|
||||
}
|
||||
is Event.ToolResult -> {
|
||||
state.postToolResult(toolName = "", result = ev.result ?: "")
|
||||
}
|
||||
is Event.Error -> {
|
||||
state.setStreaming(false)
|
||||
state.postSystem("ошибка: ${ev.message}")
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,17 @@
|
||||
package pw.binom.agentik.tui.ui
|
||||
|
||||
import androidx.compose.runtime.Composable
|
||||
import com.jakewharton.mosaic.ui.Text
|
||||
import com.jakewharton.mosaic.ui.TextStyle
|
||||
|
||||
/**
|
||||
* Нижняя подсказка с текущим набором горячих клавиш.
|
||||
*
|
||||
* При открытом help-оверлее показывает заглушку с указателем «наверху».
|
||||
*/
|
||||
@Composable
|
||||
internal fun Footer(showHelp: Boolean) {
|
||||
val hint = if (showHelp) " ↑ наверху help-оверлей ↑ "
|
||||
else " Tab focus ↑↓ scroll Enter send Esc clear F1 help Ctrl-D exit "
|
||||
Text(value = hint, textStyle = TextStyle.Italic)
|
||||
}
|
||||
@@ -0,0 +1,26 @@
|
||||
package pw.binom.agentik.tui.ui
|
||||
|
||||
import androidx.compose.runtime.Composable
|
||||
import androidx.compose.runtime.collectAsState
|
||||
import androidx.compose.runtime.getValue
|
||||
import com.jakewharton.mosaic.ui.Text
|
||||
import com.jakewharton.mosaic.ui.TextStyle
|
||||
import pw.binom.agentik.tui.AppState
|
||||
|
||||
/**
|
||||
* Верхняя инвертированная полоса с идентификатором и текущим фокусом.
|
||||
*
|
||||
* Пример: ` agentik · cli-tui:root · a1b2c3d4… · мой чат · focus=input `
|
||||
*/
|
||||
@Composable
|
||||
internal fun Header(state: AppState, focusIndex: Int) {
|
||||
val title by state.currentTitle.collectAsState()
|
||||
val convId by state.currentConversationId.collectAsState()
|
||||
val focusLabel = when (focusIndex) { 0 -> "input"; 1 -> "history"; 2 -> "sidebar"; else -> "?" }
|
||||
val convStr = convId?.let { " · ${it.take(8)}…" } ?: ""
|
||||
val titleStr = title ?: "(нет диалога)"
|
||||
Text(
|
||||
value = " agentik · ${state.config.id}$convStr · $titleStr · focus=$focusLabel ",
|
||||
textStyle = TextStyle.Bold + TextStyle.Invert,
|
||||
)
|
||||
}
|
||||
@@ -0,0 +1,26 @@
|
||||
package pw.binom.agentik.tui.ui
|
||||
|
||||
import androidx.compose.runtime.Composable
|
||||
import com.jakewharton.mosaic.modifier.Modifier
|
||||
import com.jakewharton.mosaic.ui.Column
|
||||
import com.jakewharton.mosaic.ui.Text
|
||||
import com.jakewharton.mosaic.ui.TextStyle
|
||||
|
||||
/**
|
||||
* Полноэкранный оверлей со списком горячих клавиш.
|
||||
* Включается/выключается по F1 (см. [pw.binom.agentik.tui.App]).
|
||||
*/
|
||||
@Composable
|
||||
internal fun HelpOverlay() {
|
||||
Column(modifier = Modifier) {
|
||||
Text(value = " --- HELP ---", textStyle = TextStyle.Bold + TextStyle.Invert)
|
||||
Text(value = " Tab / Shift-Tab переключить фокус (history / input / sidebar)")
|
||||
Text(value = " ↑ / ↓ скролл истории / курсор в input")
|
||||
Text(value = " ← / → курсор в input")
|
||||
Text(value = " Enter отправить сообщение")
|
||||
Text(value = " Backspace / Del удалить символ")
|
||||
Text(value = " Esc очистить input")
|
||||
Text(value = " Ctrl-D / Ctrl-C выход")
|
||||
Text(value = " F1 toggle help", textStyle = TextStyle.Italic)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,34 @@
|
||||
package pw.binom.agentik.tui.ui
|
||||
|
||||
import androidx.compose.runtime.Composable
|
||||
import androidx.compose.runtime.collectAsState
|
||||
import androidx.compose.runtime.getValue
|
||||
import com.jakewharton.mosaic.ui.Text
|
||||
import pw.binom.agentik.tui.AppState
|
||||
import pw.binom.agentik.tui.TuiMessage
|
||||
|
||||
/**
|
||||
* Прокручиваемый (через клавиатуру) лог диалога.
|
||||
* Каждое сообщение рендерится отдельной строкой с префиксом (см. [renderMessage]).
|
||||
* При пустом списке показывается подсказка.
|
||||
*/
|
||||
@Composable
|
||||
internal fun HistoryPanel(state: AppState) {
|
||||
val messages by state.messages.collectAsState()
|
||||
val rendered = if (messages.isEmpty()) {
|
||||
" (пока пусто)\n Tab — переключить фокус, F1 — подсказки.\n"
|
||||
} else {
|
||||
messages.joinToString("") { renderMessage(it) }
|
||||
}
|
||||
Text(value = rendered)
|
||||
}
|
||||
|
||||
/** Превращает [TuiMessage] в одну строку с префиксом. Потоковые чанки получают курсор `▍`. */
|
||||
internal fun renderMessage(m: TuiMessage): String = when (m) {
|
||||
is TuiMessage.System -> " ── ${m.text}\n"
|
||||
is TuiMessage.User -> " > ${m.text}\n"
|
||||
is TuiMessage.Assistant -> " ╰ ${m.text}\n"
|
||||
is TuiMessage.AssistantStreaming -> " ╰ ${m.text} ▍\n"
|
||||
is TuiMessage.ToolCall -> " ⚙ ${m.toolName}${if (!m.title.isNullOrEmpty()) ": ${m.title}" else ""}\n"
|
||||
is TuiMessage.ToolResult -> " ↳ ${m.result.take(200)}${if (m.result.length > 200) "…" else ""}\n"
|
||||
}
|
||||
@@ -0,0 +1,67 @@
|
||||
package pw.binom.agentik.tui.ui
|
||||
|
||||
import androidx.compose.runtime.Composable
|
||||
import androidx.compose.runtime.collectAsState
|
||||
import androidx.compose.runtime.getValue
|
||||
import com.jakewharton.mosaic.layout.KeyEvent
|
||||
import com.jakewharton.mosaic.layout.drawBehind
|
||||
import com.jakewharton.mosaic.layout.onPreviewKeyEvent
|
||||
import com.jakewharton.mosaic.modifier.Modifier
|
||||
import com.jakewharton.mosaic.ui.Text
|
||||
import pw.binom.agentik.tui.AppState
|
||||
|
||||
/**
|
||||
* Нижняя строка ввода с курсором.
|
||||
* При активном стриме ассистента показывает ` ⋯`, иначе ` >`.
|
||||
*
|
||||
* Содержимое строки: `<prompt> <before>|<cursorChar>|<after>` —
|
||||
* `cursorChar` — это символ, на котором стоит курсор (или пробел в конце).
|
||||
*
|
||||
* Клавиши обрабатываются через [handleInputKey] внутри `onPreviewKeyEvent`.
|
||||
*/
|
||||
@Composable
|
||||
internal fun InputLine(state: AppState) {
|
||||
val text by state.input.collectAsState()
|
||||
val cursor by state.cursor.collectAsState()
|
||||
val streaming by state.streaming.collectAsState()
|
||||
val cursorPos = cursor.coerceIn(0, text.length)
|
||||
val before = text.substring(0, cursorPos)
|
||||
val cursorChar = if (cursorPos < text.length) text[cursorPos].toString() else " "
|
||||
val afterStart = if (cursorPos < text.length) cursorPos + 1 else cursorPos
|
||||
val after = text.substring(afterStart.coerceAtMost(text.length))
|
||||
val prompt = if (streaming) " ⋯" else " >"
|
||||
|
||||
Text(
|
||||
value = "$prompt $before|$cursorChar|${after}",
|
||||
modifier = Modifier
|
||||
.onPreviewKeyEvent { ev -> handleInputKey(state, ev) }
|
||||
.drawBehind {
|
||||
// Snapshot-read state в drawBehind чтобы changes триггерили redraw.
|
||||
state.input.let { /* touch */ }
|
||||
},
|
||||
)
|
||||
}
|
||||
|
||||
/**
|
||||
* Обработка клавиш в [InputLine]. `true` = событие поглощено.
|
||||
*
|
||||
* Не перехватывает клавиши с `alt`/`ctrl` — они идут дальше
|
||||
* на корневой обработчик ([pw.binom.agentik.tui.App]).
|
||||
*/
|
||||
internal fun handleInputKey(state: AppState, ev: KeyEvent): Boolean {
|
||||
if (ev.alt || ev.ctrl) return false
|
||||
return when (ev.key) {
|
||||
"Enter" -> state.submitInput() != null
|
||||
"Backspace" -> { state.inputBackspace(); true }
|
||||
"Delete" -> { state.inputDelete(); true }
|
||||
"Left", "ArrowLeft" -> { state.inputMoveCursor(-1); true }
|
||||
"Right", "ArrowRight" -> { state.inputMoveCursor(+1); true }
|
||||
"Home" -> { state.inputCursorHome(); true }
|
||||
"End" -> { state.inputCursorEnd(); true }
|
||||
else -> {
|
||||
val s = ev.key
|
||||
if (s.length == 1) { state.inputInsert(s); true }
|
||||
else false
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,85 @@
|
||||
package pw.binom.agentik.tui
|
||||
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.flow.MutableSharedFlow
|
||||
import kotlinx.coroutines.flow.emptyFlow
|
||||
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 pw.binom.agentik.proto.MessageContext
|
||||
import kotlin.time.Instant
|
||||
|
||||
/**
|
||||
* Минимальный fake [Agent] для тестов [TuiBackend]: считает, сколько раз
|
||||
* вызвали [createConversation], и отдаёт заранее сконструированные
|
||||
* [FakeConversation].
|
||||
*/
|
||||
internal class FakeAgent(
|
||||
private val conversationFactory: () -> FakeConversation = { FakeConversation() },
|
||||
) : Agent {
|
||||
override val id: String = "fake"
|
||||
var createCount: Int = 0
|
||||
private set
|
||||
val conversations = mutableListOf<FakeConversation>()
|
||||
|
||||
override fun createConversation(temp: Boolean): Conversation {
|
||||
createCount++
|
||||
val c = conversationFactory()
|
||||
conversations += c
|
||||
return c
|
||||
}
|
||||
|
||||
override suspend fun getConversation(id: String): Conversation? =
|
||||
conversations.firstOrNull { it.id == id }
|
||||
|
||||
override suspend fun deleteConversation(id: String): Boolean =
|
||||
conversations.removeAll { it.id == id }
|
||||
|
||||
override suspend fun getConversations(offset: Int, limit: Int): List<Conversation> =
|
||||
conversations.toList()
|
||||
|
||||
override fun events(after: Instant): Flow<AgentEvent> = emptyFlow()
|
||||
}
|
||||
|
||||
/**
|
||||
* [Conversation], запоминающий все вызовы [send] и эмитящий управляемые
|
||||
* [Event] через общий [MutableSharedFlow]. Используется в тестах
|
||||
* [TuiBackend] для проверки маршрутизации событий в UI.
|
||||
*/
|
||||
internal class FakeConversation(
|
||||
override val id: String = "fake-conv",
|
||||
override val title: String? = null,
|
||||
) : Conversation {
|
||||
override val isSupportImageInput: Boolean = false
|
||||
override val isSupportImageOutput: Boolean = false
|
||||
override val isTemporal: Boolean = false
|
||||
override val updatedAt: Instant = Instant.DISTANT_PAST
|
||||
|
||||
val sent = mutableListOf<List<Content>>()
|
||||
val sentContexts = mutableListOf<MessageContext?>()
|
||||
var closed: Boolean = false
|
||||
private set
|
||||
var interrupted: Boolean = false
|
||||
private set
|
||||
|
||||
private val eventsFlow = MutableSharedFlow<Event>(extraBufferCapacity = 64)
|
||||
fun emit(e: Event) { eventsFlow.tryEmit(e) }
|
||||
|
||||
override suspend fun rename(title: String) = Unit
|
||||
|
||||
override suspend fun send(content: List<Content>, context: MessageContext?) {
|
||||
sent += content
|
||||
sentContexts += context
|
||||
}
|
||||
|
||||
override suspend fun interrupt() { interrupted = true }
|
||||
|
||||
override fun events(after: Instant): Flow<Event> = eventsFlow
|
||||
|
||||
override suspend fun getMessages(after: Instant, offset: Int, limit: Int): List<Message> = emptyList()
|
||||
|
||||
override fun close() { closed = true }
|
||||
}
|
||||
@@ -0,0 +1,249 @@
|
||||
package pw.binom.agentik.tui
|
||||
|
||||
import kotlinx.coroutines.ExperimentalCoroutinesApi
|
||||
import kotlinx.coroutines.test.runCurrent
|
||||
import kotlinx.coroutines.test.runTest
|
||||
import pw.binom.agentik.proto.Content
|
||||
import pw.binom.agentik.proto.Event
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertFalse
|
||||
import kotlin.test.assertTrue
|
||||
|
||||
/**
|
||||
* Тесты [TuiBackend]. Используем [runTest.backgroundScope] (а не TestScope)
|
||||
* для передачи в `start` — фоновые подписки должны жить параллельно с
|
||||
* телом теста и автоматически отменяться по его завершении. Иначе
|
||||
* бесконечный collect на `agent.events()` завешивает runTest на 60s
|
||||
* `UncompletedCoroutinesError`.
|
||||
*
|
||||
* [runCurrent] нужен после каждого `onUserMessage` и каждого `emit`,
|
||||
* потому что `backgroundScope` использует свой диспетчер, который не
|
||||
* продвигается через `advanceUntilIdle` — `runCurrent` прогоняет ровно
|
||||
* те задачи, что готовы к запуску сейчас.
|
||||
*/
|
||||
@OptIn(ExperimentalCoroutinesApi::class)
|
||||
class TuiBackendTest {
|
||||
|
||||
private fun fixtureConfig(server: String = "http://localhost:8080/agentik") =
|
||||
TuiConfig(server = server, id = "cli-tui:tester", historyEnabled = true)
|
||||
|
||||
@Test
|
||||
fun `start posts connected system message`() = runTest {
|
||||
val cfg = fixtureConfig()
|
||||
val state = AppState(cfg)
|
||||
val agent = FakeAgent()
|
||||
val backend = TuiBackend(state = state, agent = agent)
|
||||
backend.start(backgroundScope)
|
||||
runCurrent()
|
||||
|
||||
val sysMsgs = state.messages.value.filterIsInstance<TuiMessage.System>()
|
||||
assertTrue(
|
||||
sysMsgs.any { it.text.contains(cfg.server) },
|
||||
"ожидалось системное 'подключено к ${cfg.server}', было: ${sysMsgs.map { it.text }}",
|
||||
)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `first onUserMessage auto-creates conversation with temp=false`() = runTest {
|
||||
val cfg = fixtureConfig()
|
||||
val state = AppState(cfg)
|
||||
val agent = FakeAgent()
|
||||
val backend = TuiBackend(state = state, agent = agent)
|
||||
backend.start(backgroundScope)
|
||||
runCurrent()
|
||||
|
||||
backend.onUserMessage("привет")
|
||||
runCurrent()
|
||||
|
||||
assertEquals(1, agent.createCount, "должен быть один createConversation")
|
||||
assertEquals(listOf("привет"), agent.conversations.first().sent.flattenText())
|
||||
// temp=false — обычный (не временный) диалог: персистится на сервере
|
||||
assertFalse(agent.conversations.first().isTemporal, "диалог не должен быть временным")
|
||||
// state знает id и title нового диалога
|
||||
assertEquals("fake-conv", state.currentConversationId.value)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `second onUserMessage reuses same conversation`() = runTest {
|
||||
val cfg = fixtureConfig()
|
||||
val state = AppState(cfg)
|
||||
val agent = FakeAgent()
|
||||
val backend = TuiBackend(state = state, agent = agent)
|
||||
backend.start(backgroundScope)
|
||||
runCurrent()
|
||||
|
||||
backend.onUserMessage("раз")
|
||||
runCurrent()
|
||||
backend.onUserMessage("два")
|
||||
runCurrent()
|
||||
|
||||
assertEquals(1, agent.createCount, "новый диалог создавать не должны — переиспользуем старый")
|
||||
assertEquals(2, agent.conversations.first().sent.size)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `AppendText appends to current assistant streaming chunk`() = runTest {
|
||||
val cfg = fixtureConfig()
|
||||
val state = AppState(cfg)
|
||||
val conv = FakeConversation()
|
||||
val agent = FakeAgent(conversationFactory = { conv })
|
||||
val backend = TuiBackend(state = state, agent = agent)
|
||||
backend.start(backgroundScope)
|
||||
runCurrent()
|
||||
|
||||
backend.onUserMessage("hi")
|
||||
runCurrent()
|
||||
val now = kotlin.time.Clock.System.now()
|
||||
conv.emit(Event.StartResponse(now, Event.ResponseType.TEXT))
|
||||
conv.emit(Event.AppendText(now, "Привет"))
|
||||
conv.emit(Event.AppendText(now, ", мир"))
|
||||
runCurrent()
|
||||
|
||||
val assistantMsgs = state.messages.value.filterIsInstance<TuiMessage.AssistantStreaming>()
|
||||
assertEquals(1, assistantMsgs.size, "должен быть один streaming-чанк, не два")
|
||||
assertEquals("Привет, мир", assistantMsgs.single().text)
|
||||
assertTrue(state.streaming.value)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `End event finalizes assistant and stops streaming`() = runTest {
|
||||
val cfg = fixtureConfig()
|
||||
val state = AppState(cfg)
|
||||
val conv = FakeConversation()
|
||||
val agent = FakeAgent(conversationFactory = { conv })
|
||||
val backend = TuiBackend(state = state, agent = agent)
|
||||
backend.start(backgroundScope)
|
||||
runCurrent()
|
||||
|
||||
backend.onUserMessage("hi")
|
||||
runCurrent()
|
||||
val now = kotlin.time.Clock.System.now()
|
||||
conv.emit(Event.StartResponse(now, Event.ResponseType.TEXT))
|
||||
conv.emit(Event.AppendText(now, "ответ"))
|
||||
conv.emit(Event.End(now))
|
||||
runCurrent()
|
||||
|
||||
val last = state.messages.value.last()
|
||||
assertTrue(last is TuiMessage.Assistant, "после End последнее сообщение должно стать финальным Assistant, было: ${last::class.simpleName}")
|
||||
assertEquals("ответ", (last as TuiMessage.Assistant).text)
|
||||
assertFalse(state.streaming.value)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `Interrupted event clears streaming and posts system message`() = runTest {
|
||||
val cfg = fixtureConfig()
|
||||
val state = AppState(cfg)
|
||||
val conv = FakeConversation()
|
||||
val agent = FakeAgent(conversationFactory = { conv })
|
||||
val backend = TuiBackend(state = state, agent = agent)
|
||||
backend.start(backgroundScope)
|
||||
runCurrent()
|
||||
|
||||
backend.onUserMessage("hi")
|
||||
runCurrent()
|
||||
val now = kotlin.time.Clock.System.now()
|
||||
conv.emit(Event.StartResponse(now, Event.ResponseType.TEXT))
|
||||
conv.emit(Event.AppendText(now, "часть ответа"))
|
||||
conv.emit(Event.Interrupted(now))
|
||||
runCurrent()
|
||||
|
||||
assertFalse(state.streaming.value)
|
||||
val sysMsgs = state.messages.value.filterIsInstance<TuiMessage.System>()
|
||||
assertTrue(
|
||||
sysMsgs.any { it.text.contains("прервано") },
|
||||
"ожидалось 'прервано' в системных сообщениях, было: ${sysMsgs.map { it.text }}",
|
||||
)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `ToolCall and ToolResult events become visible tool messages`() = runTest {
|
||||
val cfg = fixtureConfig()
|
||||
val state = AppState(cfg)
|
||||
val conv = FakeConversation()
|
||||
val agent = FakeAgent(conversationFactory = { conv })
|
||||
val backend = TuiBackend(state = state, agent = agent)
|
||||
backend.start(backgroundScope)
|
||||
runCurrent()
|
||||
|
||||
backend.onUserMessage("hi")
|
||||
runCurrent()
|
||||
val now = kotlin.time.Clock.System.now()
|
||||
conv.emit(Event.ToolCall(date = now, id = "1", title = null, toolName = "echo", toolArgs = """{"x":1}"""))
|
||||
conv.emit(Event.ToolResult(date = now, id = "1", result = "ok"))
|
||||
runCurrent()
|
||||
|
||||
val toolMsgs = state.messages.value.filterIsInstance<TuiMessage.ToolCall>()
|
||||
val resultMsgs = state.messages.value.filterIsInstance<TuiMessage.ToolResult>()
|
||||
assertEquals(1, toolMsgs.size)
|
||||
assertEquals("echo", toolMsgs.single().toolName)
|
||||
assertEquals("""{"x":1}""", toolMsgs.single().args)
|
||||
assertEquals(1, resultMsgs.size)
|
||||
assertEquals("ok", resultMsgs.single().result)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `Error event posts system message and clears streaming`() = runTest {
|
||||
val cfg = fixtureConfig()
|
||||
val state = AppState(cfg)
|
||||
val conv = FakeConversation()
|
||||
val agent = FakeAgent(conversationFactory = { conv })
|
||||
val backend = TuiBackend(state = state, agent = agent)
|
||||
backend.start(backgroundScope)
|
||||
runCurrent()
|
||||
|
||||
backend.onUserMessage("hi")
|
||||
runCurrent()
|
||||
val now = kotlin.time.Clock.System.now()
|
||||
conv.emit(Event.StartResponse(now, Event.ResponseType.TEXT))
|
||||
conv.emit(Event.Error(date = now, message = "boom"))
|
||||
runCurrent()
|
||||
|
||||
val sysMsgs = state.messages.value.filterIsInstance<TuiMessage.System>()
|
||||
assertTrue(sysMsgs.any { it.text.contains("boom") }, "должно быть 'ошибка: boom'")
|
||||
assertFalse(state.streaming.value)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `onUserMessage does not swallow exceptions - state stays consistent`() = runTest {
|
||||
val cfg = fixtureConfig()
|
||||
val state = AppState(cfg)
|
||||
val agent = FakeAgent(conversationFactory = { error("server kaboom") })
|
||||
val backend = TuiBackend(state = state, agent = agent)
|
||||
backend.start(backgroundScope)
|
||||
runCurrent()
|
||||
|
||||
backend.onUserMessage("hi")
|
||||
runCurrent()
|
||||
|
||||
val sysMsgs = state.messages.value.filterIsInstance<TuiMessage.System>()
|
||||
assertTrue(
|
||||
sysMsgs.any { it.text.contains("ошибка отправки") || it.text.contains("server kaboom") },
|
||||
"должна быть системная ошибка, было: ${sysMsgs.map { it.text }}",
|
||||
)
|
||||
assertFalse(state.streaming.value, "стриминг должен быть выключен в catch-ветке")
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `StartReasoning posts thinking system message`() = runTest {
|
||||
val cfg = fixtureConfig()
|
||||
val state = AppState(cfg)
|
||||
val conv = FakeConversation()
|
||||
val agent = FakeAgent(conversationFactory = { conv })
|
||||
val backend = TuiBackend(state = state, agent = agent)
|
||||
backend.start(backgroundScope)
|
||||
runCurrent()
|
||||
|
||||
backend.onUserMessage("hi")
|
||||
runCurrent()
|
||||
val now = kotlin.time.Clock.System.now()
|
||||
conv.emit(Event.StartReasoning(now))
|
||||
runCurrent()
|
||||
|
||||
val sysMsgs = state.messages.value.filterIsInstance<TuiMessage.System>()
|
||||
assertTrue(sysMsgs.any { it.text.contains("думаю") })
|
||||
}
|
||||
}
|
||||
|
||||
private fun List<List<Content>>.flattenText(): List<String> =
|
||||
map { cs -> cs.filterIsInstance<Content.Text>().joinToString("") { it.body } }
|
||||
@@ -4,11 +4,9 @@ import pw.binom.agentik.client.AgentikAgent
|
||||
import pw.binom.agentik.proto.Agent
|
||||
|
||||
/**
|
||||
* Платформенная фабрика [Agent]. JVM-only пока: native не подключали ktor-движки.
|
||||
* Платформенные actual'ы для JVM. Используется `:client` поверх Ktor CIO.
|
||||
*/
|
||||
internal actual fun platformEnv(key: String): String? = System.getenv(key)
|
||||
|
||||
/**
|
||||
* Реализация [TuiApp.createAgent] для JVM — обычный ktor-cio через `:client`.
|
||||
*/
|
||||
internal fun jvmCreateAgent(baseUrl: String, id: String): Agent = AgentikAgent(id = id, baseUrl = baseUrl)
|
||||
internal actual fun platformCreateAgent(baseUrl: String, id: String): Agent =
|
||||
AgentikAgent(id = id, baseUrl = baseUrl)
|
||||
|
||||
@@ -9,5 +9,5 @@ import pw.binom.agentik.proto.Agent
|
||||
*/
|
||||
internal actual fun platformEnv(key: String): String? = null
|
||||
|
||||
internal fun nativeCreateAgent(baseUrl: String, id: String): Agent =
|
||||
internal actual fun platformCreateAgent(baseUrl: String, id: String): Agent =
|
||||
error("agentik-tui native target is not implemented yet (baseUrl=$baseUrl)")
|
||||
|
||||
+6
-4
@@ -11,14 +11,16 @@ group = "pw.binom.agentik"
|
||||
// прочитать через rootProject.extra["projectVersion"].
|
||||
|
||||
// Publication version: -Pversion=<tag> (CICD publishes by release tag) с
|
||||
// fallback в gradle.properties ("version=0.1.0"). Без версии maven-publish падает
|
||||
// с "Invalid publication 'kotlinMultiplatform': version cannot be empty" —
|
||||
// fallback в gradle.properties (ключ `agentik.version.default`, не `version`
|
||||
// — иначе Gradle-мерж gradle.properties и -Pversion= отдаёт приоритет
|
||||
// gradle.properties). Без версии maven-publish падает с
|
||||
// "Invalid publication 'kotlinMultiplatform': version cannot be empty" —
|
||||
// это известный gotcha: subprojects читают rootProject.version ДО того, как
|
||||
// if-блок ниже успевает его установить. Фикс: provider+orElse вычисляется
|
||||
// eagerly, и subprojects получают готовую строку.
|
||||
val projectVersion: String = providers.gradleProperty("version")
|
||||
.map { it.trimStart('v', 'V') } // strip optional "v" prefix from tag
|
||||
.getOrElse("0.1.0")
|
||||
.getOrElse(providers.gradleProperty("agentik.version.default").orElse("0.1.0-SNAPSHOT").get())
|
||||
version = projectVersion
|
||||
extra["projectVersion"] = projectVersion
|
||||
|
||||
@@ -49,7 +51,7 @@ val moduleDescriptions: Map<String, String> = mapOf(
|
||||
"storage-sqlite" to "agentik :storage-sqlite — SQLDelight реализация всех сторов на SQLite (прод-бэкенд).",
|
||||
"agent-toolsets" to "agentik :agent-toolsets — реестр инструментов + диспетчер тулов (enable_toolset/disable_toolset); переиспользуемое ядро.",
|
||||
"agentik-cli" to "agentik :agentik-cli — JVM CLI-клиент (JLine) к /agentik: REPL + slash-команды + стрим SSE.",
|
||||
"agentik-tui" to "agentik :agentik-tui — Compose-for-Mosaic TUI-клиент (desktop, без iOS) с клавиатурной навигацией без ':'-префиксов.",
|
||||
// "agentik-tui" to "agentik :agentik-tui — Compose-for-Mosaic TUI-клиент (отключён 2026-09-17)."
|
||||
"standalone" to "agentik :standalone — single-jar HTTP-сервер со всеми транспортами (AG-UI/A2A/:proto), SQLite, памятью, скилами и SOUL.",
|
||||
)
|
||||
rootProject.extra.set("moduleDescriptions", moduleDescriptions)
|
||||
|
||||
+1
-1
@@ -19,7 +19,7 @@ JVM, iOS, macOS, Linux, Windows.
|
||||
## Где используется
|
||||
|
||||
- `:agentik-cli` — REPL.
|
||||
- `:agentik-tui` — Compose-for-Mosaic клиент.
|
||||
- `:agentik-cli` — JVM/native CLI-клиент поверх `:client`.
|
||||
- Любой внешний KMP-проект, который хочет встроить агента в свой UI.
|
||||
|
||||
## Как подключить
|
||||
|
||||
+43
-18
@@ -1,25 +1,50 @@
|
||||
import org.jetbrains.kotlin.gradle.dsl.JvmTarget
|
||||
|
||||
plugins {
|
||||
alias(libs.plugins.kotlin.jvm)
|
||||
alias(libs.plugins.kotlin.multiplatform)
|
||||
alias(libs.plugins.kotlin.serialization)
|
||||
}
|
||||
|
||||
kotlin {
|
||||
compilerOptions {
|
||||
jvmTarget.set(JvmTarget.JVM_21)
|
||||
jvmToolchain(21)
|
||||
|
||||
// Только то, что нам реально нужно: JVM + 5 desktop-native. iOS не входит —
|
||||
// :client не имеет смысла на iOS, а :agentik-cli использует :client и тоже
|
||||
// без iOS. См. agentik-cli/build.gradle.kts.
|
||||
jvm()
|
||||
listOf(
|
||||
macosX64(),
|
||||
macosArm64(),
|
||||
linuxX64(),
|
||||
linuxArm64(),
|
||||
mingwX64(),
|
||||
)
|
||||
|
||||
sourceSets {
|
||||
commonMain.dependencies {
|
||||
api(project(":proto"))
|
||||
|
||||
api(libs.ktor.client.core)
|
||||
implementation(libs.ktor.client.content.negotiation)
|
||||
implementation(libs.ktor.serialization.kotlinx.json)
|
||||
|
||||
implementation(libs.kotlinx.coroutines.core)
|
||||
implementation(libs.kotlinx.serialization.core)
|
||||
implementation(libs.kotlinx.serialization.json)
|
||||
}
|
||||
commonTest.dependencies {
|
||||
implementation(libs.kotlin.test)
|
||||
implementation(libs.kotlinx.coroutines.test)
|
||||
implementation(libs.ktor.server.core)
|
||||
implementation(libs.ktor.server.test.host)
|
||||
implementation(libs.ktor.client.content.negotiation)
|
||||
implementation(libs.ktor.client.cio)
|
||||
implementation(libs.ktor.server.cio)
|
||||
implementation(libs.ktor.server.sse)
|
||||
}
|
||||
jvmTest.dependencies {
|
||||
implementation("junit:junit:4.13.2")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
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.core)
|
||||
implementation(libs.kotlinx.serialization.json)
|
||||
|
||||
// :client — это библиотека, не executable. Native-бинари объявляются
|
||||
// в :agentik-cli (он зависит от :client и реально предоставляет main).
|
||||
}
|
||||
|
||||
+10
-7
@@ -4,6 +4,7 @@ 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.prepareGet
|
||||
import io.ktor.client.request.parameter
|
||||
import io.ktor.client.request.post
|
||||
import io.ktor.client.request.setBody
|
||||
@@ -66,13 +67,15 @@ internal class AgentClient(
|
||||
}
|
||||
|
||||
override fun events(after: Instant): Flow<AgentEvent> = 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))
|
||||
httpClient.prepareGet("$agentUrl/events?after=$after") { noSseReadTimeout() }
|
||||
.execute { response ->
|
||||
check(response.status == HttpStatusCode.OK) {
|
||||
"events: server returned ${response.status}"
|
||||
}
|
||||
readSse(response.bodyAsChannel())
|
||||
.collect { payload ->
|
||||
emit(agentikJson.decodeFromString(AgentEvent.serializer(), payload))
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
+8
-16
@@ -1,9 +1,6 @@
|
||||
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
|
||||
|
||||
/**
|
||||
@@ -11,9 +8,11 @@ import pw.binom.agentik.proto.Agent
|
||||
* (модуль `:server`).
|
||||
*
|
||||
* ```
|
||||
* val http = HttpClient(CIO) { applyAgentikDefaults(token = "s3cret") }
|
||||
* val client = AgentikAgent(
|
||||
* id = "my-agent",
|
||||
* baseUrl = "http://localhost:8080/agentik",
|
||||
* httpClient = http,
|
||||
* )
|
||||
* val conv = client.createConversation(temp = false)
|
||||
* conv.send(listOf(Content.Text("hi")))
|
||||
@@ -24,20 +23,13 @@ import pw.binom.agentik.proto.Agent
|
||||
* агента не знает, поэтому клиент должен её знать сам (или взять из
|
||||
* конфига).
|
||||
*
|
||||
* [httpClient] по умолчанию — [defaultAgentikHttpClient] (CIO + JSON +
|
||||
* SSE). Можно передать свой, если нужен свой engine/логирование/аутентификация.
|
||||
* Клиент приходит снаружи: `:client` не выбирает движок. Собрать [HttpClient]
|
||||
* можно через [agentikHttpClient] (фабрика движка + опциональные движковые
|
||||
* настройки) или вручную, применив к блоку конфигурации [applyAgentikDefaults]
|
||||
* (JSON + опциональный Bearer-токен).
|
||||
*/
|
||||
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) }
|
||||
}
|
||||
httpClient: HttpClient,
|
||||
): Agent = AgentClient(httpClient = httpClient, baseUrl = baseUrl, id = id)
|
||||
+13
-7
@@ -6,6 +6,7 @@ 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.prepareGet
|
||||
import io.ktor.client.request.setBody
|
||||
import io.ktor.client.statement.bodyAsChannel
|
||||
import io.ktor.http.ContentType
|
||||
@@ -72,13 +73,18 @@ internal class ConversationClient(
|
||||
}
|
||||
|
||||
override fun events(after: Instant): Flow<Event> = 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))
|
||||
// prepareGet + execute (а не get) обязателен: `get` дожидается полного
|
||||
// тела ответа, а SSE-поток не заканчивается никогда — вызов висел бы
|
||||
// вечно. `execute` отдаёт HttpResponse со стриминговым bodyAsChannel.
|
||||
httpClient.prepareGet("$convUrl/events?after=$after") { noSseReadTimeout() }
|
||||
.execute { response ->
|
||||
check(response.status == HttpStatusCode.OK) {
|
||||
"events: server returned ${response.status}"
|
||||
}
|
||||
readSse(response.bodyAsChannel())
|
||||
.collect { payload ->
|
||||
emit(agentikJson.decodeFromString(Event.serializer(), payload))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,60 @@
|
||||
package pw.binom.agentik.client
|
||||
|
||||
import io.ktor.client.HttpClient
|
||||
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.request.header
|
||||
import io.ktor.http.HttpHeaders
|
||||
import io.ktor.serialization.kotlinx.json.json
|
||||
|
||||
/**
|
||||
* Общая конфигурация HTTP-клиента agentik — платформо-независимая часть.
|
||||
*
|
||||
* `:client` НЕ выбирает движок: его приносит потребитель. Здесь живёт только то,
|
||||
* без чего клиент несовместим с `/agentik`:
|
||||
* - JSON-конфиг [agentikJson] (обязан совпадать с серверным);
|
||||
* - при заданном [token] — `Authorization: Bearer <token>` на ВСЕ запросы
|
||||
* через [DefaultRequest] (накрывает 10 REST-вызовов и оба SSE-потока;
|
||||
* заголовок живёт на клиенте, а не в отдельных запросах).
|
||||
*
|
||||
* `null` — авторизация выключена, заголовок не отправляется.
|
||||
*
|
||||
* Потребитель, знающий свой движок, добавляет к этому движковые настройки, напр.:
|
||||
* ```
|
||||
* val http = HttpClient(CIO) {
|
||||
* engine { requestTimeout = 0 } // CIO-специфика, живёт у потребителя
|
||||
* applyAgentikDefaults(token)
|
||||
* }
|
||||
* ```
|
||||
*/
|
||||
fun HttpClientConfig<*>.applyAgentikDefaults(token: String? = null) {
|
||||
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,31 @@
|
||||
package pw.binom.agentik.client
|
||||
|
||||
import io.ktor.client.plugins.HttpTimeoutConfig
|
||||
import io.ktor.client.plugins.HttpTimeoutCapability
|
||||
import io.ktor.client.request.HttpRequestBuilder
|
||||
|
||||
/**
|
||||
* Отключает request/connect/socket-таймауты для конкретного запроса через
|
||||
* [HttpTimeoutCapability] со всеми таймаутами = [HttpTimeoutConfig.INFINITE_TIMEOUT_MS].
|
||||
*
|
||||
* Зачем: наш SSE-ридер ([readSse]) читает `bodyAsChannel()` руками и не
|
||||
* использует плагин `SSE`, поэтому движок не считает запрос SSE-шным
|
||||
* (`HttpRequestBuilder.supportsRequestTimeout` проверяет
|
||||
* `body is SSEClientContent`, а у нас тело — обычный GET без тела).
|
||||
* Без capability встроенный `CIOEngineConfig.requestTimeout` (по умолчанию
|
||||
* **15000 мс**) молча убивает долгий idle-стрим через 15 секунд.
|
||||
*
|
||||
* Конфиг создаётся заново на каждый вызов — плагин `HttpTimeout` при
|
||||
* установленном capability мутирует его поля через `?:`, так что шаренный
|
||||
* инстанс мог бы утечь между запросами.
|
||||
*/
|
||||
internal fun HttpRequestBuilder.noSseReadTimeout() {
|
||||
setCapability(
|
||||
HttpTimeoutCapability,
|
||||
HttpTimeoutConfig(
|
||||
requestTimeoutMillis = HttpTimeoutConfig.INFINITE_TIMEOUT_MS,
|
||||
connectTimeoutMillis = HttpTimeoutConfig.INFINITE_TIMEOUT_MS,
|
||||
socketTimeoutMillis = HttpTimeoutConfig.INFINITE_TIMEOUT_MS,
|
||||
),
|
||||
)
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,148 @@
|
||||
package pw.binom.agentik.client
|
||||
|
||||
import io.ktor.client.HttpClient
|
||||
import io.ktor.client.engine.cio.CIO
|
||||
import io.ktor.client.plugins.HttpRequestTimeoutException
|
||||
import io.ktor.client.request.header
|
||||
import io.ktor.client.request.prepareGet
|
||||
import io.ktor.client.statement.bodyAsChannel
|
||||
import io.ktor.server.application.call
|
||||
import io.ktor.server.engine.embeddedServer
|
||||
import io.ktor.server.response.respondBytesWriter
|
||||
import io.ktor.server.routing.get
|
||||
import io.ktor.server.routing.routing
|
||||
import io.ktor.http.ContentType
|
||||
import io.ktor.utils.io.writeStringUtf8
|
||||
import io.ktor.utils.io.readUTF8Line
|
||||
import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import kotlinx.coroutines.withTimeout
|
||||
import java.net.ServerSocket
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertFalse
|
||||
import kotlin.test.assertNotNull
|
||||
import kotlin.test.assertTrue
|
||||
import kotlin.test.fail
|
||||
|
||||
/**
|
||||
* Репродукция бага Ktor CIO: дефолтный [io.ktor.client.engine.cio.CIOEngineConfig.requestTimeout]
|
||||
* = 15 с убивает SSE read. Наш fix — [noSseReadTimeout] ставит capability
|
||||
* [io.ktor.client.plugins.HttpTimeoutCapability] со всеми таймаутами = INFINITE
|
||||
* перед каждым read-стримом.
|
||||
*
|
||||
* Тест запускает встроенный Ktor CIO-сервер на свободном порту. Сервер шлёт
|
||||
* "hello", ждёт 20 с (дольше дефолтного requestTimeout = 15 с), затем шлёт
|
||||
* "done". Без capability клиент отвалился бы на ~15 с; с capability — второе
|
||||
* сообщение доходит.
|
||||
*
|
||||
* Читаем строки пока не найдём "data: done" или пока не сработает
|
||||
* [withTimeout] (18 с — запас над server delay 20 с).
|
||||
*/
|
||||
class SseTimeoutTest {
|
||||
|
||||
private fun freePort(): Int = ServerSocket(0).use { it.localPort }
|
||||
|
||||
@Test
|
||||
fun `sse read survives past default cio timeout with noSseReadTimeout`(): Unit = runBlocking {
|
||||
val port = freePort()
|
||||
val server = embeddedServer(io.ktor.server.cio.CIO, port = port) {
|
||||
routing {
|
||||
get("/sse") {
|
||||
call.respondBytesWriter(contentType = ContentType.Text.EventStream) {
|
||||
writeStringUtf8("data: hello\n\n")
|
||||
flush()
|
||||
// 17 с — чуть больше дефолтного CIO requestTimeout = 15 с.
|
||||
// Если capability сломана, клиент упадёт на 15 с и не получит "done".
|
||||
delay(17_000)
|
||||
writeStringUtf8("data: done\n\n")
|
||||
}
|
||||
}
|
||||
}
|
||||
}.start(wait = false)
|
||||
|
||||
try {
|
||||
val client = HttpClient(CIO)
|
||||
val received = mutableListOf<String>()
|
||||
client.prepareGet("http://127.0.0.1:$port/sse") {
|
||||
header("Accept", "text/event-stream")
|
||||
noSseReadTimeout()
|
||||
}.execute { resp ->
|
||||
val ch = resp.bodyAsChannel()
|
||||
// 19 с запас: ждём, пока сервер пошлёт "done" после 17 с.
|
||||
// Если capability сломана, клиент упадёт на 15 с и мы словим исключение.
|
||||
val deadline = 19_000L
|
||||
val start = System.currentTimeMillis()
|
||||
while (System.currentTimeMillis() - start < deadline) {
|
||||
val line = withTimeout<String?>(deadline) { ch.readUTF8Line() } ?: break
|
||||
if (line.startsWith("data: ")) {
|
||||
received.add(line)
|
||||
}
|
||||
if (line == "data: done") break
|
||||
}
|
||||
}
|
||||
assertTrue(received.contains("data: hello"), "должно получить hello: $received")
|
||||
assertTrue(
|
||||
received.contains("data: done"),
|
||||
"должно получить done (SSE read не должен падать на 15 с): $received",
|
||||
)
|
||||
assertFalse(
|
||||
received.any { it == "<timeout>" },
|
||||
"SSE read упал в timeout (capability не сработал): $received",
|
||||
)
|
||||
} finally {
|
||||
server.stop(100, 200)
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Контр-тест: убеждаемся что БЕЗ [noSseReadTimeout] дефолтный
|
||||
* CIO requestTimeout = 15 с действительно убивает SSE-стрим.
|
||||
* Сервер держит stream 17 с; если клиент не выставил capability —
|
||||
* мы должны получить [HttpRequestTimeoutException] на ~15 с, не
|
||||
* дожидаясь "done".
|
||||
*/
|
||||
@Test
|
||||
fun `without noSseReadTimeout default cio requestTimeout kills the stream`(): Unit = runBlocking {
|
||||
val port = freePort()
|
||||
val server = embeddedServer(io.ktor.server.cio.CIO, port = port) {
|
||||
routing {
|
||||
get("/sse") {
|
||||
call.respondBytesWriter(contentType = ContentType.Text.EventStream) {
|
||||
writeStringUtf8("data: hello\n\n")
|
||||
flush()
|
||||
delay(17_000)
|
||||
writeStringUtf8("data: done\n\n")
|
||||
}
|
||||
}
|
||||
}
|
||||
}.start(wait = false)
|
||||
|
||||
try {
|
||||
val client = HttpClient(CIO)
|
||||
val start = System.currentTimeMillis()
|
||||
try {
|
||||
client.prepareGet("http://127.0.0.1:$port/sse") {
|
||||
header("Accept", "text/event-stream")
|
||||
// НАМЕРЕННО без noSseReadTimeout.
|
||||
}.execute { resp ->
|
||||
val ch = resp.bodyAsChannel()
|
||||
// Читаем строки, пока не придёт "data: done" — без capability
|
||||
// клиент упадёт на ~15 с до того, как сервер пошлёт done.
|
||||
while (true) {
|
||||
val line = ch.readUTF8Line() ?: break
|
||||
if (line == "data: done") break
|
||||
}
|
||||
}
|
||||
fail("без capability клиент должен словить HttpRequestTimeoutException")
|
||||
} catch (e: HttpRequestTimeoutException) {
|
||||
val elapsed = System.currentTimeMillis() - start
|
||||
assertTrue(
|
||||
elapsed in 14_000..17_000,
|
||||
"timeout должен сработать в районе 15 с (default), elapsed=$elapsed",
|
||||
)
|
||||
}
|
||||
} finally {
|
||||
server.stop(100, 200)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,158 @@
|
||||
# 01 — Слои модулей (целевое состояние)
|
||||
|
||||
Целевая модульная структура agentik. Снизу вверх:
|
||||
**приложения → runtime → домен → абстракции → платформенные impl**.
|
||||
|
||||

|
||||
|
||||
PlantUML source (для редактирования; требует Graphviz `dot` для рендеринга):
|
||||
|
||||
```plantuml
|
||||
@startuml agentik-module-layers
|
||||
skinparam componentStyle rectangle
|
||||
skinparam ranksep 60
|
||||
skinparam nodesep 30
|
||||
skinparam packageStyle rectangle
|
||||
|
||||
title agentik — слои модулей (целевое состояние)
|
||||
|
||||
' --- Applications: entry points (thin wrappers) ---
|
||||
package "Applications\n(entry points, тонкие)" {
|
||||
[Standalone\nHTTP+AG-UI+A2A] as Standalone
|
||||
[AgentikCli\nREPL] as Cli
|
||||
[AgentikAndroid\nCompose UI] as Android
|
||||
}
|
||||
|
||||
' --- Agent runtime ---
|
||||
package "Agent Runtime\n(композиция, lifecycle)" {
|
||||
[AgentCore\nBaseAgent] as AgentCore
|
||||
[AgentBuilder\nDSL] as Builder
|
||||
}
|
||||
|
||||
' --- Background work ---
|
||||
package "Background Work\n(event-driven triggers)" {
|
||||
[BackgroundEvents\nbus + events] as Ev
|
||||
[BackgroundScheduler\npolicy] as Sched
|
||||
}
|
||||
|
||||
' --- Domain logic (generic, переиспользуется) ---
|
||||
package "Domain Logic\n(generic tools)" {
|
||||
[LlmTools\nReflector/Reviewer/Miner] as LlmT
|
||||
[McpBridge\nMCP-SDK → LiteTool] as Mcp
|
||||
[Skills\nparse + store] as Skills
|
||||
}
|
||||
|
||||
' --- Storage abstractions + impls ---
|
||||
package "Storage\n(abstractions)" as StoragePkg {
|
||||
[StorageCore\ninterfaces] as StorageCore
|
||||
}
|
||||
|
||||
package "Storage\n(JVM impls)" {
|
||||
[StorageSqlite\nJDBC] as StorageSql
|
||||
[StorageInmemory\ntests] as StorageInmem
|
||||
}
|
||||
|
||||
package "Storage\n(Android impl)" {
|
||||
[StorageSqliteAndroid\nRoom/sqlite] as StorageSqlA
|
||||
}
|
||||
|
||||
' --- Memory backends ---
|
||||
package "Memory\n(abstractions)" {
|
||||
[MemoryApi\nMemorySystem/MemoryTools] as MemApi
|
||||
}
|
||||
|
||||
package "Memory\n(impls)" {
|
||||
[MemoryMd\nHermes §-files] as MemMd
|
||||
[MemoryVector\nJVector+JVM] as MemVec
|
||||
[MemoryVectorAndroid\nONNX+ANN] as MemVecA
|
||||
}
|
||||
|
||||
' --- LLM backends ---
|
||||
package "LLM\n(abstractions)" {
|
||||
[LitertApi\nLiteLlm контракт] as Litert
|
||||
}
|
||||
|
||||
package "LLM\n(impls)" {
|
||||
[LitertOpenai\nHTTP] as LitertO
|
||||
[LitertGoogle\nLiteRT JVM] as LitertG
|
||||
[LitertAndroid\nLiteRT Android] as LitertA
|
||||
}
|
||||
|
||||
' --- Inter-app protocol ---
|
||||
package "Inter-app" {
|
||||
[Proto\nAgent/Conversation] as Proto
|
||||
[A2AServer] as A2A
|
||||
}
|
||||
|
||||
' --- Зависимости (приложения → runtime → домен → абстракции → платформенные импл) ---
|
||||
Standalone ..> Builder
|
||||
Cli ..> Builder
|
||||
Android ..> Builder
|
||||
|
||||
Builder ..> AgentCore
|
||||
AgentCore ..> Proto
|
||||
AgentCore ..> StorageCore
|
||||
AgentCore ..> MemApi
|
||||
AgentCore ..> Litert
|
||||
AgentCore ..> Mcp
|
||||
AgentCore ..> Skills
|
||||
|
||||
Sched ..> Ev
|
||||
AgentCore ..> Sched
|
||||
AgentCore ..> Ev
|
||||
|
||||
Mcp ..> Litert
|
||||
LlmT ..> Litert
|
||||
|
||||
MemMd ..> MemApi
|
||||
MemVec ..> MemApi
|
||||
MemVecA ..> MemApi
|
||||
|
||||
StorageSql ..> StorageCore
|
||||
StorageInmem ..> StorageCore
|
||||
StorageSqlA ..> StorageCore
|
||||
|
||||
LitertO ..> Litert
|
||||
LitertG ..> Litert
|
||||
LitertA ..> Litert
|
||||
|
||||
Standalone ..> A2A
|
||||
Standalone ..> LitertO
|
||||
Standalone ..> LitertG
|
||||
Standalone ..> StorageSql
|
||||
Standalone ..> MemMd
|
||||
Standalone ..> MemVec
|
||||
Standalone ..> Mcp
|
||||
|
||||
Android ..> LitertA
|
||||
Android ..> StorageSqlA
|
||||
Android ..> MemMd
|
||||
Android ..> MemVecA
|
||||
|
||||
@enduml
|
||||
```
|
||||
|
||||
## Что показывает
|
||||
|
||||
- **Applications** — три точки входа: web-сервер, CLI REPL, Android-приложение. Каждое тонкое, не содержит бизнес-логики.
|
||||
- **Agent Runtime** — `BaseAgent` + `AgentBuilder` DSL. Вся композиция и lifecycle.
|
||||
- **Background Work** — `BackgroundEvents` (event-bus) + `BackgroundScheduler` (policy подписки). Event-driven, не interval-polling.
|
||||
- **Domain Logic** — generic переиспользуемые модули (`:llm-tools`, `:mcp-bridge`, `:skills`).
|
||||
- **Storage / Memory / LLM** — каждая с абстракцией и одним или несколькими impl (JVM-only или Android-only).
|
||||
- **Inter-app** — `:proto` контракты + `:a2a-server` для межагентного общения.
|
||||
|
||||
## Текущее состояние vs целевое
|
||||
|
||||
✅ Уже сделано (в этом цикле правок):
|
||||
- `:llm-tools` extracted
|
||||
- `:mcp-bridge` extracted
|
||||
- `BackgroundScheduler` стал event-driven
|
||||
- `ConversationLoop` стал отдельным компонентом (typealias `ChatConversation`)
|
||||
|
||||
⏳ Не сделано:
|
||||
- `:agent-core` (выделить `BaseAgent` + builder в отдельный KMP-модуль)
|
||||
- `:background-events` (выделить events + scheduler — пока в `:standalone`)
|
||||
- `:storage-sqlite-android`
|
||||
- `:memory-vector-android`
|
||||
- `:litert-android`
|
||||
- `:agentik-android` (само приложение)
|
||||
File diff suppressed because one or more lines are too long
|
After Width: | Height: | Size: 38 KiB |
@@ -0,0 +1,127 @@
|
||||
# 02 — Agent Builder: композиция (целевое API)
|
||||
|
||||
Как `AgentBuilder` собирает `BaseAgent` из компонентов. **Memory backend сам объявляет свои tools** — builder их авто-мержит. BackgroundScheduler подписан на события, не interval-poll.
|
||||
|
||||

|
||||
|
||||
PlantUML source (для редактирования; требует Graphviz `dot` для рендеринга):
|
||||
|
||||
```plantuml
|
||||
@startuml agent-composition
|
||||
skinparam componentStyle rectangle
|
||||
|
||||
title Agent Builder — композиция (целевое API)
|
||||
|
||||
' --- Builder ---
|
||||
rectangle "AgentBuilder" as Builder {
|
||||
rectangle "llm: LiteLlm (обязательно)" as Llm
|
||||
rectangle "storage: StorageBundle (обязательно)" as Storage
|
||||
rectangle "memory: MemorySystem (обязательно)" as Mem
|
||||
rectangle "soul: SoulProvider (default NoopSoul)" as Soul
|
||||
rectangle "tools: List<NamedTool> (авто-сборка из backends)" as Tools
|
||||
rectangle "background: BackgroundConfig (default EmptyBg)" as Bg
|
||||
rectangle "toolset: List<ToolsetContribution> (default empty)" as Ts
|
||||
}
|
||||
|
||||
' --- Backends with their tool side-effects ---
|
||||
rectangle "MemoryMd" as MdMem {
|
||||
interface "MemorySystem" as MemSys
|
||||
interface "List<NamedTool>" as MdTools
|
||||
note right
|
||||
MemoryMd.exposesTools() →
|
||||
memory_save / memory_read /
|
||||
memory_list / memory_delete
|
||||
end note
|
||||
}
|
||||
|
||||
rectangle "McpRegistry" as McpReg {
|
||||
interface "List<NamedTool>" as McpTools
|
||||
note right
|
||||
McpRegistry.namedTools →
|
||||
server__tool1, server__tool2,
|
||||
...
|
||||
end note
|
||||
}
|
||||
|
||||
rectangle "BackgroundEvents" as Events {
|
||||
interface "MutableSharedFlow<CompactionEvent|ToolCallEvent|LifecycleEvent>" as Flow
|
||||
note right
|
||||
Эмитится из:
|
||||
- CompactionCoordinator
|
||||
- ToolDispatcher
|
||||
- ConversationLoop.close()
|
||||
end note
|
||||
}
|
||||
|
||||
rectangle "BackgroundScheduler" as Sched {
|
||||
interface "policy: trigger + debounce" as Policy
|
||||
note right
|
||||
Подписан на Events.
|
||||
НИКАКОГО interval-polling.
|
||||
end note
|
||||
}
|
||||
|
||||
' --- Получаемый Agent ---
|
||||
rectangle "BaseAgent\n(impl: ConversationLoop)" as Agent {
|
||||
rectangle "send / interrupt / events" as API
|
||||
rectangle "BackgroundScheduler\nподписка" as Sub
|
||||
}
|
||||
|
||||
' --- Стрелки зависимостей ---
|
||||
Builder --> Llm
|
||||
Builder --> Storage
|
||||
Builder --> Mem
|
||||
Builder --> Soul
|
||||
Builder --> Tools
|
||||
Builder --> Bg
|
||||
Builder --> Ts
|
||||
|
||||
MdMem --> Mem : implements
|
||||
MdMem --> MdTools : exposes
|
||||
|
||||
McpReg --> McpTools : exposes
|
||||
|
||||
Bg --> Events : subscribes-to
|
||||
Bg --> Sched : holds
|
||||
|
||||
Tools <-- MdTools : auto-merge
|
||||
Tools <-- McpTools : auto-merge
|
||||
|
||||
Builder --> Agent : build()
|
||||
Agent --> API
|
||||
Agent --> Sub
|
||||
|
||||
@enduml
|
||||
```
|
||||
|
||||
## Целевой Kotlin DSL
|
||||
|
||||
```kotlin
|
||||
val agent = agentBuilder {
|
||||
// Обязательные
|
||||
llm(OpenAiLlm.fromEnv()) // или LitertAndroid.onDevice(context)
|
||||
storage(SqliteStorage(path)) // или SqliteStorage.android(context)
|
||||
memory(MemoryMd(root = "~/memory")) // или MemoryVector(embedding = HttpEmbedding(...))
|
||||
|
||||
// Опциональные
|
||||
soul(FileSoul("~/SOUL.md")) // или HttpSoul(url), NoopSoul()
|
||||
background {
|
||||
// triggers: OnClosing (reflection+mining), OnCompaction(minTurns=10, mining=true)
|
||||
// event-driven, не interval
|
||||
}
|
||||
tools {
|
||||
// memoryMd.exposesTools() + mcpRegistry.namedTools авто-подцепляются
|
||||
+FileReadTool(root = "/data")
|
||||
}
|
||||
toolset {
|
||||
+MemoryToolsToolset(memoryMd)
|
||||
}
|
||||
}.build()
|
||||
```
|
||||
|
||||
## Ключевые решения
|
||||
|
||||
- **`MemoryBackend.exposesTools()`** — backend декларирует свои tools. Не «подставить любой backend», а «backend сообщает что он умеет». Это убирает coupling «какие tools совместимы с какими backends».
|
||||
- **Builder требует только `llm + storage + memory`** как обязательные. Всё остальное — опционально с разумными default'ами (`NoopSoul`, `EmptyBackground`, `empty toolset`).
|
||||
- **`BaseAgent`** — реализация `ConversationLoop` через builder. Конструктор принимает все нужные компоненты. **Один и тот же `BaseAgent` в `:standalone`, `:agentik-cli`, `:agentik-android`** — отличается только wiring через builder.
|
||||
- **BackgroundScheduler подписан на `BackgroundEvents`** — это даёт event-driven по умолчанию. `OnEvery(n)` interval-режим — опциональный fallback (не default).
|
||||
File diff suppressed because one or more lines are too long
|
After Width: | Height: | Size: 24 KiB |
@@ -0,0 +1,129 @@
|
||||
# 03 — Multi-user chat с mention-detection
|
||||
|
||||
Точка 1 из планов. Один `BaseAgent` обслуживает N пользователей. Отвечает только когда addressed (mention или admin-команда).
|
||||
|
||||

|
||||
|
||||
PlantUML source (для редактирования; требует Graphviz `dot` для рендеринга):
|
||||
|
||||
```plantuml
|
||||
@startuml multi-user-chat
|
||||
skinparam componentStyle rectangle
|
||||
skinparam participantPadding 15
|
||||
skinparam boxPadding 10
|
||||
|
||||
title Multi-user chat с mention-detection (точка 1 из планов)
|
||||
|
||||
' --- Участники ---
|
||||
actor "User A" as UA
|
||||
actor "User B" as UB
|
||||
actor "User C\n(админ)" as UC
|
||||
participant "Telegram /\nSlack /\nMatrix" as Channel
|
||||
participant "AgentRuntime\n(BaseAgent)" as Runtime
|
||||
participant "MentionDetector" as Detector
|
||||
participant "SoulProvider" as Soul
|
||||
participant "MemorySystem\n(MdMemory)" as Memory
|
||||
participant "LlmBackend\n(LiteLlm)" as Llm
|
||||
|
||||
' --- Сценарий ---
|
||||
UA -> Channel : "@bot, что нового?"
|
||||
UB -> Channel : "люблю котов"
|
||||
UC -> Channel : "/bot status"
|
||||
Channel -> Runtime : событие чата
|
||||
|
||||
' --- Внутри Runtime ---
|
||||
Runtime -> Detector : isMentioned(message, botName)
|
||||
|
||||
note right of Detector
|
||||
variants:
|
||||
- SimpleMentionDetector (regex: @bot)
|
||||
- LlmMentionDetector (mini-classifier)
|
||||
- AdminCommandDetector (/command)
|
||||
end note
|
||||
|
||||
Detector --> Runtime : MatchResult{isMentioned, isCommand}
|
||||
|
||||
alt isMentioned или isCommand
|
||||
Runtime -> Soul : read()
|
||||
Runtime -> Memory : prefetch(query, topK)
|
||||
Runtime -> Llm : send(system + history + memory + user)
|
||||
Llm --> Runtime : response + tool_calls
|
||||
Runtime -> Memory : save(decision)
|
||||
Runtime --> Channel : ответ в нужный канал/thread
|
||||
else NOT mentioned и NOT command
|
||||
Runtime -> Runtime : drop (если не админ)
|
||||
note right
|
||||
Не отвечаем, но возможно:
|
||||
- запоминаем факт (memory-only update)
|
||||
- summary на long conversation
|
||||
end note
|
||||
end
|
||||
|
||||
@enduml
|
||||
```
|
||||
|
||||
## Ключевые модули (что нужно будет добавить)
|
||||
|
||||
### `MentionDetector` interface
|
||||
|
||||
```kotlin
|
||||
interface MentionDetector {
|
||||
data class Result(
|
||||
val isMentioned: Boolean,
|
||||
val isAdminCommand: Boolean,
|
||||
val isPrivateMessage: Boolean, // DM — всегда отвечаем
|
||||
)
|
||||
|
||||
fun detect(message: ChatMessage, botName: String): Result
|
||||
}
|
||||
```
|
||||
|
||||
Имплементации:
|
||||
- `SimpleMentionDetector` — regex `@bot`, `/command` (дешёво, latency ~0)
|
||||
- `LlmMentionDetector` — маленькая классификация через тот же LLM (точнее, но +1 LLM-вызов на каждое сообщение)
|
||||
- `HybridMentionDetector` — fast regex → fallback на LLM только если ambiguous
|
||||
|
||||
### `ChatAdapter` interface
|
||||
|
||||
```kotlin
|
||||
interface ChatAdapter {
|
||||
val channel: String // "telegram" / "slack" / "matrix"
|
||||
suspend fun listen(onMessage: (ChatMessage) -> Unit): Job
|
||||
suspend fun reply(messageId: String, text: String, threadId: String? = null)
|
||||
suspend fun isAdmin(userId: String): Boolean
|
||||
}
|
||||
```
|
||||
|
||||
Имплементации per platform. Каждая адаптирует формат platform → `ChatMessage`.
|
||||
|
||||
### Конфигурация builder'а
|
||||
|
||||
```kotlin
|
||||
agentBuilder {
|
||||
llm(...)
|
||||
storage(...)
|
||||
memory(...)
|
||||
soul(...)
|
||||
background { ... }
|
||||
chat {
|
||||
mentionDetector = HybridMentionDetector(regex = "@bot|@Agent", llmClassifier = false)
|
||||
chatAdapter = TelegramChatAdapter(token = "...")
|
||||
// На каждое сообщение:
|
||||
// 1. mentionDetector.detect()
|
||||
// 2. если isMentioned || isAdminCommand || isPrivate → process
|
||||
// 3. иначе — опционально memory-only save (тихий режим)
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
## Что это даёт
|
||||
|
||||
- Один `BaseAgent` обслуживает чат целиком (один LLM, одна память — общий контекст команды)
|
||||
- `@bot` — explicit invocation, не «agent отвечает на всё подряд»
|
||||
- `/bot status` / `/bot clear-memory` — admin-команды (отдельный канал, без LLM)
|
||||
- DM — всегда отвечает (это личное обращение)
|
||||
- В групповом чате без mention — agent может **молча учить** (memory update без ответа). Полезно для «запомнил что Вася любит котов».
|
||||
|
||||
## Текущее состояние vs целевое
|
||||
|
||||
⏳ Ничего из этого нет. Сейчас `:standalone` — это HTTP API, к которому подключаются clients. Для multi-user chat нужен новый `:chat-adapter-telegram` (или -slack / -matrix) модуль + `MentionDetector` interface.
|
||||
File diff suppressed because one or more lines are too long
|
After Width: | Height: | Size: 18 KiB |
@@ -0,0 +1,119 @@
|
||||
# 04 — Sub-agents + A2A между агентами
|
||||
|
||||
Точки 2 и 3 из планов. Orchestrator-агент spawn'ит sub-агентов с изолированным контекстом. Независимые агенты общаются через A2A.
|
||||
|
||||

|
||||
|
||||
PlantUML source (для редактирования; требует Graphviz `dot` для рендеринга):
|
||||
|
||||
```plantuml
|
||||
@startuml sub-agents-and-a2a
|
||||
skinparam componentStyle rectangle
|
||||
|
||||
title Sub-agents + A2A между агентами (точка 2+3 из планов)
|
||||
|
||||
' --- Orchestrator ---
|
||||
rectangle "OrchestratorAgent\n(BaseAgent + tools)" as Orch {
|
||||
rectangle "ConversationLoop\n(main user)" as MainConv
|
||||
}
|
||||
|
||||
' --- Sub-agent spawn ---
|
||||
rectangle "subAgent(\n task: String,\n config: AgentConfig\n): Flow<SubAgentEvent>" as SpawnAPI
|
||||
note right of SpawnAPI
|
||||
Spawn API — НЕ отдельный модуль,
|
||||
а convenience поверх BaseAgent:
|
||||
val sub = agent.spawnChild(config) {
|
||||
systemPrompt = "..."
|
||||
tools = [ReadTool, WriteTool]
|
||||
memory = EmptyMemory // изолированно
|
||||
}
|
||||
sub.events.collect { ... }
|
||||
end note
|
||||
|
||||
' --- Дочерний агент (изолированный контекст) ---
|
||||
rectangle "SubAgent\n(изолированный scope)" as Sub {
|
||||
rectangle "ConversationLoop\n(child)" as SubConv
|
||||
rectangle "backgroundScope\n(lifecycle scoped)" as SubBg
|
||||
}
|
||||
|
||||
' --- A2A между независимыми агентами ---
|
||||
rectangle "Agent A\n(BaseAgent)" as AgentA
|
||||
rectangle "Agent B\n(BaseAgent)" as AgentB
|
||||
rectangle "A2A Server\n(:a2a-server)" as A2ASrv
|
||||
|
||||
AgentA -> A2ASrv : POST /\n(application/json)
|
||||
A2ASrv -> AgentB : dispatch(message)
|
||||
AgentB --> A2ASrv : response
|
||||
A2ASrv --> AgentA : SSE / JSON-RPC
|
||||
|
||||
' --- Стрелки ---
|
||||
Orch -> SpawnAPI : calls
|
||||
SpawnAPI -> Sub : creates with custom config
|
||||
Sub -> SubBg : has its own
|
||||
Orch -> Orch : main flow continues
|
||||
Sub --> Orch : Flow<SubAgentEvent> emits\n(Started / ToolCalled / ToolResult /\nAssistantMessage / Done / Failed)
|
||||
|
||||
Orch -> A2ASrv : can also delegate to remote agent
|
||||
|
||||
@enduml
|
||||
```
|
||||
|
||||
## Sub-agents API
|
||||
|
||||
```kotlin
|
||||
sealed interface SubAgentEvent {
|
||||
data class Started(val taskId: String) : SubAgentEvent
|
||||
data class AssistantMessage(val text: String) : SubAgentEvent
|
||||
data class ToolCalled(val toolName: String, val args: JsonObject) : SubAgentEvent
|
||||
data class ToolResult(val toolName: String, val result: String) : SubAgentEvent
|
||||
data class Done(val taskId: String, val finalResult: String) : SubAgentEvent
|
||||
data class Failed(val taskId: String, val error: String) : SubAgentEvent
|
||||
}
|
||||
|
||||
interface BaseAgent {
|
||||
// ... existing methods ...
|
||||
|
||||
/**
|
||||
* Spawn дочерний агент с изолированным контекстом (memory, system prompt,
|
||||
* tools). Возвращает Flow событий жизненного цикла + результата.
|
||||
* Cancellation родителя НЕ отменяет sub-agent — sub-agent живёт до Done/Failed.
|
||||
*/
|
||||
fun spawnChild(config: SubAgentConfig): Flow<SubAgentEvent>
|
||||
}
|
||||
|
||||
data class SubAgentConfig(
|
||||
val systemPrompt: String,
|
||||
val tools: List<NamedTool> = emptyList(),
|
||||
val memory: MemorySystem = EmptyMemory(),
|
||||
val model: LiteLlm? = null, // если null — делит LLM родителя
|
||||
val maxTurns: Int = 10,
|
||||
val timeoutMs: Long = 60_000,
|
||||
)
|
||||
```
|
||||
|
||||
## Зачем изолированный scope
|
||||
|
||||
Sub-agent получает **свою копию контекста**, не делит memory с родителем. Это критично:
|
||||
- `research_subagent` — должен исследовать тему, не отвечать на основные сообщения пользователя
|
||||
- `summarize_subagent` — суммаризировать документ, не трогать основной диалог
|
||||
- `code_review_subagent` — ревьюить PR, не видеть разговор
|
||||
|
||||
Если нужно расшарить контекст — это explicit через `sharedMemory: SharedMemoryHandle` параметр, не default.
|
||||
|
||||
## A2A между независимыми агентами
|
||||
|
||||
Уже есть `:a2a-server` модуль (см. `standalone/build.gradle.kts` — `implementation(libs.a2a.server)`). Использовался для AG-UI/A2A протокола в `:standalone`. Можно переиспользовать для межагентного общения.
|
||||
|
||||
Сценарий: orchestrator-agent не может сам решить задачу → делегирует remote-агенту через A2A → получает response → продолжает. Это уже работающая инфраструктура.
|
||||
|
||||
## Текущее состояние vs целевое
|
||||
|
||||
✅ Уже есть:
|
||||
- `:a2a-server` подключён
|
||||
- `BaseAgent.spawnChild` — **не существует**, но `ConversationLoop` уже умеет создавать изолированный scope через свой `agentScope` — нужна только обёртка
|
||||
|
||||
⏳ Не сделано:
|
||||
- `SubAgentConfig` + `Flow<SubAgentEvent>` API
|
||||
- `EmptyMemory` (null-object для изолированного scope)
|
||||
- Lifecycle management (parent dies → child должен complete or be cancelled?)
|
||||
- Сериализация sub-agent state для отладки (event log)
|
||||
File diff suppressed because one or more lines are too long
|
After Width: | Height: | Size: 14 KiB |
@@ -0,0 +1,171 @@
|
||||
# 05 — Android Agent Stack
|
||||
|
||||
Что меняется vs `:standalone`. Цель: `BaseAgent` тот же самый, но platform-impl разные (Storage, LLM, Vector Memory, MCP).
|
||||
|
||||

|
||||
|
||||
PlantUML source (для редактирования; требует Graphviz `dot` для рендеринга):
|
||||
|
||||
```plantuml
|
||||
@startuml android-agent-stack
|
||||
skinparam componentStyle rectangle
|
||||
|
||||
title Android Agent Stack — что меняется vs Standalone
|
||||
|
||||
' --- Android side ---
|
||||
package "Android Application" {
|
||||
[MainActivity\n(Compose)] as Activity
|
||||
[AndroidAgentRunner\n(workmanager / service)] as Runner
|
||||
[AndroidAgentBuilder] as AndroidBuilder
|
||||
}
|
||||
|
||||
package "Android-specific impls" {
|
||||
[StorageSqliteAndroid\n(Room/sqlite)] as StorageA
|
||||
[MemoryVectorAndroid\n(ONNX runtime + ANN)] as MemVecA
|
||||
[LitertAndroid\n(NNAPI delegate)] as LitertA
|
||||
[SoulFileAndroid\n(context.filesDir)] as SoulA
|
||||
[McpRegistry\nstdio: ProcessBuilder] as McpA
|
||||
}
|
||||
|
||||
' --- Shared (KMP) ---
|
||||
package "Agent Runtime (shared)" {
|
||||
[AgentCore\nBaseAgent] as AgentCore
|
||||
[AgentBuilder] as Builder
|
||||
}
|
||||
|
||||
package "Domain (shared)" {
|
||||
[LlmTools\ncommonMain] as LlmT
|
||||
[BackgroundEvents\ncommonMain] as Ev
|
||||
[McpBridge\njvmMain] as McpB
|
||||
[Skills\ncommonMain] as Skills
|
||||
}
|
||||
|
||||
package "Memory (shared impl)" {
|
||||
[MemoryMd\n(commonMain)] as MemMd
|
||||
[MemoryApi\ninterfaces] as MemApi
|
||||
}
|
||||
|
||||
' --- Зависимости ---
|
||||
Activity --> Runner
|
||||
Runner --> AndroidBuilder
|
||||
AndroidBuilder --> AgentCore
|
||||
|
||||
AndroidBuilder --> StorageA
|
||||
AndroidBuilder --> MemVecA
|
||||
AndroidBuilder --> LitertA
|
||||
AndroidBuilder --> SoulA
|
||||
AndroidBuilder --> McpA
|
||||
|
||||
AgentCore --> LlmT
|
||||
AgentCore --> Ev
|
||||
AgentCore --> McpB
|
||||
AgentCore --> Skills
|
||||
AgentCore --> MemMd
|
||||
|
||||
' --- Главные отличия от Standalone ---
|
||||
note right of LitertA
|
||||
On-device inference.
|
||||
LiteRT с NNAPI delegate →
|
||||
работает на CPU/GPU/NPU
|
||||
прямо на устройстве, без сети.
|
||||
|
||||
vs Standalone: HTTP-only
|
||||
(OpenAI-compatible).
|
||||
end note
|
||||
|
||||
note right of StorageA
|
||||
android.database.sqlite
|
||||
через Room или сырой API.
|
||||
|
||||
vs Standalone: JDBC +
|
||||
Sqlite-JDBC driver
|
||||
(только JVM).
|
||||
end note
|
||||
|
||||
note right of MemVecA
|
||||
JVector JVM-only. На Android
|
||||
нужна альтернатива —
|
||||
ONNX Runtime + какой-нибудь
|
||||
ANN (Annoy/HNSW).
|
||||
|
||||
Или пока без vector memory,
|
||||
только MemoryMd.
|
||||
end note
|
||||
|
||||
note right of McpA
|
||||
MCP через ProcessBuilder
|
||||
на Android работает, но
|
||||
subprocess lifecycle
|
||||
сложнее (foreground service
|
||||
нужен для долгого subprocess).
|
||||
end note
|
||||
|
||||
@enduml
|
||||
```
|
||||
|
||||
## Что общего с `:standalone`
|
||||
|
||||
**`BaseAgent`, `BackgroundScheduler`, `LlmTools`, `McpBridge`, `Skills`, `MemoryMd` — всё KMP (commonMain).** Android-agent = `:standalone` с другим wiring'ом. Не нужно переписывать agent logic.
|
||||
|
||||
## Что другое
|
||||
|
||||
| Компонент | `:standalone` (JVM) | `:agentik-android` (Android) | Сложность |
|
||||
|---|---|---|---|
|
||||
| Storage | `:storage-sqlite` (JDBC + Sqlite-JDBC) | `:storage-sqlite-android` (Room или raw) | Низкая — тот же `StorageBundle` interface |
|
||||
| LLM | `:litert-openai` (HTTP), `:litert-google` (LiteRT JVM) | `:litert-android` (LiteRT Android, NNAPI delegate) | Средняя — нужен новый модуль |
|
||||
| Vector memory | `:memory-vector` (JVector) | `:memory-vector-android` (ONNX Runtime + HNSW/Annoy) | Высокая — JVector JVM-only, нужна альтернатива |
|
||||
| SOUL provider | `FileSoulProvider` (path) | `SoulFileAndroid` (`context.filesDir`) | Низкая |
|
||||
| MCP | `McpRegistry` (ProcessBuilder, stdio subprocess) | Тот же `McpRegistry`, но subprocess в foreground service | Средняя — нужен Android service |
|
||||
| Embedding | `HttpEmbeddingClient` (HTTP) | Тот же ИЛИ on-device (ONNX) | Средняя |
|
||||
|
||||
## Минимальный Android agent (v1)
|
||||
|
||||
Если не нужны все фичи сразу — минимум:
|
||||
|
||||
```kotlin
|
||||
val agent = androidAgentBuilder(context) {
|
||||
llm(LitertAndroid.onDevice(context, modelPath = "/data/local/tmp/model.litertlm"))
|
||||
storage(SqliteStorage.android(context, "agent.db"))
|
||||
memory(MemoryMd.root(context.filesDir.resolve("memory")))
|
||||
soul(FileSoul(context.filesDir.resolve("SOUL.md")))
|
||||
background {
|
||||
// OnClosing + OnCompaction работают так же как на JVM
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
Без MCP, без vector memory (только MemoryMd на файлах), только on-device LLM. Достаточно для off-line агента.
|
||||
|
||||
## Foreground service для MCP
|
||||
|
||||
Если нужны MCP-серверы (например, локальный file-system MCP) — subprocess нужен foreground service чтобы Android не убил его при выключении экрана. Это добавляет сложности:
|
||||
|
||||
```kotlin
|
||||
class McpForegroundService : Service() {
|
||||
override fun onStartCommand(intent: Intent?, flags: Int, startId: Int): Int {
|
||||
startForeground(NOTIFICATION_ID, notification)
|
||||
val proc = ProcessBuilder(command, args).start()
|
||||
// ... route stdio to McpLiteToolAdapter ...
|
||||
return START_STICKY
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
Пока можно без этого (только если MCP нужен на Android).
|
||||
|
||||
## Текущее состояние vs целевое
|
||||
|
||||
✅ KMP-ready:
|
||||
- `:llm-tools` (commonMain, платформо-агностик)
|
||||
- `:mcp-bridge` (jvmMain — Android-вариант через `:mcp-bridge-android`)
|
||||
- `:skills` (commonMain)
|
||||
- `:memory-md` (commonMain)
|
||||
- `:proto` (commonMain)
|
||||
|
||||
⏳ Не существует:
|
||||
- `:storage-sqlite-android`
|
||||
- `:memory-vector-android`
|
||||
- `:litert-android`
|
||||
- `:agentik-android` (само приложение)
|
||||
- `:agent-core` (выделить BaseAgent + builder)
|
||||
- `:background-events` (выделить events + scheduler)
|
||||
File diff suppressed because one or more lines are too long
|
After Width: | Height: | Size: 25 KiB |
@@ -0,0 +1,64 @@
|
||||
# agentik — диаграммы архитектуры
|
||||
|
||||
PlantUML-схемы для обсуждения будущей структуры (Android agent, multi-user chat, sub-agents, A2A). Это **целевое состояние**, не текущее.
|
||||
|
||||
## Файлы
|
||||
|
||||
Каждый `.md` содержит:
|
||||
- Краткое описание (что показывает)
|
||||
- **Пред-рендеренный SVG** (``) — гарантированно показывается **везде**
|
||||
- PlantUML source в ` ```plantuml ` блоке — для редактирования (требует Graphviz `dot` для рендеринга)
|
||||
- Дополнительный markdown-текст (что нужно сделать, текущее vs целевое)
|
||||
|
||||
| Файл | Что показывает |
|
||||
|---|---|
|
||||
| [01-module-layers.md](./01-module-layers.md) | Целевая модульная структура (приложения → runtime → домен → абстракции → платформенные impl). Что в каком слое и кто от кого зависит. |
|
||||
| [02-agent-composition.md](./02-agent-composition.md) | Как `AgentBuilder` собирает `BaseAgent` из компонентов. Memory backend сам объявляет свои tools. BackgroundScheduler подписан на события (НЕ interval-poll). |
|
||||
| [03-multi-user-chat.md](./03-multi-user-chat.md) | Сценарий: чат с N пользователями, mention-detection, админ-команды, agent отвечает только когда addressed. |
|
||||
| [04-sub-agents.md](./04-sub-agents.md) | Orchestrator spawn'ит sub-agent с изолированным контекстом, получает `Flow<SubAgentEvent>`. A2A между независимыми агентами через `:a2a-server`. |
|
||||
| [05-android-stack.md](./05-android-stack.md) | Что меняется на Android: on-device LLM (NNAPI), Room/sqlite, ONNX-based vector memory, foreground-service для MCP subprocess. |
|
||||
|
||||
## Почему SVG + PlantUML source
|
||||
|
||||
PlantUML требует Java + (для component/class/deployment диаграмм) Graphviz `dot`. Если `dot` не установлен — рендерер падает с ошибкой "Executable dot does not exist".
|
||||
|
||||
Решение: **пре-рендерим в SVG один раз** и вставляем как `<img>`. Диаграмма гарантированно показывается в любом markdown-viewer (GitHub, IntelliJ, VSCode, GitLab) без зависимостей. PlantUML source в code block остаётся для редактирования.
|
||||
|
||||
## Как редактировать диаграмму
|
||||
|
||||
1. Меняешь PlantUML-source в ` ```plantuml ` блоке `.md` файла.
|
||||
2. Ре-рендеришь SVG:
|
||||
```bash
|
||||
mkdir -p /tmp/plantuml-work && chmod 777 /tmp/plantuml-work
|
||||
cp docs/diagrams/*.md /tmp/plantuml-work/
|
||||
docker run --rm -v /tmp/plantuml-work:/work plantuml/plantuml -tsvg /work/*.md
|
||||
cp /tmp/plantuml-work/*.svg docs/diagrams/
|
||||
```
|
||||
3. Проверяешь что SVG обновился:
|
||||
```bash
|
||||
ls -la docs/diagrams/*.svg
|
||||
```
|
||||
4. Коммитишь оба файла: `.md` (source) и `.svg` (rendered).
|
||||
|
||||
Требует Docker (или локального PlantUML+Graphviz). `apt install graphviz` для Arch/Manjaro.
|
||||
|
||||
## Контекст
|
||||
|
||||
Текущий код движется в эту сторону:
|
||||
- `:llm-tools` extracted ✅
|
||||
- `:mcp-bridge` extracted ✅
|
||||
- `BackgroundScheduler` стал event-driven ✅
|
||||
- `ConversationLoop` стал отдельным компонентом ✅
|
||||
|
||||
Не сделано (см. детали в каждом .md):
|
||||
- `:agent-core` (выделить `BaseAgent` + builder)
|
||||
- `:background-events` (выделить events + scheduler)
|
||||
- `:storage-sqlite-android`, `:memory-vector-android`, `:litert-android`
|
||||
- `:agentik-android` (само приложение)
|
||||
- `MentionDetector` interface + adapters для multi-user chat
|
||||
- `BaseAgent.spawnChild` + `Flow<SubAgentEvent>`
|
||||
|
||||
Подробнее:
|
||||
- `STANDALONE-REVIEW.md` — что плохо в текущем коде
|
||||
- `MEMORY-DESIGN.md` — детали memory архитектуры
|
||||
- `STANDALONE.md` — текущий standalone
|
||||
+6
-2
@@ -1,5 +1,9 @@
|
||||
# Default version for local builds; overridden by `-Pversion=<tag>` from CI/CD.
|
||||
version=0.1.0
|
||||
# Default version for local builds only (когда CI/CD не передал -Pversion=<tag>).
|
||||
# Имя ключа специально НЕ 'version' — иначе Gradle-мерж gradle.properties и
|
||||
# -Pversion= возьмёт default из gradle.properties. Передавай через CICD:
|
||||
# ./gradlew ... -Pversion=$(git describe --tags)
|
||||
# см. .gitea/workflows/release.yml (использует -Pversion=$GITHUB_REF_NAME).
|
||||
agentik.version.default=0.1.0-SNAPSHOT
|
||||
|
||||
# KMP jvm target uses JDK 21 for both compilation and toolchain.
|
||||
org.gradle.jvmargs=-Xmx4096M -XX:+UseG1GC
|
||||
|
||||
@@ -13,8 +13,9 @@ jvector = "3.0.6"
|
||||
text-embedding-kmp = "3.0.0-SNAPSHOT"
|
||||
kotlin-logging = "3.0.5"
|
||||
logback = "1.5.18"
|
||||
jline = "3.30.0"
|
||||
mosaic = "0.18.0"
|
||||
clikt = "5.0.3"
|
||||
kotlinx-cli = "0.3.6"
|
||||
|
||||
[plugins]
|
||||
kotlin-multiplatform = { id = "org.jetbrains.kotlin.multiplatform", version.ref = "kotlin" }
|
||||
@@ -54,6 +55,7 @@ ktor-server-content-negotiation = { module = "io.ktor:ktor-server-content-negoti
|
||||
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-curl = { module = "io.ktor:ktor-client-curl", version.ref = "ktor" }
|
||||
ktor-client-content-negotiation = { module = "io.ktor:ktor-client-content-negotiation", version.ref = "ktor" }
|
||||
ktor-server-test-host = { module = "io.ktor:ktor-server-test-host", version.ref = "ktor" }
|
||||
ktor-client-sse = { module = "io.ktor:ktor-client-sse", version.ref = "ktor" }
|
||||
@@ -61,8 +63,12 @@ ktor-client-sse = { module = "io.ktor:ktor-client-sse", version.ref = "ktor" }
|
||||
# --- Model Context Protocol (MCP) ---
|
||||
mcp-sdk-client = { module = "io.modelcontextprotocol:kotlin-sdk-client", version = "0.15.0" }
|
||||
|
||||
# --- CLI: JLine (readline для JVM-таргета) ---
|
||||
jline = { module = "org.jline:jline", version.ref = "jline" }
|
||||
# --- CLI: clikt (ajalt). KMP, native Linux/macOS/Windows включая linuxArm64. ---
|
||||
# https://ajalt.github.io/clikt/
|
||||
# Артефакт один и тот же — `com.github.ajalt.clikt:clikt` — Gradle module
|
||||
# metadata резолвит per-target variant (clikt-jvm / clikt-linuxarm64 / ...).
|
||||
clikt = { module = "com.github.ajalt.clikt:clikt-core", version.ref = "clikt" }
|
||||
kotlinx-cli = { module = "org.jetbrains.kotlinx:kotlinx-cli", version.ref = "kotlinx-cli" }
|
||||
|
||||
# --- TUI: Mosaic (Jetpack Compose → ANSI-терминал), jvm + desktop-native. ---
|
||||
# https://github.com/JakeWharton/mosaic
|
||||
|
||||
@@ -0,0 +1,37 @@
|
||||
// Generic LLM-side tools: LlmReflector, SkillMiner, LlmMemoryReviewer,
|
||||
// ContextCompactor + парсеры/промпты. Вынесены из :standalone (god class)
|
||||
// — переиспользуемы в :agentik-cli / :agentik-tui и любых других клиентах.
|
||||
//
|
||||
// Зависимости — все JVM-only контракты: litert.api JVM-only для LiteLlm
|
||||
// (он и так JVM-only), :memory-api / :storage-core / :skills — commonMain,
|
||||
// доступные JVM target'у.
|
||||
@file:OptIn(org.jetbrains.kotlin.gradle.ExperimentalKotlinGradlePluginApi::class)
|
||||
|
||||
plugins {
|
||||
alias(libs.plugins.kotlin.multiplatform)
|
||||
}
|
||||
|
||||
kotlin {
|
||||
jvmToolchain(21)
|
||||
|
||||
jvm()
|
||||
|
||||
sourceSets {
|
||||
commonMain.dependencies {
|
||||
api(project(":memory-api"))
|
||||
api(project(":storage-core"))
|
||||
api(project(":skills"))
|
||||
api(libs.litert.api)
|
||||
implementation(libs.kotlinx.coroutines.core)
|
||||
implementation(libs.kotlinx.serialization.json)
|
||||
}
|
||||
commonTest.dependencies {
|
||||
implementation(kotlin("test"))
|
||||
implementation(libs.kotlinx.coroutines.core)
|
||||
}
|
||||
jvmMain.dependencies {
|
||||
// mu.KotlinLogging — JVM-only, для SkillMiner'а
|
||||
implementation(libs.kotlin.logging)
|
||||
}
|
||||
}
|
||||
}
|
||||
+1
-1
@@ -1,4 +1,4 @@
|
||||
package pw.binom.agentik.standalone.agent
|
||||
package pw.binom.agentik.llm.tools
|
||||
|
||||
import pw.binom.litert.LiteContentPart
|
||||
import pw.binom.litert.LiteConversation
|
||||
+1
-1
@@ -1,4 +1,4 @@
|
||||
package pw.binom.agentik.standalone.agent.memory
|
||||
package pw.binom.agentik.llm.tools
|
||||
|
||||
import kotlinx.coroutines.CoroutineDispatcher
|
||||
import kotlinx.coroutines.withContext
|
||||
+2
-2
@@ -1,4 +1,4 @@
|
||||
package pw.binom.agentik.standalone.agent
|
||||
package pw.binom.agentik.llm.tools
|
||||
|
||||
import kotlinx.coroutines.CoroutineDispatcher
|
||||
import kotlinx.coroutines.withContext
|
||||
@@ -27,7 +27,7 @@ import kotlin.time.Clock
|
||||
*/
|
||||
class LlmReflector(
|
||||
private val llm: LiteLlm,
|
||||
private val maxTurns: Int = 6,
|
||||
val maxTurns: Int = 6,
|
||||
private val maxTokens: Int = 512,
|
||||
private val dispatcher: CoroutineDispatcher = kotlinx.coroutines.Dispatchers.IO,
|
||||
private val clock: Clock = Clock.System,
|
||||
+1
-1
@@ -1,4 +1,4 @@
|
||||
package pw.binom.agentik.standalone.agent
|
||||
package pw.binom.agentik.llm.tools
|
||||
|
||||
/**
|
||||
* Минимальный парсер JSON-ответа от [LlmReflector].
|
||||
+1
-1
@@ -1,4 +1,4 @@
|
||||
package pw.binom.agentik.standalone.agent
|
||||
package pw.binom.agentik.llm.tools
|
||||
|
||||
import pw.binom.agentik.memory.ConversationTurn
|
||||
|
||||
+1
-1
@@ -1,4 +1,4 @@
|
||||
package pw.binom.agentik.standalone.agent.memory
|
||||
package pw.binom.agentik.llm.tools
|
||||
|
||||
import pw.binom.agentik.memory.MemoryCategory
|
||||
import pw.binom.agentik.memory.MemoryReviewDecision
|
||||
+1
-1
@@ -1,4 +1,4 @@
|
||||
package pw.binom.agentik.standalone.agent.memory
|
||||
package pw.binom.agentik.llm.tools
|
||||
|
||||
import pw.binom.agentik.memory.ReviewedTurn
|
||||
|
||||
+1
-1
@@ -1,4 +1,4 @@
|
||||
package pw.binom.agentik.standalone.agent
|
||||
package pw.binom.agentik.llm.tools
|
||||
|
||||
import kotlinx.coroutines.CoroutineDispatcher
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
+1
-1
@@ -1,4 +1,4 @@
|
||||
package pw.binom.agentik.standalone.agent
|
||||
package pw.binom.agentik.llm.tools
|
||||
|
||||
import pw.binom.agentik.skills.SkillFile
|
||||
|
||||
+1
-1
@@ -1,4 +1,4 @@
|
||||
package pw.binom.agentik.standalone.agent
|
||||
package pw.binom.agentik.llm.tools
|
||||
|
||||
import pw.binom.agentik.memory.ConversationTurn
|
||||
import pw.binom.agentik.skills.SkillFile
|
||||
@@ -0,0 +1,39 @@
|
||||
// Generic MCP (Model Context Protocol) bridge — переиспользуемый модуль,
|
||||
// который превращает любой MCP-сервер (stdio subprocess или HTTP endpoint)
|
||||
// в набор [LiteTool]-адаптеров.
|
||||
//
|
||||
// Вынесен из :standalone — MCP не специфичен для standalone'а, это generic
|
||||
// мост между MCP-SDK и litert-kmp. Может переиспользоваться в :agentik-cli
|
||||
// или :agentik-tui когда те снова включатся.
|
||||
//
|
||||
// Зависимости:
|
||||
// - :agent-toolsets для NamedTool (обёртка для LiteTool + имя-как-видит-модель)
|
||||
// - litert.api для LiteTool контракта
|
||||
// - MCP SDK (JVM-only)
|
||||
// - Ktor client (для StreamableHttpClientTransport)
|
||||
// - kotlinx-serialization для парсинга конфига
|
||||
plugins {
|
||||
alias(libs.plugins.kotlin.jvm)
|
||||
alias(libs.plugins.kotlin.serialization)
|
||||
}
|
||||
|
||||
kotlin {
|
||||
jvmToolchain(21)
|
||||
}
|
||||
|
||||
dependencies {
|
||||
implementation(project(":agent-toolsets"))
|
||||
|
||||
api(libs.litert.api)
|
||||
|
||||
implementation(libs.mcp.sdk.client)
|
||||
implementation(libs.ktor.client.core)
|
||||
implementation(libs.ktor.client.cio)
|
||||
implementation(libs.ktor.client.content.negotiation)
|
||||
implementation(libs.ktor.serialization.kotlinx.json)
|
||||
|
||||
implementation(libs.kotlinx.coroutines.core)
|
||||
implementation(libs.kotlinx.serialization.json)
|
||||
|
||||
implementation(libs.kotlin.logging)
|
||||
}
|
||||
+1
-1
@@ -1,4 +1,4 @@
|
||||
package pw.binom.agentik.standalone.mcp
|
||||
package pw.binom.agentik.mcp.bridge
|
||||
|
||||
|
||||
import kotlinx.serialization.SerialName
|
||||
+2
-2
@@ -1,4 +1,4 @@
|
||||
package pw.binom.agentik.standalone.mcp
|
||||
package pw.binom.agentik.mcp.bridge
|
||||
|
||||
import mu.KotlinLogging
|
||||
|
||||
@@ -30,9 +30,9 @@ import kotlinx.serialization.json.doubleOrNull
|
||||
import kotlinx.serialization.json.intOrNull
|
||||
import kotlinx.serialization.json.longOrNull
|
||||
import kotlinx.serialization.json.put
|
||||
import pw.binom.agentik.standalone.agent.NamedTool
|
||||
import pw.binom.litert.LiteTool
|
||||
import java.util.concurrent.ConcurrentHashMap
|
||||
import pw.binom.agentik.toolsets.NamedTool
|
||||
|
||||
/**
|
||||
* Реестр подключённых MCP-серверов.
|
||||
+2
-2
@@ -22,8 +22,8 @@
|
||||
|
||||
- `:server` — Ktor-фасад, маппит `Agent` ↔ HTTP/SSE.
|
||||
- `:client` — Ktor-клиент, маппит HTTP/SSE ↔ `Agent/Conversation`.
|
||||
- `:agentik-cli`, `:agentik-tui` — оба работают поверх `:client`,
|
||||
а следовательно поверх `:proto`.
|
||||
- `:agentik-cli` — работает поверх `:client`, а следовательно поверх `:proto`.
|
||||
*(`:agentik-tui` был исключён из сборки 2026-09-17.)*
|
||||
- `:standalone` — реализует `Agent` (через `ChatAgent`) и пишет/читает
|
||||
`Message`/`Event` напрямую через storage.
|
||||
|
||||
|
||||
+1
-1
@@ -20,7 +20,7 @@ IRC/MCP) общаясь с одним сервером по стабильном
|
||||
|
||||
- `:standalone` подключает `Route.agentikAgent(agent)` в свой
|
||||
embedded Netty engine.
|
||||
- Любые клиенты (наши `:client`, `:agentik-cli`, `:agentik-tui`, или
|
||||
- Любые клиенты (наши `:client`, `:agentik-cli`, или
|
||||
внешние web-фронтенды) идут через этот контракт.
|
||||
|
||||
## Как подключить
|
||||
|
||||
@@ -37,6 +37,8 @@ kotlin {
|
||||
commonTest.dependencies {
|
||||
implementation(libs.kotlinx.coroutines.core)
|
||||
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 /health` — `"ok"`
|
||||
*/
|
||||
fun Route.agentikAgent(agent: Agent, path: String = "/agentik") {
|
||||
fun Route.agentikAgent(agent: Agent, path: String = "/agentik", token: String? = null) {
|
||||
route(path) {
|
||||
install(ContentNegotiation) {
|
||||
json(agentikJson)
|
||||
}
|
||||
if (token != null) {
|
||||
install(BearerTokenPlugin) {
|
||||
this.token = token
|
||||
}
|
||||
}
|
||||
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)
|
||||
}
|
||||
}
|
||||
}
|
||||
+11
-1
@@ -28,6 +28,14 @@ dependencyResolutionManagement {
|
||||
rootProject.name = "agentik"
|
||||
|
||||
include(":standalone")
|
||||
// Generic LLM-side tools: LlmReflector, SkillMiner, LlmMemoryReviewer,
|
||||
// ContextCompactor + парсеры/промпты. Вынесены из :standalone (god class)
|
||||
// — переиспользуемы в :agentik-cli / :agentik-tui и любых других клиентах.
|
||||
include(":llm-tools")
|
||||
// Generic MCP-bridge: McpConfig, McpRegistry, McpLiteToolAdapter. Вынесены
|
||||
// из :standalone — MCP не специфичен для standalone'а, это generic мост
|
||||
// между MCP-SDK и LiteTool. Содержит :agent-toolsets (NamedTool).
|
||||
include(":mcp-bridge")
|
||||
// Собственный протокол agentik. Пока в нём пилим, потом вынесем.
|
||||
include(":proto")
|
||||
// Парсер скилов (YAML-frontmatter + markdown body, opencode-style).
|
||||
@@ -41,7 +49,9 @@ include(":client")
|
||||
include(":agentik-cli")
|
||||
// TUI-клиент поверх :client — Compose-style UI (Mosaic от Jake Wharton),
|
||||
// рендерится в ANSI-терминал. KMP со всеми desktop-целями (без ios).
|
||||
include(":agentik-tui")
|
||||
// include(":agentik-tui") — отключено 2026-09-17: пользователь признал TUI-подход неудачным.
|
||||
// Папка agentik-tui/ оставлена на диске для возможного возврата; из сборки исключена.
|
||||
|
||||
// Встраиваемая долговременная память агента. `:memory-api` — интерфейсы,
|
||||
// `:memory-md` — реализация на базе §-файлов (Hermes-style).
|
||||
include(":memory-api")
|
||||
|
||||
@@ -69,6 +69,8 @@ AGENTIK_GOOGLE_MODEL_PATH=/root/gemma-4-E2B-it.litertlm \
|
||||
|---|---|---|
|
||||
| `AGENTIK_PORT` | `8080` | Порт HTTP-сервера |
|
||||
| `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_LLM_BACKEND` | `openai` | `openai` или `google` |
|
||||
| `AGENTIK_LLM_MODEL` | (выбирается по backend) | Имя модели |
|
||||
|
||||
@@ -59,6 +59,12 @@ kotlin {
|
||||
implementation(project(":storage-core"))
|
||||
implementation(project(":storage-sqlite"))
|
||||
implementation(project(":agent-toolsets"))
|
||||
// Generic LLM-side tools (LlmReflector, SkillMiner, LlmMemoryReviewer,
|
||||
// ContextCompactor, парсеры/промпты). Вынесены из :standalone.
|
||||
implementation(project(":llm-tools"))
|
||||
// Generic MCP-bridge (McpConfig, McpRegistry, McpLiteToolAdapter).
|
||||
// Вынесен из :standalone — generic мост между MCP-SDK и LiteTool.
|
||||
implementation(project(":mcp-bridge"))
|
||||
|
||||
// litert-google: встроенный LiteRT-LM движок, нужен только на runtime
|
||||
runtimeOnly(libs.litert.google)
|
||||
|
||||
@@ -11,8 +11,8 @@ import kotlinx.serialization.json.put
|
||||
import pw.binom.agentik.memory.ConversationTurn
|
||||
import pw.binom.agentik.proto.Agent
|
||||
import pw.binom.agentik.standalone.agent.ChatConversation
|
||||
import pw.binom.agentik.standalone.agent.LlmReflector
|
||||
import pw.binom.agentik.standalone.agent.SkillMiner
|
||||
import pw.binom.agentik.llm.tools.LlmReflector
|
||||
import pw.binom.agentik.llm.tools.SkillMiner
|
||||
import pw.binom.agentik.standalone.agent.memory.Curator
|
||||
import pw.binom.agentik.storage.StorageBundle
|
||||
import pw.binom.agentik.skills.SkillStore
|
||||
|
||||
@@ -21,16 +21,17 @@ import pw.binom.agentik.memory.vector.embedding.SiglipEmbeddingProvider
|
||||
import pw.binom.agentik.server.agentikAgent
|
||||
import pw.binom.agentik.skills.SkillCatalog
|
||||
import pw.binom.agentik.standalone.agent.ChatAgent
|
||||
import pw.binom.agentik.standalone.agent.LiteLlmContextCompactor
|
||||
import pw.binom.agentik.standalone.agent.LlmReflector
|
||||
import pw.binom.agentik.standalone.agent.memory.LlmMemoryReviewer
|
||||
import pw.binom.agentik.standalone.config.AgentikConfig
|
||||
import pw.binom.agentik.standalone.config.AgentikConfig.MemoryBackend
|
||||
import pw.binom.agentik.llm.tools.LiteLlmContextCompactor
|
||||
import pw.binom.agentik.llm.tools.LlmReflector
|
||||
import pw.binom.agentik.llm.tools.LlmMemoryReviewer
|
||||
import pw.binom.agentik.standalone.config.AppConfig
|
||||
import pw.binom.agentik.standalone.config.AppConfig.MemoryBackend
|
||||
import pw.binom.agentik.standalone.llm.LlmBackend
|
||||
import pw.binom.agentik.standalone.llm.ModelDownloader
|
||||
import pw.binom.agentik.standalone.mcp.McpRegistry
|
||||
import pw.binom.agentik.mcp.bridge.McpRegistry
|
||||
import pw.binom.agentik.storage.sqlite.SqliteStores
|
||||
import java.io.File
|
||||
import pw.binom.agentik.llm.tools.SkillMiner
|
||||
/**
|
||||
* standalone-контейнер agentik:
|
||||
* - :server (proto): встраиваемый Ktor (CIO), порт AGENTIK_PORT (default 8080)
|
||||
@@ -48,8 +49,10 @@ import java.io.File
|
||||
* GET /a2a/.well-known/agent-card.json -> AgentCard
|
||||
* GET /health -> "ok"
|
||||
*
|
||||
* Вся конфигурация — [AgentikConfig.fromEnv] (см. [AgentikConfig]). Источники:
|
||||
* Вся конфигурация — [AppConfig.fromEnv] (см. [AppConfig]). Источники:
|
||||
* - AGENTIK_PORT / AGENTIK_DB_PATH
|
||||
* - AGENTIK_TOKEN — Bearer-токен для HTTP-фасада /agentik (пусто — авторизация выключена)
|
||||
* - AGENTIK_A2A_TOKEN — Bearer-токен для A2A-фасада /a2a (пусто — авторизация выключена)
|
||||
* - LLM: AGENTIK_LLM_BACKEND, OPENAI_* либо AGENTIK_GOOGLE_*
|
||||
* - MCP: AGENTIK_MCP_CONFIG=<path>.json (формат Claude Desktop)
|
||||
* - Skills: AGENTIK_SKILLS_DIR=<path> (папка с SKILL.md / *.yaml)
|
||||
@@ -98,7 +101,7 @@ private fun printHelp() {
|
||||
* Если файл по PATH уже есть и совпадает по размеру с HEAD — no-op (exit 0).
|
||||
*/
|
||||
private fun runPullModel(args: List<String>) {
|
||||
val config = AgentikConfig.fromEnv()
|
||||
val config = AppConfig.fromEnv()
|
||||
val google = config.llm.google
|
||||
?: error("pull-model: требуется AGENTIK_LLM_BACKEND=google (сейчас ${config.llm.backend})")
|
||||
|
||||
@@ -148,7 +151,7 @@ private fun formatBytes(b: Long): String = when {
|
||||
}
|
||||
|
||||
private fun runServer() {
|
||||
val config = AgentikConfig.fromEnv()
|
||||
val config = AppConfig.fromEnv()
|
||||
|
||||
// Перед созданием LLM: если backend=google и файл по AGENTIK_GOOGLE_MODEL_PATH
|
||||
// отсутствует — качаем автоматически (только при AGENTIK_AUTO_DOWNLOAD_MODEL=1),
|
||||
@@ -195,14 +198,14 @@ private fun runServer() {
|
||||
}
|
||||
|
||||
val llm = config.llm.createLlm()
|
||||
val storage = SqliteStores.open(dbPath = config.dbPath).asBundle()
|
||||
val storage = SqliteStores.open(dbPath = config.agent.dbPath).asBundle()
|
||||
val mcpRegistry = McpRegistry.fromConfig(config.mcp)
|
||||
|
||||
// Хранилище скилов: если skillsDir задан, читаем каталог + создаём
|
||||
// DiskSkillStore для self-improvement (`skill_save`/`skill_delete`).
|
||||
// Один и тот же файл-каталог используется и для чтения (read_skill),
|
||||
// и для записи — никаких рассинхронов.
|
||||
val skillStore: pw.binom.agentik.skills.SkillStore? = config.skillsDir?.let { dir ->
|
||||
val skillStore: pw.binom.agentik.skills.SkillStore? = config.agent.skillsDir?.let { dir ->
|
||||
pw.binom.agentik.skills.DiskSkillStore(File(dir))
|
||||
}
|
||||
val skills = skillStore?.catalog ?: SkillCatalog.EMPTY
|
||||
@@ -210,18 +213,18 @@ private fun runServer() {
|
||||
// Skill mining: фоновый LLM-прогон, который находит переиспользуемые скилы,
|
||||
// которые модель забыла сохранить через `skill_save`. Работает только когда
|
||||
// есть куда писать (skillStore) и интервал > 0.
|
||||
val skillMiner: pw.binom.agentik.standalone.agent.SkillMiner? =
|
||||
if (skillStore != null && config.skillMiningInterval > 0) {
|
||||
pw.binom.agentik.standalone.agent.SkillMiner(
|
||||
val skillMiner: pw.binom.agentik.llm.tools.SkillMiner? =
|
||||
if (skillStore != null && config.skillMining.interval > 0) {
|
||||
pw.binom.agentik.llm.tools.SkillMiner(
|
||||
llm = llm,
|
||||
maxTurns = config.skillMiningMaxTurns,
|
||||
maxTurns = config.skillMining.maxTurns,
|
||||
)
|
||||
} else null
|
||||
|
||||
// SOUL.md — файл персоны. Если задан — читается как plain text/markdown,
|
||||
// вставляется в самое начало systemInstruction. Если отсутствует — exit-code != 0
|
||||
// (на старте агента это фатально: нечего показывать LLM).
|
||||
val soulBody = config.soulPath?.let { path ->
|
||||
val soulBody = config.agent.soulPath?.let { path ->
|
||||
val file = File(path)
|
||||
if (!file.exists() || !file.isFile) {
|
||||
log.warn { "SOUL file not found: $path" }
|
||||
@@ -235,8 +238,8 @@ private fun runServer() {
|
||||
// - MD (дефолт) — Hermes-style §-файлы в AGENTIK_MEMORY_DIR (~/.agentik/memory)
|
||||
// - VECTOR — SQLite + JVector + LLM-эмбеддинги (тот же agentik.db для metadata)
|
||||
// - OFF — память выключена (memoryDir="off" или memoryBackend="off")
|
||||
val rawMemory = config.memoryDir
|
||||
val memorySystem: MemorySystem? = when (config.memoryBackend) {
|
||||
val rawMemory = config.memory.dir
|
||||
val memorySystem: MemorySystem? = when (config.memory.backend) {
|
||||
MemoryBackend.OFF -> {
|
||||
println(" memory: disabled")
|
||||
null
|
||||
@@ -253,8 +256,8 @@ private fun runServer() {
|
||||
}
|
||||
}
|
||||
MemoryBackend.VECTOR -> {
|
||||
val embedding: pw.binom.agentik.memory.vector.EmbeddingProvider = when (config.embeddingBackend) {
|
||||
AgentikConfig.EmbeddingBackend.HTTP -> {
|
||||
val embedding: pw.binom.agentik.memory.vector.EmbeddingProvider = when (config.embedding.backend) {
|
||||
AppConfig.EmbeddingBackend.HTTP -> {
|
||||
val llm = config.llm
|
||||
// Берём базовый URL + API key у активного LLM-бэкенда.
|
||||
// Поддерживается только OPENAI (LiteLLM proxy тоже работает, т.к. /v1/embeddings
|
||||
@@ -266,38 +269,38 @@ private fun runServer() {
|
||||
HttpEmbeddingClient(
|
||||
apiUrl = oa.baseUrl.trimEnd('/'),
|
||||
apiKey = oa.apiKey,
|
||||
model = config.embeddingModel,
|
||||
dimension = config.embeddingDimension,
|
||||
model = config.embedding.model,
|
||||
dimension = config.embedding.dimension,
|
||||
)
|
||||
}
|
||||
AgentikConfig.EmbeddingBackend.SIGLIP -> {
|
||||
val modelPath = checkNotNull(config.embeddingModelPath) {
|
||||
AppConfig.EmbeddingBackend.SIGLIP -> {
|
||||
val modelPath = checkNotNull(config.embedding.modelPath) {
|
||||
"AGENTIK_EMBEDDING_BACKEND=siglip требует AGENTIK_EMBEDDING_MODEL_PATH"
|
||||
}
|
||||
val tokenizerPath = checkNotNull(config.embeddingTokenizerPath) {
|
||||
val tokenizerPath = checkNotNull(config.embedding.tokenizerPath) {
|
||||
"AGENTIK_EMBEDDING_BACKEND=siglip требует AGENTIK_EMBEDDING_TOKENIZER_PATH"
|
||||
}
|
||||
SiglipEmbeddingProvider(modelPath = modelPath, tokenizerPath = tokenizerPath)
|
||||
}
|
||||
}
|
||||
VectorMemorySystem.open(
|
||||
dbPath = config.dbPath,
|
||||
dbPath = config.agent.dbPath,
|
||||
embedding = embedding,
|
||||
).also {
|
||||
val backendLabel = when (config.embeddingBackend) {
|
||||
AgentikConfig.EmbeddingBackend.HTTP ->
|
||||
"model=${config.embeddingModel}, dim=${config.embeddingDimension}"
|
||||
AgentikConfig.EmbeddingBackend.SIGLIP ->
|
||||
val backendLabel = when (config.embedding.backend) {
|
||||
AppConfig.EmbeddingBackend.HTTP ->
|
||||
"model=${config.embedding.model}, dim=${config.embedding.dimension}"
|
||||
AppConfig.EmbeddingBackend.SIGLIP ->
|
||||
"model=siglip2-base (on-device), dim=${embedding.dimension}"
|
||||
}
|
||||
println(" memory: db=${config.dbPath} (vector-backend, $backendLabel)")
|
||||
println(" memory: db=${config.agent.dbPath} (vector-backend, $backendLabel)")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Контекстное окно модели (для compaction'а working memory).
|
||||
// Если null — compaction выключен. Резолвится один раз из LlmConfig/env.
|
||||
val contextWindow: Int? = config.llm.resolveContextWindow()
|
||||
val contextWindow: Int? = config.llm.contextWindow
|
||||
val contextCompactor = if (contextWindow != null) LiteLlmContextCompactor(liteLlm = llm) else null
|
||||
|
||||
// Review-loop: всегда используем LlmMemoryReviewer поверх LiteLlm, если память включена.
|
||||
@@ -312,11 +315,11 @@ private fun runServer() {
|
||||
|
||||
// Self-reflection: reflector работает только когда LLM доступен (нужен LiteLlm)
|
||||
// и interval > 0. Загружаем top-K последних рефлексий из SQLite в system prompt.
|
||||
val reflector: pw.binom.agentik.standalone.agent.LlmReflector? =
|
||||
if (config.reflectionInterval > 0) LlmReflector(llm = llm) else null
|
||||
val reflector: pw.binom.agentik.llm.tools.LlmReflector? =
|
||||
if (config.reflection.interval > 0) LlmReflector(llm = llm) else null
|
||||
val recentReflections: List<pw.binom.agentik.storage.Reflection> =
|
||||
if (config.reflectionTopK > 0) kotlinx.coroutines.runBlocking {
|
||||
storage.reflectionStore.listRecent(config.reflectionTopK)
|
||||
if (config.reflection.topK > 0) kotlinx.coroutines.runBlocking {
|
||||
storage.reflectionStore.listRecent(config.reflection.topK)
|
||||
} else emptyList()
|
||||
|
||||
val agent = ChatAgent(
|
||||
@@ -332,13 +335,11 @@ private fun runServer() {
|
||||
memoryReviewer = memoryReviewer,
|
||||
soulBody = soulBody,
|
||||
contextWindow = contextWindow,
|
||||
compressionThreshold = config.compressionThreshold,
|
||||
compressionThreshold = config.memory.compressionThreshold,
|
||||
contextCompactor = contextCompactor,
|
||||
recentReflections = recentReflections,
|
||||
reflector = reflector,
|
||||
reflectionInterval = config.reflectionInterval,
|
||||
skillMiner = skillMiner,
|
||||
skillMiningInterval = config.skillMiningInterval,
|
||||
)
|
||||
|
||||
// Куратор памяти: фоновая архивация stale-заметок. Поднимается до server'а,
|
||||
@@ -351,12 +352,17 @@ private fun runServer() {
|
||||
c
|
||||
} else null
|
||||
|
||||
val server = embeddedServer(CIO, port = config.port) {
|
||||
val server = embeddedServer(CIO, port = config.agent.port) {
|
||||
routing {
|
||||
get("/health") { call.respondText("ok") }
|
||||
agentikAgent(agent, path = "/agentik")
|
||||
a2aAgent(agentName = "agentik", handler = A2aBridge(agent), path = "/a2a")
|
||||
if (config.debugEndpoints) {
|
||||
agentikAgent(agent, path = "/agentik", token = config.agent.authToken)
|
||||
a2aAgent(
|
||||
agentName = "agentik",
|
||||
handler = A2aBridge(agent),
|
||||
path = "/a2a",
|
||||
token = config.agent.a2aToken,
|
||||
)
|
||||
if (config.debug.endpoints) {
|
||||
debugRoutes(
|
||||
agent = agent,
|
||||
storage = storage,
|
||||
@@ -368,20 +374,20 @@ private fun runServer() {
|
||||
}
|
||||
}
|
||||
}
|
||||
println("agentik standalone listening on http://localhost:${config.port}")
|
||||
println("agentik standalone listening on http://localhost:${config.agent.port}")
|
||||
println(" GET /health")
|
||||
println(" POST /agentik/conversations -> 201")
|
||||
println(" GET /agentik/conversations/{id}/events -> SSE")
|
||||
println(" POST /a2a/ -> A2A JSON-RPC (message/send, tasks/get, tasks/cancel)")
|
||||
println(" GET /a2a/.well-known/agent-card.json -> AgentCard")
|
||||
println(" storage: ${config.dbPath}")
|
||||
println(" storage: ${config.agent.dbPath}")
|
||||
println(" llm: ${config.llm.backend} ${config.llm.modelInfo()}")
|
||||
println(" mcp: ${mcpRegistry.allTools.size} tools from ${mcpRegistry.connectedServerCount} servers")
|
||||
println(" skills: ${skills.size} loaded${config.skillsDir?.let { " from $it" } ?: ""}")
|
||||
if (config.soulPath != null) println(" soul: ${config.soulPath} (${soulBody?.length ?: 0} chars)")
|
||||
println(" memory: ${if (memorySystem == null) "disabled" else "${config.memoryBackend.name.lowercase()}-backend"}")
|
||||
println(" skills: ${skills.size} loaded${config.agent.skillsDir?.let { " from $it" } ?: ""}")
|
||||
if (config.agent.soulPath != null) println(" soul: ${config.agent.soulPath} (${soulBody?.length ?: 0} chars)")
|
||||
println(" memory: ${if (memorySystem == null) "disabled" else "${config.memory.backend.name.lowercase()}-backend"}")
|
||||
if (contextWindow != null) {
|
||||
println(" compaction: enabled, threshold=${config.compressionThreshold}, window=$contextWindow tokens")
|
||||
println(" compaction: enabled, threshold=${config.memory.compressionThreshold}, window=$contextWindow tokens")
|
||||
} else {
|
||||
println(" compaction: disabled (OPENAI_CONTEXT_WINDOW not set)")
|
||||
}
|
||||
@@ -389,9 +395,9 @@ private fun runServer() {
|
||||
println(" curator: enabled (interval=${pw.binom.agentik.standalone.agent.memory.Curator.DEFAULT_INTERVAL}, maxAge=${pw.binom.agentik.standalone.agent.memory.Curator.DEFAULT_MAX_AGE})")
|
||||
}
|
||||
if (skillMiner != null) {
|
||||
println(" skill-mining: enabled (interval=${config.skillMiningInterval} turns, maxTurns=${config.skillMiningMaxTurns})")
|
||||
println(" skill-mining: enabled (interval=${config.skillMining.interval} turns, maxTurns=${config.skillMining.maxTurns})")
|
||||
}
|
||||
if (config.debugEndpoints) {
|
||||
if (config.debug.endpoints) {
|
||||
println(" debug endpoints: enabled (/debug/reflect, /debug/skill-mine, /debug/curate, /debug/compact, /debug/tokens)")
|
||||
}
|
||||
// Token stats по существующим диалогам (агрегат на старте — каждая запись
|
||||
|
||||
@@ -0,0 +1,84 @@
|
||||
package pw.binom.agentik.standalone.agent
|
||||
|
||||
import kotlinx.coroutines.channels.BufferOverflow
|
||||
import kotlinx.coroutines.flow.MutableSharedFlow
|
||||
import kotlinx.coroutines.flow.SharedFlow
|
||||
import kotlinx.coroutines.flow.asSharedFlow
|
||||
|
||||
/**
|
||||
* Internal event bus для background work — separate from [ConversationEvents]
|
||||
* (который это SSE-event stream для клиента).
|
||||
*
|
||||
* Background work fires на **structural events**, не на interval-polling:
|
||||
* - **Review** — триггерится в CompactionCoordinator ПРЯМО ПЕРЕД удалением ходов
|
||||
* из working memory (last chance вытащить факты). Это уже было сделано
|
||||
* через `reviewer.reviewPreCompaction()` — оставляем как есть.
|
||||
* - **Reflection** — на `ConversationLifecycleEvent.Closing` (финальная
|
||||
* рефлексия перед закрытием) ИЛИ накопление N tool-failures в окне
|
||||
* (что-то идёт не так).
|
||||
* - **Skill mining** — на `CompactionEvent.Triggered` если `turnsToDelete > N`
|
||||
* (есть контент для минера) ИЛИ на `ConversationLifecycleEvent.Closing`.
|
||||
*
|
||||
* Subscribers (BackgroundScheduler) решают, что делать. НЕ текстовая
|
||||
* инспекция, НЕ regex — только структурные события с явным семантическим
|
||||
* смыслом. См. STANDALONE-REVIEW раздел "event-driven background".
|
||||
*/
|
||||
|
||||
/** Эмитится из [ToolDispatcher] после каждого `runToolAndPersist` (success/failure). */
|
||||
sealed interface ToolCallEvent {
|
||||
val toolName: String
|
||||
|
||||
data class Succeeded(
|
||||
override val toolName: String,
|
||||
val durationMs: Long,
|
||||
) : ToolCallEvent
|
||||
|
||||
data class Failed(
|
||||
override val toolName: String,
|
||||
val error: String,
|
||||
) : ToolCallEvent
|
||||
}
|
||||
|
||||
/** Эмитится из [CompactionCoordinator] ПЕРЕД `workingMemory.compact(...)`. */
|
||||
sealed interface CompactionEvent {
|
||||
/**
|
||||
* Compaction сейчас удалит N ходов из working memory. Background work
|
||||
* имеет последний шанс вытащить оттуда данные.
|
||||
*/
|
||||
data class Triggered(
|
||||
val turnsToDelete: Int,
|
||||
val conversationId: String,
|
||||
) : CompactionEvent
|
||||
}
|
||||
|
||||
/** Эмитится из `ConversationLoop.close()` сразу ПЕРЕД `agentScope.cancel()`. */
|
||||
sealed interface ConversationLifecycleEvent {
|
||||
data class Closing(val conversationId: String) : ConversationLifecycleEvent
|
||||
}
|
||||
|
||||
internal class BackgroundEventBus {
|
||||
private val _toolCallEvents = MutableSharedFlow<ToolCallEvent>(
|
||||
replay = 0,
|
||||
extraBufferCapacity = 256,
|
||||
onBufferOverflow = BufferOverflow.DROP_OLDEST,
|
||||
)
|
||||
val toolCallEvents: SharedFlow<ToolCallEvent> get() = _toolCallEvents.asSharedFlow()
|
||||
|
||||
private val _compactionEvents = MutableSharedFlow<CompactionEvent>(
|
||||
replay = 0,
|
||||
extraBufferCapacity = 16,
|
||||
onBufferOverflow = BufferOverflow.DROP_OLDEST,
|
||||
)
|
||||
val compactionEvents: SharedFlow<CompactionEvent> get() = _compactionEvents.asSharedFlow()
|
||||
|
||||
private val _lifecycleEvents = MutableSharedFlow<ConversationLifecycleEvent>(
|
||||
replay = 0,
|
||||
extraBufferCapacity = 16,
|
||||
onBufferOverflow = BufferOverflow.DROP_OLDEST,
|
||||
)
|
||||
val lifecycleEvents: SharedFlow<ConversationLifecycleEvent> get() = _lifecycleEvents.asSharedFlow()
|
||||
|
||||
fun tryEmit(event: ToolCallEvent): Boolean = _toolCallEvents.tryEmit(event)
|
||||
fun tryEmit(event: CompactionEvent): Boolean = _compactionEvents.tryEmit(event)
|
||||
fun tryEmit(event: ConversationLifecycleEvent): Boolean = _lifecycleEvents.tryEmit(event)
|
||||
}
|
||||
+226
@@ -0,0 +1,226 @@
|
||||
package pw.binom.agentik.standalone.agent
|
||||
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Job
|
||||
import kotlinx.coroutines.flow.filterIsInstance
|
||||
import kotlinx.coroutines.flow.launchIn
|
||||
import kotlinx.coroutines.flow.merge
|
||||
import kotlinx.coroutines.flow.onEach
|
||||
import kotlinx.coroutines.launch
|
||||
import mu.KotlinLogging
|
||||
import pw.binom.agentik.memory.ConversationTurn
|
||||
import pw.binom.agentik.memory.MemoryReviewer
|
||||
import pw.binom.agentik.memory.MemoryStore
|
||||
import pw.binom.agentik.skills.SkillStore
|
||||
import pw.binom.agentik.storage.Content
|
||||
import pw.binom.agentik.storage.ReflectionStore
|
||||
import pw.binom.agentik.storage.WorkingMemoryEntry
|
||||
import pw.binom.agentik.storage.WorkingMemoryStore
|
||||
import pw.binom.agentik.llm.tools.LlmReflector
|
||||
import pw.binom.agentik.llm.tools.SkillMiner
|
||||
import java.util.concurrent.atomic.AtomicLong
|
||||
|
||||
/**
|
||||
* Background work подписчик на [BackgroundEventBus]. Заменяет старую interval-based
|
||||
* логику (`maybeScheduleReview/Reflection/SkillMining` с `userTurnCount % N == 0`).
|
||||
*
|
||||
* Подписки:
|
||||
* - [ConversationLifecycleEvent.Closing] → финальный reflection + skill mining
|
||||
* перед закрытием conversation (last chance вытащить insights).
|
||||
* - [CompactionEvent.Triggered] → skill mining если `turnsToDelete > MIN_COMPACTION_FOR_MINING`.
|
||||
* Review уже сделан внутри CompactionCoordinator (`reviewPreCompaction`) — не дублируем.
|
||||
* - [ToolCallEvent.Failed] (накопительно) → reflection если 2+ фейлов в окне 60 сек
|
||||
* (что-то пошло не так — самоанализ полезен).
|
||||
*
|
||||
* НЕ текстовая инспекция, НЕ regex, НЕ interval-polling. Только структурные
|
||||
* события с явным семантическим смыслом.
|
||||
*/
|
||||
internal data class BackgroundConfig(
|
||||
val memoryReviewer: MemoryReviewer?,
|
||||
val memoryStore: MemoryStore?,
|
||||
val reflectionStore: ReflectionStore?,
|
||||
val reflector: LlmReflector?,
|
||||
val skillMiner: SkillMiner?,
|
||||
val skillMiningStore: SkillStore?,
|
||||
)
|
||||
|
||||
internal class BackgroundScheduler(
|
||||
private val state: ConversationState,
|
||||
private val workingMemory: WorkingMemoryStore,
|
||||
private val config: BackgroundConfig,
|
||||
private val backgroundEvents: BackgroundEventBus,
|
||||
) {
|
||||
private val log = KotlinLogging.logger {}
|
||||
|
||||
private val lastSkillMiningAt = AtomicLong(0)
|
||||
private val lastReflectionAt = AtomicLong(0)
|
||||
|
||||
/** Recent tool failure timestamps (ms). Trimmed to [FAILURE_WINDOW_MS]. */
|
||||
private val toolFailures = mutableListOf<Long>()
|
||||
private val toolFailuresLock = Any()
|
||||
|
||||
private var subscriptionJob: Job? = null
|
||||
|
||||
/** Запустить подписки. Вызывать один раз после конструктора. */
|
||||
fun start(scope: CoroutineScope) {
|
||||
if (subscriptionJob?.isActive == true) return
|
||||
subscriptionJob = scope.launch {
|
||||
// Merge all three event flows into one subscription scope. Each onEach
|
||||
// returns Unit, so launchIn merges them as cold flows.
|
||||
merge(
|
||||
backgroundEvents.lifecycleEvents
|
||||
.filterIsInstance<ConversationLifecycleEvent.Closing>()
|
||||
.onEach { onClosing() },
|
||||
backgroundEvents.compactionEvents
|
||||
.filterIsInstance<CompactionEvent.Triggered>()
|
||||
.onEach { onCompaction(it) },
|
||||
backgroundEvents.toolCallEvents
|
||||
.filterIsInstance<ToolCallEvent.Failed>()
|
||||
.onEach { onToolFailure() },
|
||||
).collect {}
|
||||
}
|
||||
}
|
||||
|
||||
private suspend fun onClosing() {
|
||||
if (state.isTemporal) return
|
||||
val convId = state.id
|
||||
runFinalReflection(convId)
|
||||
runFinalSkillMining(convId)
|
||||
}
|
||||
|
||||
private fun runFinalReflection(convId: String) {
|
||||
val reflector = config.reflector ?: return
|
||||
val store = config.reflectionStore ?: return
|
||||
state.agentScope.launch {
|
||||
try {
|
||||
val turns = recentTurns(reflector.maxTurns)
|
||||
if (turns.isEmpty()) return@launch
|
||||
val reflection = reflector.reflect(turns) ?: return@launch
|
||||
runCatching { store.insert(reflection.copy(conversationId = convId)) }
|
||||
.onFailure { log.warn(it) { "final reflection insert failed: ${it.message}" } }
|
||||
log.info { "final reflection on close: conv=$convId score=${reflection.score}/5" }
|
||||
} catch (e: Throwable) {
|
||||
log.warn(e) { "final reflection failed for $convId: ${e.message}" }
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private fun runFinalSkillMining(convId: String) {
|
||||
val miner = config.skillMiner ?: return
|
||||
val store = config.skillMiningStore ?: return
|
||||
state.agentScope.launch {
|
||||
try {
|
||||
val turns = recentTurns(miner.maxTurns)
|
||||
if (turns.isEmpty()) return@launch
|
||||
val mined = miner.mine(turns, store.catalog.skills)
|
||||
for (s in mined) {
|
||||
runCatching { store.upsert(s) }
|
||||
.onFailure { log.warn(it) { "final skill mining upsert '${s.name}' failed: ${it.message}" } }
|
||||
}
|
||||
log.info { "final skill mining on close: conv=$convId turns=${turns.size} existing=${store.catalog.skills.size} mined=${mined.size}" }
|
||||
} catch (e: Throwable) {
|
||||
log.warn(e) { "final skill mining failed for $convId: ${e.message}" }
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private fun onCompaction(event: CompactionEvent.Triggered) {
|
||||
if (state.isTemporal) return
|
||||
if (event.turnsToDelete < MIN_COMPACTION_FOR_MINING) return
|
||||
val miner = config.skillMiner ?: return
|
||||
val store = config.skillMiningStore ?: return
|
||||
val now = System.currentTimeMillis()
|
||||
// Debounce: не чаще раза в минуту
|
||||
if (now - lastSkillMiningAt.get() < MINING_DEBOUNCE_MS) return
|
||||
lastSkillMiningAt.set(now)
|
||||
val convId = event.conversationId
|
||||
state.agentScope.launch {
|
||||
try {
|
||||
val turns = recentTurns(miner.maxTurns)
|
||||
if (turns.isEmpty()) return@launch
|
||||
val mined = miner.mine(turns, store.catalog.skills)
|
||||
for (s in mined) {
|
||||
runCatching { store.upsert(s) }
|
||||
.onFailure { log.warn(it) { "compaction skill mining upsert '${s.name}' failed: ${it.message}" } }
|
||||
}
|
||||
log.info { "compaction skill mining: conv=$convId turnsToDelete=${event.turnsToDelete} mined=${mined.size}" }
|
||||
} catch (e: Throwable) {
|
||||
log.warn(e) { "compaction skill mining failed for $convId: ${e.message}" }
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private fun onToolFailure() {
|
||||
val reflector = config.reflector ?: return
|
||||
if (state.isTemporal) return
|
||||
val store = config.reflectionStore ?: return
|
||||
|
||||
val now = System.currentTimeMillis()
|
||||
val shouldReflect = synchronized(toolFailuresLock) {
|
||||
toolFailures.add(now)
|
||||
// Trim old failures outside the window
|
||||
val cutoff = now - FAILURE_WINDOW_MS
|
||||
toolFailures.removeAll { it < cutoff }
|
||||
toolFailures.size >= FAILURE_THRESHOLD
|
||||
}
|
||||
if (!shouldReflect) return
|
||||
|
||||
// Debounce reflection globally (не чаще раза в 5 мин)
|
||||
if (now - lastReflectionAt.get() < REFLECTION_DEBOUNCE_MS) {
|
||||
log.debug { "reflection debounced: ${toolFailures.size} failures accumulated but reflection fired recently" }
|
||||
return
|
||||
}
|
||||
lastReflectionAt.set(now)
|
||||
|
||||
// Clear failure window — fresh accounting period
|
||||
synchronized(toolFailuresLock) { toolFailures.clear() }
|
||||
|
||||
val convId = state.id
|
||||
state.agentScope.launch {
|
||||
try {
|
||||
val turns = recentTurns(reflector.maxTurns)
|
||||
if (turns.isEmpty()) return@launch
|
||||
val reflection = reflector.reflect(turns) ?: return@launch
|
||||
runCatching { store.insert(reflection.copy(conversationId = convId)) }
|
||||
.onFailure { log.warn(it) { "reflection after failures insert failed: ${it.message}" } }
|
||||
log.info { "reflection triggered by tool failures: conv=$convId score=${reflection.score}/5" }
|
||||
} catch (e: Throwable) {
|
||||
log.warn(e) { "reflection after tool failures failed for $convId: ${e.message}" }
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private suspend fun recentTurns(maxTurns: Int): List<ConversationTurn> {
|
||||
val rows = workingMemory.list(state.id)
|
||||
val pairs = mutableListOf<ConversationTurn>()
|
||||
var pendingUser: String? = null
|
||||
for (row in rows) {
|
||||
when (val e = row.entry) {
|
||||
is WorkingMemoryEntry.User -> pendingUser = e.content.text()
|
||||
is WorkingMemoryEntry.Assistant -> {
|
||||
val user = pendingUser ?: ""
|
||||
pendingUser = null
|
||||
pairs += ConversationTurn(userMessage = user, assistantMessage = e.content.text())
|
||||
}
|
||||
else -> {}
|
||||
}
|
||||
}
|
||||
return pairs.takeLast(maxTurns)
|
||||
}
|
||||
|
||||
private fun List<Content>.text(): String =
|
||||
filterIsInstance<Content.Text>().joinToString("\n") { it.body }
|
||||
|
||||
companion object {
|
||||
/** Минимум ходов, удаляемых compaction'ом, чтобы trigger'ить skill mining. */
|
||||
private const val MIN_COMPACTION_FOR_MINING = 10
|
||||
/** Дебаунс skill mining между запусками. */
|
||||
private const val MINING_DEBOUNCE_MS = 60_000L
|
||||
/** Debounce reflection между запусками (накопительный, не per-failure). */
|
||||
private const val REFLECTION_DEBOUNCE_MS = 300_000L
|
||||
/** Сколько tool-failures в окне должно накопиться чтобы trigger reflection. */
|
||||
private const val FAILURE_THRESHOLD = 2
|
||||
/** Окно для accumulation tool-failures. */
|
||||
private const val FAILURE_WINDOW_MS = 60_000L
|
||||
}
|
||||
}
|
||||
Some files were not shown because too many files have changed in this diff Show More
Reference in New Issue
Block a user