Ассистент Порфирий: тула времени, сессии, чат-UI, голосовой цикл; BT-транспорт очки↔телефон
- BT (RFCOMM SPP 00001105) вместо WiFi: BtGlassesTransport/BtServerTransport, реконнект 3с - PlaybackService (foreground, WAKE_LOCK) против морозов HONOR - Ассистент: AgentTools (TimeTool ISO 8601), TOOL_CALL-протокол, сессии SQLite (UUID) - Чат-UI: создание/сброс/переключение сессий, автообновление ленты 2с - LlmClient на прямом OkHttp (Koog давал пустой content на длинных историях) - Голосовой цикл: тройной тап на очках → STT → активная сессия → ShowText - Тесты: AgentToolsTest (3), GlassesServerTest переписан (6) — зелёные
This commit is contained in:
@@ -1,23 +1,18 @@
|
||||
package pw.binom.viewmate.phone
|
||||
|
||||
import io.ktor.server.application.Application
|
||||
import io.ktor.server.application.install
|
||||
import io.ktor.server.cio.CIO
|
||||
import io.ktor.server.engine.EmbeddedServer
|
||||
import io.ktor.server.engine.embeddedServer
|
||||
import io.ktor.server.routing.routing
|
||||
import io.ktor.server.websocket.DefaultWebSocketServerSession
|
||||
import io.ktor.server.websocket.WebSockets
|
||||
import io.ktor.server.websocket.webSocket
|
||||
import io.ktor.websocket.Frame
|
||||
import io.ktor.websocket.readText
|
||||
import io.ktor.websocket.send
|
||||
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.GlassesMode
|
||||
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
|
||||
@@ -41,14 +36,18 @@ const val GLASSES_WS_PATH = "/ws/glasses"
|
||||
|
||||
/**
|
||||
* Менеджер подключённых очков (телефон). Потокобезопасен (ConcurrentHashMap).
|
||||
* Раздаёт исходящие сообщения всем очкам и обрабатывает входящие
|
||||
* (Hello → Welcome, жесты/статус/позицию → лог, GlassesOff → пауза,
|
||||
* DownloadProgress → карта скачивания на очки).
|
||||
* Сессии адресуются **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,
|
||||
) {
|
||||
private val sessions = ConcurrentHashMap<DefaultWebSocketServerSession, Unit>()
|
||||
/** Сессии: connId → Unit. */
|
||||
private val sessions = ConcurrentHashMap<String, Unit>()
|
||||
|
||||
private val _connected = MutableStateFlow(0)
|
||||
val connected: StateFlow<Int> = _connected.asStateFlow()
|
||||
@@ -110,32 +109,29 @@ class GlassesHub(
|
||||
@Volatile
|
||||
var sttFactory: (() -> SttStreamer?)? = null
|
||||
|
||||
fun add(session: DefaultWebSocketServerSession) {
|
||||
sessions[session] = Unit
|
||||
fun add(connId: String) {
|
||||
sessions[connId] = Unit
|
||||
_connected.value = sessions.size
|
||||
log("hub", "очки подключились, всего: ${sessions.size}")
|
||||
log("hub", "очки подключились (connId $connId, всего: ${sessions.size})")
|
||||
}
|
||||
|
||||
fun remove(session: DefaultWebSocketServerSession) {
|
||||
sessions.remove(session)
|
||||
fun remove(connId: String) {
|
||||
sessions.remove(connId)
|
||||
_connected.value = sessions.size
|
||||
log("hub", "очки отключились, всего: ${sessions.size}")
|
||||
log("hub", "очки отключились (connId $connId, всего: ${sessions.size})")
|
||||
}
|
||||
|
||||
/** Отправить сообщение всем подключённым очкам. */
|
||||
/** Отправить сообщение всем подключённым очкам (JSON-строкой через [transport]). */
|
||||
suspend fun broadcast(msg: HostToGlasses) {
|
||||
val text = protocolJson.encodeToString(HostToGlasses.serializer(), msg)
|
||||
for (session in sessions.keys) {
|
||||
try {
|
||||
session.send(text)
|
||||
} catch (e: Exception) {
|
||||
log("hub", "не удалось отправить ${msg::class.simpleName}: ${e.message}")
|
||||
}
|
||||
for (connId in sessions.keys) {
|
||||
runCatching { transport.send(connId, text) }
|
||||
.onFailure { log("hub", "не удалось отправить ${msg::class.simpleName} на $connId: ${it.message}") }
|
||||
}
|
||||
}
|
||||
|
||||
/** Обработать входящее текстовое сообщение от очков. */
|
||||
suspend fun handle(session: DefaultWebSocketServerSession, text: String) {
|
||||
/** Обработать входящее текстовое сообщение от очков (JSON-строкой из [transport]). */
|
||||
suspend fun handle(connId: String, text: String) {
|
||||
val msg = try {
|
||||
protocolJson.decodeFromString(GlassesToHost.serializer(), text)
|
||||
} catch (e: Exception) {
|
||||
@@ -144,8 +140,8 @@ class GlassesHub(
|
||||
}
|
||||
when (msg) {
|
||||
is Hello -> {
|
||||
log("glasses", "hello (app=${msg.appVersion}) → welcome")
|
||||
send(session, Welcome(mode = state.mode, sessionId = state.activeSessionId))
|
||||
log("glasses", "hello (app=${msg.appVersion}, connId $connId) → welcome")
|
||||
send(connId, Welcome(mode = state.mode, sessionId = state.activeSessionId))
|
||||
}
|
||||
|
||||
is Gesture -> {
|
||||
@@ -199,80 +195,67 @@ class GlassesHub(
|
||||
is StopStt -> {
|
||||
if (msg.cancel) {
|
||||
log("stt", "распознавание отменено (очки)")
|
||||
send(session, SttCancel())
|
||||
send(connId, SttCancel())
|
||||
stt?.reset()
|
||||
} else {
|
||||
val full = stt?.finish() ?: ""
|
||||
log("stt", "ВЕСЬ ТЕКСТ: $full")
|
||||
send(session, SttDone(full))
|
||||
send(connId, SttDone(full))
|
||||
stt?.reset()
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private suspend fun send(session: DefaultWebSocketServerSession, msg: HostToGlasses) {
|
||||
private suspend fun send(connId: String, msg: HostToGlasses) {
|
||||
val text = protocolJson.encodeToString(HostToGlasses.serializer(), msg)
|
||||
session.send(text)
|
||||
}
|
||||
}
|
||||
|
||||
/** Ktor-модуль телефона: WebSocket-роут для подключения очков. */
|
||||
fun Application.glassesServerModule(hub: GlassesHub) {
|
||||
install(WebSockets) {
|
||||
pingPeriodMillis = 5_000
|
||||
timeoutMillis = 5_000
|
||||
maxFrameSize = Long.MAX_VALUE
|
||||
masking = false
|
||||
}
|
||||
routing {
|
||||
webSocket(GLASSES_WS_PATH) {
|
||||
hub.add(this)
|
||||
try {
|
||||
for (frame in incoming) {
|
||||
if (frame is Frame.Text) {
|
||||
val text = frame.readText()
|
||||
hub.handle(this, text)
|
||||
}
|
||||
}
|
||||
} catch (e: Exception) {
|
||||
log("server", "WS-ошибка: ${e.message}")
|
||||
e.printStackTrace()
|
||||
} finally {
|
||||
hub.remove(this)
|
||||
}
|
||||
}
|
||||
runCatching { transport.send(connId, text) }
|
||||
.onFailure { log("hub", "не удалось отправить ${msg::class.simpleName} на $connId: ${it.message}") }
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* WS-сервер телефона (CIO-движок — работает на Android, в отличие от Netty).
|
||||
* Живёт в приложении (Application), слушает 0.0.0.0:[port]/ws/glasses.
|
||||
* Сервер связи очков (запускается в 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,
|
||||
) {
|
||||
val hub = GlassesHub(state)
|
||||
/** Активный транспорт: [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 var server: EmbeddedServer<*, *>? = null
|
||||
private val scope = CoroutineScope(SupervisorJob() + Dispatchers.IO)
|
||||
private var acceptJob: Job? = null
|
||||
|
||||
/** Запустить приём подключений очков ([transport.accept] — бесконечный цикл). */
|
||||
fun start() {
|
||||
if (server != null) return
|
||||
server = embeddedServer(CIO, port = port, host = "0.0.0.0") {
|
||||
glassesServerModule(hub)
|
||||
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) },
|
||||
)
|
||||
}
|
||||
server!!.start(wait = false)
|
||||
log("server", "WS-сервер запущен на 0.0.0.0:$port$GLASSES_WS_PATH")
|
||||
}
|
||||
|
||||
fun stop() {
|
||||
server?.stop(1_000, 3_000)
|
||||
server = null
|
||||
transport.close()
|
||||
scope.cancel()
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user