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.Dispatchers 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.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.ShowText import pw.binom.viewmate.core.protocol.SttCancel import pw.binom.viewmate.core.protocol.SttPhrase 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.JellyfinMediaSource import pw.binom.viewmate.phone.agent.LlmClient import pw.binom.viewmate.phone.agent.MEDIA_TOOLSET import pw.binom.viewmate.phone.agent.MediaToolExecutor 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.SttStreamer import pw.binom.viewmate.phone.stt.WhisperStt /** * 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) } /** История общения с ассистентом (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 /** Ленивый ассистент (паттерн ensureStt): один экземпляр на процесс, не на каждую фразу. */ private fun ensureAssistant(): Assistant? { synchronized(assistantLock) { assistant?.let { return it } if (BuildConfig.LLM_API_KEY.isBlank()) return null val llm = LlmClient(apiKey = BuildConfig.LLM_API_KEY, model = BuildConfig.LLM_MODEL) 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, system -> llm.chat(history, system) }, toolkit = tk, ).also { assistant = it } } } /** Тулсет 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-ключ не задан — фраза пропущена") 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(AssistantStateMsg(AssistantState.IDLE)) } } /** * Отправка текста в ассистента из чата на телефоне (не с очков). * Та же очередь, что и голосовые фразы — сериализация 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) } /** * Лениво создаёт стриминговое распознавание (модели во filesDir/models — * кладутся через adb push + run-as cp, как в этапе 1). Вызывается из хаба * на первый SttAudio; тяжёлое создание Whisper — один раз на процесс. */ private fun ensureStt(): SttStreamer? { synchronized(sttLock) { sttStreamer?.let { return it } val modelsDir = File(filesDir, "models") val vadModel = File(modelsDir, "silero_vad.onnx") val encoder = File(modelsDir, "small-encoder.int8.onnx") val decoder = File(modelsDir, "small-decoder.int8.onnx") val tokens = File(modelsDir, "small-tokens.txt") if (!vadModel.isFile || !encoder.isFile || !decoder.isFile || !tokens.isFile) { log("stt", "модели не найдены в ${modelsDir.absolutePath} — STT недоступен") return null } val whisper = WhisperStt( encoderPath = encoder.absolutePath, decoderPath = decoder.absolutePath, tokensPath = tokens.absolutePath, numThreads = 4, ) val s = SttStreamer( stt = whisper, vadModelPath = vadModel.absolutePath, onPhrase = { phrase, full -> log("stt", "фраза: $phrase") scope.launch { server.hub.broadcast(SttPhrase(phrase, full)) } scope.launch { assistantChannel.send(phrase) } }, onSilence30s = { log("stt", "тишина 30 с — автo-отмена") scope.launch { server.hub.broadcast(SttCancel()) } }, ) sttStreamer = s log("stt", "SttStreamer готов (Whisper-small int8, VAD fp32)") 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)") server.start() log("server", "WS-сервер запущен на ${PhoneConfig.SERVER_PORT}") nsd.publish() server.hub.sttFactory = { ensureStt() } scope.launch { assistantLoop() } val lastPositionAt = AtomicLong(0) val lastWasPlaying = AtomicBoolean(false) // каждый PlaybackPosition от очков (5с) → синхронизация аудио-плеера // + автоподхват звука (очки играют после рестарта телефона, audioSync не активен) server.hub.onPlaybackPosition = { positionMs, playing, itemId, audioIndex -> lastPositionAt.set(System.currentTimeMillis()) lastWasPlaying.set(playing) audioSync.sync(positionMs, playing) 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 } }