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.
 
 
 
 
 
 

944 lines
36 KiB

import 'dart:async';
import 'dart:typed_data';
import 'package:device_plugin_interface/device_plugin_interface.dart';
import 'package:get/get.dart';
import '../../core/utils/logger.dart';
import '../../devices/device_hub.dart';
import '../models/agent_module.dart';
import '../models/emai_message.dart';
import '../models/user_Info.dart' as user_info;
import '../utils/background_keepalive.dart';
import '../utils/cue_tone.dart';
import '../utils/peer_speech_gate.dart' show rmsOfPcm16;
import 'agent_module_defaults.dart';
import 'assistant_directive_service.dart';
import 'bailian_multimodal_service.dart';
import 'bailian_native_dialog.dart';
import 'emai_chat_store.dart';
import 'mcp_token_service.dart';
/// 设备按键唤醒的 AI 实时对话会话——**厂商无关**,要求设备具备
/// [DeviceCapability.aiWake]。
///
/// 把「设备按键唤醒」接到阿里百炼多模态上:设备麦的音频由插件解码成 16k PCM
/// (`openMic(ai)`)推给百炼,百炼回的 16k PCM 交给 `openSpeaker(ai)` 由插件/原生
/// 编码后按节奏发回设备播放。整条链路上这里只搬运 PCM。
///
/// **握手是设备主导的**,App 只能应答(恒玄时序,其它厂商在插件里对齐到同一语义):
/// ```
/// 设备按键 ──► wakeEvents(ptt) 恒玄: BB 64
/// App ──► ai.accept 恒玄: AA 6A,必须 3 秒内,否则设备停止送音频
/// …会话进行中…
/// App ──► ai.exit 恒玄: AA 65
/// 设备 ──► wakeEvents(hangup) 恒玄: BB E5;2 秒不回就本地强制结束
/// ```
///
/// 常驻服务:按键在任意页面(含 App 在后台、EMAI 页没打开)都要能唤醒。
class DeviceAiSessionService extends GetxService {
static const String _tag = 'DeviceAiSession';
static DeviceAiSessionService get to => Get.find<DeviceAiSessionService>();
/// 设备侧的会话与手机 EMAI 页各用一条独立的 WebSocket。
/// 共用一个单例的话两边会互相 stop 掉对方、下行音频也会串。
///
/// ⚠️ **接哪个百炼应用要看是哪台设备唤醒的**:耳机/支架接 EMAI,车载香薰接
/// 「Smartcar 车载助手」(各自的 app_id 不同)。所以这里不能是 final 单例,
/// 由 [_ensureService] 按厂商换实例;app_id 没变就复用,不白开 WebSocket。
BailianMultimodalService _svc = BailianMultimodalService();
/// 原生百炼客户端。非空 = 本次会话走原生直连:耳机麦 PCM 与合成音频都不经过
/// Dart,这边只收文本、状态、下行字节数。为 null 时整条链路退回 [_svc]。
///
/// ⚠️ 两条路**不会同时开**:[_useNative] 是本次会话的唯一开关,
/// 两边都起来的话上行会被推两份,服务端 VAD 会把它当成一直有人在说话。
BailianNativeDialog? _native;
bool get _useNative => _native != null;
StreamSubscription? _nativeDownSub;
StreamSubscription? _nativeMicSub;
/// 原生上报的是**累计**下行字节数,这里存上一次的值好算增量
int _nativeDownSeen = 0;
/// 下次开原生会话要用的百炼应用(跟着 [_ensureService] 换)
String? _wantNativeAppId;
String? _svcAppId;
/// 厂商 → 该接哪个 agent 的 app_id。走的是「厂商 → 品类 → 默认编排第一个模块
/// → 模块声明的 bailianAppId」这条链,**设备层不认识 agent、agent 层不认识厂商**。
/// 查不到返回 null = 用全局默认(EMAI)。
String? _resolveAgentAppId(String vendor) {
try {
return bailianAppIdForModule(AgentDefaults.primaryModuleForVendor(vendor));
} catch (e) {
Logger.w(_tag, '解析 agent app_id 失败(回退默认): $e');
return null;
}
}
/// 按需换一个指向目标 app_id 的服务实例。
void _ensureService(String? appId) {
final want = (appId ?? '').trim();
final cur = (_svcAppId ?? '').trim();
if (want == cur) return;
Logger.i(_tag,
'切换设备侧百炼应用: ${cur.isEmpty ? '(默认)' : cur} → ${want.isEmpty ? '(默认)' : want}');
// 旧实例可能还挂着 WebSocket,先停掉再丢
unawaited(_svc.stop());
_svc = BailianMultimodalService(appId: want.isEmpty ? null : want);
// 原生那份也跟着换:两个 agent 用的是百炼里不同的应用,共用会串会话
unawaited(_native?.dispose());
_native = null;
_wantNativeAppId = want.isEmpty ? null : want;
_svcAppId = want;
}
/// 设备上行 PCM 就是 16k/16bit/mono,与百炼上行要求一致
static const int _deviceRate = 16000;
/// 上行攒够这么多字节再发一包:3200B = 100ms @16k。设备每 20ms 上来一包,
/// 逐包转发就是 50 条 WS 消息/秒,攒成 100ms 一条降到 1/5。
static const int _micBatchBytes = 3200;
final BytesBuilder _micAccum = BytesBuilder(copy: false);
/// 收到唤醒后必须在这个时间内应答。我们是立即回,这个常量只是把固件约束写在代码里。
static const Duration _ackDeadline = Duration(seconds: 3);
/// 发出 ai.exit 后等设备回 hangup 的时间,超时就本地强制结束
static const Duration _stopAckTimeout = Duration(seconds: 2);
StreamSubscription? _wakeSub;
StreamSubscription? _hubSub;
StreamSubscription? _micSub;
StreamSubscription? _eventSub;
StreamSubscription? _audioSub;
DeviceSession? _session;
DeviceAudioSource? _mic;
DeviceAudioSink? _spk;
Timer? _stopAckTimer;
Timer? _playoutTimer;
/// 会话是否在跑(已应答,还没收到/发出结束)
final RxBool isActive = false.obs;
/// 手机 EMAI 页(打开时由它自己 attach 上来)。
/// 不用 `Get.put<EmaiSessionOwner>(controller)`:GetX 的 put 会对实例调
/// onInit,而页面正是在自己的 onInit 里注册的,会递归到 Stack Overflow。
EmaiSessionOwner? _phonePage;
int _respondingBytes = 0;
DateTime? _respondingStartedAt;
/// 上行是否被门控住(AI 正在设备里播报期间)。
/// 设备没有把下发的音频接进自己的 AEC,播出去的 AI 声音会被麦克风原样收回,
/// duplex 模式下服务端 VAD 把它当成用户插话。没有 AEC 就只能做半双工。
bool _uplinkGated = false;
Timer? _gateTimer;
/// 门控的兜底计时:开门后无论如何在这个时限内解除。
///
/// ⚠️ 原先门控**只**在 respondingEnded → 播完估算 → _gateTail 这一条链上解除。
/// 任何一轮 AI 出了声却没等到 respondingEnded(非致命 error、WS 抖一下、纯指令帧
/// 收尾),门就永远开着:上行从此全是静音(_onMicPcm 里替换成零),空闲计时也不再
/// 重新武装(_touchActivity 见 _uplinkGated 直接 return),会话活着但是聋的 ——
/// 就是「对着耳机说话,AI 再也不回」。
Timer? _gateSafety;
static const Duration _gateMax = Duration(seconds: 90);
bool _tearingDown = false;
static const Duration _gateTail = Duration(milliseconds: 500);
/// 双方多久没动静就自动结束。
///
/// ⚠️ 原来是 5s,而且只有**百炼事件**才刷新它——设备上行 PCM 到了不算:AI 说完、
/// 门控解除之后用户必须在 5 秒内让服务端 VAD 判到说话,否则会话被判「双方无声」
/// 直接结束。锁屏 + 蓝牙上行延迟 + 蜂窝 RTT,这 5 秒非常容易超,
/// 这就是「锁屏后用耳机对话很容易中断」的直接原因(2026-09-20 苹果测试)。
/// 现在放到 15s,并且设备上行里**有声音**(RMS 过阈值)也算活动。
static const Duration _idleTimeout = Duration(seconds: 15);
Timer? _idleTimer;
/// 设备上行判「有人在说」的 RMS 阈值(与 PeerSpeechGate 同一档:静音 0~50,
/// 正常说话 1000~8000),以及刷新空闲计时的节流间隔。
static const int _micActivityRms = 600;
static const Duration _micActivityThrottle = Duration(milliseconds: 500);
DateTime? _lastMicActivityAt;
/// 百炼 WebSocket 掉了以后的自动重连(蜂窝网切换、弱网抖动)。
/// 原来 `closed` 一到就整个会话拆掉——用户还戴着耳机在说话。
static const int _wsReconnectMax = 2;
static const List<Duration> _wsReconnectDelays = [
Duration(seconds: 1),
Duration(seconds: 3),
];
int _wsReconnectAttempts = 0;
Timer? _wsReconnectTimer;
/// 设备断连后的观察期:iOS 上 Echo-one 的系统级连接本来就会自行断续,
/// 一抖就拆会话太激进;期间设备又回到 ready 就继续,到期没回来才拆。
static const Duration _deviceDropGrace = Duration(seconds: 3);
Timer? _deviceDropTimer;
DateTime? _firstAudioAt;
int _micPackets = 0;
int _downPackets = 0;
EmaiChatStore get _store => EmaiChatStore.to;
EmaiMessage? _streamingUser;
EmaiMessage? _streamingAssistant;
@override
void onInit() {
super.onInit();
final hub = DeviceHub.to;
_wakeSub = hub.wakeEvents.listen(_onWake, onError: (e) {
Logger.e(_tag, '唤醒事件流错误: $e');
});
_hubSub = hub.sessionEvents.listen((e) {
if (e.type != DeviceSessionEventType.connectionStateChanged) return;
if (!isActive.value || _session?.deviceId != e.deviceId) return;
if (e.connectionState == DeviceConnectionState.disconnected) {
if (_deviceDropTimer != null) return;
Logger.w(_tag, '设备断连,观察 ${_deviceDropGrace.inSeconds}s 看它回不回来');
_deviceDropTimer = Timer(_deviceDropGrace, () {
_deviceDropTimer = null;
if (!isActive.value) return;
Logger.w(_tag, '设备 ${_deviceDropGrace.inSeconds}s 内没回来,结束 AI 会话');
_teardown(notifyDevice: false, playEndTone: false);
});
} else if (e.connectionState == DeviceConnectionState.ready &&
_deviceDropTimer != null) {
_deviceDropTimer!.cancel();
_deviceDropTimer = null;
// 新的 session 对象(插件断开重连会重建),把引用换过来,之后的 ai.state /
// ai.exit 才发得出去;上下行流是插件级的,不用重开。
final fresh = hub.sessions.firstWhereOrNull((s) => s.deviceId == e.deviceId);
if (fresh != null) _session = fresh;
Logger.w(_tag, '设备在观察期内回来了,AI 会话继续');
_touchActivity();
}
});
Logger.i(_tag, '设备 AI 会话服务已就绪,等待按键唤醒');
}
@override
void onClose() {
_wakeSub?.cancel();
_hubSub?.cancel();
_teardown(notifyDevice: true, playEndTone: false);
super.onClose();
}
// ---------------- 设备事件 ----------------
void _onWake(DeviceWakeEvent e) {
switch (e.reason) {
case WakeReason.ptt:
case WakeReason.voiceWake:
_handleStart(e.deviceId);
break;
case WakeReason.hangup:
Logger.i(_tag, '设备通知结束 AI');
_teardown(notifyDevice: false);
break;
default:
break;
}
}
/// 起本次对话。两条路径的参数**必须一致**,所以只在这一处写。
Future<bool> _startDialog() async {
final token = await McpTokenService.get();
if (_useNative) {
return _native!.start(
mode: 'duplex',
downstreamRate: _deviceRate,
userId: _currentUserId(),
source: 'earphone',
mcpToken: token,
);
}
return _svc.start(
mode: 'duplex',
downstreamRate: _deviceRate,
userId: _currentUserId(),
source: 'earphone',
mcpToken: token,
);
}
/// 试着开原生直连。两端都要开:
/// - azure_speech 那端([BailianNativeDialog])负责 WS 与攒包;
/// - bluetooth_manager 那端(`DeviceFeatures.aiPcmNativeBridge`)负责把耳机麦
/// PCM 从 Dart 的 EventChannel 上摘下来、并接住回灌的合成音频。
///
/// 任一端不支持就整条退回 Dart,**不能只开一半**:只开翻译端的话上行一帧都收不到,
/// 只开设备端的话 PCM 被摘走却没人接,两种都表现为「唤醒了但 AI 不说话」。
Future<BailianNativeDialog?> _tryOpenNativeDialog(DeviceSession session) async {
if (!await BailianNativeDialog.isAvailable()) {
Logger.i(_tag, '原生百炼客户端不可用,走 Dart 链路');
return null;
}
try {
final r = await session.invokeFeature(
DeviceFeatures.aiPcmNativeBridge, {'enabled': true});
if (r['enabled'] != true) {
Logger.i(_tag, '设备端不支持 AI 原生桥,走 Dart 链路');
return null;
}
} on DeviceException catch (e) {
Logger.i(_tag, '设备端不支持 AI 原生桥(${e.code}),走 Dart 链路');
return null;
} catch (e) {
Logger.w(_tag, '开设备端 AI 原生桥失败,走 Dart 链路: $e');
return null;
}
Logger.i(_tag, 'EMAI 走原生直连(耳机 PCM 不经过 Dart)');
return BailianNativeDialog(appId: _wantNativeAppId);
}
/// 告诉服务端「这一轮本地已经播完了」,它据此把对话状态推回 Listening/Idle。
///
/// ⚠️ 这里**不要**去动耳机的下行队列。原生 `AiAudioDownlink.flush()` 的语义是
/// **清空**(`outQueue.clear()`),不是"把尾包发出去"——名字容易读反,
/// 在这里调它等于把每一句回复的尾音切掉。真正要丢弃待播的地方只有
/// 「真人插话」,那条走 `_spk?.discardPending()`,两种链路都有效
/// (下行发送器是同一个原生对象,跟谁喂的 PCM 无关)。
void _notifyLocalRespondingEnded() {
if (_useNative) {
unawaited(_native!.notifyLocalRespondingEnded());
} else {
_svc.notifyLocalRespondingEnded();
}
}
/// 原生只报**累计**下行字节数(约 100ms 一条):PCM 已经直接回了耳机,
/// Dart 这边要做的只有开门控、计时、计数,三件事都只依赖长度。
void _onNativeDownstreamTotal(int total) {
final delta = total - _nativeDownSeen;
_nativeDownSeen = total;
if (delta > 0) _noteDownstream(delta);
}
void _onMicPcm(Uint8List pcm) {
if (!isActive.value || pcm.isEmpty) return;
_micPackets++;
// 用户在说话 = 有活动。只看能量、带节流,静音帧不续命
if (!_uplinkGated && rmsOfPcm16(pcm) >= _micActivityRms) {
final now = DateTime.now();
if (_lastMicActivityAt == null ||
now.difference(_lastMicActivityAt!) >= _micActivityThrottle) {
_lastMicActivityAt = now;
_touchActivity();
}
}
if (_micPackets == 1 || _micPackets % 50 == 0) {
Logger.i(_tag,
'设备上行 #$_micPackets (${pcm.length}B/包)${_uplinkGated ? ' [门控中]' : ''}');
}
// 门控期间推等长静音而不是断流:上行保持连续,服务端 VAD 收到的是
// "真的没人说话",而不是一段数据缺口
_micAccum.add(_uplinkGated ? Uint8List(pcm.length) : pcm);
if (_micAccum.length >= _micBatchBytes) {
_svc.pushAudio(_micAccum.takeBytes());
}
}
Future<void> _handleStart(String deviceId) async {
if (isActive.value) {
Logger.w(_tag, '已有会话在跑,忽略重复唤醒');
return;
}
final hub = DeviceHub.to;
final session =
hub.sessions.firstWhereOrNull((s) => s.deviceId == deviceId) ??
hub.sessionWith(DeviceCapability.aiWake);
if (session == null || session.state != DeviceConnectionState.ready) {
Logger.e(_tag, '设备未处于已连接状态,拒绝本次唤醒');
return;
}
// ★ 先按厂商选中该说话的那个 agent(车载香薰 → Smartcar 车载助手),
// 再去看 isConfigured —— 反过来的话检查的是上一台设备的应用
_ensureService(_resolveAgentAppId(session.vendor));
Logger.i(_tag,
'设备按键唤醒 AI (${session.vendor},应用=${_svcAppId?.isNotEmpty == true ? _svcAppId : '默认'})');
// 起不来就别应答,直接让设备退出 AI 模式——应答之后设备会一直送音频,
// 而我们这边没有会话接,它会白等到超时
if (!_svc.isConfigured) {
Logger.e(_tag, '百炼未配置(${_svc.configHint}),拒绝本次唤醒');
await _sendExit(session);
return;
}
// 固件要求 3 秒内确认,先应答再做后面的初始化。
// 反过来的话,百炼握手(一次 WebSocket 往返)可能就吃掉大半个预算。
final ackAt = DateTime.now();
_session = session;
try {
await session.invokeFeature(DeviceFeatures.aiAccept);
} catch (e) {
Logger.e(_tag, '应答唤醒失败: $e');
_session = null;
return;
}
isActive.value = true;
_wsReconnectAttempts = 0;
// 换了设备就重新试一次 ai.state(上一台不支持不代表这台也不支持)
_aiStateUnsupported = false;
_quitAfterPlayout = false;
_micPackets = 0;
_downPackets = 0;
_uplinkGated = false;
_micAccum.clear();
await _yieldPhonePage();
// iOS:会话期间保持一个静音的后台音频播放,否则锁屏几秒后进程被挂起,
// BLE 上行与 WebSocket 一起停摆(见 BackgroundKeepalive)。Android 空操作。
await BackgroundKeepalive.acquire(_tag);
// ★ 先试原生直连:上行 PCM 与合成音频都不经过 Dart。
// 插件老 / 设备不是恒玄 → 返回 null,整条链路退回下面的 Dart 路径。
_native = await _tryOpenNativeDialog(session);
try {
if (_useNative) {
_eventSub = _native!.events.listen(_onBailianEvent);
_nativeDownSub = _native!.downstreamBytes.listen(_onNativeDownstreamTotal);
// 设备上行有声只用来续空闲计时,不能当 speechStarted(那会丢弃待播队列)
_nativeMicSub = _native!.micActivity.listen((_) => _touchActivity());
} else {
_mic = await session.openMic(route: AudioRoute.ai);
_micSub = _mic!.pcm.listen(_onMicPcm);
_eventSub = _svc.events.listen(_onBailianEvent);
// 下行 16k:G.722 编码器就是 16k,让服务端直接出 16k 省掉一次重采样
_audioSub = _svc.audioOut.listen(_onDownstreamPcm);
}
// 两条路都要开扬声器:它会启动耳机的 AI 下行发送器(G.722 编码与节流)。
// 原生直连时我们不往里写 PCM,但那个发送器必须在跑,否则原生推下去的
// 合成音频没人发。提示音(叮/嘟)两种模式下都还是由 Dart 写。
_spk = await session.openSpeaker(route: AudioRoute.ai);
} catch (e) {
Logger.e(_tag, '打开设备音频通路失败: $e');
await _teardown(notifyDevice: true, playEndTone: false);
return;
}
// 先给一声「叮」,让用户立刻知道已经在听了——百炼握手还要一次往返
await _playCue(CueTone.ding(sampleRate: _deviceRate));
_touchActivity();
final ok = await _startDialog();
final elapsed = DateTime.now().difference(ackAt);
if (elapsed > _ackDeadline) {
Logger.w(
_tag,
'会话初始化耗时 ${elapsed.inMilliseconds}ms,已超过固件 '
'${_ackDeadline.inSeconds}s 预算(应答已提前回,通常无碍)');
}
if (!ok) {
Logger.e(_tag, '百炼会话启动失败,结束');
await _teardown(notifyDevice: true);
}
}
// ---------------- 百炼事件 ----------------
void _onBailianEvent(BailianEvent e) {
Logger.i(_tag, '百炼事件 $e');
switch (e.type) {
case BailianEventType.speechStarted:
case BailianEventType.speechContent:
case BailianEventType.speechEnded:
case BailianEventType.respondingStarted:
case BailianEventType.respondingContent:
case BailianEventType.respondingEnded:
_touchActivity();
break;
default:
break;
}
switch (e.type) {
case BailianEventType.speechStarted:
// 只有真人插话才丢队列。门控期间的 speechStarted 是 AI 自己的回声触发的
if (!_uplinkGated) {
_spk?.discardPending();
_cancelPlayoutTimer();
}
_streamingUser = null;
break;
case BailianEventType.speechContent:
_streamingUser ??= _store
.append(EmaiMessage(isUser: true, text: '', fromEarphone: true));
_store.applyText(_streamingUser!, e.text);
if (e.isFinal) {
_store.finalize(_streamingUser);
_streamingUser = null;
}
break;
case BailianEventType.speechEnded:
if (e.text.isNotEmpty && _streamingUser != null) {
_store.applyText(_streamingUser!, e.text);
}
_store.finalize(_streamingUser);
_streamingUser = null;
break;
case BailianEventType.respondingStarted:
_respondingBytes = 0;
_firstAudioAt = null;
_respondingStartedAt = DateTime.now();
_cancelPlayoutTimer();
_streamingAssistant = null;
break;
case BailianEventType.respondingContent:
if (e.directives.isNotEmpty &&
Get.isRegistered<AssistantDirectiveService>()) {
AssistantDirectiveService.to.handle(e.directives);
}
// 「关闭对话」:模型会下发 quit 指令并同时回一句告别语。
// 不能立刻结束——那样告别语还没播完就被掐了;标记下来,等播完再收尾。
if (e.directives.any((d) => d.name == _quitDirective)) {
Logger.i(_tag, '收到 quit 指令,本轮播完后结束会话');
_quitAfterPlayout = true;
}
// 纯指令帧没有文本,别建空气泡
if (e.text.isEmpty && e.finishReason == 'command_calls') break;
_streamingAssistant ??= _store
.append(EmaiMessage(isUser: false, text: '', fromEarphone: true));
_store.applyText(_streamingAssistant!, e.text);
break;
case BailianEventType.respondingEnded:
if (e.text.isNotEmpty && _streamingAssistant != null) {
_store.applyText(_streamingAssistant!, e.text);
}
_store.finalize(_streamingAssistant);
_streamingAssistant = null;
_logResponseAudioStats();
_schedulePlayoutDone();
break;
case BailianEventType.error:
Logger.e(_tag, '会话错误: ${e.code} ${e.text} fatal=${e.fatal}');
// 网络类的致命错误(socket/连接被掐)先当掉线重连,不当「功能结束」
if (e.fatal && _looksLikeNetworkError(e.text) && _scheduleWsReconnect(e.text)) {
_store.dropIfEmpty(_streamingUser);
_streamingUser = null;
_store.dropIfEmpty(_streamingAssistant);
_streamingAssistant = null;
break;
}
_store.dropIfEmpty(_streamingUser);
_streamingUser = null;
_store.dropIfEmpty(_streamingAssistant);
_streamingAssistant = null;
if (e.text.isNotEmpty) {
_store.append(EmaiMessage(
isUser: false,
text: e.text,
isFinal: true,
isError: true,
fromEarphone: true));
_store.persist();
}
if (e.fatal) {
_teardown(notifyDevice: true);
} else {
// 非致命错误这一轮不会再有 respondingEnded 了,门必须在这里关,
// 否则上行永远是静音(见 _gateSafety 的说明)。
_cancelPlayoutTimer();
_closeGate();
}
break;
case BailianEventType.stateChanged:
// 把会话状态同步给设备,驱动它的指示灯/表情动画。
// ⚠️ 只有声明了这个 feature 的厂商会真发字节,其余抛 not_supported 被吞掉。
switch (e.state) {
case BailianDialogState.listening:
_pushAiState(DeviceAiState.listening);
break;
case BailianDialogState.thinking:
_pushAiState(DeviceAiState.thinking);
break;
case BailianDialogState.responding:
_pushAiState(DeviceAiState.speaking);
break;
default:
break;
}
break;
case BailianEventType.closed:
if (isActive.value && !_tearingDown) {
if (_scheduleWsReconnect('closed')) break;
Logger.w(_tag, '百炼会话关闭,结束设备 AI');
_teardown(notifyDevice: true);
}
break;
default:
break;
}
}
/// 退回 Dart 链路时才会跑(原生直连时 PCM 不经过这里)
void _onDownstreamPcm(Uint8List pcm) {
if (pcm.isEmpty) return;
_noteDownstream(pcm.length);
_spk?.write(pcm);
}
/// 收到 [n] 字节下行音频时要做的记账。两条路径共用:
/// 走 Dart 时由 [_onDownstreamPcm] 调,走原生时由 [_onNativeDownstreamTotal] 调。
/// 开门控、估播完时长、计数——都只依赖字节数,不需要 PCM 本身。
void _noteDownstream(int n) {
if (!isActive.value || n <= 0) return;
if (_firstAudioAt == null) {
_firstAudioAt = DateTime.now();
_openGate();
}
_respondingBytes += n;
_downPackets++;
if (_downPackets == 1 || _downPackets % 50 == 0) {
Logger.i(_tag, '下行给设备 #$_downPackets (${n}B)');
}
}
/// 有动静就把静默计时重新拨回 5 秒。播报期间不计时。
/// 设备是否支持 [DeviceFeatures.aiState]。不支持就别每次都试、别每次都打日志。
bool _aiStateUnsupported = false;
/// 模型下发的「结束对话」指令名(`extra_info.commands` 里的 `name`)
static const String _quitDirective = 'quit';
/// 收到 quit 指令,等这一轮告别语播完就收尾
bool _quitAfterPlayout = false;
/// 把 AI 会话状态推给设备(驱动指示灯/表情)。
///
/// **fire-and-forget**:设备不支持、或这条命令失败,都不该影响对话本身,
/// 所以既不 await 也不往上抛。
void _pushAiState(DeviceAiState state) {
if (_aiStateUnsupported) return;
final s = _session;
if (s == null || s.state != DeviceConnectionState.ready) return;
s.invokeFeature(DeviceFeatures.aiState, {'state': state}).catchError(
(Object e) {
if (e is DeviceException && e.code == DeviceErrorCode.notSupported) {
// 这台设备没有表情/指示灯,记一次就够了
_aiStateUnsupported = true;
Logger.i(_tag, '设备不支持 ai.state,后续不再推送会话状态');
} else {
Logger.w(_tag, '推送 ai.state($state) 失败: $e');
}
return <String, Object?>{};
},
);
}
static bool _looksLikeNetworkError(String text) {
final t = text.toLowerCase();
const markers = [
'socket', 'connection', 'network', 'timed out', 'timeout', 'reset',
'broken pipe', 'unreachable', 'host lookup', 'websocket', 'closed',
];
return markers.any(t.contains);
}
/// 百炼 WS 掉了:会话还活着、设备还连着,就换条连接接着聊。
/// 返回 false 表示预算用完,调用方按原来的方式结束会话。
bool _scheduleWsReconnect(String reason) {
if (!isActive.value || _tearingDown) return false;
if (_wsReconnectAttempts >= _wsReconnectMax) {
Logger.w(_tag, '百炼 WS 重连 $_wsReconnectMax 次仍失败($reason),放弃');
return false;
}
if (_wsReconnectTimer != null) return true; // 已在排队
final delay = _wsReconnectDelays[_wsReconnectAttempts];
_wsReconnectAttempts++;
Logger.w(_tag, '百炼 WS 掉线($reason),${delay.inSeconds}s 后第 $_wsReconnectAttempts 次重连');
// 重连期间别让空闲计时把会话判死
_idleTimer?.cancel();
_idleTimer = null;
_wsReconnectTimer = Timer(delay, () async {
_wsReconnectTimer = null;
if (!isActive.value || _tearingDown) return;
_cancelPlayoutTimer();
_closeGate();
final ok = await _startDialog();
if (!isActive.value) return;
if (ok) {
Logger.w(_tag, '百炼 WS 第 $_wsReconnectAttempts 次重连成功');
_touchActivity();
} else if (!_scheduleWsReconnect('start failed')) {
_teardown(notifyDevice: true);
}
});
return true;
}
void _touchActivity() {
_idleTimer?.cancel();
_idleTimer = null;
if (!isActive.value || _uplinkGated) return;
_idleTimer = Timer(_idleTimeout, () {
if (!isActive.value) return;
Logger.i(_tag, '双方 ${_idleTimeout.inSeconds}s 无声,自动结束 AI 会话');
requestStop();
});
}
/// 在设备里放一段提示音(走 AI 下行通道)。播放期间开门控。
Future<void> _playCue(Uint8List pcm, {bool waitForPlayout = false}) async {
final spk = _spk;
if (spk == null || _session?.state != DeviceConnectionState.ready) return;
_openGate();
await spk.write(pcm);
final ms = pcm.length * 1000 ~/ (_deviceRate * 2);
if (waitForPlayout) {
// 结束音要等它真的发完+播完,否则紧接着关通道会把它清掉
await Future.delayed(Duration(milliseconds: ms + 400));
} else {
_gateTimer?.cancel();
_gateTimer = Timer(Duration(milliseconds: ms + 300), _closeGate);
}
}
void _openGate() {
_gateTimer?.cancel();
_gateTimer = null;
_gateSafety?.cancel();
_gateSafety = Timer(_gateMax, () {
if (!_uplinkGated) return;
Logger.w(_tag, '门控开了 ${_gateMax.inSeconds}s 还没等到播完信号,强制解除');
_cancelPlayoutTimer();
_notifyLocalRespondingEnded();
_closeGate();
});
if (_uplinkGated) return;
_uplinkGated = true;
// 原生直连时逐帧替换静音在原生做,这里只把结论推下去
unawaited(_native?.setGated(true));
_idleTimer?.cancel();
_idleTimer = null;
Logger.i(_tag, '上行门控开启(AI 播报中)');
}
void _closeGate() {
_gateTimer?.cancel();
_gateTimer = null;
_gateSafety?.cancel();
_gateSafety = null;
if (!_uplinkGated) return;
_uplinkGated = false;
unawaited(_native?.setGated(false));
Logger.i(_tag, '上行门控解除');
_touchActivity();
}
void _logResponseAudioStats() {
final start = _firstAudioAt;
if (start == null || _respondingBytes == 0) return;
final ms = DateTime.now().difference(start).inMilliseconds;
final audioMs = _respondingBytes * 1000 ~/ (_deviceRate * 2);
Logger.i(_tag,
'本轮下行音频 ${_respondingBytes}B ≈ ${audioMs}ms(按16k算),实际用时 ${ms}ms');
}
/// 估算设备什么时候把这一轮播完,然后回 `LocalRespondingEnded`。
void _schedulePlayoutDone() {
_cancelPlayoutTimer();
const bytesPerSecond = _deviceRate * 2;
final totalMs = _respondingBytes * 1000 ~/ bytesPerSecond;
final elapsedMs = DateTime.now()
.difference(_respondingStartedAt ?? DateTime.now())
.inMilliseconds;
// 上限 10 分钟,跟原生下行队列的容量对齐
final remainMs = (totalMs - elapsedMs).clamp(0, 600000);
Logger.i(_tag, '预计还需 ${remainMs}ms 播完(已推 ${totalMs}ms 音频)');
_playoutTimer = Timer(Duration(milliseconds: remainMs), () {
if (!isActive.value) return;
_notifyLocalRespondingEnded();
_gateTimer = Timer(_gateTail, _closeGate);
if (_quitAfterPlayout) {
_quitAfterPlayout = false;
Logger.i(_tag, '告别语已播完,按 quit 指令结束会话');
_teardown(notifyDevice: true, playEndTone: false);
}
});
}
void _cancelPlayoutTimer() {
_playoutTimer?.cancel();
_playoutTimer = null;
}
// ---------------- 结束 ----------------
/// 主动结束:发 ai.exit,等设备回 hangup;超时就本地强制收尾。
Future<void> requestStop() async {
if (!isActive.value) return;
Logger.i(_tag, '请求设备退出 AI 模式');
final s = _session;
if (s != null) await _sendExit(s);
_stopAckTimer?.cancel();
_stopAckTimer = Timer(_stopAckTimeout, () {
if (!isActive.value) return;
Logger.w(_tag, '${_stopAckTimeout.inSeconds}s 未收到设备的结束应答,本地强制结束');
_teardown(notifyDevice: false);
});
}
Future<void> _sendExit(DeviceSession s) async {
if (s.state != DeviceConnectionState.ready) return;
try {
await s.invokeFeature(DeviceFeatures.aiExit);
} catch (e) {
Logger.w(_tag, '发送 ai.exit 失败: $e');
}
}
/// [notifyDevice] 为 true 时补发 ai.exit:不是设备主导的结束,必须告诉设备退出。
/// [playEndTone] 为 false 用于设备已经断连的场景。
Future<void> _teardown({
required bool notifyDevice,
bool playEndTone = true,
}) async {
if (_tearingDown) return;
_tearingDown = true;
try {
await _doTeardown(notifyDevice: notifyDevice, playEndTone: playEndTone);
} finally {
_tearingDown = false;
}
}
Future<void> _doTeardown({
required bool notifyDevice,
bool playEndTone = true,
}) async {
_stopAckTimer?.cancel();
_stopAckTimer = null;
_idleTimer?.cancel();
_idleTimer = null;
_gateSafety?.cancel();
_gateSafety = null;
_wsReconnectTimer?.cancel();
_wsReconnectTimer = null;
_wsReconnectAttempts = 0;
_deviceDropTimer?.cancel();
_deviceDropTimer = null;
_lastMicActivityAt = null;
_cancelPlayoutTimer();
_micAccum.clear();
if (isActive.value && playEndTone) {
// 先把还没播完的 AI 语音丢掉,不然「咚」要排在几十秒的队尾
await _spk?.discardPending();
await _playCue(CueTone.dong(sampleRate: _deviceRate), waitForPlayout: true);
}
_closeGate();
final wasActive = isActive.value;
isActive.value = false;
_store.dropIfEmpty(_streamingUser);
_streamingUser = null;
_store.dropIfEmpty(_streamingAssistant);
_streamingAssistant = null;
if (wasActive) _store.persist();
await _eventSub?.cancel();
_eventSub = null;
await _audioSub?.cancel();
_audioSub = null;
await _micSub?.cancel();
_micSub = null;
await _nativeDownSub?.cancel();
_nativeDownSub = null;
await _nativeMicSub?.cancel();
_nativeMicSub = null;
_nativeDownSeen = 0;
await _spk?.close();
_spk = null;
await _mic?.close();
_mic = null;
// 原生直连时两端都要关:只关一端会把耳机麦的 PCM 一直摘在原生里没人收
final nat = _native;
_native = null;
if (nat != null) {
await nat.stop();
final s0 = _session;
if (s0 != null) {
try {
await s0.invokeFeature(
DeviceFeatures.aiPcmNativeBridge, {'enabled': false});
} catch (e) {
Logger.w(_tag, '关设备端 AI 原生桥失败: $e');
}
}
}
await _svc.stop();
final s = _session;
if (wasActive && notifyDevice && s != null) await _sendExit(s);
_session = null;
await BackgroundKeepalive.release(_tag);
if (wasActive) {
Logger.i(_tag, '设备 AI 会话已结束 (上行 $_micPackets 包 / 下行 $_downPackets 包)');
await _resumePhonePage();
}
}
// ---------------- 与手机 EMAI 页互斥 ----------------
void attachPhonePage(EmaiSessionOwner page) => _phonePage = page;
void detachPhonePage(EmaiSessionOwner page) {
if (identical(_phonePage, page)) _phonePage = null;
}
Future<void> _yieldPhonePage() async {
final ctrl = _phonePage;
if (ctrl == null) return;
Logger.i(_tag, '手机 EMAI 页正开着,让它让出会话');
try {
await ctrl.yieldToEarphone();
} catch (e) {
Logger.w(_tag, '通知 EMAI 页让出失败: $e');
}
}
Future<void> _resumePhonePage() async {
final ctrl = _phonePage;
if (ctrl == null) return;
try {
await ctrl.resumeFromEarphone();
} catch (e) {
Logger.w(_tag, '通知 EMAI 页恢复失败: $e');
}
}
String _currentUserId() {
try {
return user_info.User.isLoggedIn()
? user_info.User.instance.uid
: 'eaimar-guest';
} catch (_) {
return 'eaimar-guest';
}
}
}
/// 手机 EMAI 页实现这个接口,让设备会话能在不反向依赖 modules 层的前提下
/// 通知它让出/恢复。EmaiController 在 onInit/onClose 里
/// 调 [DeviceAiSessionService.attachPhonePage] / [DeviceAiSessionService.detachPhonePage]。
abstract class EmaiSessionOwner {
/// 设备接管:停掉本页的播放与会话,并且不要自动重连
Future<void> yieldToEarphone();
/// 设备结束:把本页的会话接回来
Future<void> resumeFromEarphone();
}