agent: связь с сервером + UI-экраны + ASR Qwen3 + settings

This commit is contained in:
2026-09-23 11:53:57 +03:00
parent d5aa844fa0
commit b923bc6196
29 changed files with 3845 additions and 32 deletions
+27 -1
View File
@@ -41,8 +41,34 @@ android {
}
dependencies {
// Ядро агента (Nexus caffeine). Тянет :proto, ktor-client-core и т.д.
// Ядро агента (Nexus caffeine). Тянет :proto, :journal-api (транзитивно в v9), ktor-client-core и т.д.
implementation(libs.agentik.client)
// journal-api объявлен отдельно, чтобы IDE/компилятор видел MessageRecord в AgentMappers
// (транзитивная зависимость в v9, но прямой импорт делаем явным).
implementation(libs.agentik.journal.api)
// InMemoryJournalStore — RAM-кэш истории для HttpChatSession. Персистентный
// кэш (SQLite) подключим отдельной задачей.
implementation(libs.agentik.journal.inmemory)
// Persistent journal-кеш через ksqlite. Подменяет InMemoryJournalStore —
// история сообщений переживает kill приложения. Начиная с journal-ksqlite:12
// конструктор `KsqliteJournalStore(path)` публичный — рефлексия и
// явная зависимость `pw.binom.db:ksqlite` больше не нужны.
implementation(libs.agentik.journal.ksqlite)
// ksqlite опубликован в Maven Central (0.1.2). Наш ConversationMetaRepository
// работает с SQLiteConnection напрямую — общая абстракция
// :journal-ksqlite (MutableJournalStore) про метаданные диалогов не знает.
// Декларируем явно, чтобы транзитивная compile-api за journal-ksqlite не
// подсунула несовместимую версию.
implementation(libs.ksqlite)
// ASR-стек: модель Qwen3-ASR на устройстве (auto-download через QwenModelProvider).
implementation(libs.asr.qwen3)
implementation(libs.mic.api)
implementation(libs.androidx.work.runtime.ktx)
// Markdown-парсер (ru.otpbank.ai:markdown) — нужно для рендеринга
// сообщений ассистента. Хостится в нашем Nexus (см. settings.gradle.kts).
implementation(libs.ru.otpbank.markdown)
implementation(libs.androidx.core.ktx)
implementation(libs.androidx.activity.compose)
+37
View File
@@ -0,0 +1,37 @@
{
"version": "06",
"source": "http://static.binom.pw/models/qwen3_06/",
"totalBytes": 987015347,
"files": [
{
"path": "conv_frontend.onnx",
"size": 44148281,
"sha256": "d22dc4423e0940e49884e903d2ea2f7e5567c14fc1aed97e4e26d6b8f208ef9e"
},
{
"path": "encoder.int8.onnx",
"size": 182491662,
"sha256": "60748d3e6744a57c9c91e1b17424a6c2990567e8adceb0783940c03ed98fa9d9"
},
{
"path": "decoder.int8.onnx",
"size": 755914231,
"sha256": "4f6885be5959ae26af3089d38ee7972c5fafbeeb1cf8d5e76eab6d8b61ca5771"
},
{
"path": "tokenizer/merges.txt",
"size": 1671853,
"sha256": "8831e4f1a044471340f7c0a83d7bd71306a5b867e95fd870f74d0c5308a904d5"
},
{
"path": "tokenizer/tokenizer_config.json",
"size": 12487,
"sha256": "4942d005604266809309cabc9f4e9cb89ce855d59b14681fdc0e1cc62ea26c4c"
},
{
"path": "tokenizer/vocab.json",
"size": 2776833,
"sha256": "ca10d7e9fb3ed18575dd1e277a2579c16d108e32f27439684afa0e10b1440910"
}
]
}
@@ -1,9 +1,77 @@
package pw.binom.agentik.android
import android.app.Application
import androidx.work.BackoffPolicy
import androidx.work.Constraints
import androidx.work.ExistingPeriodicWorkPolicy
import androidx.work.NetworkType
import androidx.work.PeriodicWorkRequestBuilder
import androidx.work.WorkManager
import pw.binom.agentik.android.agent.AgentConnection
import pw.binom.agentik.android.agent.real.HttpAgentConnection
import pw.binom.agentik.android.asr.QwenDownloadWorker
import pw.binom.agentik.android.settings.SettingsRepository
import java.util.concurrent.TimeUnit
/**
* Точка входа приложения. Пока пустая: сюда будут подниматься настройки
* (адрес агента и токен) и хранилище.
* Точка входа приложения. Поднимает фоновые сервисы:
* - периодический [QwenDownloadWorker]: раз в 6 часов проверяет/докачивает голосовую модель
* Qwen3-ASR по Wi-Fi (UNMETERED). Зарегистрирован через [ExistingPeriodicWorkPolicy.KEEP],
* чтобы перезапуск приложения не сбрасывал уже запланированный запуск и не дёргал воркер
* лишний раз.
* - singleton-фабрика [agentConnection] — единая точка подключения к агенту,
* живёт до вызова `close()` (никогда в этой PR — приложение короткоживущее).
* В фазах тестирования можно подменить на stub: `(application as AgentikApp).agentConnection = StubAgentConnection(...)`.
* - [settingsRepository] — JSON-персистентность [pw.binom.agentik.android.settings.Settings]
* в `filesDir/settings.json`. Передаётся в UI через [MainActivity].
*/
class AgentikApp : Application()
class AgentikApp : Application() {
/**
* Singleton-фабрика агента-фасада. Создаётся лениво при первом обращении,
* живёт до вызова `close()` (никогда в этой PR — приложение короткоживущее).
*
* NB: используем `lateinit` (а не `by Delegates.notNull()`), чтобы
* работал `::agentConnection.isInitialized` — `Delegates.notNull()`
* делегат не поддерживает эту проверку.
*/
lateinit var agentConnection: AgentConnection
private set
/**
* Singleton-репозиторий для JSON-персистентности настроек в `filesDir/settings.json`.
* Создаётся лениво при первом обращении (после `onCreate()`) — `Application`
* сам передаёт себя как `Context`.
*/
val settingsRepository: SettingsRepository by lazy { SettingsRepository(this) }
override fun onCreate() {
super.onCreate()
if (!::agentConnection.isInitialized) {
agentConnection = HttpAgentConnection(context = this)
}
scheduleQwenDownload()
}
private fun scheduleQwenDownload() {
val request = PeriodicWorkRequestBuilder<QwenDownloadWorker>(6, TimeUnit.HOURS)
.setConstraints(
Constraints.Builder()
// Wi-Fi only — мобильная сеть не качает 940 МБ.
.setRequiredNetworkType(NetworkType.UNMETERED)
.build(),
)
.setBackoffCriteria(BackoffPolicy.EXPONENTIAL, 30, TimeUnit.MINUTES)
.build()
WorkManager.getInstance(this).enqueueUniquePeriodicWork(
UNIQUE_WORK_NAME,
ExistingPeriodicWorkPolicy.KEEP,
request,
)
}
private companion object {
const val UNIQUE_WORK_NAME = "qwen-model-download"
}
}
@@ -4,42 +4,21 @@ import android.os.Bundle
import androidx.activity.ComponentActivity
import androidx.activity.compose.setContent
import androidx.activity.enableEdgeToEdge
import androidx.compose.foundation.layout.Column
import androidx.compose.foundation.layout.fillMaxSize
import androidx.compose.foundation.layout.padding
import androidx.compose.material3.MaterialTheme
import androidx.compose.material3.Scaffold
import androidx.compose.material3.Text
import androidx.compose.runtime.Composable
import androidx.compose.ui.Modifier
import androidx.compose.ui.unit.dp
import pw.binom.agentik.android.ui.screen.AppRoot
import pw.binom.agentik.android.ui.theme.AgentikTheme
class MainActivity : ComponentActivity() {
override fun onCreate(savedInstanceState: Bundle?) {
super.onCreate(savedInstanceState)
enableEdgeToEdge()
val app = application as AgentikApp
setContent {
AgentikTheme {
Scaffold(modifier = Modifier.fillMaxSize()) { innerPadding ->
StartScreen(Modifier.padding(innerPadding))
}
AppRoot(
connection = app.agentConnection,
settingsRepository = app.settingsRepository,
)
}
}
}
}
@Composable
private fun StartScreen(modifier: Modifier = Modifier) {
Column(modifier = modifier.padding(24.dp)) {
Text(
text = "Agentik",
style = MaterialTheme.typography.headlineMedium,
)
Text(
text = "Клиент готов к работе. Дальше — список диалогов и чат.",
style = MaterialTheme.typography.bodyMedium,
modifier = Modifier.padding(top = 12.dp),
)
}
}
@@ -0,0 +1,80 @@
package pw.binom.agentik.android.agent
import kotlinx.coroutines.flow.StateFlow
import pw.binom.agentik.journal.ConversationStore
/**
* Фасад над `pw.binom.agentik.proto.Agent`. Скрывает proto-типы от UI.
*
* Реализации: [stub.StubAgentConnection] (текущая фаза) и HTTP-импл на
* `AgentikAgent(...)` (следующая фаза). Сам контракт lifecycle-а:
* один [AgentConnection] на приложение, на каждый диалог — отдельный
* [ChatSession], живущий до явного [close].
*/
interface AgentConnection : AutoCloseable {
/** Текущее состояние подключения к агенту. */
val state: StateFlow<ConnectionState>
/**
* Read-only view на список диалогов. В v12 клиентский [AgentikAgent]
* оборачивает [pw.binom.agentik.proto.Agent.conversationStore] в локальный
* in-memory кэш, синхронизируемый через outbox-события — UI через
* [ConversationStore.listFlow] получает auto-refresh без ручного refetch.
*
* Доступ к [conversationStore] валиден только после успешного [configure].
*/
val conversationStore: ConversationStore
/**
* Подключиться/переподключиться к агенту. Бросает на ошибке конфигурации
* (например, недоступный URL). На стабе — no-op после валидации URL.
*/
suspend fun configure(id: String, name: String, baseUrl: String, token: String?)
/** Список диалогов из локального кэша. */
suspend fun listConversations(): List<ConversationSummary>
/**
* Открыть существующий диалог (если [convId] != null) или создать новый.
* Возвращённый [ChatSession] живёт, пока не закрыт через `close()`.
*/
fun openConversation(convId: String? = null): ChatSession
override fun close()
}
sealed interface ConnectionState {
/** Не сконфигурирован. */
data object Disconnected : ConnectionState
/** Подключаемся / проверяем связь. */
data class Connecting(val agentId: String) : ConnectionState
/** Связь есть, агент отвечает. [agentName] — как его зовёт сервер/клиент. */
data class Connected(val agentId: String, val agentName: String) : ConnectionState
/** Связи нет / конфиг плохой. */
data class Failed(val reason: String) : ConnectionState
}
/**
* Локально-кэшированная «карточка» диалога. Превью последнего сообщения и
* счётчик непрочитанных берутся из локального кэша, не из протокола —
* десктопные требования R24/R26 (нет данных — нет цифры).
*
* [unreadCount] теперь **всегда известен** (v15+):
* - диалог ни разу не открывали (`last_seen_message_at` отсутствует) →
* `journal.count(convId)` (все сообщения — непрочитанные);
* - открывали → `journal.count(convId, lastSeen)` (только новые после lastSeen).
*
* `0` — открывали, но новых сообщений нет. UI бейдж скрывается.
*/
data class ConversationSummary(
val id: String,
val title: String?,
val updatedAt: kotlin.time.Instant,
val lastMessagePreview: String?,
val unreadCount: Int,
)
@@ -0,0 +1,94 @@
package pw.binom.agentik.android.agent
import pw.binom.agentik.journal.MessageRecord
import pw.binom.agentik.outbox.Event
/**
* Маппинг `outbox.Event` → `UiEvent` и `journal.MessageRecord` → `UiMessage`.
*
* Сигнатуры сверялись по `proto-jvm-9-sources.jar`, `outbox-api-jvm-9-sources.jar`
* и `journal-api-jvm-9-sources.jar`.
*
* **NB**: `Event` в v9 живёт в `pw.binom.agentik.outbox.*`. В `pw.binom.agentik.proto.*`
* есть `typealias Event = pw.binom.agentik.outbox.Event`, но Kotlin не
* резолвит nested-классы через typealias при `is`-проверках, поэтому импортируем
* `outbox.Event` напрямую.
*
* - `Event` — расширяемый протокольный sealed-тип; `when` НЕ exhaustive
* (`else -> null`): клиент должен пережить добавление новых вариантов на
* сервере без падений маппинга.
* - `MessageRecord` — стабильный (наш кэш); `when` exhaustive, без `else`.
* Если journal-api получит новый вариант, маппер упадёт на этапе компиляции —
* это сигнал разработчику добавить ветку.
*/
/**
* Маппинг `outbox.Event` → `UiEvent`. Возвращает `null` для вариантов, которые
* UI сейчас игнорирует (`StartResponse`/`StartReasoning`/`Interrupted`).
*
* **NB по сигнатурам v9** (сверялся по `outbox-api-jvm-9-sources.jar`):
* - `Event.AppendText(date, body)` → [UiEvent.AppendText]
* - `Event.AppendImage(date, body, mime)` → [UiEvent.AppendImage] (mime игнорируем на UI-слое)
* - `Event.ToolCall(date, id, title, toolName, toolArgs)` → [UiEvent.ToolCall]
* - `Event.ToolResult(date, id, result)` → [UiEvent.ToolResult]. В v9
* `Event.ToolResult` НЕ несёт `toolName` — UI связывает по `id`.
* - `Event.Error(date, message, code)` → [UiEvent.Error]
*/
@Suppress("REDUNDANT_ELSE_IN_WHEN") // `else` оставлен намеренно: когда протокол добавит вариант, наш маппер не сломается.
fun Event.toUiEvent(): UiEvent? = when (this) {
is Event.StartResponse -> null // UI на это не реагирует
is Event.StartReasoning -> null // рассуждения не рендерим (пока)
is Event.AppendText -> UiEvent.AppendText(body)
is Event.AppendImage -> UiEvent.AppendImage(body)
is Event.ToolCall -> UiEvent.ToolCall(id, toolName, title, toolArgs)
is Event.ToolResult -> UiEvent.ToolResult(toolCallId = toolCallId, toolName = toolName, result = result)
is Event.End -> UiEvent.End
is Event.Interrupted -> null // ход «закрывается» без сообщения
is Event.Error -> UiEvent.Error(message, code)
else -> null // неизвестный вариант — пропустить
}
/**
* Маппинг `pw.binom.agentik.journal.MessageRecord` → `UiMessage`.
*
* `MessageRecord` — sealed, стабильный: `when` exhaustive без `else`.
* Если journal-api получит новый вариант — компилятор скажет добавить ветку.
*
* `interrupted` у [UiMessage.Assistant] всегда `false`: interrupted живёт
* только в live-`Event.Interrupted`, в журнал не пишется.
*/
fun MessageRecord.toUiMessage(): UiMessage = when (this) {
is MessageRecord.UserMessage -> UiMessage.User(id, createdAt, content.toPlainText())
is MessageRecord.AssistantMessage -> UiMessage.Assistant(
id, createdAt,
content.toPlainText(),
interrupted = false,
)
is MessageRecord.ToolCall -> UiMessage.ToolCall(
id, createdAt,
toolName,
title = toolTitle,
args = toolArgsJson,
)
is MessageRecord.ToolResult -> UiMessage.ToolResult(
id, createdAt,
toolCallId,
toolName = toolName,
result = result,
)
is MessageRecord.Error -> UiMessage.Error(
id, createdAt,
message,
code = code,
)
}
/**
* Склеивает текстовые куски [pw.binom.agentik.journal.Content.Text] в
* одну строку — для плоского [UiMessage.User.text]/[UiMessage.Assistant.text].
* Изображения и прочие не-текстовые content-ы игнорируются (для показа
* картинок в истории будет отдельный UI-элемент).
*/
private fun List<pw.binom.agentik.journal.Content>.toPlainText(): String =
filterIsInstance<pw.binom.agentik.journal.Content.Text>()
.joinToString(separator = "") { it.body }
@@ -0,0 +1,54 @@
package pw.binom.agentik.android.agent
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.StateFlow
import kotlin.time.Instant
/**
* Сессия одного диалога. Хранит локальный кэш истории (для backfill),
* подписку на live-события текущего хода и снимок метаданных диалога.
*
* Создаётся через [AgentConnection.openConversation], закрывается
* владельцем (или через [AgentConnection.close]).
*/
interface ChatSession : AutoCloseable {
val conversationId: String
/** Снимок диалога (метаданные). Обновляется на rename / send. */
val snapshot: StateFlow<ConversationSnapshot?>
/** История из локального кэша. Не содержит ещё не «закрытые» ходы (нет `End`). */
val messages: StateFlow<List<UiMessage>>
/** Live-события текущего хода. Используется для streaming-отрисовки. */
val liveEvents: Flow<UiEvent>
/** Отправить текстовое сообщение. */
suspend fun send(text: String)
/** Прервать текущий ход агента (R11). */
suspend fun interrupt()
/** Переименовать диалог. */
suspend fun rename(newTitle: String)
/** Явно докачать новые записи с сервера (после End или ручного refresh). */
suspend fun refresh()
override fun close()
}
/**
* Метаданные диалога. Зеркалит `pw.binom.agentik.proto.ConversationSnapshot`,
* но без зависимости на proto-тип. Имена полей соответствуют REST DTO
* `:client` (`isSupportImageInput`/`isSupportImageOutput`/`isTemporal`).
*/
data class ConversationSnapshot(
val id: String,
val title: String?,
val updatedAt: Instant,
val supportsImageInput: Boolean,
val supportsImageOutput: Boolean,
val isTemporal: Boolean,
)
@@ -0,0 +1,108 @@
package pw.binom.agentik.android.agent
import kotlin.time.Instant
/**
* UI-проекция записи в истории диалога. Полностью отвязана от протокольных
* типов `:client`/`proto`/`journal` — маппинг делается в [AgentMappers].
*
* Все варианты несут стабильный [id] и [timestamp] для ключей Compose-списков
* и виртуализации.
*/
sealed interface UiMessage {
val id: String
val timestamp: Instant
data class User(
override val id: String,
override val timestamp: Instant,
val text: String,
) : UiMessage
/**
* Ассистентский ход. [interrupted] на исторических записях **всегда**
* `false`: interrupted живёт в `Event.Interrupted` в live-стриме, в журнал
* не попадает (см. [AgentMappers]).
*/
data class Assistant(
override val id: String,
override val timestamp: Instant,
val text: String,
val interrupted: Boolean = false, // прерван кнопкой «Стоп» — R39
) : UiMessage
data class ToolCall(
override val id: String,
override val timestamp: Instant,
val toolName: String,
val title: String? = null,
val args: String? = null,
) : UiMessage
/**
* Результат инструмента. [toolCallId] — FK на [ToolCall.id], для UI-связи.
* [toolName] nullable для обратной совместимости: в v10 добавлено, в
* исторических записях (до v10) может быть `null`.
*/
data class ToolResult(
override val id: String,
override val timestamp: Instant,
val toolCallId: String,
val toolName: String? = null,
val result: String? = null,
) : UiMessage
/**
* Ошибочный ход. [code] опционален — в v9 у `MessageRecord.Error` это
* `String?`, исторически могло быть `null`.
*/
data class Error(
override val id: String,
override val timestamp: Instant,
val message: String,
val code: String? = null,
) : UiMessage
}
/**
* Live-события текущего хода. Эмитятся в [ChatSession.liveEvents] по мере
* того, как агент стримит ответ. В UI накапливаются в текущий «бабл»
* ассистента до прихода [End] / [Error] / `interrupted`.
*
* **NB**: `Interrupted` намеренно отсутствует — на стабе ход просто
* завершается, а реальный протокол эмитит его отдельным
* `proto.Event.Interrupted`. В фазе «жир» добавим отдельный вариант,
* когда UI начнёт его рендерить.
*/
sealed interface UiEvent {
data class AppendText(val body: String) : UiEvent
data class AppendImage(val body: ByteArray) : UiEvent
/**
* Tool-call из live-стрима. [id] — стабильный идентификатор вызова
* (из `Event.ToolCall.id`); по нему UI свяжет последующий
* [ToolResult.id] с этим вызовом.
*/
data class ToolCall(
val id: String,
val toolName: String,
val title: String? = null,
val args: String? = null,
) : UiEvent
/**
* Tool-result из live-стрима. [toolCallId] — FK на [ToolCall.id]
* (бывшее поле `id` в v9, переименовано в v10 для консистентности с
* [MessageRecord.ToolResult.toolCallId] и [UiMessage.ToolResult.toolCallId]).
* [toolName] nullable, доступно с v10 — раньше UI держал локальный
* `Map<toolCallId, toolName>`, теперь не нужен.
*/
data class ToolResult(
val toolCallId: String,
val toolName: String? = null,
val result: String? = null,
) : UiEvent
data object End : UiEvent
data class Error(val message: String, val code: String? = null) : UiEvent
}
@@ -0,0 +1,163 @@
package pw.binom.agentik.android.agent.real
import android.util.Log
import kotlinx.coroutines.CancellationException
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Job
import kotlinx.coroutines.launch
import pw.binom.agentik.journal.MessageRecord
import pw.binom.agentik.journal.MutableConversationStore
import pw.binom.agentik.journal.MutableJournalStore
import pw.binom.agentik.outbox.AgentEvent
import pw.binom.agentik.outbox.CommonEvent
import pw.binom.agentik.outbox.Event
import pw.binom.agentik.proto.Agent
import kotlin.time.Instant
/**
* Background-sync из `agent.outbox.events()` в локальный SQLite-кэш.
*
* Запускается в [HttpAgentConnection.configure] сразу после успешного
* `AgentikAgent(...)` + [ReconnectingOutbox]. Не требует открытия чата —
* слушает глобальный поток событий агента и переносит их в:
*
* - [conversationStore] (таблица `conversation`) — `Created` /
* `Deleted` / `Renamed` / `Touched`;
* - [journal] (таблица `message`) — полный refresh на каждый `Event.End`;
* - [metaRepo] (таблица `conversation_meta`) — `markSeen` после End для
* корректного счётчика unread в [HttpAgentConnection.listConversations].
*
* **Не пишет tool-call/tool-result/error в фоне** — они придут через
* полный refresh после `Event.End` (там же захватятся streamed
* `UserMessage`/`AssistantMessage`, которые через outbox передаются только
* инкрементально как `AppendText`, а в журнал фиксируются целиком при
* завершении хода). Это зеркалит поведение `HttpChatSession.Event.End`.
*
* **Один subscribe на поток**: используем тот же `ReconnectingOutbox`,
* который уже создан в [HttpAgentConnection]; повторный consumer у общего
* `MutableSharedFlow` безопасен (broadcast).
*
* **Закрытие**: [close] отменяет [syncJob] — потребители `journal` /
* `conversationStore` / `metaRepo` не закрываются (они переиспользуются
* на уровне соединения и закрываются в [HttpAgentConnection.close]).
*/
internal class AgentEventSync(
private val agent: Agent,
private val conversationStore: MutableConversationStore,
private val journal: MutableJournalStore,
private val metaRepo: ConversationMetaRepository,
private val parentScope: CoroutineScope,
private val withSqliteLock: suspend ((suspend () -> Unit) -> Unit),
) : AutoCloseable {
private val syncJob: Job = parentScope.launch {
try {
agent.outbox.events(Instant.DISTANT_PAST).collect { commonEvent ->
// Сериализуем с HttpChatSession и listConversations — без мьютекса
// параллельный writer может вклиниться между cache.append и
// markSeen, оставив несогласованный счётчик unread.
withSqliteLock {
when (commonEvent) {
is CommonEvent.Agent -> handleAgent(commonEvent.event)
is CommonEvent.Conversation -> handleConversation(commonEvent)
}
}
}
} catch (e: CancellationException) {
throw e
} catch (e: Exception) {
// outbox-стрим закрывается штатно после Failed; логируем на debug,
// чтобы не засорять logcat на ожидаемых кейсах (disconnect/long offline).
Log.d(TAG, "events stream terminated: ${e.message}")
}
}
private suspend fun handleAgent(ev: AgentEvent) {
when (ev) {
is AgentEvent.Created -> {
// Достаём полный record через локальный in-memory кэш Agent'а
// (см. `wrapWithLocalConversationCache` в `:client`) — там же
// происходит merge из этого же outbox. Если seed в кэше ещё не
// случился (rare: connect при offline), get() вернёт null —
// следующее Touched/refresh догонит.
agent.conversationStore.get(ev.conversationId)
?.let { conversationStore.upsert(it) }
}
is AgentEvent.Deleted -> {
conversationStore.delete(ev.id)
}
is AgentEvent.Renamed -> {
val existing = conversationStore.get(ev.id) ?: return
conversationStore.upsert(existing.copy(title = ev.title))
}
is AgentEvent.Touched -> {
conversationStore.touch(ev.id, ev.updatedAt)
}
}
}
private suspend fun handleConversation(ev: CommonEvent.Conversation) {
val convId = ev.conversationId
when (val event = ev.event) {
is Event.End -> refreshAfterEnd(convId, event.date)
// StartResponse / StartReasoning / AppendText / AppendImage /
// Interrupted / ToolCall / ToolResult / ToolFailed / Error /
// ConversationClosing / CompactionTriggered — промежуточные или
// системные события; полная транскрипция восстановится через
// refreshAfterEnd на ближайшем End (или при следующем открытии чата).
else -> Unit
}
}
private suspend fun refreshAfterEnd(convId: String, endDate: Instant) {
// Cursor bug fix (как в HttpChatSession): читаем из metaRepo,
// а не из cache.list(limit=1).lastOrNull() (последний возвращает OLDEST
// из-за ASC-сортировки — см. memory `ui-wireup-state`).
val lastSeen = metaRepo.getLastSeen(convId) ?: Instant.DISTANT_PAST
try {
agent.journal.list(convId, lastSeen, offset = 0, limit = REFRESH_LIMIT)
.forEach { record -> appendIfAbsent(record) }
} catch (e: Exception) {
// Оффлайн / 5xx — пропускаем refresh, остальные writers (markSeen
// / upsert ConversationRecord) тоже не выполняем, чтобы не оставлять
// несогласованного состояния.
Log.w(TAG, "refresh-on-End for $convId failed: ${e.message}", e)
return
}
// markSeen ДО upsert ConversationRecord — порядок не критичен, но так
// логично: сначала фиксируем прогресс чтения, потом обновляем карточку диалога.
metaRepo.markSeen(convId, endDate)
// updatedAt диалога мог измениться в этом ходе — синхронизируем snapshot.
agent.conversationStore.get(convId)?.let { conversationStore.upsert(it) }
}
/**
* Append с защитой от повторной доставки (reconnect/retry): `message.id` —
* PRIMARY KEY, повторный INSERT бросает constraint violation. Ловим и идём
* дальше — естественный dedup без модификации `KsqliteJournalStore` и без
* отдельного SQL-connection.
*
* Альтернатива — INSERT OR IGNORE через прямой SQL потребовала бы
* собственного SQLitePreparedStatement и отдельного connection'а;
* в сценарии с правильной cursor-дисциплиной (`lastSeen`-based fetch +
* `markSeen`) дубликаты редки, поэтому простой try/catch достаточен.
*/
private suspend fun appendIfAbsent(record: MessageRecord) {
try {
journal.append(record)
} catch (e: CancellationException) {
throw e
} catch (e: Exception) {
Log.d(TAG, "skip duplicate ${record.id}: ${e.message}")
}
}
override fun close() {
syncJob.cancel()
}
private companion object {
const val TAG = "AgentEventSync"
const val REFRESH_LIMIT = 200
}
}
@@ -0,0 +1,137 @@
package pw.binom.agentik.android.agent.real
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.withContext
import pw.binom.db.ksqlite.SQLiteConnection
import pw.binom.db.ksqlite.SQLitePreparedStatement
import kotlin.time.Instant
/**
* Хранилище метаданных диалогов: `last_seen_message_at` per-conversation.
*
* Живёт в той же SQLite-БД, что и [KsqliteJournalStore]
* (`filesDir/journal/agentik.db`), но **отдельная таблица** `conversation_meta`.
* Три отдельные абстракции per desktop-требований R41.4:
* messages / settings / conversation-list snapshot — каждая со своим owner'ом.
*
* **Не блокирует**: `markSeen` / `getLastSeen` — O(1) по primary key.
*
* ## Schema
*
* ```sql
* CREATE TABLE IF NOT EXISTS conversation_meta (
* conversation_id TEXT PRIMARY KEY,
* last_seen_message_at INTEGER NOT NULL -- epoch millis (Instant.toEpochMilliseconds)
* )
* ```
*
* **Версионирование** — `Schema.CURRENT_VERSION` инкрементируется при ЛЮБОМ
* изменении DDL. Сейчас 1.
*
* ## API-источник
*
* Реальные сигнатуры ksqlite взяты из upstream `KsqliteJournalStore`:
* - `SQLiteConnection.open(path)` / `SQLiteConnection.memory(name)`
* - `connection.prepare(sql)` → `SQLitePreparedStatement`
* - `connection.exec(sql)` — raw DDL
* - `stmt.bindText(idx, value)` — для строк
* - `stmt.bindLong(idx, value)` — для Long
* - `stmt.executeUpdate()` — для INSERT/UPDATE/DELETE
* - `stmt.executeQuery()` → `SQLiteResultSet` (closeable)
* - `rs.next()`: Boolean, `rs.getLong(idx)`, `rs.getText(idx)` — nullable!
*/
class ConversationMetaRepository private constructor(
private val connection: SQLiteConnection,
private val ownsConnection: Boolean,
) : AutoCloseable {
/** Открывает файловое соединение и берёт на себя его закрытие в [close]. */
constructor(path: String) : this(
connection = SQLiteConnection.open(path = path),
ownsConnection = true,
)
/** Внешнее соединение — репо НЕ закрывает его в [close]. Для shared-connection bundles. */
constructor(connection: SQLiteConnection) : this(
connection = connection,
ownsConnection = false,
)
init {
Schema.migrate(connection)
}
private val markSeenStmt: SQLitePreparedStatement = connection.prepare(
"INSERT OR REPLACE INTO conversation_meta (conversation_id, last_seen_message_at) VALUES (?, ?)"
)
private val getLastSeenStmt: SQLitePreparedStatement = connection.prepare(
"SELECT last_seen_message_at FROM conversation_meta WHERE conversation_id = ?"
)
private val getAllLastSeenStmt: SQLitePreparedStatement = connection.prepare(
"SELECT conversation_id, last_seen_message_at FROM conversation_meta"
)
/** UPSERT: запоминает, что пользователь видел диалог до момента [at]. */
suspend fun markSeen(convId: String, at: Instant): Unit = withContext(Dispatchers.IO) {
val ts = at.toEpochMilliseconds()
markSeenStmt.reset()
markSeenStmt.clearBindings()
markSeenStmt.bindText(1, convId)
markSeenStmt.bindLong(2, ts)
markSeenStmt.executeUpdate()
}
/** Возвращает `last_seen` или null (никогда не открывали). */
suspend fun getLastSeen(convId: String): Instant? = withContext(Dispatchers.IO) {
getLastSeenStmt.reset()
getLastSeenStmt.clearBindings()
getLastSeenStmt.bindText(1, convId)
getLastSeenStmt.executeQuery().use { rs ->
if (!rs.next()) return@use null
val ts = rs.getLong(1) ?: return@use null
Instant.fromEpochMilliseconds(ts)
}
}
/** Bulk-чтение для списка диалогов (избегаем N+1 запросов). */
suspend fun getAllLastSeen(): Map<String, Instant> = withContext(Dispatchers.IO) {
getAllLastSeenStmt.reset()
getAllLastSeenStmt.clearBindings()
getAllLastSeenStmt.executeQuery().use { rs ->
val result = mutableMapOf<String, Instant>()
while (rs.next()) {
val convId = rs.getText(1) ?: continue
val ts = rs.getLong(2) ?: continue
result[convId] = Instant.fromEpochMilliseconds(ts)
}
result.toMap()
}
}
override fun close() {
runCatching { markSeenStmt.close() }
runCatching { getLastSeenStmt.close() }
runCatching { getAllLastSeenStmt.close() }
if (ownsConnection) connection.close()
}
/**
* DDL для таблицы `conversation_meta`. По аналогии с
* `pw.binom.agentik.journal.ksqlite.Schema`. Идемпотентна
* (`IF NOT EXISTS`), можно звать повторно.
*/
object Schema {
const val CURRENT_VERSION: Int = 1
fun migrate(conn: SQLiteConnection) {
// DDL без prepared statements — однократный CREATE TABLE,
// exec() достаточно.
conn.exec(
"CREATE TABLE IF NOT EXISTS conversation_meta (" +
"conversation_id TEXT PRIMARY KEY," +
"last_seen_message_at INTEGER NOT NULL" +
")"
)
}
}
}
@@ -0,0 +1,383 @@
package pw.binom.agentik.android.agent.real
import android.content.Context
import io.ktor.client.engine.HttpClientEngineFactory
import io.ktor.client.engine.okhttp.OkHttp
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.cancel
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.MutableSharedFlow
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asSharedFlow
import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.flow.toList
import kotlinx.coroutines.launch
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
import pw.binom.agentik.android.agent.AgentConnection
import pw.binom.agentik.android.agent.ChatSession
import pw.binom.agentik.android.agent.ConnectionState
import pw.binom.agentik.android.agent.ConversationSummary
import pw.binom.agentik.client.AgentikAgent
import pw.binom.agentik.client.ConnectionStatus
import pw.binom.agentik.client.ReconnectingOutbox
import pw.binom.agentik.journal.ConversationStore
import pw.binom.agentik.journal.MutableJournalStore
import pw.binom.agentik.journal.ksqlite.KsqliteJournalStore
import pw.binom.agentik.outbox.CommonEvent
import pw.binom.agentik.proto.Agent
import java.io.File
import java.util.concurrent.CopyOnWriteArraySet
import kotlin.time.Instant
/**
* Реальная реализация [AgentConnection] поверх `pw.binom.agentik:client:12`.
*
* Lifecycle: создаётся один раз, держит [Agent] (а через него — `HttpClient`)
* до [close]. На каждый [configure] закрывает старый [Agent] и создаёт новый.
*
* **Не использует Compose, не использует DI.**
* Настройки хранятся в памяти в локальном поле configure-аргументов;
* их JSON-персистентность — за [SettingsRepository] (отдельный слой).
*
* **Persistent journal cache.** В конструктор требуется [Context] —
* используется для пути к SQLite-файлу кэша (`context.filesDir/journal/agentik.db`).
* Один [MutableJournalStore] живёт на всё соединение и переиспользуется
* всеми [HttpChatSession] (изоляция — по `conversation_id` в схеме).
* В v12 secondary constructor `KsqliteJournalStore(path)` стал public —
* reflection-workaround больше не нужен (см. memory `journal-ksqlite-reflection-workaround`).
*
* **Auto-reconnect.** В v10 поверх `agent.outbox` оборачивается
* [ReconnectingOutbox] — он обеспечивает два независимых потока:
* `events()` (live-события с курсором между реконнектами) и
* `connectionStatus()` (lifecycle Connecting/Connected/Disconnected/Failed).
* Это критично на мобильной сети — обрыв SSE не теряет события и не требует
* ручного переподключения.
*
* **EngineFactory.** По умолчанию — OkHttp. Для тестов / платформ, где
* OkHttp недоступен, в конструктор можно передать другую фабрику (например,
* CIO в unit-тестах на JVM). В production Android это OkHttp.
*/
class HttpAgentConnection(
private val context: Context,
private val engineFactory: HttpClientEngineFactory<*> = OkHttp,
private val scope: CoroutineScope = CoroutineScope(SupervisorJob() + Dispatchers.IO),
) : AgentConnection {
@Volatile private var agent: Agent? = null
@Volatile private var reconnectingOutbox: ReconnectingOutbox? = null
@Volatile private var currentAgentId: String? = null
@Volatile private var currentAgentName: String? = null
private val _state = MutableStateFlow<ConnectionState>(ConnectionState.Disconnected)
override val state: StateFlow<ConnectionState> = _state.asStateFlow()
/**
* Lifecycle-эмиты диалогов агента (Created / Deleted / Renamed / Touched).
* Для UI-подписчиков списка диалогов: на каждый эмит перезапрашиваем
* `conversationStore.listFlow()` (см. ConversationListScreen). Под капотом
* фильтрует [CommonEvent.Agent] из `reconnectingOutbox.events()`.
*/
private val _agentEvents = MutableSharedFlow<Unit>(extraBufferCapacity = 8)
val agentEvents: Flow<Unit> = _agentEvents.asSharedFlow()
/**
* Read-only view на список диалогов. В v12 [AgentikAgent] оборачивает
* свой [Agent.conversationStore] в локальный in-memory кэш,
* синхронизируемый через outbox-события.
* UI читает `conversationStore.listFlow()` — auto-refresh без ручного refetch.
*
* Доступ валиден только после успешного [configure] — до этого бросает
* [IllegalStateException] (экран списка диалогов показывается только при
* [ConnectionState.Connected]).
*/
override val conversationStore: ConversationStore
get() = agent?.conversationStore
?: throw IllegalStateException("AgentConnection not configured (state=${_state.value})")
/**
* Путь к shared SQLite-файлу journal-кэша. Создаётся лениво при первом
* [openConversation], один на всё соединение.
*/
private val journalDbPath: String = run {
val dir = File(context.filesDir, "journal")
dir.mkdirs()
File(dir, "agentik.db").absolutePath
}
/**
* Single shared [MutableJournalStore] на всё соединение. Создаётся лениво
* в [openConversation], закрывается в [close]. Все [HttpChatSession]
* используют один экземпляр (single-DB подход — изоляция по
* `conversation_id` в схеме).
*/
@Volatile private var sharedCache: MutableJournalStore? = null
/**
* Репозиторий метаданных диалогов (per-conversation `last_seen_message_at`).
* Таблица `conversation_meta` живёт в той же БД, что и [sharedCache]
* (отдельная абстракция — desktop R41.4 требует три отдельных слоя:
* messages / settings / conversation-list snapshot). Доступ к metaRepo
* есть всегда (создаётся лениво при первом обращении), даже если ни один
* чат ещё не открывался — нужен для `listConversations()`, который
* вызывается ДО любого [openConversation].
*
* `internal` — [HttpChatSession] дёргает `markSeen` после backfill,
* на каждом `Event.End` и при close().
*/
internal val metaRepo: ConversationMetaRepository by lazy {
ConversationMetaRepository(journalDbPath)
}
/**
* Persistent SQLite-реестр диалогов агента (таблица `conversation`).
* Один экземпляр на всё соединение — обновляется в фоне через
* [AgentEventSync] из `outbox.events()`, читается напрямую из
* [listConversations] без обращения к серверу.
*
* Тонкая переупаковка над [KsqliteMutableConversationStore] из
* `:journal-ksqlite:15` (см. typealias [LocalConversationRepository]):
* эта либа уже владеет таблицей `conversation` — дублировать DDL
* не нужно.
*/
private val localConversations: LocalConversationRepository by lazy {
LocalConversationRepository(journalDbPath)
}
/**
* Background-sync из `outbox.events()` в локальный кэш
* ([localConversations] + [sharedCache] + [metaRepo]).
* Не требует открытия чата. Создаётся в [configure] после успешного
* `ReconnectingOutbox`, отменяется в [close] ДО закрытия outbox.
*/
private var eventSync: AgentEventSync? = null
/**
* Открытые [HttpChatSession] для безопасного закрытия в [close] —
* иначе принудительное отключение во время активного чата оставит
* session-корутину и кэш живыми.
*/
private val openSessions = CopyOnWriteArraySet<HttpChatSession>()
/**
* Сериализует доступ к SQLite (`KsqliteJournalStore` + [metaRepo],
* оба пишут в `filesDir/journal/agentik.db`). Без этого при конкурентных
* записях возможен SQLITE_BUSY — у SQLite file-level locking, и второй
* writer может ждать или получить ошибку. Использовать через
* [withSqliteLock] из всех мест, где дёргаются [sharedCache] и/или [metaRepo].
*/
private val sqliteMutex = Mutex()
/**
* Запускает [block] под [sqliteMutex]. Используется для атомарных
* операций над SQLite-кэшем (backfill + publishMessages, Event.End
* fetch + publish + markSeen, refresh-merge и т.д.) — чтобы между
* этими шагами не вклинился другой writer.
*/
suspend fun <T> withSqliteLock(block: suspend () -> T): T = sqliteMutex.withLock { block() }
override suspend fun configure(id: String, name: String, baseUrl: String, token: String?) {
// Валидация URL — пустой или не http(s) сразу бросает IllegalArgumentException,
// не выставляя state.
require(baseUrl.isNotBlank() && (baseUrl.startsWith("http://") || baseUrl.startsWith("https://"))) {
"baseUrl must be non-empty and start with http:// or https://"
}
// Закрыть старые chat-sessions + agent + outbox, если были.
openSessions.forEach { it.close() }
openSessions.clear()
reconnectingOutbox?.close()
reconnectingOutbox = null
agent?.close()
agent = null
currentAgentId = null
currentAgentName = null
_state.value = ConnectionState.Connecting(id)
val newAgent: Agent = try {
AgentikAgent(
id = id,
baseUrl = baseUrl,
engineFactory = engineFactory,
token = token,
)
} catch (e: Exception) {
val reason = e.message ?: e::class.simpleName ?: "unknown"
_state.value = ConnectionState.Failed(reason)
return
}
agent = newAgent
currentAgentId = id
currentAgentName = name
// v10+: ReconnectingOutbox оборачивает agent.outbox и даёт два потока —
// events() (с курсором между реконнектами) и connectionStatus()
// (lifecycle Connecting/Connected/Disconnected/Failed). Probe через
// getConversations() больше не нужен — первый эмитт
// ConnectionStatus.Connected сам индикатор успешной связи.
val outbox = ReconnectingOutbox(newAgent.outbox, scope)
reconnectingOutbox = outbox
// ConnectionStatus → наш ConnectionState. Connecting/Disconnected
// мапятся в Connecting (transient — клиент переподключится сам через
// BackoffPolicy). Failed — терминальная ошибка после maxAttempts.
scope.launch {
outbox.connectionStatus().collect { status ->
_state.value = when (status) {
is ConnectionStatus.Connecting -> ConnectionState.Connecting(currentAgentId ?: id)
is ConnectionStatus.Connected -> ConnectionState.Connected(
currentAgentId ?: id,
currentAgentName ?: name,
)
is ConnectionStatus.Disconnected -> ConnectionState.Connecting(currentAgentId ?: id)
is ConnectionStatus.Failed -> ConnectionState.Failed(
status.cause.message
?: status.cause::class.simpleName
?: "outbox failed",
)
}
}
}
// Lifecycle-эмиты диалогов (Created / Deleted / Renamed / Touched)
// — единичный тик в [agentEvents]. Подписчик UI (ConversationListScreen)
// перезапрашивает список через `conversationStore.listFlow()`.
scope.launch {
outbox.events(Instant.DISTANT_PAST).collect { commonEvent ->
if (commonEvent is CommonEvent.Agent) {
_agentEvents.tryEmit(Unit)
}
}
}
// Background-sync `outbox.events()` → локальный SQLite-кэш.
// Делает listConversations() мгновенным (читает локальную таблицу
// `conversation` вместо обращения к серверу). Запускаем ПОСЛЕ outbox —
// иначе первый эмитт событий может уйти до того, как sync готов.
// `sharedCache` lazy — к моменту первого refresh события уже инициализирован
// (см. acquireJournalStore); metaRepo / localConversations — ленивые,
// init при первом touch в AgentEventSync.
eventSync = AgentEventSync(
agent = newAgent,
conversationStore = localConversations,
journal = sharedCache ?: acquireJournalStore(),
metaRepo = metaRepo,
parentScope = scope,
withSqliteLock = { block -> withSqliteLock<Unit>(block) },
)
}
override suspend fun listConversations(): List<ConversationSummary> {
// Инициализация sharedCache нужна до listConversations — без него unread-цифры
// были бы нулевыми (раньше этот путь считал через `agent.journal.count`,
// теперь — через локальный SQLite-кэш).
val cache = acquireJournalStore()
return withSqliteLock {
try {
// Локальный snapshot бесед — без HTTP. AgentEventSync поддерживает
// эту таблицу в актуальном состоянии в фоне (Created/Deleted/Renamed/Touched),
// плюс подтягивает записи из `journal.list(...)` на каждый `Event.End`.
val convs = localConversations.listFlow().toList()
val meta = metaRepo.getAllLastSeen()
convs.map { rec ->
val lastSeen = meta[rec.id]
// Unread считаем из локального кэша. cache.list(...) грузит
// страницу в память (в т.ч. Int.MAX_VALUE для never-opened);
// при текущих масштабах это приемлемо. Оптимизация через SQL
// COUNT — будущая задача.
val unread: Int = if (lastSeen == null) {
cache.list(rec.id, Instant.DISTANT_PAST, 0, Int.MAX_VALUE).size
} else {
cache.list(rec.id, lastSeen, 0, Int.MAX_VALUE).size
}
ConversationSummary(
id = rec.id,
title = rec.title,
updatedAt = rec.updatedAt,
lastMessagePreview = null, // preview из кэша — отдельная задача
unreadCount = unread,
)
}
} catch (e: Exception) {
emptyList() // SQLite недоступен — пусть UI покажет старый список
}
}
}
override fun openConversation(convId: String?): ChatSession {
val a = agent
?: throw IllegalStateException("AgentConnection not configured (state=${_state.value})")
val ob = reconnectingOutbox
?: throw IllegalStateException("AgentConnection outbox not initialized (state=${_state.value})")
val cache = acquireJournalStore()
val session = HttpChatSession(a, ob, convId, cache, scope, this)
openSessions += session
return session
}
override fun close() {
// Background-sync отменяем ДО outbox — иначе pending emissions после
// reconnect-cancel могут дёрнуть writer'ы уже закрытого SQLite-стора.
runCatching { eventSync?.close() }
eventSync = null
// Закрываем chat-sessions первыми — они ссылаются на sharedCache и
// могут держать активные корутины в scope.
openSessions.forEach { it.close() }
openSessions.clear()
// reconnectingOutbox первым — он запускает background-loop в scope,
// должен быть отменён до того как мы отменим сам scope.
reconnectingOutbox?.close()
reconnectingOutbox = null
agent?.close()
agent = null
currentAgentId = null
currentAgentName = null
_state.value = ConnectionState.Disconnected
// Local conversation registry — persistent SQLite-таблица в том же
// файле, что и sharedCache. Закрываем перед sharedCache, чтобы
// backfill-end race не дёрнул writer во время teardown.
runCatching { localConversations.close() }
// Shared journal cache — единственный владелец SQLite connection
// для всего соединения. Закрываем здесь; все session'ы уже
// отменены и больше не дернут кэш.
runCatching { sharedCache?.close() }
sharedCache = null
// ConversationMetaRepository — отдельная SQLite-таблица в том же файле,
// держит свой SQLitePreparedStatement (markSeen / getAllLastSeen). Закрываем
// после sharedCache (порядок не критичен — обе таблицы в одном файле, но
// sharedCache мог быть lazy-неинициализирован, тогда metaRepo остался
// единственным открывателем соединения).
runCatching { metaRepo.close() }
scope.cancel()
}
/**
* Удаление сессии из tracking-множества. Вызывается из [HttpChatSession.close]
* — иначе повторный close сессии или накопление мёртвых ссылок.
*/
internal fun onSessionClosed(session: HttpChatSession) {
openSessions.remove(session)
}
/**
* Ленивая инициализация кэш-стора. Все [HttpChatSession] делят один
* экземпляр (single-DB подход). Double-checked locking: contention только
* на самом первом вызове.
*/
private fun acquireJournalStore(): MutableJournalStore {
sharedCache?.let { return it }
return synchronized(this) {
sharedCache ?: KsqliteJournalStore(path = journalDbPath).also { sharedCache = it }
}
}
}
@@ -0,0 +1,252 @@
package pw.binom.agentik.android.agent.real
import android.util.Log
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.Job
import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.MutableSharedFlow
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asSharedFlow
import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.launch
import pw.binom.agentik.android.agent.ChatSession
import pw.binom.agentik.android.agent.ConversationSnapshot
import pw.binom.agentik.android.agent.UiEvent
import pw.binom.agentik.android.agent.UiMessage
import pw.binom.agentik.android.agent.toUiEvent
import pw.binom.agentik.android.agent.toUiMessage
import pw.binom.agentik.client.ReconnectingOutbox
import pw.binom.agentik.journal.MutableJournalStore
import pw.binom.agentik.outbox.CommonEvent
import pw.binom.agentik.outbox.Event
import pw.binom.agentik.proto.Agent
import pw.binom.agentik.proto.Content
import pw.binom.agentik.proto.Conversation
import kotlin.time.Instant
/**
* Реальная реализация [ChatSession] поверх [Conversation] + persistent
* [MutableJournalStore] (SQLite через `:journal-ksqlite`).
*
* **Persistent cache.** Переданный снаружи [cache] живёт дольше сессии
* — это общий кэш на уровне [HttpAgentConnection], переиспользуется
* между диалогами (изоляция — по `conversation_id` в схеме). История
* переживает kill приложения.
*
* Init flow:
* 1. get-or-create conversation (внутри [parentScope].launch).
* 2. Backfill через `agent.journal.listFlow(c.id, Instant.DISTANT_PAST)` —
* авто-пагинация, Flow завершается когда сервер отдал пустую страницу.
* Каждая запись пишется в [cache], потом публикуется в [messages]
* (через `MessageRecord.toUiMessage()`).
* 3. Live через `outbox.events(Instant.DISTANT_PAST)` — единый поток всех
* событий агента с авто-reconnect (см. [ReconnectingOutbox]). Фильтруем
* по `CommonEvent.Conversation` с нашим `convId`. На каждый [Event.End]
* догоняем последние записи через `agent.journal.list(c.id, latest, 0, 100)`
* и обновляем [messages]. Другие события маппятся в [UiEvent] и эмитятся
* в [liveEvents].
*
* **Закрытие сессии НЕ закрывает [cache]** — он переиспользуется. Закрытие
* делает [HttpAgentConnection.close]. На [close] уведомляем [owner], чтобы
* тот убрал сессию из tracking-множества и не пытался закрыть её повторно.
*/
internal class HttpChatSession(
private val agent: Agent,
private val outbox: ReconnectingOutbox,
private val initialConvId: String?,
private val cache: MutableJournalStore,
private val parentScope: CoroutineScope,
private val owner: HttpAgentConnection,
) : ChatSession {
private val sessionJob: Job = SupervisorJob(parentScope.coroutineContext[Job])
private val sessionScope = CoroutineScope(parentScope.coroutineContext + sessionJob + Dispatchers.IO)
private var conv: Conversation? = null
private val _snapshot = MutableStateFlow<ConversationSnapshot?>(null)
private val _messages = MutableStateFlow<List<UiMessage>>(emptyList())
private val _liveEvents = MutableSharedFlow<UiEvent>(extraBufferCapacity = 16)
override val conversationId: String
get() = conv?.id ?: initialConvId.orEmpty()
override val snapshot: StateFlow<ConversationSnapshot?> = _snapshot.asStateFlow()
override val messages: StateFlow<List<UiMessage>> = _messages.asStateFlow()
override val liveEvents: Flow<UiEvent> = _liveEvents.asSharedFlow()
init {
sessionScope.launch { initialize() }
}
private suspend fun initialize() {
val c: Conversation = try {
initialConvId?.let { agent.getConversation(it) }
?: agent.createConversation(temp = false)
} catch (e: Exception) {
// Логировать (Log.e) и закрыть сессию — UI увидит snapshot=null и поймёт.
Log.e(TAG, "failed to open conversation: ${e.message}", e)
return
} ?: run {
Log.w(TAG, "conversation $initialConvId not found")
return
}
conv = c
_snapshot.value = ConversationSnapshot(
id = c.id,
title = c.title,
updatedAt = c.updatedAt,
supportsImageInput = c.isSupportImageInput,
supportsImageOutput = c.isSupportImageOutput,
isTemporal = c.isTemporal,
)
// Backfill. listFlow — auto-paginating Flow<MessageRecord>, завершается
// когда сервер отдал пустую страницу. Атомарно под sqliteMutex:
// append + publishMessages + markSeen должны видеть консистентное
// состояние SQLite (cache.append пишет в ту же БД, что и metaRepo).
try {
agent.journal.listFlow(c.id, Instant.DISTANT_PAST)
.collect { rec ->
owner.withSqliteLock { cache.append(rec) }
}
owner.withSqliteLock {
val records = cache.list(c.id, Instant.DISTANT_PAST, 0, Int.MAX_VALUE)
_messages.value = records.map { it.toUiMessage() }
// Mark seen: после полного backfill'а диалог считается прочитанным.
// Берём свежайший timestamp из кэша (Int.MAX_VALUE limit + .lastOrNull()
// — ASC-сортировка cache.list означает, что .lastOrNull() вернёт
// именно максимальный createdAt). null → кэш пуст → markSeen не зовём.
val lastCached = records.lastOrNull()?.createdAt
if (lastCached != null) {
owner.metaRepo.markSeen(c.id, lastCached)
}
}
} catch (e: Exception) {
Log.e(TAG, "backfill failed: ${e.message}", e)
// Продолжаем — live всё равно работает
}
// Live events + refresh-on-End. v10: подписываемся на
// `ReconnectingOutbox.events()` (единый поток всех событий агента,
// с авто-reconnect и курсором) и фильтруем по нашему convId —
// берём только `CommonEvent.Conversation` с нужным conversationId.
try {
outbox.events(Instant.DISTANT_PAST).collect { commonEvent ->
val convEv = commonEvent as? CommonEvent.Conversation ?: return@collect
if (convEv.conversationId != c.id) return@collect
val ev = convEv.event
ev.toUiEvent()?.let { uiEvent ->
_liveEvents.tryEmit(uiEvent)
}
if (ev is Event.End) {
// **Cursor bug fix.** Раньше здесь читался
// свежайший timestamp из кэша через вызов
// `cache.list(... limit = 1).lastOrNull()`.
// Это ВСЕГДА возвращало OLDEST-сообщение (journal-api
// сортирует ASC; limit=1 + .lastOrNull() = единственная
// запись = самая старая) → journal.list(after=oldest)
// тащил почти всё заново → дубликаты в кэше.
// Теперь читаем `last_seen_message_at` из metaRepo —
// именно тот cursor, который обновляет markSeen.
val lastSeen = owner.metaRepo.getLastSeen(c.id) ?: Instant.DISTANT_PAST
// Атомарно под sqliteMutex: fetch → cache.append →
// cache.list → publishMessages → markSeen. Без мьютекса
// параллельный refresh (или новый backfill после reconnect)
// может вклиниться между append и publishMessages →
// рассинхрон UI vs SQLite.
owner.withSqliteLock {
try {
agent.journal.list(c.id, lastSeen, offset = 0, limit = 100)
.forEach { cache.append(it) }
} catch (e: Exception) {
Log.w(TAG, "refresh on End failed: ${e.message}", e)
}
val records = cache.list(c.id, Instant.DISTANT_PAST, 0, Int.MAX_VALUE)
_messages.value = records.map { it.toUiMessage() }
// Mark seen: после очередного завершённого хода ассистента
// пользователь всё, что пришло в этот End, уже увидел.
// ev.date — серверный timestamp момента завершения хода.
owner.metaRepo.markSeen(c.id, ev.date)
}
}
}
} catch (e: Exception) {
Log.e(TAG, "live events failed: ${e.message}", e)
}
}
override suspend fun send(text: String) {
val c = conv ?: throw IllegalStateException("ChatSession not yet initialized")
c.send(listOf(Content.Text(text)), context = null)
}
override suspend fun interrupt() {
conv?.interrupt()
}
override suspend fun rename(newTitle: String) {
val c = conv ?: throw IllegalStateException("ChatSession not yet initialized")
c.rename(newTitle)
// snapshot обновится по событию outbox (AgentEvent.Renamed) при следующем
// events-цикле или ручном refresh — для UX приемлемо.
}
override suspend fun refresh() {
val c = conv ?: return
// **Cursor bug fix.** Тот же баг, что и в Event.End блоке: см. KDoc
// там. Заменено на owner.metaRepo.getLastSeen() вместо чтения
// из cache через limit=1 + lastOrNull (что возвращало OLDEST).
val lastSeen = owner.metaRepo.getLastSeen(c.id) ?: Instant.DISTANT_PAST
try {
owner.withSqliteLock {
agent.journal.list(c.id, lastSeen, offset = 0, limit = 200)
.forEach { cache.append(it) }
val records = cache.list(c.id, Instant.DISTANT_PAST, 0, Int.MAX_VALUE)
_messages.value = records.map { it.toUiMessage() }
}
} catch (e: Exception) {
Log.w(TAG, "manual refresh failed: ${e.message}", e)
}
}
override fun close() {
// cancel только наш sessionJob, не parentScope (connection живёт дальше).
// cache не закрываем — он переиспользуется между сессиями и закроется
// в HttpAgentConnection.close().
sessionJob.cancel()
// Финальный markSeen — best-effort, на случай если End-событие не
// пришло (стрим оборвался, пользователь нажал «назад»).
// Запускаем на parentScope (переживает sessionJob.cancel()): close()
// не-suspend, поэтому блокирующий markSeen здесь — анти-паттерн.
// parentScope в свою очередь переживёт до HttpAgentConnection.close() —
// если он будет вызван совсем сразу, markSeen отвалится в runCatching
// (acceptable: best-effort).
conv?.let { c ->
parentScope.launch {
try {
owner.withSqliteLock {
val lastCached = cache.list(c.id, Instant.DISTANT_PAST, 0, Int.MAX_VALUE)
.lastOrNull()?.createdAt
if (lastCached != null) {
owner.metaRepo.markSeen(c.id, lastCached)
}
}
} catch (e: Exception) {
Log.w(TAG, "markSeen on close failed: ${e.message}", e)
}
}
}
// Убираем себя из owner.openSessions, чтобы не копиться мёртвыми
// ссылками и не получить повторный close через owner.close().
owner.onSessionClosed(this)
}
private companion object {
const val TAG = "HttpChatSession"
}
}
@@ -0,0 +1,26 @@
package pw.binom.agentik.android.agent.real
import pw.binom.agentik.journal.ksqlite.KsqliteMutableConversationStore
/**
* Persistent SQLite-backed реестр диалогов на Android.
*
* `:journal-ksqlite:15` уже владеет таблицей `conversation` (см.
* `Schema.migrate` / `Schema.TABLE_CONVERSATION`). Тип
* [KsqliteMutableConversationStore] — реализация `MutableConversationStore`
* с двумя формами конструктора: `(path)` — открывает файловое соединение
* и владеет им, `(connection, ownsConnection)` — внешнее соединение
* (для shared-connection bundles).
*
* На Android держим ОДНО файловое соединение на таблицу `conversation`
* (`filesDir/journal/agentik.db`) — изоляция от других таблиц в той же БД
* (`message`, `conversation_meta`) обеспечивается отдельными connections.
* Конкуренция сериализуется через [HttpAgentConnection.sqliteMutex].
*
* Почему не свой класс по образцу [KsqliteJournalStore]: явное дублирование
* `Schema.migrate` / набора prepared-statement'ов — лишний код, который
* придётся поддерживать в двух местах. Typealias даёт имя уровня
* Android-слоя (отражает роль в lifecycle [HttpAgentConnection]) без
* потери связи с upstream-реализацией.
*/
typealias LocalConversationRepository = KsqliteMutableConversationStore
@@ -0,0 +1,201 @@
package pw.binom.agentik.android.agent.stub
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.cancel
import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.runBlocking
import pw.binom.agentik.android.agent.AgentConnection
import pw.binom.agentik.android.agent.ChatSession
import pw.binom.agentik.android.agent.ConnectionState
import pw.binom.agentik.android.agent.ConversationSummary
import pw.binom.agentik.android.agent.ConversationSnapshot
import pw.binom.agentik.android.agent.UiMessage
import pw.binom.agentik.journal.ConversationRecord
import pw.binom.agentik.journal.ConversationStore
import pw.binom.agentik.journal.MutableConversationStore
import pw.binom.agentik.journal.inmemory.InMemoryMutableConversationStore
import java.util.UUID
import java.util.concurrent.CopyOnWriteArraySet
import kotlin.time.Clock
import kotlin.time.Duration.Companion.days
import kotlin.time.Duration.Companion.hours
import kotlin.time.Duration.Companion.minutes
/**
* Заглушка [AgentConnection]: возвращает фейковые диалоги и валидирует
* конфиг (URL). Эмулирует handshake `Disconnected → Connecting → Connected`.
*
* Используется до подключения HTTP-импла на `AgentikAgent(...)`.
*
* [conversationStore] — [InMemoryMutableConversationStore], предзаполненный
* тремя фейковыми [ConversationRecord] (аналогично [listConversations]).
* Сидинг через [runBlocking] в init — список из трёх записей, инициализация
* занимает микросекунды, поэтому main-thread блокировка приемлема для stub'а.
* Это нужно, чтобы UI через `listFlow()` сразу видел данные без видимой
* вспышки "empty state" при первом запуске.
*/
internal class StubAgentConnection : AgentConnection {
private val scope = CoroutineScope(SupervisorJob() + Dispatchers.Default)
private val _state = MutableStateFlow<ConnectionState>(ConnectionState.Disconnected)
override val state: StateFlow<ConnectionState> = _state.asStateFlow()
private val openSessions = CopyOnWriteArraySet<StubChatSession>()
/**
* Mutable-импл для сидинга; наружу торчит как read-only [ConversationStore].
* Доступен сразу после конструирования — список видим через `listFlow()`.
*/
private val mutableConversationStore: MutableConversationStore = InMemoryMutableConversationStore()
override val conversationStore: ConversationStore = mutableConversationStore
init {
// Предзаполняем кэш фейковыми диалогами — UI через `listFlow()`
// сразу получает данные, без вспышки "empty state".
runBlocking { seedStubConversations(mutableConversationStore) }
}
private suspend fun seedStubConversations(store: MutableConversationStore) {
val now = Clock.System.now()
store.upsert(
ConversationRecord(
id = "stub-conv-1",
title = "Ремонт на кухне",
isTemporal = false,
createdAt = now - 1.hours,
updatedAt = now - 1.hours,
),
)
store.upsert(
ConversationRecord(
id = "stub-conv-2",
title = "Парсер логов",
isTemporal = false,
createdAt = now - 3.hours,
updatedAt = now - 3.hours,
),
)
store.upsert(
ConversationRecord(
id = "stub-conv-3",
title = null, // новый диалог без названия — REQUIREMENTS.md
isTemporal = false,
createdAt = now - 1.days,
updatedAt = now - 1.days,
),
)
}
override suspend fun configure(id: String, name: String, baseUrl: String, token: String?) {
require(baseUrl.isNotBlank()) { "baseUrl must not be blank" }
require(baseUrl.startsWith("http://") || baseUrl.startsWith("https://")) {
"baseUrl must start with http:// or https:// (got '$baseUrl')"
}
_state.value = ConnectionState.Connecting(id)
// Имитация round-trip'а к серверу.
delay(HANDSHAKE_DELAY_MS)
_state.value = ConnectionState.Connected(agentId = id, agentName = name)
}
override suspend fun listConversations(): List<ConversationSummary> {
val now = Clock.System.now()
return listOf(
ConversationSummary(
id = "stub-conv-1",
title = "Ремонт на кухне",
updatedAt = now - 1.hours,
lastMessagePreview = "А если столешница из массива?",
unreadCount = 2,
),
ConversationSummary(
id = "stub-conv-2",
title = "Парсер логов",
updatedAt = now - 3.hours,
lastMessagePreview = "Закончили, отчёт готов",
unreadCount = 0,
),
ConversationSummary(
id = "stub-conv-3",
title = null, // новый диалог без названия — REQUIREMENTS.md
updatedAt = now - 1.days,
lastMessagePreview = null,
unreadCount = 0, // v15: non-null, never-opened показываем «0»
),
)
}
override fun openConversation(convId: String?): ChatSession {
val id = convId ?: "stub-conv-${UUID.randomUUID()}"
val now = Clock.System.now()
val session = StubChatSession(
conversationId = id,
initialSnapshot = ConversationSnapshot(
id = id,
title = convId?.let { titleForKnownConv(it) },
updatedAt = now,
supportsImageInput = false,
supportsImageOutput = false,
isTemporal = false,
),
initialMessages = if (convId != null && convId in KNOWN_CONV_IDS) {
historyForKnownConv(convId, now)
} else {
emptyList()
},
)
openSessions += session
return session
}
override fun close() {
openSessions.forEach { it.close() }
openSessions.clear()
scope.cancel()
}
private fun titleForKnownConv(id: String): String? = when (id) {
"stub-conv-1" -> "Ремонт на кухне"
"stub-conv-2" -> "Парсер логов"
"stub-conv-3" -> null
else -> "Stub ${id.takeLast(6)}"
}
/**
* 5–7 фейковых сообщений: U/A/TC/TR с реальной хронологией, последнее
* по времени — `Assistant`. Используется как стартовая история для
* известных диалогов.
*/
private fun historyForKnownConv(id: String, now: kotlin.time.Instant): List<UiMessage> {
val base = now - 1.hours
return when (id) {
"stub-conv-1" -> listOf(
UiMessage.User("u1", base, "Привет, планирую ремонт на кухне."),
UiMessage.Assistant("a1", base + 1.minutes, "Здравствуйте! С чего начнём — столешница или фартук?"),
UiMessage.User("u2", base + 2.minutes, "Столешница. Хочу кварц, но боюсь, что дорого."),
UiMessage.ToolCall("tc1", base + 2.minutes, toolName = "calc_price", args = "{material:quartz}"),
UiMessage.ToolResult("tr1", base + 3.minutes, toolCallId = "tc1", result = "{price: 45000}"),
UiMessage.Assistant("a2", base + 4.minutes, "Кварц — от 45 000 ₽/п.м. Альтернатива — искусственный камень, дешевле на 30%."),
UiMessage.User("u3", base + 50.minutes, "А если столешница из массива?"),
)
"stub-conv-2" -> listOf(
UiMessage.User("u1", base, "Распарси access.log за вчера"),
UiMessage.ToolCall("tc1", base + 1.minutes, toolName = "parse_log", args = "{date:2025-09-19}"),
UiMessage.ToolResult("tr1", base + 2.minutes, toolCallId = "tc1", result = "{lines:12834}"),
UiMessage.Assistant("a1", base + 3.minutes, "Готово, 12 834 строк. 4 ошибки 5xx, см. отчёт."),
)
else -> emptyList()
}
}
private companion object {
const val HANDSHAKE_DELAY_MS = 200L
val KNOWN_CONV_IDS: Set<String> = setOf("stub-conv-1", "stub-conv-2", "stub-conv-3")
}
}
@@ -0,0 +1,115 @@
package pw.binom.agentik.android.agent.stub
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.Job
import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.cancel
import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.MutableSharedFlow
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asSharedFlow
import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.launch
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
import pw.binom.agentik.android.agent.ChatSession
import pw.binom.agentik.android.agent.ConversationSnapshot
import pw.binom.agentik.android.agent.UiEvent
import pw.binom.agentik.android.agent.UiMessage
import kotlin.time.Clock
/**
* Заглушка [ChatSession]: держит в памяти фейковую историю и эмулирует
* стрим ответа с задержками. Используется, пока HTTP-фасад не подключён.
*
* **Lifecycle**: `openConversation()` создаёт сессию и регистрирует её у
* [StubAgentConnection]; [close] отменяет свой scope и снимает регистрацию.
* Конкурентные вызовы `send`/`interrupt` сериализуются мьютексом, чтобы
* предыдущая эмуляция корректно отменялась до старта новой.
*/
internal class StubChatSession(
override val conversationId: String,
initialSnapshot: ConversationSnapshot,
initialMessages: List<UiMessage>,
) : ChatSession {
private val scope = CoroutineScope(SupervisorJob() + Dispatchers.Default)
private val mutex = Mutex()
private val _snapshot = MutableStateFlow(initialSnapshot)
override val snapshot: StateFlow<ConversationSnapshot?> = _snapshot.asStateFlow()
private val _messages = MutableStateFlow(initialMessages)
override val messages: StateFlow<List<UiMessage>> = _messages.asStateFlow()
private val _liveEvents = MutableSharedFlow<UiEvent>(extraBufferCapacity = 16)
override val liveEvents: Flow<UiEvent> = _liveEvents.asSharedFlow()
/** Текущая «идущая» эмуляция хода. Новая отменяет старую. */
private var runningJob: Job? = null
override suspend fun send(text: String) = mutex.withLock {
if (text.isBlank()) return@withLock
val sendTime = Clock.System.now()
_messages.value = _messages.value + UiMessage.User(
id = "user-${System.currentTimeMillis()}",
timestamp = sendTime,
text = text,
)
runningJob?.cancel()
runningJob = scope.launch { runFakeTurn(text) }
}
override suspend fun interrupt() = mutex.withLock {
runningJob?.cancel()
runningJob = null
_messages.value = _messages.value + UiMessage.Assistant(
id = "assistant-${System.currentTimeMillis()}",
timestamp = Clock.System.now(),
text = "[прервано пользователем]",
interrupted = true,
)
}
override suspend fun rename(newTitle: String) {
_snapshot.value = _snapshot.value.copy(title = newTitle)
}
/** No-op: стаб ничего не догоняет. */
override suspend fun refresh() = Unit
override fun close() {
scope.cancel()
}
/**
* Эмулирует один ход: эмитит чанки текста + tool-call/tool-result + `End`,
* затем добавляет финальный [UiMessage.Assistant] в историю. Задержка
* 150мс между чанками — чтобы UI мог порисовать streaming-эффект.
*/
private suspend fun runFakeTurn(userText: String) {
_liveEvents.emit(UiEvent.AppendText("Думаю над "))
delay(STREAM_DELAY_MS)
_liveEvents.emit(UiEvent.AppendText("вопросом \"$userText\"..."))
delay(STREAM_DELAY_MS)
_liveEvents.emit(UiEvent.ToolCall(id = STUB_TOOL_ID, toolName = "echo"))
delay(STREAM_DELAY_MS)
_liveEvents.emit(UiEvent.ToolResult(toolCallId = STUB_TOOL_ID, toolName = "echo", result = ""))
delay(STREAM_DELAY_MS)
_liveEvents.emit(UiEvent.End)
_messages.value = _messages.value + UiMessage.Assistant(
id = "assistant-${System.currentTimeMillis()}",
timestamp = Clock.System.now(),
text = "Думаю над вопросом \"$userText\"...",
)
}
private companion object {
const val STREAM_DELAY_MS = 150L
/** Стабильный id для связи stub-ToolCall ↔ stub-ToolResult в live-стриме. */
const val STUB_TOOL_ID = "stub-tool-1"
}
}
@@ -0,0 +1,208 @@
package pw.binom.agentik.android.asr
import android.app.Notification
import android.app.NotificationChannel
import android.app.NotificationManager
import android.app.PendingIntent
import android.content.Context
import android.content.Intent
import android.content.pm.ServiceInfo
import android.os.Build
import androidx.core.app.NotificationCompat
import androidx.core.content.getSystemService
import androidx.work.CoroutineWorker
import androidx.work.ForegroundInfo
import androidx.work.ListenableWorker
import androidx.work.WorkerParameters
import kotlinx.coroutines.CancellationException
import kotlinx.coroutines.async
import kotlinx.coroutines.coroutineScope
import kotlinx.coroutines.flow.collect
import kotlin.math.roundToInt
/**
* Periodic-воркер, который доводит голосовую модель Qwen3-ASR до [ModelState.Ready]
* при наличии Wi-Fi (UNMETERED constraint задаётся на стороне [androidx.work.WorkRequest]).
*
* Поведение:
* 1. Сразу поднимает foreground-сервис с нотификацией, чтобы ОС не убила процесс во время
* скачивания ~940 МБ. Без `setForeground(...)` Android почти гарантированно убьёт процесс
* при сворачивании приложения.
* 2. Параллельно: запускает [QwenModelProvider.ensureReady] и слушает [QwenModelProvider.state],
* обновляя нотификацию по мере прогресса.
* 3. По завершении — отменяет коллектор и возвращает [ListenableWorker.Result.success] /
* [ListenableWorker.Result.retry].
*
* Воркер живёт недолго и держит собственный экземпляр [QwenModelProvider] (у того свой
* `CoroutineScope(SupervisorJob + Dispatchers.IO)`), поэтому гонки с другими вызовами провайдера
* из UI не возникает.
*/
class QwenDownloadWorker(
appContext: Context,
params: WorkerParameters,
) : CoroutineWorker(appContext, params) {
override suspend fun doWork(): ListenableWorker.Result {
ensureChannel(applicationContext)
val provider = QwenModelProvider(applicationContext)
// Шаг 1: сразу поднимаем foreground-сервис с текущим состоянием модели.
// Это критично для 940 МБ — без foreground-сервиса Android убьёт процесс.
setForeground(makeForegroundInfo(provider.state.value))
return try {
coroutineScope {
val downloadJob = async {
try {
val outcome = provider.ensureReady()
if (outcome.isSuccess) {
ListenableWorker.Result.success()
} else {
// Сетевая/IO-ошибка — WorkManager перезапланирует по constraints + backoff.
ListenableWorker.Result.retry()
}
} catch (e: CancellationException) {
throw e
}
}
val collectorJob = async {
try {
provider.state.collect { state ->
setForeground(makeForegroundInfo(state))
}
} catch (e: CancellationException) {
// Коллектор отменяется по завершении загрузки — это норма.
throw e
}
}
val result = downloadJob.await()
// ensureReady() завершилась — останавливаем коллектор, чтобы он не висел вечно.
// Foreground-сервис сам остановится при возврате из doWork().
collectorJob.cancel()
result
}
} catch (e: CancellationException) {
// WorkManager остановил воркер (constraints изменились, пользователь отменил и т.п.).
// Пробрасываем — WorkManager сам разберётся со статусом.
throw e
}
}
// --- нотификация -------------------------------------------------------
private fun makeForegroundInfo(state: ModelState): ForegroundInfo {
val notification = buildNotification(state)
// Android 14+ требует явный foreground service type.
return if (Build.VERSION.SDK_INT >= Build.VERSION_CODES.UPSIDE_DOWN_CAKE) {
ForegroundInfo(
NOTIFICATION_ID,
notification,
ServiceInfo.FOREGROUND_SERVICE_TYPE_DATA_SYNC,
)
} else {
ForegroundInfo(NOTIFICATION_ID, notification)
}
}
private fun buildNotification(state: ModelState): Notification {
val builder = NotificationCompat.Builder(applicationContext, CHANNEL_ID)
.setSmallIcon(android.R.drawable.stat_sys_download)
.setOnlyAlertOnce(true)
.setOngoing(true)
.setContentIntent(makeContentIntent())
.setPriority(NotificationCompat.PRIORITY_LOW)
.setCategory(NotificationCompat.CATEGORY_PROGRESS)
when (state) {
ModelState.NotDownloaded, ModelState.Ready -> {
builder.setContentTitle("Проверка голосовой модели")
if (state is ModelState.Ready) {
builder.setContentText("Голосовая модель готова")
}
// Без прогресс-бара.
builder.setProgress(0, 0, false)
}
is ModelState.Downloading -> {
val percent = if (state.bytesTotal > 0L) {
((state.bytesDone.toDouble() / state.bytesTotal) * 100.0).roundToInt()
.coerceIn(0, 100)
} else {
0
}
builder.setContentTitle("Скачивание голосовой модели")
builder.setContentText(
"${toMb(state.bytesDone)} / ${toMb(state.bytesTotal)} МБ",
)
builder.setProgress(100, percent, false)
}
is ModelState.Failed -> {
builder.setContentTitle("Ошибка скачивания голосовой модели")
builder.setContentText(state.reason)
builder.setProgress(0, 0, false)
// Не выкидывается автоматически — пользователь сам закроет.
builder.setOngoing(false)
}
}
return builder.build()
}
private fun makeContentIntent(): PendingIntent {
val intent: Intent = applicationContext.packageManager
.getLaunchIntentForPackage(applicationContext.packageName)
?: Intent().apply {
setClassName(
applicationContext,
"pw.binom.agentik.android.MainActivity",
)
}
intent.addFlags(Intent.FLAG_ACTIVITY_CLEAR_TOP or Intent.FLAG_ACTIVITY_SINGLE_TOP)
return PendingIntent.getActivity(
applicationContext,
0,
intent,
PendingIntent.FLAG_UPDATE_CURRENT or PendingIntent.FLAG_IMMUTABLE,
)
}
private fun toMb(bytes: Long): Long = bytes / (1024L * 1024L)
companion object {
private const val CHANNEL_ID = "qwen-model-download"
private const val NOTIFICATION_ID = 1001
private const val CHANNEL_NAME = "Голосовая модель"
private const val CHANNEL_DESCRIPTION =
"Скачивание и проверка голосовой модели Qwen3-ASR"
@Volatile
private var channelEnsured = false
/**
* Создаёт канал нотификации один раз за жизнь процесса. На pre-O (API < 26) каналов нет —
* no-op. Канал создаётся здесь, а не в `Application.onCreate`, чтобы не зависеть от
* инициализации WorkManager: воркер может запуститься и без явной регистрации канала,
* и канал всё равно будет готов.
*/
fun ensureChannel(context: Context) {
if (channelEnsured) return
synchronized(this) {
if (channelEnsured) return
val nm = context.getSystemService<NotificationManager>() ?: return
if (Build.VERSION.SDK_INT >= Build.VERSION_CODES.O) {
val channel = NotificationChannel(
CHANNEL_ID,
CHANNEL_NAME,
NotificationManager.IMPORTANCE_LOW,
).apply {
description = CHANNEL_DESCRIPTION
setShowBadge(false)
}
nm.createNotificationChannel(channel)
}
channelEnsured = true
}
}
}
}
@@ -0,0 +1,308 @@
package pw.binom.agentik.android.asr
import android.content.Context
import android.util.Log
import java.io.File
import java.io.FileOutputStream
import java.io.IOException
import java.io.OutputStream
import java.security.DigestOutputStream
import java.security.MessageDigest
import kotlinx.coroutines.CancellationException
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Deferred
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.async
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.launch
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
import kotlinx.coroutines.withContext
import kotlinx.serialization.Serializable
import kotlinx.serialization.json.Json
import okhttp3.OkHttpClient
import okhttp3.Request
/**
* Манифест модели Qwen3-ASR: версия, базовый URL и список файлов с контрольными суммами.
* Лежит в `assets/qwen3_manifest.json` и парсится по требованию.
*/
@Serializable
data class QwenManifest(
val version: String,
val source: String,
val totalBytes: Long,
val files: List<ManifestFile>,
)
/**
* Один файл модели. `path` может содержать `/` (например `tokenizer/vocab.json`) —
* подкаталог при скачивании сохраняется.
*/
@Serializable
data class ManifestFile(
val path: String,
val size: Long,
val sha256: String,
)
/**
* Состояние модели на устройстве. В `Downloading.bytesDone / bytesTotal`
* отдаётся глобальный прогресс по всем файлам из манифеста.
*/
sealed interface ModelState {
/** Модель ещё не скачана (или была удалена / загрузка отменена). */
data object NotDownloaded : ModelState
/** Идёт скачивание или проверка уже существующих файлов. */
data class Downloading(val bytesDone: Long, val bytesTotal: Long) : ModelState
/** Все файлы из манифеста на месте и прошли проверку SHA-256. */
data object Ready : ModelState
/** Скачивание упало. `reason` — короткое описание для логов и UI. */
data class Failed(val reason: String) : ModelState
}
/**
* Менеджер загрузки модели Qwen3-ASR. Сам ASR-движок модель не качает — это наша обёртка:
* читает манифест из assets, проверяет файлы по SHA-256, докачивает недостающие через OkHttp.
*
* Потокобезопасен: внутренний [Mutex] сериализует вызовы [ensureReady], а [cancelDownload]
* прерывает текущую загрузку снаружи. Состояние наружу — только через [state] (StateFlow).
*/
class QwenModelProvider(context: Context) {
private val appContext: Context = context.applicationContext
/** Куда кладём файлы модели: `filesDir/models/qwen3-asr/`. */
val modelDir: File = File(appContext.filesDir, "models/qwen3-asr")
private val client: OkHttpClient = OkHttpClient()
private val scope: CoroutineScope = CoroutineScope(SupervisorJob() + Dispatchers.IO)
private val mutex: Mutex = Mutex()
private val _state: MutableStateFlow<ModelState> = MutableStateFlow(ModelState.NotDownloaded)
/** Состояние модели. Горячий поток, текущее значение всегда доступно через `.value`. */
val state: StateFlow<ModelState> = _state.asStateFlow()
/** Текущая загрузка — чтобы [cancelDownload] мог её прервать. */
@Volatile
private var currentJob: Deferred<Unit>? = null
init {
// На старте один раз проверяем файлы: если всё уже на месте — сразу Ready.
scope.launch {
try {
if (checkAllReady()) {
_state.value = ModelState.Ready
Log.i(TAG, "init: все файлы модели уже на месте, состояние Ready")
} else {
Log.i(TAG, "init: модель не готова, ждём вызова ensureReady()")
}
} catch (e: CancellationException) {
throw e
} catch (e: Throwable) {
Log.w(TAG, "init: не удалось проверить готовность: ${e.message}")
}
}
}
/**
* Идемпотентно доводит модель до [ModelState.Ready]: проверяет все файлы по SHA-256,
* докачивает недостающие или повреждённые. Прогресс пишется в [state].
*
* @return `Result.success(modelDir)` если все файлы валидны;
* `Result.failure(e)` при сетевой/IO-ошибке.
* При отмене (CancellationException) состояние становится [ModelState.NotDownloaded],
* исключение пробрасывается вызывающему.
*/
suspend fun ensureReady(): Result<File> = mutex.withLock {
val job = scope.async<Unit> {
try {
val manifest = readManifest()
val total = manifest.totalBytes
var done: Long = 0L
for (entry in manifest.files) {
ensureFile(entry, manifest) { currentBytes ->
_state.value = ModelState.Downloading(done + currentBytes, total)
}
done += entry.size
}
_state.value = ModelState.Ready
Log.i(
TAG,
"модель готова: версия ${manifest.version}, " +
"${manifest.files.size} файлов, $total байт",
)
} catch (e: CancellationException) {
_state.value = ModelState.NotDownloaded
Log.i(TAG, "загрузка модели отменена")
throw e
} catch (e: Throwable) {
val reason = e.message ?: e::class.simpleName ?: "unknown error"
Log.e(TAG, "загрузка модели не удалась: $reason", e)
_state.value = ModelState.Failed(reason)
throw e
}
}
currentJob = job
try {
job.await()
Result.success(modelDir)
} catch (e: CancellationException) {
// cancelDownload() или отмена скоупа вызывающего — гасим job и пробрасываем.
job.cancel()
throw e
} catch (e: Throwable) {
Result.failure(e)
} finally {
currentJob = null
}
}
/** Отменяет текущую загрузку, если она идёт. Безопасно вызывать, когда загрузки нет. */
fun cancelDownload() {
currentJob?.cancel()
}
/** Удаляет каталог модели и сбрасывает состояние в [ModelState.NotDownloaded]. */
suspend fun delete() {
withContext(Dispatchers.IO) {
if (modelDir.exists()) {
val deleted = modelDir.deleteRecursively()
Log.i(TAG, "каталог модели удалён: success=$deleted")
}
_state.value = ModelState.NotDownloaded
}
}
// --- внутреннее --------------------------------------------------------
private fun readManifest(): QwenManifest {
val text = appContext.assets.open(MANIFEST_ASSET)
.bufferedReader()
.use { it.readText() }
return JSON.decodeFromString(QwenManifest.serializer(), text)
}
/** Проверяет, что все файлы из манифеста существуют и их SHA-256 совпадает. */
private suspend fun checkAllReady(): Boolean = withContext(Dispatchers.IO) {
val manifest = readManifest()
manifest.files.all { entry ->
val target = File(modelDir, entry.path)
target.exists() && sha256Of(target) == entry.sha256
}
}
/**
* Гарантирует, что [entry] присутствует в [modelDir] и совпадает по SHA-256.
* Если файл уже валиден — пропускает. Если повреждён — удаляет и качает заново.
*/
private suspend fun ensureFile(
entry: ManifestFile,
manifest: QwenManifest,
onProgress: (bytesInCurrentFile: Long) -> Unit,
) = withContext(Dispatchers.IO) {
val target = File(modelDir, entry.path)
target.parentFile?.let { parent ->
if (!parent.exists() && !parent.mkdirs()) {
throw IOException("не удалось создать каталог ${parent.absolutePath}")
}
}
if (target.exists()) {
val existing = sha256Of(target)
if (existing == entry.sha256) {
Log.i(TAG, "[${entry.path}] уже на месте, SHA совпал")
onProgress(entry.size)
return@withContext
}
Log.w(TAG, "[${entry.path}] SHA не совпал (есть ${existing.take(12)}…), перекачиваем")
if (!target.delete()) {
throw IOException("не удалось удалить повреждённый ${target.absolutePath}")
}
}
download(entry, manifest, target, onProgress)
}
private fun download(
entry: ManifestFile,
manifest: QwenManifest,
target: File,
onProgress: (Long) -> Unit,
) {
val url = manifest.source + entry.path
Log.i(TAG, "[${entry.path}] скачиваем $url")
val request = Request.Builder().url(url).build()
val digest = MessageDigest.getInstance("SHA-256")
client.newCall(request).execute().use { response ->
if (!response.isSuccessful) {
throw IOException("HTTP ${response.code} при скачивании $url")
}
val body = response.body ?: throw IOException("пустое тело ответа для $url")
FileOutputStream(target).use { fos ->
DigestOutputStream(fos, digest).use { dos ->
val input = body.byteStream()
val buf = ByteArray(READ_BUFFER)
var written: Long = 0L
while (true) {
val n = input.read(buf)
if (n < 0) break
dos.write(buf, 0, n)
written += n
onProgress(written)
}
dos.flush()
}
}
}
val hash = digest.digest().toHex()
if (hash != entry.sha256) {
target.delete()
throw IOException(
"[${entry.path}] SHA не совпал: ожидался ${entry.sha256}, получили $hash",
)
}
Log.i(TAG, "[${entry.path}] скачан и проверен, ${entry.size} байт")
}
private fun sha256Of(file: File): String {
val digest = MessageDigest.getInstance("SHA-256")
file.inputStream().use { input ->
DigestOutputStream(NULL_OUTPUT_STREAM, digest).use { dos ->
input.copyTo(dos)
}
}
return digest.digest().toHex()
}
private fun ByteArray.toHex(): String =
joinToString(separator = "") { byte -> "%02x".format(byte) }
companion object {
private const val TAG = "QwenModelProvider"
private const val MANIFEST_ASSET = "qwen3_manifest.json"
private const val READ_BUFFER = 64 * 1024
private val JSON: Json = Json { ignoreUnknownKeys = true }
// Свой null-приёмник вместо java.io.OutputStream.nullOutputStream(),
// который доступен в Android только с API 33.
private val NULL_OUTPUT_STREAM: OutputStream = object : OutputStream() {
override fun write(b: Int) {}
override fun write(b: ByteArray, off: Int, len: Int) {}
}
}
}
@@ -0,0 +1,22 @@
package pw.binom.agentik.android.settings
import kotlinx.serialization.Serializable
/**
* Персистируемые настройки подключения к агенту. Хранятся в
* `filesDir/settings.json` через [SettingsRepository] — соответствует
* desktop-варианту (R41.5), где settings живут в JSON-файле, а не в БД.
*
* @property agentId стабильный client-id (UUID-строка, генерируется
* один раз при первом подключении и переиспользуется).
* @property agentName отображаемое имя клиента.
* @property baseUrl базовый URL агента (например, `https://agent.example.com`).
* @property token bearer-токен для авторизации (опционально, может быть null).
*/
@Serializable
data class Settings(
val agentId: String,
val agentName: String,
val baseUrl: String,
val token: String? = null,
)
@@ -0,0 +1,55 @@
package pw.binom.agentik.android.settings
import android.content.Context
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.withContext
import kotlinx.serialization.json.Json
import java.io.File
/**
* Repository для JSON-персистентности [Settings] в `filesDir/settings.json`.
*
* По аналогии с desktop-вариантом (R41.5): настройки в JSON-файле, не в БД.
* Простая read/write — без шифрования, без миграций, без бэкапов.
*
* Все операции файлового ввода/вывода уходят в [Dispatchers.IO] —
* вызовы из UI-потока не блокируют рекомпозицию.
*
* @param context Android context; используется только `context.filesDir`.
*/
class SettingsRepository(private val context: Context) {
private val file: File
get() = File(context.filesDir, FILE_NAME)
private val json = Json {
ignoreUnknownKeys = true
prettyPrint = false
}
/**
* Прочитать настройки. Возвращает `null`, если файла нет или он
* нечитаем/повреждён (повреждённый JSON не должен ронять приложение —
* пользователю проще ввести заново, чем ловить крэш при старте).
*/
suspend fun load(): Settings? = withContext(Dispatchers.IO) {
val f = file
if (!f.exists()) return@withContext null
runCatching { json.decodeFromString<Settings>(f.readText()) }.getOrNull()
}
/**
* Сохранить настройки. Перезаписывает файл целиком (atomic write —
* `writeText` под Android эквивалентен truncate+write для маленького
* файла, что приемлемо для настроек).
*/
suspend fun save(settings: Settings): Unit = withContext(Dispatchers.IO) {
val f = file
f.parentFile?.mkdirs()
f.writeText(json.encodeToString(settings))
}
private companion object {
const val FILE_NAME = "settings.json"
}
}
@@ -0,0 +1,196 @@
package pw.binom.agentik.android.ui.common
import androidx.compose.foundation.BorderStroke
import androidx.compose.foundation.background
import androidx.compose.foundation.border
import androidx.compose.foundation.layout.*
import androidx.compose.foundation.shape.RoundedCornerShape
import androidx.compose.foundation.text.BasicTextField
import androidx.compose.material.icons.Icons
import androidx.compose.material.icons.automirrored.filled.Send
import androidx.compose.material.icons.filled.Stop
import androidx.compose.material3.Icon
import androidx.compose.material3.LocalTextStyle
import androidx.compose.material3.MaterialTheme
import androidx.compose.material3.Surface
import androidx.compose.material3.Text
import androidx.compose.runtime.Composable
import androidx.compose.ui.Alignment
import androidx.compose.ui.Modifier
import androidx.compose.ui.draw.clip
import androidx.compose.ui.graphics.SolidColor
import androidx.compose.ui.text.font.FontWeight
import androidx.compose.ui.text.style.TextAlign
import androidx.compose.ui.tooling.preview.Preview
import androidx.compose.ui.unit.dp
import pw.binom.agentik.android.ui.theme.AgentikTheme
/** Состояние строки ввода; маппится [InputBar]'ом на один из трёх макетов. */
sealed interface InputBarState {
data class Idle(
val text: String = "",
val canSend: Boolean = false,
val micState: MicState = MicState.Idle,
) : InputBarState
data class Recording(val elapsedMillis: Long = 0L) : InputBarState
data object Streaming : InputBarState
}
@Composable
fun InputBar(
state: InputBarState,
onTextChange: (String) -> Unit = {},
onStartRecording: () -> Unit = {},
onStopRecording: () -> Unit = {},
onCancelRecording: () -> Unit = {},
onSend: () -> Unit = {},
onStopAgent: () -> Unit = {},
onMicDisabledClick: () -> Unit = {},
modifier: Modifier = Modifier,
) {
Column(
modifier = modifier
.background(MaterialTheme.colorScheme.surface)
.padding(start = 12.dp, end = 12.dp, top = 10.dp, bottom = 14.dp),
) {
when (state) {
is InputBarState.Idle -> IdleLayout(
state, onTextChange, onStartRecording, onStopRecording, onSend, onMicDisabledClick,
)
is InputBarState.Recording -> RecordingLayout(state, onStopRecording, onCancelRecording)
is InputBarState.Streaming -> StreamingLayout(onStopAgent)
}
}
}
@Composable
private fun IdleLayout(
state: InputBarState.Idle,
onTextChange: (String) -> Unit,
onStartRecording: () -> Unit,
onStopRecording: () -> Unit,
onSend: () -> Unit,
onMicDisabledClick: () -> Unit,
) {
val alpha = if (state.canSend) 1f else 0.35f
Row(verticalAlignment = Alignment.CenterVertically, horizontalArrangement = Arrangement.spacedBy(9.dp)) {
FieldBox(
text = state.text, placeholder = "Сообщение…", onValueChange = onTextChange,
modifier = Modifier.weight(1f),
)
MicButton(
state = state.micState, onStart = onStartRecording, onStop = onStopRecording,
onDisabledClick = onMicDisabledClick,
)
Surface(
onClick = onSend, enabled = state.canSend, shape = RoundedCornerShape(14.dp),
color = MaterialTheme.colorScheme.primary.copy(alpha = alpha),
modifier = Modifier.size(44.dp),
) {
Box(Modifier.fillMaxSize(), Alignment.Center) {
Icon(
imageVector = Icons.AutoMirrored.Filled.Send, contentDescription = "Отправить",
tint = MaterialTheme.colorScheme.onPrimary.copy(alpha = alpha),
modifier = Modifier.size(20.dp),
)
}
}
}
}
@Composable
private fun RecordingLayout(
state: InputBarState.Recording,
onStopRecording: () -> Unit,
onCancelRecording: () -> Unit,
) {
RecordingPanel(
elapsedMillis = state.elapsedMillis, onCancel = onCancelRecording, onSend = onStopRecording,
modifier = Modifier.padding(bottom = 10.dp),
)
Row(verticalAlignment = Alignment.CenterVertically, horizontalArrangement = Arrangement.spacedBy(9.dp)) {
FieldBox(
text = "", placeholder = "Идёт запись…", onValueChange = {},
modifier = Modifier.weight(1f),
)
MicButton(state = MicState.Idle, onStart = {}, onStop = onStopRecording, recording = true)
}
RecNote("Говорите — текст появится в поле. ✕ отменяет запись целиком.")
}
@Composable
private fun StreamingLayout(onStopAgent: () -> Unit) {
Surface(
onClick = onStopAgent, shape = RoundedCornerShape(14.dp),
color = MaterialTheme.colorScheme.error.copy(alpha = 0.10f),
border = BorderStroke(1.dp, MaterialTheme.colorScheme.error),
modifier = Modifier.fillMaxWidth().height(44.dp),
) {
Row(
Modifier.fillMaxSize(),
verticalAlignment = Alignment.CenterVertically,
horizontalArrangement = Arrangement.spacedBy(9.dp, Alignment.CenterHorizontally),
) {
Icon(
imageVector = Icons.Default.Stop, contentDescription = null,
tint = MaterialTheme.colorScheme.error, modifier = Modifier.size(16.dp),
)
Text("Стоп", color = MaterialTheme.colorScheme.error, fontWeight = FontWeight.SemiBold)
}
}
RecNote("Агент отвечает — можно прервать")
}
@Composable
private fun FieldBox(
text: String, placeholder: String, onValueChange: (String) -> Unit, modifier: Modifier = Modifier,
) {
BasicTextField(
value = text, onValueChange = onValueChange,
modifier = modifier
.height(44.dp)
.clip(RoundedCornerShape(14.dp))
.background(MaterialTheme.colorScheme.background)
.border(1.dp, MaterialTheme.colorScheme.outline, RoundedCornerShape(14.dp))
.padding(horizontal = 14.dp),
textStyle = LocalTextStyle.current.copy(color = MaterialTheme.colorScheme.onSurface),
cursorBrush = SolidColor(MaterialTheme.colorScheme.primary),
singleLine = true,
decorationBox = { inner ->
Box(Modifier.fillMaxSize(), Alignment.CenterStart) {
if (text.isEmpty()) Text(placeholder, color = MaterialTheme.colorScheme.onSurfaceVariant)
inner()
}
},
)
}
@Composable
private fun RecNote(text: String) {
Text(
text = text, color = MaterialTheme.colorScheme.onSurfaceVariant,
style = MaterialTheme.typography.bodySmall,
modifier = Modifier.fillMaxWidth().padding(top = 8.dp),
textAlign = TextAlign.Center,
)
}
@Preview
@Composable private fun InputBarPreviewIdle() = AgentikTheme {
InputBar(state = InputBarState.Idle(text = "Привет", canSend = true))
}
@Preview
@Composable private fun InputBarPreviewIdleMicDisabled() = AgentikTheme {
InputBar(state = InputBarState.Idle(micState = MicState.Disabled("Нет доступа к микрофону")))
}
@Preview
@Composable private fun InputBarPreviewRecording() = AgentikTheme {
InputBar(state = InputBarState.Recording(elapsedMillis = 7_000))
}
@Preview
@Composable private fun InputBarPreviewStreaming() = AgentikTheme {
InputBar(state = InputBarState.Streaming)
}
@@ -0,0 +1,139 @@
package pw.binom.agentik.android.ui.common
import androidx.compose.foundation.BorderStroke
import androidx.compose.foundation.background
import androidx.compose.foundation.layout.Box
import androidx.compose.foundation.layout.fillMaxSize
import androidx.compose.foundation.layout.size
import androidx.compose.foundation.shape.CircleShape
import androidx.compose.material.icons.Icons
import androidx.compose.material.icons.filled.Mic
import androidx.compose.material.icons.filled.Stop
import androidx.compose.material3.ExperimentalMaterial3Api
import androidx.compose.material3.Icon
import androidx.compose.material3.MaterialTheme
import androidx.compose.material3.PlainTooltip
import androidx.compose.material3.Surface
import androidx.compose.material3.Text
import androidx.compose.material3.TooltipBox
import androidx.compose.material3.TooltipDefaults
import androidx.compose.material3.rememberTooltipState
import androidx.compose.runtime.Composable
import androidx.compose.ui.Alignment
import androidx.compose.ui.Modifier
import androidx.compose.ui.graphics.Color
import androidx.compose.ui.tooling.preview.Preview
import androidx.compose.ui.unit.dp
import pw.binom.agentik.android.ui.theme.AgentikTheme
/**
* Состояние кнопки микрофона в [InputBarState.Idle].
* - [Idle] — готова к записи;
* - [Disabled] — недоступна (например, нет разрешения на микрофон); при нажатии
* показывается [tooltip].
*/
sealed interface MicState {
data object Idle : MicState
data class Disabled(val tooltip: String) : MicState
}
/**
* Круглая 44×44 кнопка микрофона. Только визуал — наружу колбэки.
*
* Визуальные режимы:
* - `state == MicState.Idle` и `recording == false` — серая, иконка микрофона,
* клик вызывает [onStart];
* - `state is MicState.Disabled` — серая, иконка микрофона, клик вызывает
* [onDisabledClick] и показывает тултип с [MicState.Disabled.tooltip];
* - `recording == true` — красная, иконка «стоп», клик вызывает [onStop].
*/
@OptIn(ExperimentalMaterial3Api::class)
@Composable
fun MicButton(
state: MicState,
onStart: () -> Unit,
onStop: () -> Unit,
onDisabledClick: () -> Unit = {},
recording: Boolean = false,
modifier: Modifier = Modifier,
) {
val isDisabled = !recording && state is MicState.Disabled
val tooltip = (state as? MicState.Disabled)?.tooltip
val containerColor =
if (recording) MaterialTheme.colorScheme.error
else MaterialTheme.colorScheme.surface
val borderColor =
if (recording) MaterialTheme.colorScheme.error
else MaterialTheme.colorScheme.outline
val iconTint =
if (recording) Color.White
else MaterialTheme.colorScheme.onSurfaceVariant
val iconVector = if (recording) Icons.Default.Stop else Icons.Default.Mic
val description = when {
recording -> "Остановить запись"
isDisabled -> tooltip.orEmpty()
else -> "Начать запись"
}
val onClick = when {
recording -> onStop
isDisabled -> onDisabledClick
else -> onStart
}
val button: @Composable () -> Unit = {
Surface(
onClick = onClick,
shape = CircleShape,
color = containerColor,
border = BorderStroke(1.dp, borderColor),
modifier = modifier.size(44.dp),
) {
Box(
modifier = Modifier.fillMaxSize(),
contentAlignment = Alignment.Center,
) {
Icon(
imageVector = iconVector,
contentDescription = description,
tint = iconTint,
modifier = Modifier.size(20.dp),
)
}
}
}
if (tooltip != null && !recording) {
TooltipBox(
positionProvider = TooltipDefaults.rememberPlainTooltipPositionProvider(),
tooltip = { PlainTooltip { Text(tooltip) } },
state = rememberTooltipState(isPersistent = false),
) {
button()
}
} else {
button()
}
}
@Preview
@Composable
private fun MicButtonPreviewIdle() = AgentikTheme {
MicButton(state = MicState.Idle, onStart = {}, onStop = {})
}
@Preview
@Composable
private fun MicButtonPreviewRecording() = AgentikTheme {
MicButton(state = MicState.Idle, onStart = {}, onStop = {}, recording = true)
}
@Preview
@Composable
private fun MicButtonPreviewDisabled() = AgentikTheme {
MicButton(
state = MicState.Disabled("Нет доступа к микрофону"),
onStart = {},
onStop = {},
)
}
@@ -0,0 +1,138 @@
package pw.binom.agentik.android.ui.common
import androidx.compose.foundation.BorderStroke
import androidx.compose.foundation.background
import androidx.compose.foundation.border
import androidx.compose.foundation.layout.*
import androidx.compose.foundation.shape.CircleShape
import androidx.compose.foundation.shape.RoundedCornerShape
import androidx.compose.material.icons.Icons
import androidx.compose.material.icons.automirrored.filled.Send
import androidx.compose.material.icons.filled.Close
import androidx.compose.material3.Icon
import androidx.compose.material3.LocalTextStyle
import androidx.compose.material3.MaterialTheme
import androidx.compose.material3.Surface
import androidx.compose.material3.Text
import androidx.compose.runtime.Composable
import androidx.compose.ui.Alignment
import androidx.compose.ui.Modifier
import androidx.compose.ui.draw.clip
import androidx.compose.ui.text.font.FontWeight
import androidx.compose.ui.tooling.preview.Preview
import androidx.compose.ui.unit.dp
import pw.binom.agentik.android.ui.theme.AgentikTheme
/**
* Панель, которая показывается над строкой ввода во время записи голоса
* ([InputBarState.Recording]). Содержит декоративную «волну» из 9 полосок,
* таймер записи в формате m:ss, кнопки «отменить» и «отправить».
*
* Только визуал: таймер тикает извне через [elapsedMillis], за отправку/отмену
* отвечает родитель.
*/
@Composable
fun RecordingPanel(
elapsedMillis: Long,
onCancel: () -> Unit,
onSend: () -> Unit,
modifier: Modifier = Modifier,
) {
val seconds = elapsedMillis / 1000L
val timeText = "%d:%02d".format(seconds / 60, seconds % 60)
Surface(
onClick = {},
shape = RoundedCornerShape(14.dp),
color = MaterialTheme.colorScheme.background,
border = BorderStroke(1.dp, MaterialTheme.colorScheme.error),
modifier = modifier,
) {
Row(
modifier = Modifier
.fillMaxWidth()
.padding(horizontal = 13.dp, vertical = 10.dp),
verticalAlignment = Alignment.CenterVertically,
horizontalArrangement = Arrangement.spacedBy(11.dp),
) {
Wave(modifier = Modifier.weight(1f))
Text(
text = timeText,
color = MaterialTheme.colorScheme.error,
fontWeight = FontWeight.Bold,
style = LocalTextStyle.current.copy(fontFeatureSettings = "tnum"),
)
CancelAction(onClick = onCancel)
SendAction(onClick = onSend)
}
}
}
private val BarHeights = listOf(0.30f, 0.70f, 0.45f, 1.00f, 0.55f, 0.80f, 0.35f, 0.65f, 0.25f)
@Composable
private fun Wave(modifier: Modifier = Modifier) {
Row(
modifier = modifier
.height(22.dp)
.fillMaxWidth(),
verticalAlignment = Alignment.CenterVertically,
horizontalArrangement = Arrangement.spacedBy(3.dp),
) {
BarHeights.forEach { ratio ->
Box(
modifier = Modifier
.width(3.dp)
.fillMaxHeight(ratio)
.clip(RoundedCornerShape(2.dp))
.background(MaterialTheme.colorScheme.error),
)
}
}
}
@Composable
private fun CancelAction(onClick: () -> Unit) {
Surface(
onClick = onClick,
shape = CircleShape,
color = MaterialTheme.colorScheme.background,
modifier = Modifier.size(34.dp),
) {
Box(contentAlignment = Alignment.Center, modifier = Modifier.fillMaxSize()) {
Icon(
imageVector = Icons.Default.Close,
contentDescription = "Отменить запись",
tint = MaterialTheme.colorScheme.error,
modifier = Modifier.size(18.dp),
)
}
}
}
@Composable
private fun SendAction(onClick: () -> Unit) {
Surface(
onClick = onClick,
shape = RoundedCornerShape(14.dp),
color = MaterialTheme.colorScheme.primary,
modifier = Modifier.size(38.dp),
) {
Box(contentAlignment = Alignment.Center, modifier = Modifier.fillMaxSize()) {
Icon(
imageVector = Icons.AutoMirrored.Filled.Send,
contentDescription = "Отправить запись",
tint = MaterialTheme.colorScheme.onPrimary,
modifier = Modifier.size(20.dp),
)
}
}
}
@Preview
@Composable
private fun RecordingPanelPreview() = AgentikTheme {
Box(modifier = Modifier.padding(12.dp)) {
RecordingPanel(elapsedMillis = 7_000, onCancel = {}, onSend = {})
}
}
@@ -0,0 +1,41 @@
package pw.binom.agentik.android.ui.markdown
import androidx.compose.foundation.layout.Box
import androidx.compose.runtime.Composable
import androidx.compose.ui.Modifier
import androidx.compose.ui.graphics.Color
import androidx.compose.ui.platform.LocalContext
import androidx.compose.ui.unit.TextUnit
import androidx.compose.ui.unit.sp
/**
* Единая точка входа для UI: парсит `text` как Markdown и рендерит через
* скопированный [renderMarkdownSegments]. Цвет/размер берутся из
* [androidx.compose.material3.MaterialTheme.colorScheme.onSurface] и
* 14.sp, если вызывающий не передал явно — этого достаточно для большинства
* пузырей чата; тонкая настройка — через `textColor`/`fontSize`.
*
* Клик по ссылке открывает её системным браузером через [MarkdownRenderer.openUrl]
* (Intent.ACTION_VIEW + FLAG_ACTIVITY_NEW_TASK).
*/
@Composable
fun MarkdownText(
text: String,
modifier: Modifier = Modifier,
textColor: Color = Color.Unspecified,
fontSize: TextUnit = 14.sp,
isBubbleScrolled: Boolean = false,
) {
val context = LocalContext.current
// [renderMarkdownSegments] не принимает Modifier — оборачиваем в Box,
// чтобы вызывающий мог влиять на размеры/отступы снаружи.
Box(modifier = modifier) {
renderMarkdownSegments(
content = text,
textColor = textColor,
fontSize = fontSize,
isBubbleScrolled = isBubbleScrolled,
onUrlClick = { url -> MarkdownRenderer.openUrl(context, url) },
)
}
}
@@ -0,0 +1,61 @@
package pw.binom.agentik.android.ui.screen
import androidx.compose.runtime.Composable
import androidx.compose.runtime.LaunchedEffect
import androidx.compose.runtime.collectAsState
import androidx.compose.runtime.getValue
import androidx.compose.runtime.mutableStateOf
import androidx.compose.runtime.saveable.rememberSaveable
import androidx.compose.runtime.setValue
import pw.binom.agentik.android.agent.AgentConnection
import pw.binom.agentik.android.agent.ConnectionState
import pw.binom.agentik.android.settings.SettingsRepository
/**
* Корневой композабл: переключает SettingsScreen / ConversationListScreen /
* ChatScreen по текущему [ConnectionState] и локальному `selectedConvId`.
*
* Без навигационной библиотеки — простой `when` + rememberSaveable.
*
* - state ≠ Connected → [SettingsScreen] (форма подключения).
* - state = Connected, convId = null → [ConversationListScreen] (список диалогов).
* - state = Connected, convId != null → [ChatScreen] (открытый диалог).
*
* Сброс [selectedConvId] в null на любой не-Connected state — чтобы при
* разрыве соединения пользователь не оказался внутри ChatScreen с мёртвой
* сессией.
*/
@Composable
fun AppRoot(
connection: AgentConnection,
settingsRepository: SettingsRepository,
) {
val state by connection.state.collectAsState()
var selectedConvId by rememberSaveable { mutableStateOf<String?>(null) }
LaunchedEffect(state) {
if (state !is ConnectionState.Connected) selectedConvId = null
}
when {
state !is ConnectionState.Connected -> {
SettingsScreen(
connection = connection,
settingsRepository = settingsRepository,
)
}
selectedConvId == null -> {
ConversationListScreen(
connection = connection,
onOpenChat = { id -> selectedConvId = id },
)
}
else -> {
ChatScreen(
connection = connection,
convId = selectedConvId!!,
onBack = { selectedConvId = null },
)
}
}
}
@@ -0,0 +1,399 @@
package pw.binom.agentik.android.ui.screen
import androidx.compose.foundation.layout.Arrangement
import androidx.compose.foundation.layout.Box
import androidx.compose.foundation.layout.Column
import androidx.compose.foundation.layout.PaddingValues
import androidx.compose.foundation.layout.Row
import androidx.compose.foundation.layout.fillMaxSize
import androidx.compose.foundation.layout.fillMaxWidth
import androidx.compose.foundation.layout.padding
import androidx.compose.foundation.layout.widthIn
import androidx.compose.foundation.lazy.LazyColumn
import androidx.compose.foundation.lazy.items
import androidx.compose.foundation.lazy.rememberLazyListState
import androidx.compose.foundation.shape.RoundedCornerShape
import androidx.compose.material.icons.Icons
import androidx.compose.material.icons.automirrored.filled.ArrowBack
import androidx.compose.material.icons.automirrored.filled.Send
import androidx.compose.material3.CircularProgressIndicator
import androidx.compose.material3.ExperimentalMaterial3Api
import androidx.compose.material3.Icon
import androidx.compose.material3.IconButton
import androidx.compose.material3.MaterialTheme
import androidx.compose.material3.OutlinedTextField
import androidx.compose.material3.Scaffold
import androidx.compose.material3.Surface
import androidx.compose.material3.Text
import androidx.compose.material3.TopAppBar
import androidx.compose.material3.Button
import androidx.compose.runtime.Composable
import androidx.compose.runtime.DisposableEffect
import androidx.compose.runtime.LaunchedEffect
import androidx.compose.runtime.collectAsState
import androidx.compose.runtime.getValue
import androidx.compose.runtime.mutableStateOf
import androidx.compose.runtime.remember
import androidx.compose.runtime.rememberCoroutineScope
import androidx.compose.runtime.setValue
import androidx.compose.ui.Alignment
import androidx.compose.ui.Modifier
import androidx.compose.ui.graphics.Color
import androidx.compose.ui.text.font.FontFamily
import androidx.compose.ui.unit.dp
import kotlinx.coroutines.launch
import pw.binom.agentik.android.agent.AgentConnection
import pw.binom.agentik.android.agent.ChatSession
import pw.binom.agentik.android.agent.ConversationSnapshot
import pw.binom.agentik.android.agent.UiEvent
import pw.binom.agentik.android.agent.UiMessage
import pw.binom.agentik.android.ui.markdown.MarkdownText
/**
* Экран одного открытого диалога. Жизненный цикл:
* 1. `LaunchedEffect(convId)` → `connection.openConversation(convId)`. Если
* convId не найден на сервере — `session` остаётся null, после таймаута
* показываем "Диалог не найден" + кнопку "Назад".
* 2. `LaunchedEffect(session)` → собираем `session.liveEvents` в локальный
* `liveBuffer` (буфер «в процессе»). На `UiEvent.End` буфер очищается —
* при следующем `End` к этому моменту [ChatSession.messages] уже
* содержит закрытый ход.
* 3. `DisposableEffect(Unit) { onDispose { session?.close() } }` —
* закрываем сессию при уходе с экрана (AppRoot переключает обратно на
* список диалогов при `onBack()`, либо при выходе из приложения).
*
* UI:
* - TopAppBar: заголовок = `snapshot.title` или "Новый диалог", кнопка "Назад"
* вызывает `onBack()` → AppRoot рекомпозится в ConversationListScreen.
* - Тело: LazyColumn с историей + «живой» бабл внизу, пока
* `liveBuffer` не пуст.
* - Низ: простой input bar (`OutlinedTextField` + Send). Voice (MicButton /
* RecordingPanel / InputBar из ui/common) — отдельная задача.
*/
@OptIn(ExperimentalMaterial3Api::class)
@Composable
fun ChatScreen(
connection: AgentConnection,
convId: String,
onBack: () -> Unit,
modifier: Modifier = Modifier,
) {
var session by remember { mutableStateOf<ChatSession?>(null) }
var sessionError by remember { mutableStateOf<String?>(null) }
var notFound by remember { mutableStateOf(false) }
var liveBuffer by remember { mutableStateOf<List<UiEvent>>(emptyList()) }
var inputText by remember { mutableStateOf("") }
val scope = rememberCoroutineScope()
val listState = rememberLazyListState()
LaunchedEffect(convId) {
try {
session = connection.openConversation(convId)
} catch (e: Exception) {
sessionError = e.message ?: e::class.simpleName ?: "Не удалось открыть диалог"
}
// Если после секунды session всё ещё null и ошибки не было — диалог
// не найден на сервере. openConversation у нас не различает "не найден"
// от "создаю новый с convId", поэтому явный флаг через секунду.
kotlinx.coroutines.delay(1000)
if (session == null && sessionError == null) {
notFound = true
}
}
LaunchedEffect(session) {
val current = session ?: return@LaunchedEffect
current.liveEvents.collect { ev ->
liveBuffer = if (ev is UiEvent.End) emptyList() else liveBuffer + ev
}
}
DisposableEffect(Unit) {
onDispose {
session?.close()
}
}
// Подписки на сообщения и метаданные. Когда session ещё null,
// collectAsState не вызывается (короткое замыкание по `?.`).
val messages: List<UiMessage> =
session?.messages?.collectAsState(initial = emptyList())?.value ?: emptyList()
val snapshot: ConversationSnapshot? =
session?.snapshot?.collectAsState()?.value
Scaffold(
modifier = modifier,
topBar = {
TopAppBar(
title = { Text(snapshot?.title ?: "Новый диалог") },
navigationIcon = {
IconButton(onClick = onBack) {
Icon(
imageVector = Icons.AutoMirrored.Filled.ArrowBack,
contentDescription = "Назад",
)
}
},
)
},
bottomBar = {
ChatInputBar(
text = inputText,
onTextChange = { inputText = it },
enabled = session != null,
onSend = {
val s = session ?: return@ChatInputBar
val text = inputText.trim()
if (text.isEmpty()) return@ChatInputBar
inputText = ""
scope.launch {
try {
s.send(text)
} catch (_: Exception) {
// Подавляем: на следующем End в messages появится
// UiMessage.Error (если сервер пришлёт), а Failed
// уже обрабатывается в AppRoot. Здесь просто не
// теряем набранное.
}
}
},
)
},
) { innerPadding ->
Box(
modifier = Modifier
.padding(innerPadding)
.fillMaxSize(),
contentAlignment = Alignment.Center,
) {
val err = sessionError
when {
notFound -> {
Column(
horizontalAlignment = Alignment.CenterHorizontally,
verticalArrangement = Arrangement.spacedBy(12.dp),
modifier = Modifier.padding(24.dp),
) {
Text(text = "Диалог не найден")
Button(onClick = onBack) { Text("Назад") }
}
}
err != null -> {
Text(
text = "Ошибка: $err",
color = MaterialTheme.colorScheme.error,
modifier = Modifier.padding(24.dp),
)
}
session == null -> {
CircularProgressIndicator()
}
else -> {
// Авто-скролл вниз при появлении нового сообщения или live-блока.
LaunchedEffect(messages.size, liveBuffer.size) {
if (messages.isNotEmpty() || liveBuffer.isNotEmpty()) {
listState.animateScrollToItem(messages.size)
}
}
MessagesList(
messages = messages,
liveBuffer = liveBuffer,
listState = listState,
)
}
}
}
}
}
@Composable
private fun MessagesList(
messages: List<UiMessage>,
liveBuffer: List<UiEvent>,
listState: androidx.compose.foundation.lazy.LazyListState,
) {
LazyColumn(
state = listState,
modifier = Modifier.fillMaxSize(),
contentPadding = PaddingValues(horizontal = 12.dp, vertical = 8.dp),
verticalArrangement = Arrangement.spacedBy(8.dp),
) {
items(messages, key = { it.id }) { msg ->
MessageBubble(msg)
}
if (liveBuffer.isNotEmpty()) {
item(key = "live-bubble") {
LiveBubble(liveBuffer)
}
}
}
}
@OptIn(ExperimentalMaterial3Api::class)
@Composable
private fun ChatInputBar(
text: String,
onTextChange: (String) -> Unit,
enabled: Boolean,
onSend: () -> Unit,
) {
Surface(
tonalElevation = 3.dp,
modifier = Modifier.fillMaxWidth(),
) {
Row(
modifier = Modifier
.fillMaxWidth()
.padding(horizontal = 12.dp, vertical = 8.dp),
verticalAlignment = Alignment.CenterVertically,
horizontalArrangement = Arrangement.spacedBy(8.dp),
) {
OutlinedTextField(
value = text,
onValueChange = onTextChange,
placeholder = { Text("Сообщение…") },
enabled = enabled,
modifier = Modifier.weight(1f),
maxLines = 4,
)
IconButton(
onClick = onSend,
enabled = enabled && text.isNotBlank(),
) {
Icon(
imageVector = Icons.AutoMirrored.Filled.Send,
contentDescription = "Отправить",
)
}
}
}
}
/* ----------------------------- bubble composables ----------------------------- */
@Composable
private fun MessageBubble(msg: UiMessage) {
when (msg) {
is UiMessage.User -> UserBubble(msg)
is UiMessage.Assistant -> AssistantBubble(msg)
is UiMessage.ToolCall -> ToolCallBubble(msg)
is UiMessage.ToolResult -> ToolResultBubble(msg)
is UiMessage.Error -> ErrorBubble(msg)
}
}
@Composable
private fun BubbleSurface(
color: Color,
content: @Composable () -> Unit,
) {
Surface(
color = color,
shape = RoundedCornerShape(12.dp),
modifier = Modifier.widthIn(max = 320.dp),
) {
Box(modifier = Modifier.padding(horizontal = 12.dp, vertical = 8.dp)) {
content()
}
}
}
@Composable
private fun UserBubble(msg: UiMessage.User) {
Row(
modifier = Modifier.fillMaxWidth(),
horizontalArrangement = Arrangement.End,
) {
BubbleSurface(color = MaterialTheme.colorScheme.primaryContainer) {
Text(msg.text, style = MaterialTheme.typography.bodyLarge)
}
}
}
@Composable
private fun AssistantBubble(msg: UiMessage.Assistant) {
Row(
modifier = Modifier.fillMaxWidth(),
horizontalArrangement = Arrangement.Start,
) {
BubbleSurface(color = MaterialTheme.colorScheme.surfaceVariant) {
// Markdown-рендеринг вместо plain Text: кодовые блоки, таблицы,
// ссылки, инлайн-стили (жирный/курсив/strikethrough). Цвет —
// стандартный onSurface от MaterialTheme (MarkdownText сам
// подберёт дефолты; тут bubble color — surfaceVariant).
MarkdownText(text = msg.text)
}
}
}
@Composable
private fun ToolCallBubble(msg: UiMessage.ToolCall) {
BubbleSurface(color = MaterialTheme.colorScheme.surface) {
Column {
Text(
text = "🔧 ${msg.toolName}",
style = MaterialTheme.typography.labelMedium,
)
msg.title?.takeIf { it.isNotBlank() }?.let {
Text(it, style = MaterialTheme.typography.bodySmall)
}
msg.args?.takeIf { it.isNotBlank() }?.let {
Text(
text = it,
style = MaterialTheme.typography.bodySmall.copy(fontFamily = FontFamily.Monospace),
)
}
}
}
}
@Composable
private fun ToolResultBubble(msg: UiMessage.ToolResult) {
BubbleSurface(color = MaterialTheme.colorScheme.surface) {
Text(
text = "✓ ${msg.toolName ?: "tool"}: ${msg.result.orEmpty()}",
style = MaterialTheme.typography.bodySmall,
)
}
}
@Composable
private fun ErrorBubble(msg: UiMessage.Error) {
Text(
text = "⚠ ${msg.message}",
color = MaterialTheme.colorScheme.error,
style = MaterialTheme.typography.bodyMedium,
modifier = Modifier.fillMaxWidth(),
)
}
@Composable
private fun LiveBubble(events: List<UiEvent>) {
val text = events.filterIsInstance<UiEvent.AppendText>().joinToString("") { it.body }
val toolCalls = events.filterIsInstance<UiEvent.ToolCall>()
val toolResults = events.filterIsInstance<UiEvent.ToolResult>()
Row(
modifier = Modifier.fillMaxWidth(),
horizontalArrangement = Arrangement.Start,
) {
BubbleSurface(color = MaterialTheme.colorScheme.surfaceVariant) {
Column(verticalArrangement = Arrangement.spacedBy(4.dp)) {
if (text.isNotEmpty()) {
Text(text, style = MaterialTheme.typography.bodyLarge)
}
toolCalls.forEach { tc ->
Text(
text = "🔧 ${tc.toolName}(${tc.args.orEmpty()})",
style = MaterialTheme.typography.labelMedium,
)
}
toolResults.forEach { tr ->
Text(
text = "✓ ${tr.toolName ?: "tool"}: ${tr.result.orEmpty()}",
style = MaterialTheme.typography.labelMedium,
)
}
}
}
}
}
@@ -0,0 +1,290 @@
package pw.binom.agentik.android.ui.screen
import androidx.compose.foundation.background
import androidx.compose.foundation.clickable
import androidx.compose.foundation.layout.Arrangement
import androidx.compose.foundation.layout.Box
import androidx.compose.foundation.layout.Column
import androidx.compose.foundation.layout.PaddingValues
import androidx.compose.foundation.layout.Row
import androidx.compose.foundation.layout.Spacer
import androidx.compose.foundation.layout.fillMaxSize
import androidx.compose.foundation.layout.fillMaxWidth
import androidx.compose.foundation.layout.padding
import androidx.compose.foundation.layout.size
import androidx.compose.foundation.lazy.LazyColumn
import androidx.compose.foundation.lazy.items
import androidx.compose.foundation.shape.CircleShape
import androidx.compose.material.icons.Icons
import androidx.compose.material.icons.automirrored.filled.Logout
import androidx.compose.material.icons.filled.Add
import androidx.compose.material.icons.filled.ChatBubbleOutline
import androidx.compose.material3.Button
import androidx.compose.material3.ExperimentalMaterial3Api
import androidx.compose.material3.ExtendedFloatingActionButton
import androidx.compose.material3.Icon
import androidx.compose.material3.IconButton
import androidx.compose.material3.MaterialTheme
import androidx.compose.material3.Scaffold
import androidx.compose.material3.Surface
import androidx.compose.material3.Text
import androidx.compose.material3.TopAppBar
import androidx.compose.runtime.Composable
import androidx.compose.runtime.LaunchedEffect
import androidx.compose.runtime.getValue
import androidx.compose.runtime.mutableStateOf
import androidx.compose.runtime.remember
import androidx.compose.runtime.rememberCoroutineScope
import androidx.compose.runtime.setValue
import androidx.compose.ui.Alignment
import androidx.compose.ui.Modifier
import androidx.compose.ui.text.style.TextOverflow
import androidx.compose.ui.unit.dp
import kotlinx.coroutines.launch
import pw.binom.agentik.android.agent.AgentConnection
import pw.binom.agentik.android.agent.ConversationSummary
import pw.binom.agentik.android.agent.real.HttpAgentConnection
import java.time.LocalDateTime
import java.time.ZoneId
import kotlin.time.Clock
import kotlin.time.Instant
/**
* Экран списка диалогов. Три состояния:
* - загрузка (пустой list + "Загрузка…" — короткая фаза до первого эмита);
* - пустой (EmptyState с CTA "Новый диалог");
* - список (`LazyColumn` с `ConversationRow`).
*
* **Источник данных (v12):** [AgentConnection.listConversations] — обогащённый
* список из connection layer (id / title / updatedAt / lastMessagePreview /
* unreadCount). Внутри `HttpAgentConnection.listConversations()` ходит в
* `conversationStore.listFlow()` + `metaRepo.getAllLastSeen()` + локальный
* кэш для подсчёта `unreadCount`. Stub-импл возвращает статический список.
*
* **Live-refresh:** [HttpAgentConnection.agentEvents] фильтрует
* `CommonEvent.Agent` (Created / Deleted / Renamed / Touched) из outbox — на
* каждый эмит перезапрашиваем список через `listConversations()`. Это
* сохраняет v10-семантику "UI обновляется на lifecycle диалога" поверх
* нового API. Для Stub-режима signal недоступен — refresh только при mount.
*
* Действия:
* - FAB "+" / EmptyState-кнопка → `connection.openConversation(null)` →
* `onOpenChat(newId)`.
* - TopAppBar "Disconnect" → `connection.close()` → AppRoot рекомпозится
* в SettingsScreen (state → Disconnected).
* - Row-tap → `onOpenChat(convId)`.
*/
@OptIn(ExperimentalMaterial3Api::class)
@Composable
fun ConversationListScreen(
connection: AgentConnection,
onOpenChat: (String) -> Unit,
modifier: Modifier = Modifier,
) {
var summaries by remember { mutableStateOf<List<ConversationSummary>>(emptyList()) }
val scope = rememberCoroutineScope()
suspend fun reload() {
try {
// HttpAgentConnection.listConversations() внутри делает
// listFlow().toList() + metaRepo.getAllLastSeen() + cache lookup.
// StubAgentConnection возвращает статику.
summaries = connection.listConversations()
} catch (_: Exception) {
// Сервер недоступен / seed ещё не прошёл — оставляем предыдущий список.
}
}
// Первичная загрузка.
LaunchedEffect(connection) { reload() }
// Live-refresh на lifecycle-события диалога (только HTTP). Каст к
// HttpAgentConnection: интерфейс AgentConnection этот signal не
// экспонирует.
val httpConnection = connection as? HttpAgentConnection
LaunchedEffect(httpConnection) {
httpConnection?.agentEvents?.collect { reload() }
}
Scaffold(
modifier = modifier,
topBar = {
TopAppBar(
title = { Text("Диалоги") },
actions = {
IconButton(onClick = { connection.close() }) {
Icon(
imageVector = Icons.AutoMirrored.Filled.Logout,
contentDescription = "Отключиться",
)
}
},
)
},
floatingActionButton = {
ExtendedFloatingActionButton(
onClick = {
scope.launch {
try {
val session = connection.openConversation(convId = null)
onOpenChat(session.conversationId)
} catch (_: Exception) {
// Не критично — пользователь попробует ещё раз
}
}
},
icon = {
Icon(
imageVector = Icons.Filled.Add,
contentDescription = null,
)
},
text = { Text("Новый диалог") },
)
},
) { innerPadding ->
Box(
modifier = Modifier
.padding(innerPadding)
.fillMaxSize(),
contentAlignment = Alignment.Center,
) {
when {
summaries.isEmpty() -> {
// Первичная загрузка — короткая фаза, listFlow пейджит
// через list() пока не упрётся в пустую страницу.
Text(
text = "Загрузка…",
color = MaterialTheme.colorScheme.onSurfaceVariant,
)
}
else -> {
LazyColumn(
modifier = Modifier.fillMaxSize(),
contentPadding = PaddingValues(vertical = 8.dp),
) {
items(summaries, key = { it.id }) { summary ->
ConversationRow(
summary = summary,
onClick = { onOpenChat(summary.id) },
)
}
}
}
}
}
}
}
@Composable
private fun ConversationRow(summary: ConversationSummary, onClick: () -> Unit) {
Surface(
modifier = Modifier
.fillMaxWidth()
.clickable(onClick = onClick),
color = MaterialTheme.colorScheme.surface,
) {
Row(
modifier = Modifier
.fillMaxWidth()
.padding(horizontal = 16.dp, vertical = 12.dp),
verticalAlignment = Alignment.CenterVertically,
horizontalArrangement = Arrangement.spacedBy(12.dp),
) {
// Avatar-кружок — placeholder, цвет из темы (R1.1).
Box(
modifier = Modifier
.size(40.dp)
.background(
color = MaterialTheme.colorScheme.primaryContainer,
shape = CircleShape,
),
contentAlignment = Alignment.Center,
) {
Text(
text = (summary.title ?: "Б").take(1).uppercase(),
style = MaterialTheme.typography.titleMedium,
color = MaterialTheme.colorScheme.onPrimaryContainer,
)
}
Column(modifier = Modifier.weight(1f)) {
Row(verticalAlignment = Alignment.CenterVertically) {
Text(
text = summary.title ?: "Без названия",
style = MaterialTheme.typography.titleSmall,
maxLines = 1,
overflow = TextOverflow.Ellipsis,
modifier = Modifier.weight(1f),
)
Text(
text = formatTimestamp(summary.updatedAt),
style = MaterialTheme.typography.labelSmall,
color = MaterialTheme.colorScheme.onSurfaceVariant,
)
}
// v12: lastMessagePreview ещё не подключён (отдельная задача).
// Пока показываем placeholder.
Spacer(Modifier.size(2.dp))
Row(verticalAlignment = Alignment.CenterVertically) {
Text(
text = summary.lastMessagePreview ?: "—",
style = MaterialTheme.typography.bodySmall,
color = MaterialTheme.colorScheme.onSurfaceVariant,
maxLines = 1,
overflow = TextOverflow.Ellipsis,
modifier = Modifier.weight(1f),
)
// Unread badge: показывается только при unreadCount > 0.
// В v15 значение всегда известно (server-side count()).
UnreadBadge(unread = summary.unreadCount)
}
}
}
}
}
/**
* Бейдж непрочитанных сообщений. Маленький кружок с цифрой
* (или «99+» если > 99). Прячется при `unread == 0` (всё прочитано,
* либо never-opened — в v15 server-side count() даёт 0 для пустых
* never-opened диалогов).
*/
@Composable
private fun UnreadBadge(unread: Int) {
if (unread <= 0) return
Surface(
shape = CircleShape,
color = MaterialTheme.colorScheme.primary,
modifier = Modifier.size(24.dp),
) {
Box(contentAlignment = Alignment.Center) {
Text(
text = if (unread > 99) "99+" else unread.toString(),
color = MaterialTheme.colorScheme.onPrimary,
style = MaterialTheme.typography.labelSmall,
)
}
}
}
/**
* Форматирует [Instant] как `HH:mm` для сегодняшних дат, иначе `dd.MM`.
* Использует `java.time` (встроен в Android начиная с API 26 — наш `minSdk`).
* Сравнение с "сегодня" — через системный [Clock] (тестируемо).
*
* Конвертация `kotlin.time.Instant` → `java.time.Instant` через
* `toEpochMilliseconds()` + `Instant.ofEpochMilli(...)` — единственный
* способ пересечь границу типов.
*/
private fun formatTimestamp(instant: Instant): String {
val zone = ZoneId.systemDefault()
val javaNow = java.time.Instant.ofEpochMilli(Clock.System.now().toEpochMilliseconds())
val javaInstant = java.time.Instant.ofEpochMilli(instant.toEpochMilliseconds())
val now = LocalDateTime.ofInstant(javaNow, zone).toLocalDate()
val dt = LocalDateTime.ofInstant(javaInstant, zone)
return if (dt.toLocalDate() == now) {
"%02d:%02d".format(dt.hour, dt.minute)
} else {
"%02d.%02d".format(dt.dayOfMonth, dt.monthValue)
}
}
@@ -0,0 +1,204 @@
package pw.binom.agentik.android.ui.screen
import android.util.Log
import androidx.compose.foundation.layout.Arrangement
import androidx.compose.foundation.layout.Column
import androidx.compose.foundation.layout.fillMaxSize
import androidx.compose.foundation.layout.fillMaxWidth
import androidx.compose.foundation.layout.padding
import androidx.compose.foundation.rememberScrollState
import androidx.compose.foundation.verticalScroll
import androidx.compose.material3.Button
import androidx.compose.material3.ExperimentalMaterial3Api
import androidx.compose.material3.LinearProgressIndicator
import androidx.compose.material3.MaterialTheme
import androidx.compose.material3.OutlinedTextField
import androidx.compose.material3.Scaffold
import androidx.compose.material3.Text
import androidx.compose.material3.TopAppBar
import androidx.compose.runtime.Composable
import androidx.compose.runtime.LaunchedEffect
import androidx.compose.runtime.collectAsState
import androidx.compose.runtime.getValue
import androidx.compose.runtime.mutableStateOf
import androidx.compose.runtime.remember
import androidx.compose.runtime.rememberCoroutineScope
import androidx.compose.runtime.setValue
import androidx.compose.ui.Modifier
import androidx.compose.ui.text.input.PasswordVisualTransformation
import androidx.compose.ui.unit.dp
import kotlinx.coroutines.launch
import pw.binom.agentik.android.agent.AgentConnection
import pw.binom.agentik.android.agent.ConnectionState
import pw.binom.agentik.android.settings.Settings
import pw.binom.agentik.android.settings.SettingsRepository
import java.util.UUID
/**
* Экран подключения к агенту. Показывается в трёх случаях:
* - старт приложения ([ConnectionState.Disconnected]);
* - в процессе [ConnectionState.Connecting] — показываем прогресс;
* - при [ConnectionState.Failed] — показываем причину ошибки.
*
* Persistence: настройки читаются из [SettingsRepository] на старте
* (предзаполняют форму), и сохраняются после успешного `configure(...)`
* (когда state становится Connected — см. note в Connect-handler).
* Если файл `settings.json` отсутствует или повреждён, генерируется
* новый `agentId` (UUID) и поля показываются пустыми.
*
* agentId: при первом запуске генерируется как `agentik-android-<UUID>`,
* дальше переиспользуется из persisted settings.
*/
@OptIn(ExperimentalMaterial3Api::class)
@Composable
fun SettingsScreen(
connection: AgentConnection,
settingsRepository: SettingsRepository,
modifier: Modifier = Modifier,
) {
val state by connection.state.collectAsState()
val isConnecting = state is ConnectionState.Connecting
// Дефолты для первого запуска. Если settings.json существует — будут
// перезаписаны в LaunchedEffect ниже.
val initialAgentId = remember { "agentik-android-${UUID.randomUUID()}" }
var agentId by remember { mutableStateOf(initialAgentId) }
var agentName by remember { mutableStateOf("Android") }
var baseUrl by remember { mutableStateOf("") }
var token by remember { mutableStateOf("") }
val scope = rememberCoroutineScope()
// Загрузка persisted настроек. runCatching в repository уже глотает
// любые ошибки (битый JSON, отсутствующий файл) — здесь дополнительной
// обработки не нужно.
LaunchedEffect(Unit) {
val saved = settingsRepository.load() ?: return@LaunchedEffect
agentId = saved.agentId
agentName = saved.agentName
baseUrl = saved.baseUrl
token = saved.token.orEmpty()
}
Scaffold(
modifier = modifier,
topBar = {
TopAppBar(
title = { Text("Подключение к агенту") },
)
},
) { innerPadding ->
Column(
modifier = Modifier
.padding(innerPadding)
.padding(horizontal = 24.dp, vertical = 16.dp)
.fillMaxSize()
.verticalScroll(rememberScrollState()),
verticalArrangement = Arrangement.spacedBy(16.dp),
) {
Text(
text = "Agentik",
style = MaterialTheme.typography.headlineMedium,
)
Text(
text = "Введите адрес агента для подключения. После успешного коннекта откроется чат.",
style = MaterialTheme.typography.bodyMedium,
)
OutlinedTextField(
value = agentId,
onValueChange = { agentId = it },
label = { Text("Agent ID (client)") },
singleLine = true,
modifier = Modifier.fillMaxWidth(),
)
OutlinedTextField(
value = agentName,
onValueChange = { agentName = it },
label = { Text("Имя клиента") },
singleLine = true,
modifier = Modifier.fillMaxWidth(),
)
OutlinedTextField(
value = baseUrl,
onValueChange = { baseUrl = it },
label = { Text("URL агента") },
placeholder = { Text("https://agent.example.com") },
singleLine = true,
modifier = Modifier.fillMaxWidth(),
)
OutlinedTextField(
value = token,
onValueChange = { token = it },
label = { Text("Bearer token (опционально)") },
singleLine = true,
visualTransformation = PasswordVisualTransformation(),
modifier = Modifier.fillMaxWidth(),
)
Button(
onClick = {
scope.launch {
try {
connection.configure(
id = agentId,
name = agentName,
baseUrl = baseUrl,
token = token.takeIf { it.isNotBlank() },
)
// configure() может либо бросить (валидация URL)
// либо вернуть нормально с state = Connecting /
// Failed. Сохраняем настройки в обоих не-бросковых
// случаях — пользователь явно их ввёл, при следующем
// запуске они должны быть видны (особенно если
// первая попытка коннекта провалилась из-за сети).
settingsRepository.save(
Settings(
agentId = agentId,
agentName = agentName,
baseUrl = baseUrl,
token = token.takeIf { it.isNotBlank() },
),
)
} catch (e: Exception) {
// HttpAgentConnection.configure бросает IllegalArgumentException
// ДО выставления state на некорректном URL — логируем,
// пользователь увидит стейт в Disconnected и сможет исправить.
Log.w(TAG, "configure() failed: ${e.message}", e)
}
}
},
enabled = baseUrl.isNotBlank() && !isConnecting,
modifier = Modifier.fillMaxWidth(),
) {
Text(if (isConnecting) "Подключение…" else "Подключиться")
}
// Диагностика под полями: прогресс / ошибка.
when (val s = state) {
is ConnectionState.Connecting -> {
Column(verticalArrangement = Arrangement.spacedBy(8.dp)) {
Text(
text = "Подключение к ${s.agentId}…",
style = MaterialTheme.typography.bodyMedium,
)
LinearProgressIndicator(modifier = Modifier.fillMaxWidth())
}
}
is ConnectionState.Failed -> {
Text(
text = "Ошибка: ${s.reason}",
style = MaterialTheme.typography.bodyMedium,
color = MaterialTheme.colorScheme.error,
)
}
else -> Unit
}
}
}
}
private const val TAG = "SettingsScreen"
+17 -1
View File
@@ -8,7 +8,15 @@ androidx-activity-compose = "1.10.1"
androidx-lifecycle = "2.9.0"
kotlinx-serialization = "1.11.0"
kotlinx-coroutines = "1.11.0"
agentik = "5"
agentik = "15"
journal-api = "15"
journal-inmemory = "15"
journal-ksqlite = "15"
ksqlite = "0.1.2"
asrQwen3 = "4"
micApi = "1.0.0"
work = "2.10.0"
ru-otpbank-markdown = "0.51.0"
android-compileSdk = "35"
android-minSdk = "26"
android-targetSdk = "35"
@@ -31,6 +39,14 @@ kotlinx-coroutines-core = { module = "org.jetbrains.kotlinx:kotlinx-coroutines-c
ktor-client-core = { module = "io.ktor:ktor-client-core", version.ref = "ktor" }
ktor-client-okhttp = { module = "io.ktor:ktor-client-okhttp", version.ref = "ktor" }
agentik-client = { module = "pw.binom.agentik:client", version.ref = "agentik" }
agentik-journal-api = { module = "pw.binom.agentik:journal-api", version.ref = "journal-api" }
agentik-journal-inmemory = { module = "pw.binom.agentik:journal-inmemory", version.ref = "journal-inmemory" }
agentik-journal-ksqlite = { module = "pw.binom.agentik:journal-ksqlite", version.ref = "journal-ksqlite" }
ksqlite = { module = "pw.binom.db:ksqlite", version.ref = "ksqlite" }
asr-qwen3 = { module = "pw.binom.asr:asr-qwen3", version.ref = "asrQwen3" }
mic-api = { module = "pw.binom.mic:mic-api", version.ref = "micApi" }
androidx-work-runtime-ktx = { module = "androidx.work:work-runtime-ktx", version.ref = "work" }
ru-otpbank-markdown = { module = "ru.otpbank.ai:markdown", version.ref = "ru-otpbank-markdown" }
[plugins]
android-application = { id = "com.android.application", version.ref = "agp" }
+13
View File
@@ -23,6 +23,19 @@ dependencyResolutionManagement {
// Nexus отдаёт по plain HTTP — без флага Gradle откажется.
isAllowInsecureProtocol = true
}
// Хост для ru.otpbank.ai:markdown (парсер Markdown из дизайна BORROW.md, раздел 3).
maven {
name = "developerspace-prod-mvn-hosted"
url = uri("http://192.168.76.117/repository/developerspace-prod-mvn-hosted/")
// Nexus отдаёт по plain HTTP — без флага Gradle откажется.
isAllowInsecureProtocol = true
}
// JitPack хостит onnx-asr/sherpa-onnx и другие GitHub-проекты,
// которые тянет asr-qwen3-android транзитивно.
maven {
name = "jitpack"
url = uri("https://jitpack.io")
}
}
}