You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.

208 lines
6.4 KiB

import Foundation
#if canImport(JL_BLEKit)
import JL_BLEKit
#endif
/// AI-assistant bridge — mirrors Android `feature/assistant/AssistantBridge.kt`
/// + `JieliAssistantPort.kt`.
///
/// Pipeline:
/// 1. Enter RCSP `MODE_RECORD` (=1) with `recordtype = .byDevice`
/// (equivalent to Android `STRATEGY_DEVICE_ALWAYS_RECORDING`). The
/// headset starts pushing OPUS audio frames continuously.
/// 2. `OpusStreamDecoder` (packet size 40 — critical, otherwise 80% of
/// frames get dropped, same gotcha the Android comment calls out) emits
/// 16-bit / 16 kHz / mono / 20 ms PCM.
/// 3. Each PCM frame is published on the EventDispatcher as `assistantAudio`.
/// 4. Errors / lifecycle events as `assistantError` / `assistantStart` /
/// `assistantEnd`.
///
/// TTS playback (stop direction) — for now we accept playback frames and
/// queue them but do not yet write back through A2DP locally. The Android
/// path writes via a `LocalPlayer` (AudioTrack USAGE_MEDIA) which the OS
/// re-routes through the connected A2DP earphone. iOS equivalent is
/// `AVAudioEngine`/`AudioQueue`; the orchestrator can drive this directly
/// or call back through `assistantPlayback` (TODO once needed).
class AssistantBridge {
weak var server: JieliHomeServer?
init(server: JieliHomeServer) { self.server = server }
private var running: Bool = false
private var decoder: OpusStreamDecoder?
private var sequence: Int64 = 0
#if canImport(JL_BLEKit)
private var sinkWrapper: AssistantSinkWrapper?
#endif
func isRunning() -> Bool { running }
func start() -> Bool {
if running { return true }
#if canImport(JL_BLEKit)
guard let server = server else { return false }
guard let manager = server.translationCoordinator.ensureManager() else {
server.dispatcher.send([
"type": "assistantError",
"code": "device.assistant.no_device",
"message": "no connected device",
])
return false
}
if !manager.trIsSupportTranslate() {
server.dispatcher.send([
"type": "assistantError",
"code": "device.assistant.not_supported",
"message": "device does not support translation",
])
return false
}
// Build the OPUS decoder that turns 40-byte JieLi OPUS packets into
// 16k / 16-bit / mono / 20 ms PCM and republish each frame.
let dec = OpusStreamDecoder(
channels: 1,
packetSize: 40,
sampleRate: 16000,
onPcm: { [weak self] pcm in self?.publishPcm(pcm) },
onError: { [weak self] code, msg in
self?.server?.dispatcher.send([
"type": "assistantError",
"code": "device.assistant.decoder_failed",
"message": "opus: code=\(code) msg=\(msg ?? "")",
])
}
)
dec.start()
decoder = dec
// Install ourselves as the active translation sink before triggering
// `trStartTranslate`. Reusing the call-translation runtime is overkill
// for this single-leg path; we drive the manager directly.
let wrapper = AssistantSinkWrapper(owner: self)
server.translationCoordinator.activeSink = wrapper
sinkWrapper = wrapper
let mode = JLTranslateSetMode()
mode.modeType = .onlyRecord
mode.dataType = .OPUS
mode.channel = 1
mode.sampleRate = 16000
manager.recordtype = .byDevice
// Fire-and-forget — failure surfaces through the delegate as
// `onError`, which we map to `assistantError` events.
manager.trStartTranslate(mode) { [weak self] status, err in
guard let self = self else { return }
if status != .success {
self.server?.dispatcher.send([
"type": "assistantError",
"code": "device.assistant.enter_failed",
"message": "trStartTranslate status=\(status.rawValue) err=\(err?.localizedDescription ?? "")",
])
self.cleanup()
self.running = false
}
}
running = true
server.dispatcher.send([
"type": "assistantStart",
"sampleRate": 16000,
"tsMs": Int(Date().timeIntervalSince1970 * 1000),
])
return true
#else
return false
#endif
}
func stop() {
if !running { return }
running = false
cleanup()
server?.dispatcher.send([
"type": "assistantEnd",
"tsMs": Int(Date().timeIntervalSince1970 * 1000),
])
}
func shutdown() { stop() }
// MARK: - Internal
private func publishPcm(_ pcm: Data) {
sequence &+= 1
server?.dispatcher.send([
"type": "assistantAudio",
"encoding": "pcm16",
"sampleRate": 16000,
"channels": 1,
"bitsPerSample": 16,
"sequence": sequence,
"tsMs": Int(Date().timeIntervalSince1970 * 1000),
"pcm": pcm,
])
}
private func cleanup() {
decoder?.stop(); decoder = nil
#if canImport(JL_BLEKit)
if let server = server, let wrapper = sinkWrapper,
server.translationCoordinator.activeSink === wrapper {
server.translationCoordinator.activeSink = nil
}
sinkWrapper = nil
server?.translationCoordinator.manager?.trExitMode { _, _ in }
#endif
}
// MARK: - Sink callbacks
#if canImport(JL_BLEKit)
fileprivate func handleAudio(_ audio: JLTranslateAudio) {
let payload = audio.data
if payload.isEmpty { return }
switch audio.audioType {
case .PCM:
publishPcm(payload)
case .OPUS:
decoder?.feedEncoded(payload)
default:
// Speex / MSBC / JLA_V2 — not on the assistant path. Drop quietly.
break
}
}
fileprivate func handleModeChange(_ mode: JLTranslateSetMode) {
if mode.modeType == .idle && running {
server?.dispatcher.send([
"type": "assistantError",
"code": "device.assistant.mode_exited",
"message": "headset exited record mode",
])
}
}
fileprivate func handleError(_ error: Error) {
server?.dispatcher.send([
"type": "assistantError",
"code": "device.assistant.translation_error",
"message": error.localizedDescription,
])
}
#endif
}
#if canImport(JL_BLEKit)
private final class AssistantSinkWrapper: TranslationManagerSink {
weak var owner: AssistantBridge?
init(owner: AssistantBridge) { self.owner = owner }
func onModeChange(uuid: String, mode: JLTranslateSetMode) { owner?.handleModeChange(mode) }
func onReceiveAudioData(uuid: String, audio: JLTranslateAudio) { owner?.handleAudio(audio) }
func onError(uuid: String, error: Error) { owner?.handleError(error) }
func onCallingStateChanged(uuid: String, isCalling: Bool) {}
func onSendAudioQueueOver(uuid: String) {}
}
#endif