package pw.binom.viewmate.phone import android.app.Application import java.io.File import java.util.concurrent.atomic.AtomicBoolean import java.util.concurrent.atomic.AtomicLong import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.CancellationException import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.Job import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.channels.Channel import kotlinx.coroutines.flow.MutableStateFlow import kotlinx.coroutines.flow.StateFlow import kotlinx.coroutines.flow.asStateFlow import kotlinx.coroutines.delay import kotlinx.coroutines.flow.collect import kotlinx.coroutines.launch import pw.binom.viewmate.core.GlassesMode import pw.binom.viewmate.core.log.BatchingLogCollector import pw.binom.viewmate.core.log.FileLogSpool import pw.binom.viewmate.core.log.LokiConfig import pw.binom.viewmate.core.log.SystemClock import pw.binom.viewmate.core.log.lokiLogSink import pw.binom.viewmate.core.media.JellyfinClient import pw.binom.viewmate.core.media.MirrorClient import pw.binom.viewmate.core.phone.PhoneActions import pw.binom.viewmate.core.phone.PhoneState import pw.binom.viewmate.core.protocol.AssistantState import pw.binom.viewmate.core.protocol.AssistantStateMsg import pw.binom.viewmate.core.protocol.ChatHistoryItem import pw.binom.viewmate.core.protocol.ChatHistoryMsg import pw.binom.viewmate.core.protocol.ShowText import pw.binom.viewmate.core.protocol.SttCancel import pw.binom.viewmate.core.protocol.SttDone import pw.binom.viewmate.phone.agent.AgentToolkit import pw.binom.viewmate.phone.agent.Assistant import pw.binom.viewmate.phone.agent.ChatDao import pw.binom.viewmate.phone.agent.ChatDb import pw.binom.viewmate.phone.agent.ChatRole import pw.binom.viewmate.phone.agent.Gemma4E2B import pw.binom.viewmate.phone.agent.GemmaModelDownloader import pw.binom.viewmate.phone.agent.JellyfinMediaSource import pw.binom.viewmate.phone.agent.LocalLlmClient import pw.binom.viewmate.phone.agent.LocalLlmModelState import pw.binom.viewmate.phone.agent.LlmClient import pw.binom.viewmate.phone.agent.LlmModelChoice import pw.binom.viewmate.phone.agent.LlmPrefs import pw.binom.viewmate.phone.agent.MEDIA_TOOLSET import pw.binom.viewmate.phone.agent.MediaToolExecutor import pw.binom.viewmate.phone.agent.RemoteLlmClient import pw.binom.viewmate.phone.agent.SkillRegistry import pw.binom.viewmate.phone.agent.SkillRepository import pw.binom.viewmate.phone.agent.ToolsetRegistry import pw.binom.viewmate.phone.stt.PhraseRecognizer import pw.binom.viewmate.phone.stt.Qwen3AsrStt import pw.binom.viewmate.phone.stt.SttStreamer /** * Application телефона: создаёт клиентов (Jellyfin/mirror), PhoneState, * WS-сервер (GlassesServer), DownloadManager и PhoneActions. Синглтон-доступ через companion. */ class PhoneApp : Application() { val state: PhoneState = PhoneState() /** Сервер связи очков: транспорт по [PhoneConfig.USE_BLUETOOTH] — дефолт WiFi (WS), Bluetooth SPP — флагом. */ val server: GlassesServer = GlassesServer( port = PhoneConfig.SERVER_PORT, state = state, useBluetooth = PhoneConfig.USE_BLUETOOTH, btContext = this, ) /** Публикация сервера в локальной сети (mDNS `_viewmate._tcp`) для авто-обнаружения очками. */ val nsd: NsdPublisher by lazy { NsdPublisher(this) } /** HTTP-клиенты бэкендов (Jellyfin-каталог, media-mirror). */ val jellyfin: JellyfinClient = JellyfinClient(PhoneConfig.JELLYFIN_URL, PhoneConfig.JELLYFIN_API_KEY) val mirror: MirrorClient = MirrorClient(PhoneConfig.MIRROR_URL, PhoneConfig.MIRROR_API_KEY) private val _userId = MutableStateFlow("") /** Первый пользователь Jellyfin (подтягивается один раз при старте). */ val userId: StateFlow = _userId.asStateFlow() /** Реальное скачивание зеркал на телефон (с прогрессом). */ val downloadManager: DownloadManager by lazy { DownloadManager(this, mirror) } /** Синхронизированный аудио-плеер (звук на телефоне при просмотре на очках). */ val audioSync: AudioSyncPlayer by lazy { AudioSyncPlayer(this) } /** Лог-коллектор: свои записи + пришедшие от очков → сервер (спул — filesDir/logspool). */ lateinit var logCollector: BatchingLogCollector private set /** История общения с ассистентом (SQLite, "assistant.db"). */ val chatDao: ChatDao by lazy { ChatDao(ChatDb(this)) } /** Автоматика паузы/продолжения звука при обрыве связи. */ private val connectionAutoResume = ConnectionAutoResume() var actions: PhoneActions? = null private set private val scope = CoroutineScope(SupervisorJob() + Dispatchers.IO) /** itemId, у которого автоподхват звука уже не удался — не дёргать повторно каждые 5с. */ private var lastFailedItemId: String? = null private var sttLock = Any() private var sttStreamer: SttStreamer? = null /** * Очередь распознанных фраз для ассистента. UNLIMITED — `send` неблокирующий, * STT-поток не блокируется. Единственный consumer — [assistantLoop] (сериализация * LLM-вызовов: фразы обрабатываются по одной). */ private val assistantChannel = Channel(Channel.UNLIMITED) private var assistantLock = Any() private var assistant: Assistant? = null private var toolkit: AgentToolkit? = null /** Модель, на которой создан текущий [assistant] — смена выбора пересоздаёт ассистента. */ private var assistantChoice: LlmModelChoice? = null /** Активный локальный LLM-клиент (движок закрываем при смене выбора). */ private var localLlm: LocalLlmClient? = null /** Текущий локальный LLM-клиент (для DebugHttpServer / статуса). */ fun localLlmClient(): LocalLlmClient? = localLlm /** * LLM-клиент по текущему выбору [llmChoice] для OpenAI-HTTP-сервера * ([OpenAiHttpServer]) — без тул-цикла ассистента (клиент OpenAI API сам * управляет контекстом). Локальный клиент переиспользует [localLlm] * (движок LiteRT один на процесс). null — серверный без ключа или локальная * модель ещё не скачана. */ fun openAiLlm(): LlmClient? = synchronized(assistantLock) { when (_llmChoice.value) { LlmModelChoice.REMOTE -> if (BuildConfig.LLM_API_KEY.isBlank()) null else RemoteLlmClient(apiKey = BuildConfig.LLM_API_KEY, model = BuildConfig.LLM_MODEL) LlmModelChoice.LOCAL -> if (!localModelFile().isFile) null else localLlm ?: LocalLlmClient( modelFile = localModelFile(), cacheDir = File(cacheDir, "litertlm"), ).also { localLlm = it } } } /** Выбор модели LLM (персистится, [pw.binom.viewmate.phone.agent.LlmPrefs]). */ // lazy: конструктор Application выполняется ДО attachBaseContext — контекст (и // getSharedPreferences/getExternalFilesDir) доступен только после; ленивость — // тот же паттерн, что у chatDao/downloadManager (иначе NPE при старте). val llmPrefs: LlmPrefs by lazy { LlmPrefs(this) } private val _llmChoice: MutableStateFlow by lazy { MutableStateFlow(llmPrefs.current()) } val llmChoice: StateFlow by lazy { _llmChoice.asStateFlow() } /** Состояние локальной модели Gemma 4 E2B для UI настроек. */ private val _localLlmState: MutableStateFlow by lazy { MutableStateFlow(initialLocalLlmState()) } val localLlmState: StateFlow by lazy { _localLlmState.asStateFlow() } private var localLlmDownloadJob: Job? = null private fun initialLocalLlmState(): LocalLlmModelState = if (localModelFile().isFile) LocalLlmModelState.Ready else LocalLlmModelState.Idle /** Файл локальной модели: getExternalFilesDir(null)/gemma-4-E2B-it.litertlm. */ fun localModelFile(): File = File(getExternalFilesDir(null) ?: filesDir, Gemma4E2B.FILE_NAME) /** Сменить выбор модели: персист + пересоздание ассистента при следующей фразе. */ fun setLlmModelChoice(choice: LlmModelChoice) { if (_llmChoice.value == choice) return llmPrefs.set(choice) _llmChoice.value = choice synchronized(assistantLock) { assistant = null } log("llm", "выбрана модель: ${choice.name}") } /** * Скачать локальную модель Gemma 4 E2B (файл в getExternalFilesDir(null)), * прогресс в процентах в [localLlmState], контрольная сумма SHA-256. * Повторный вызов во время скачивания — no-op. */ fun downloadLocalLlmModel() { val current = localLlmDownloadJob if (current != null && current.isActive) return localLlmDownloadJob = scope.launch { val target = localModelFile() _localLlmState.value = LocalLlmModelState.Downloading(0) try { GemmaModelDownloader(target).download { percent -> _localLlmState.value = LocalLlmModelState.Downloading(percent) } _localLlmState.value = LocalLlmModelState.Ready log("llm", "локальная модель скачана: ${target.name}, ${target.length()} байт") } catch (e: CancellationException) { throw e } catch (e: Exception) { _localLlmState.value = LocalLlmModelState.Error(e.message ?: "неизвестная ошибка") log("llm", "не удалось скачать модель: ${e.message}") } } } /** Остановить скачивание локальной модели (частичный файл удаляется). */ fun cancelLocalLlmDownload() { localLlmDownloadJob?.cancel() } /** * Ленивый ассистент (паттерн ensureStt): один экземпляр на процесс, не на каждую фразу. * LLM-клиент берётся по текущему выбору [llmChoice]: серверный или локальный * LiteRT-LM; при смене выбора ассистент создаётся заново. */ private fun ensureAssistant(): Assistant? { synchronized(assistantLock) { val choice = _llmChoice.value if (assistant != null) { if (assistantChoice == choice) return assistant assistant = null } if (choice == LlmModelChoice.REMOTE && BuildConfig.LLM_API_KEY.isBlank()) { assistantChoice = choice return null } runCatching { localLlm?.close() } localLlm = null val llm = when (choice) { LlmModelChoice.REMOTE -> RemoteLlmClient(apiKey = BuildConfig.LLM_API_KEY, model = BuildConfig.LLM_MODEL) LlmModelChoice.LOCAL -> { // cacheDir — подкаталог cache/litertlm: в 11:46 именно с ним // (без спекуляции) движок инициализировался за 12с; с корнем // cache/ initialize() виснет. GPU-делегату нужен СВОЙ каталог. val cacheDir = File(cacheDir, "litertlm") LocalLlmClient(modelFile = localModelFile(), cacheDir = cacheDir).also { localLlm = it } } } val toolsets = ToolsetRegistry() toolsets.register(MEDIA_TOOLSET, mediaToolExecutor()) val tk = toolkit ?: AgentToolkit( skills = SkillRegistry(SkillRepository(File(filesDir, "skills")).loadAll()), toolsets = toolsets, showText = { text -> scope.launch { server.hub.broadcast(ShowText(text)) } }, ).also { toolkit = it } return Assistant( chatDao = chatDao, chat = { history, toolSpecs, system -> llm.chatWithTools(history, toolSpecs, system) }, toolkit = tk, // Нативный тул-канал LiteRT: только локальная модель видит схемы // тул у движка; системный промпт тогда без текстовых TOOL_CALL-линий. nativeTools = choice == LlmModelChoice.LOCAL, // Локальной модели нужно время: инициализация движка ~17с на // CPU + генерация. Серверная укладывается в 30с, локальной // даём 5 минут (иначе withTimeoutOrNull режет генерацию). timeoutMs = if (choice == LlmModelChoice.LOCAL) 300_000L else 30_000L, ).also { assistant = it assistantChoice = choice } } } /** Тулсет media: тулы работают через [JellyfinMediaSource] (каталог) и уровня-приложения [agentPlayMedia] («включить на очках»). */ private fun mediaToolExecutor(): MediaToolExecutor = MediaToolExecutor( media = JellyfinMediaSource(jellyfin) { userId.value }, player = { id, seekMs -> agentPlayMedia(id, seekMs) }, ) /** * Тул ассистента media_play: проверки (очки на связи, видео скачано на очки, * звук скачан на телефон) + запуск через [PhoneActions.watch] (механизм DetailsScreen) * + звук на телефоне выбранной дорожкой с позиции [seekMs]. Возвращает текст для модели. */ private suspend fun agentPlayMedia(id: String, seekMs: Long): String { val hub = server.hub if (hub.connected.value == 0) return "Очки не подключены." val mediaDir = hub.mediaDir.value if (mediaDir.isNullOrBlank()) return "Очки ещё не опубликовали медиакаталог — попробуйте позже." if (id !in hub.downloadedItemIds.value) return "Видео не скачано на очки. Скачайте его вручную через телефон." if (!downloadManager.audioReady(id)) return "Звук не готов: сначала скачайте аудио на телефон." val acts = actions ?: return "Телефон ещё не готов." val audioIndex = state.currentAudioIndex log("agent", "play: $id, звук #$audioIndex, seek=${seekMs} мс") val result = acts.watch( itemId = id, audioIndex = audioIndex, localVideoPath = "file://$mediaDir/$id/video.mkv", phoneAudioReady = true, startPositionMs = seekMs, ) if (!result.contains("включил")) return result // Звук на телефоне (с позиции, как DetailsScreen): локальный ogg или S3-стрим. if (audioSync.isActive) audioSync.stop() val audio = resolveAudioSourceDetailed( mirrorFiles = mirror.mirrorByItem(id)?.files, audioIndex = audioIndex, localAudioFile = { downloadManager.audioFile(id, it) }, ) if (audio.source == null) return "Видео включено, но источники звука отсутствуют: ${audio.reason}" audioSync.play(audio.source, seekMs.coerceAtLeast(0L)) return "Включено: «${state.currentItem?.Name ?: id}»" } /** * Единственный consumer фраз: фраза → LLM (с историей из SQLite) → ShowText на очки. * Ошибка LLM (нет сети, 401 и т.п.) — коротким сообщением пользователю, не падаем. */ private suspend fun assistantLoop() { for (phrase in assistantChannel) { log("assistant", "получена фраза: $phrase") val assistant = ensureAssistant() if (assistant == null) { log("assistant", "серверный LLM без ключа (LLM_API_KEY) — фраза пропущена") continue } server.hub.broadcast(AssistantStateMsg(AssistantState.THINKING)) val answer = runCatching { val active = _activeSessionId.value if (active != null) { assistant.processInSession(active, phrase) } else { assistant.process(phrase) } } .onFailure { log("assistant", "ошибка: ${it.message}") } .getOrElse { "Ассистент недоступен: ${it.message}" } server.hub.broadcast(ShowText(answer)) // Очки показывают историю всей сессии (свайпами); ответ уже в ней. server.hub.broadcast(chatHistoryMsg()) server.hub.broadcast(AssistantStateMsg(AssistantState.IDLE)) } } /** * История активной (или последней) сессии для очков: только роли user/assistant, * tool-сообщения и toolCalls очкам не нужны (ТЗ-касяк №1). */ private fun chatHistoryMsg(): ChatHistoryMsg { val id = _activeSessionId.value ?: chatDao.latestSessionId() if (id == null) return ChatHistoryMsg(sessionId = null, messages = emptyList()) val items = chatDao.messagesForSession(id) .filter { it.role == ChatRole.USER || it.role == ChatRole.ASSISTANT } .map { ChatHistoryItem(role = it.role.wire, content = it.content) } return ChatHistoryMsg(sessionId = id.toString(), messages = items) } /** * Отправка текста в ассистента из чата на телефоне (не с очков). * Та же очередь, что и голосовые фразы — сериализация LLM сохраняется. */ fun sendToAssistant(text: String) { scope.launch { assistantChannel.send(text) } } /** Активная сессия чата (для UI телефона). null — «последняя» (дефолт). */ private val _activeSessionId = MutableStateFlow(null) val activeSessionId: StateFlow = _activeSessionId /** Создать новую сессию и сделать её активной. */ fun newChatSession(): Long { val id = chatDao.createSession() _activeSessionId.value = id return id } /** Переключить активную сессию. null — вернуться к «последней». */ fun selectChatSession(id: Long?) { _activeSessionId.value = id } /** Сбросить (очистить сообщения) активной или последней сессии. */ fun resetChatSession() { val id = _activeSessionId.value ?: chatDao.latestSessionId() ?: return chatDao.clearSession(id) } /** * Лениво создаёт стриминговое распознавание (модели Qwen3-ASR во filesDir/models — * кладутся через adb push + run-as cp, как в этапе 1). Вызывается из хаба * на первый SttAudio; тяжёлое создание Qwen3-ASR — один раз на процесс. */ private fun ensureStt(): SttStreamer? { synchronized(sttLock) { sttStreamer?.let { return it } val modelsDir = File(filesDir, "models") val qwenConv = File(modelsDir, "conv_frontend.onnx") val qwenEnc = File(modelsDir, "encoder.int8.onnx") val qwenDec = File(modelsDir, "decoder.int8.onnx") val qwenTok = File(modelsDir, "tokenizer") if (!qwenConv.isFile || !qwenEnc.isFile || !qwenDec.isFile || !qwenTok.isDirectory) { log("stt", "модели Qwen3-ASR не найдены в ${modelsDir.absolutePath} — STT недоступен") return null } val stt: PhraseRecognizer = Qwen3AsrStt( convFrontendPath = qwenConv.absolutePath, encoderPath = qwenEnc.absolutePath, decoderPath = qwenDec.absolutePath, tokenizerDir = qwenTok.absolutePath, numThreads = 4, hotwords = STT_HOTWORDS, ) val dumpPcm = if (File(filesDir, "stt_dump.marker").exists()) { File(filesDir, "stt_dump_phone.raw") } else { null } log("stt", "бэкенд: qwen3_asr (hotwords: $STT_HOTWORDS)") val s = SttStreamer( stt = stt, onSilence30s = { log("stt", "тишина 30 с — автo-отмена") scope.launch { server.hub.broadcast(SttCancel(reason = "timeout")) } }, onAutoFinished = { full -> log("stt", "тишина 30 с — ВЕСЬ ТЕКСТ: $full") // Тот же путь, что и результат по клику: в очки SttDone + в LLM. scope.launch { server.hub.broadcast(SttDone(full)) } scope.launch { assistantChannel.send(full.trim()) } }, dumpPcm = dumpPcm, ) sttStreamer = s log("stt", "PCM-дампа (телефон): ${dumpPcm?.absolutePath ?: "off"}") log("stt", "SttStreamer готов (qwen3_asr)") return s } } /** * Автоподхват звука после рестарта телефона: очки продолжают играть и шлют * PlaybackPosition с itemId+audioIndex, а audioSync не активен (телефон «не в курсе»). * Создаём аудио-плеер с позиции очков — как в DetailsScreen.onWatch. */ private suspend fun autostartAudio(itemId: String?, audioIndex: Int, positionMs: Long, playing: Boolean) { if (itemId == null || audioSync.isActive || itemId == lastFailedItemId) return val mirrorFiles = runCatching { mirror.mirrorByItem(itemId) }.getOrNull()?.files val source = resolveAudioSource(downloadManager, itemId, mirrorFiles, audioIndex) if (source != null) { lastFailedItemId = null audioSync.play(source, positionMs) log("audio", "автоподхват звука: itemId=$itemId idx=$audioIndex pos=$positionMs") val item = runCatching { jellyfin.item(userId.value, itemId) }.getOrNull() if (item != null) { state.currentItem = item state.currentAudioIndex = audioIndex state.mode = GlassesMode.MOVIE state.playing = playing } } else { log("audio", "автоподхват звука не удался: itemId=$itemId idx=$audioIndex (нет источника)") lastFailedItemId = itemId } } override fun onCreate() { super.onCreate() instance = this PlaybackService.start(this) log("server", "playback-сервис запущен (foreground, WAKE_LOCK)") MercuryBridge.init(this) server.start() log("server", "WS-сервер запущен на ${PhoneConfig.SERVER_PORT}") DebugHttpServer.start(this) OpenAiHttpServer.start(this) nsd.publish() server.hub.sttFactory = { ensureStt() } server.hub.onStopFullText = { phrase -> scope.launch { assistantChannel.send(phrase) } } // Лог-коллектор: свои записи + пачки от очков → Loki. instance в Loki // выставляется по source записи ("phone"/"glasses"), спул — filesDir/logspool. logCollector = BatchingLogCollector( device = "phone", sink = lokiLogSink(), spool = FileLogSpool(File(filesDir, "logspool").absolutePath, 64), clock = SystemClock(), flushIntervalMs = 60_000, maxBatchEntries = 500, ) logCollector.start(scope) attachLogCollector(logCollector) pw.binom.viewmate.core.logCollector = logCollector server.hub.onLogBatch = { entries -> logCollector.recordBatch(entries) } log( "server", "лог-коллектор запущен (sink: loki ${LokiConfig().url}, instance по source; " + "спул: filesDir/logspool)", ) scope.launch { assistantLoop() } val lastPositionAt = AtomicLong(0) val lastWasPlaying = AtomicBoolean(false) // каждый PlaybackPosition от очков (5с) → синхронизация аудио-плеера // + автоподхват звука (очки играют после рестарта телефона, audioSync не активен) server.hub.onPlaybackPosition = { positionMs, playing, itemId, audioIndex, sentTs -> lastPositionAt.set(System.currentTimeMillis()) lastWasPlaying.set(playing) audioSync.sync(positionMs, playing, sentTs) scope.launch { autostartAudio(itemId, audioIndex, positionMs, playing) } } // watchdog: очки молчат (экран погас/радио спит) — пауза звука по таймауту scope.launch { while (true) { delay(PhoneConfig.WATCHDOG_TICK_MS) val silentFor = System.currentTimeMillis() - lastPositionAt.get() if (silentFor > PhoneConfig.AUDIO_SILENCE_PAUSE_MS && audioSync.isActive && lastWasPlaying.get()) { log("audio", "очки молчат ${silentFor}мс — пауза (экран погас?)") audioSync.pause() lastWasPlaying.set(false) // не спамить каждые 2с } } } // автоматика обрыва связи: потеря → пауза, восстановление → продолжить с того же места scope.launch(Dispatchers.Main) { var prevConnected = server.hub.connected.value > 0 server.hub.connected.collect { count -> val nowConnected = count > 0 when (connectionEvent(prevConnected, nowConnected)) { ConnectionEvent.LOST -> { connectionAutoResume.onLost(audioSync.isActive, state.playing) if (audioSync.isActive) { log("audio", "связь потеряна — пауза") audioSync.pause() } } ConnectionEvent.RESTORED -> { if (connectionAutoResume.consumeResume()) { log("audio", "связь восстановлена — продолжаю") audioSync.play() } } ConnectionEvent.NONE -> Unit } prevConnected = nowConnected } } scope.launch { val userId = runCatching { jellyfin.users().firstOrNull()?.Id }.getOrNull() ?: "" _userId.value = userId if (userId.isBlank()) log("app", "не удалось получить пользователя Jellyfin") actions = PhoneActions( jellyfin = jellyfin, mirror = mirror, sender = server.sender, state = state, userId = userId, ) log("app", "PhoneActions готов (userId=$userId)") } } override fun onTerminate() { PlaybackService.stop(this) nsd.unpublish() sttStreamer?.close() server.stop() super.onTerminate() } companion object { lateinit var instance: PhoneApp private set /** Подсказки словаря Qwen3-ASR (запятая-разделитель). */ private const val STT_HOTWORDS = "тулсет,тулы,умеешь" } }