a434d5bb12
- agent/Assistant.kt: оркестратор диалога (SQLite-история + Koog), таймаут 30с, сессия одна живая, история — последние 20 сообщений - ChatDao.latestSessionId(); SYSTEM_PROMPT нейтральный (TODO: персона) - PhoneApp: очередь фраз (Channel UNLIMITED), assistantLoop (THINKING/ShowText/IDLE), ensureAssistant ленивый; ключ/модель из BuildConfig (local.properties) - Очки: assistantText StateFlow, ShowText → оверлей (белый текст речи + серый ответ), оверлей остаётся после Click до новой сессии - AssistantTest (androidTest, 3 теста: сессия/история/таймаут) — OK (6 tests всего) - README: строка про ассистента
278 lines
13 KiB
Kotlin
278 lines
13 KiB
Kotlin
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.Assistant
|
||
import pw.binom.viewmate.phone.agent.ChatDao
|
||
import pw.binom.viewmate.phone.agent.ChatDb
|
||
import pw.binom.viewmate.phone.agent.LlmClient
|
||
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()
|
||
|
||
/** WS-сервер для очков. */
|
||
val server: GlassesServer = GlassesServer(port = PhoneConfig.SERVER_PORT, state = state)
|
||
|
||
/** Публикация сервера в локальной сети (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<String> = _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<String>(Channel.UNLIMITED)
|
||
|
||
private var assistantLock = Any()
|
||
private var assistant: Assistant? = 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)
|
||
return Assistant(chatDao = chatDao, chat = { history, system -> llm.chat(history, system) })
|
||
.also { assistant = it }
|
||
}
|
||
}
|
||
|
||
/**
|
||
* Единственный consumer фраз: фраза → LLM (с историей из SQLite) → ShowText на очки.
|
||
* Ошибка LLM (нет сети, 401 и т.п.) — коротким сообщением пользователю, не падаем.
|
||
*/
|
||
private suspend fun assistantLoop() {
|
||
for (phrase in assistantChannel) {
|
||
val assistant = ensureAssistant()
|
||
if (assistant == null) {
|
||
log("assistant", "LLM-ключ не задан — фраза пропущена")
|
||
continue
|
||
}
|
||
server.hub.broadcast(AssistantStateMsg(AssistantState.THINKING))
|
||
val answer = runCatching { assistant.process(phrase) }
|
||
.onFailure { log("assistant", "ошибка: ${it.message}") }
|
||
.getOrElse { "Ассистент недоступен: ${it.message}" }
|
||
server.hub.broadcast(ShowText(answer))
|
||
server.hub.broadcast(AssistantStateMsg(AssistantState.IDLE))
|
||
}
|
||
}
|
||
|
||
/**
|
||
* Лениво создаёт стриминговое распознавание (модели во 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
|
||
|
||
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() {
|
||
nsd.unpublish()
|
||
sttStreamer?.close()
|
||
server.stop()
|
||
super.onTerminate()
|
||
}
|
||
|
||
companion object {
|
||
lateinit var instance: PhoneApp
|
||
private set
|
||
}
|
||
}
|