package pw.binom.viewmate.phone import android.content.Context import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.Job import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.cancel import kotlinx.coroutines.flow.MutableStateFlow import kotlinx.coroutines.flow.StateFlow import kotlinx.coroutines.flow.asStateFlow import kotlinx.coroutines.launch import kotlinx.serialization.decodeFromString import kotlinx.serialization.encodeToString import pw.binom.viewmate.core.log.LogEntry import pw.binom.viewmate.core.net.GlassesServerTransport import pw.binom.viewmate.core.phone.GlassesSender import pw.binom.viewmate.core.phone.PhoneState import pw.binom.viewmate.core.protocol.DownloadProgress import pw.binom.viewmate.core.protocol.GlassesOff import pw.binom.viewmate.core.protocol.GlassesStatus import pw.binom.viewmate.core.protocol.GlassesToHost import pw.binom.viewmate.core.protocol.Gesture import pw.binom.viewmate.core.protocol.Hello import pw.binom.viewmate.core.protocol.HostToGlasses import pw.binom.viewmate.core.protocol.LogBatchMsg import pw.binom.viewmate.core.protocol.PlaybackPosition import pw.binom.viewmate.core.protocol.SttAudio import pw.binom.viewmate.core.protocol.SttCancel import pw.binom.viewmate.core.protocol.SttDone import pw.binom.viewmate.core.protocol.StopStt import pw.binom.viewmate.core.protocol.Welcome import pw.binom.viewmate.core.protocol.protocolJson import pw.binom.viewmate.phone.stt.SttStreamer import java.util.concurrent.ConcurrentHashMap const val GLASSES_WS_PATH = "/ws/glasses" /** * Менеджер подключённых очков (телефон). Потокобезопасен (ConcurrentHashMap). * Сессии адресуются **connId-строками** (UUID — WiFi, MAC-адрес — Bluetooth), * не WS-объектами (TASK-transport.md п.1): WS-мост — в [WifiServerTransport], * SPP-сокеты — в [BtServerTransport]. Исходящие — JSON-строкой через * [GlassesServerTransport.send] (декод протокола — в хабe, не в транспорте). */ class GlassesHub( private val state: PhoneState, /** Выбранный транспорт (дефолт — [WifiServerTransport]). */ private val transport: GlassesServerTransport, ) { /** Сессии: connId → Unit. */ private val sessions = ConcurrentHashMap() private val _connected = MutableStateFlow(0) val connected: StateFlow = _connected.asStateFlow() private val _gesturesReceived = MutableStateFlow(0) val gesturesReceived: StateFlow = _gesturesReceived.asStateFlow() private val _batteryPercent = MutableStateFlow(null) val batteryPercent: StateFlow = _batteryPercent.asStateFlow() private val _storageUsedGb = MutableStateFlow(null) val storageUsedGb: StateFlow = _storageUsedGb.asStateFlow() private val _storageTotalGb = MutableStateFlow(null) val storageTotalGb: StateFlow = _storageTotalGb.asStateFlow() private val _storageTotalBytes = MutableStateFlow(null) val storageTotalBytes: StateFlow = _storageTotalBytes.asStateFlow() private val _storageFreeBytes = MutableStateFlow(null) val storageFreeBytes: StateFlow = _storageFreeBytes.asStateFlow() private val _mediaBytes = MutableStateFlow(null) val mediaBytes: StateFlow = _mediaBytes.asStateFlow() private val _appVersion = MutableStateFlow(null) val appVersion: StateFlow = _appVersion.asStateFlow() private val _mediaDir = MutableStateFlow(null) val mediaDir: StateFlow = _mediaDir.asStateFlow() private val _downloadedItemIds = MutableStateFlow>(emptySet()) /** itemId полностью скачанных на очки (из периодического GlassesStatus). */ val downloadedItemIds: StateFlow> = _downloadedItemIds.asStateFlow() private val _glassesDownloads = MutableStateFlow>>(emptyMap()) /** Прогресс скачивания на очки: itemId → имя файла → последний DownloadProgress от очков. */ val glassesDownloads: StateFlow>> = _glassesDownloads.asStateFlow() private val _positionMs = MutableStateFlow(null) val positionMs: StateFlow = _positionMs.asStateFlow() /** * Хук на позицию очков: вызывается на каждый PlaybackPosition (и GlassesOff * с playing=false). [itemId]+[audioIndex] — для автоподхвата звука после * рестарта телефона (см. PhoneApp.autostartAudio). */ var onPlaybackPosition: ((positionMs: Long, playing: Boolean, itemId: String?, audioIndex: Int, sentTs: Long, durationMs: Long) -> Unit)? = null /** * Хук на лог-пачки очков: телефон прибавляет записи в свой лог-коллектор * (BatchingLogCollector.recordBatch) — они уйдут на сервер вместе со своими. */ @Volatile var onLogBatch: ((List) -> Unit)? = null /** * Стриминговое распознавание речи с очков. Лениво создаётся через * [sttFactory] (модели во filesDir/models); onSilence30s/onAutoFinished * рассылаются через [broadcast]. */ @Volatile var stt: SttStreamer? = null /** * Единственный путь текста в LLM: вызывает PhoneApp (assistantChannel.send) * на StopStt(cancel=false) с непустым текстом. Промежуточные VAD-фразы * отсутствуют — только полный текст (по клику или по тишине 30 с). */ @Volatile var onStopFullText: ((String) -> Unit)? = null /** Фабрика SttStreamer на первый SttAudio (тяжёлая — создаётся один раз). */ @Volatile var sttFactory: (() -> SttStreamer?)? = null fun add(connId: String) { sessions[connId] = Unit _connected.value = sessions.size log("hub", "очки подключились (connId $connId, всего: ${sessions.size})") } fun remove(connId: String) { sessions.remove(connId) _connected.value = sessions.size log("hub", "очки отключились (connId $connId, всего: ${sessions.size})") } /** Отправить сообщение всем подключённым очкам (JSON-строкой через [transport]). */ suspend fun broadcast(msg: HostToGlasses) { val text = protocolJson.encodeToString(HostToGlasses.serializer(), msg) for (connId in sessions.keys) { runCatching { transport.send(connId, text) } .onFailure { log("hub", "не удалось отправить ${msg::class.simpleName} на $connId: ${it.message}") } } } /** Обработать входящее текстовое сообщение от очков (JSON-строкой из [transport]). */ suspend fun handle(connId: String, text: String) { val msg = try { protocolJson.decodeFromString(GlassesToHost.serializer(), text) } catch (e: Exception) { log("hub", "ОШИБКА декодирования входящего: ${e::class.simpleName}: ${e.message} (текст: ${text.take(120)})") return } when (msg) { is Hello -> { log("glasses", "hello (app=${msg.appVersion}, connId $connId) → welcome") send(connId, Welcome(mode = state.mode, sessionId = state.activeSessionId)) } is Gesture -> { _gesturesReceived.value++ log("glasses", "жест с тачпада: ${msg.gesture}") } is GlassesStatus -> { _batteryPercent.value = msg.batteryPercent _storageUsedGb.value = msg.storageUsedGb _storageTotalGb.value = msg.storageTotalGb if (msg.storageTotalBytes > 0) _storageTotalBytes.value = msg.storageTotalBytes if (msg.storageFreeBytes > 0) _storageFreeBytes.value = msg.storageFreeBytes if (msg.mediaBytes > 0) _mediaBytes.value = msg.mediaBytes if (msg.mediaDir.isNotBlank()) _mediaDir.value = msg.mediaDir if (msg.appVersion.isNotBlank()) _appVersion.value = msg.appVersion _downloadedItemIds.value = msg.downloadedItemIds.toSet() log( "glasses", "статус: батарея ${msg.batteryPercent}%, память ${msg.storageUsedGb}/${msg.storageTotalGb} ГБ, " + "медиа ${msg.mediaBytes} байт, скачано: ${msg.downloadedItemIds.size}, app=${msg.appVersion}", ) } is DownloadProgress -> { val prev = _glassesDownloads.value[msg.itemId] ?: emptyMap() val updated = _glassesDownloads.value + (msg.itemId to (prev + (msg.fileName to msg))) _glassesDownloads.value = updated log("glasses", "скачивание на очки: ${msg.itemId} ${msg.fileName} ${msg.percent}% (${msg.phase})") } is PlaybackPosition -> { _positionMs.value = msg.positionMs this.state.playing = msg.playing val state = if (msg.playing) "играет" else "пауза" log("glasses", "позиция: ${msg.positionMs} мс, $state, item=${msg.itemId} idx=${msg.audioIndex}") onPlaybackPosition?.invoke(msg.positionMs, msg.playing, msg.itemId, msg.audioIndex, msg.sentTs, msg.durationMs) } is GlassesOff -> { log("glasses", "очки выключились → пауза (${msg.reason})") this.state.playing = false onPlaybackPosition?.invoke(_positionMs.value ?: 0L, false, null, 0, 0L, -1L) } is LogBatchMsg -> { log("glasses", "лог-пачка с очков: ${msg.entries.size} записей") onLogBatch?.invoke(msg.entries) } is SttAudio -> { val s = stt ?: sttFactory?.invoke()?.also { stt = it } s?.accept(msg.data) } is StopStt -> { if (msg.cancel) { log("stt", "распознавание отменено (очки)") send(connId, SttCancel(reason = SttCancel.REASON_CANCEL)) stt?.reset() } else { val full = stt?.finish() ?: "" log("stt", "ВЕСЬ ТЕКСТ: $full") stt?.reset() if (full.isNotBlank()) { send(connId, SttDone(full)) val phrase = full.trim() log("stt", "фраза в LLM (по клику): $phrase") onStopFullText?.invoke(phrase) } else { // Фразы не было — очки возвращаются в «слушаю» (SttCancel). log("stt", "фразы не было — очки в «слушаю»") broadcast(SttCancel(reason = SttCancel.REASON_USER)) } } } } } private suspend fun send(connId: String, msg: HostToGlasses) { val text = protocolJson.encodeToString(HostToGlasses.serializer(), msg) runCatching { transport.send(connId, text) } .onFailure { log("hub", "не удалось отправить ${msg::class.simpleName} на $connId: ${it.message}") } } } /** * Сервер связи очков (запускается в PhoneApp). [transport] выбирается флагом * [useBluetooth] (TASK-transport.md п.4): дефолт — [WifiServerTransport] * (Ktor CIO, `0.0.0.0:$port/ws/glasses`), Bluetooth — [BtServerTransport] * (SPP). Поведение по умолчанию (WiFi) не меняется. */ class GlassesServer( private val port: Int = PhoneConfig.SERVER_PORT, private val state: PhoneState = PhoneState(), /** `false` (дефолт) — WiFi; `true` — Bluetooth SPP. */ useBluetooth: Boolean = false, /** Context для runtime-проверки BT-пермишенов (нужен только в [BtServerTransport]). */ private val btContext: Context? = null, ) { /** Активный транспорт: [WifiServerTransport] (дефолт) или [BtServerTransport]. */ val transport: GlassesServerTransport = if (useBluetooth) BtServerTransport(btContext) else WifiServerTransport(port) val hub = GlassesHub(state, transport) /** GlassesSender-адаптер поверх хаба (для PhoneActions). */ val sender: GlassesSender = object : GlassesSender { override suspend fun send(msg: HostToGlasses) = hub.broadcast(msg) } private val scope = CoroutineScope(SupervisorJob() + Dispatchers.IO) private var acceptJob: Job? = null /** Запустить приём подключений очков ([transport.accept] — бесконечный цикл). */ fun start() { if (acceptJob != null) return acceptJob = scope.launch { transport.accept( onMessage = { connId, json -> hub.handle(connId, json) }, onConnected = { connId -> hub.add(connId) }, onDisconnected = { connId -> hub.remove(connId) }, ) } } fun stop() { transport.close() scope.cancel() } }