Преглед изворни кода

feat(session): 实现房间恢复功能,增加会话服务以管理最后使用的房间

wenhongquan пре 3 недеља
родитељ
комит
ffb8cc2d1a

+ 22 - 12
AGENTS.md

@@ -177,8 +177,15 @@ sequenceDiagram
     Note over A,W: 全部离开
     A->>LK: disconnect
     LK->>W: participant_disconnected
-    W->>W: 房间空 → 退出
-    W->>LK: disconnect
+    W->>W: 房间空 → 启动 60s 退出倒计时
+
+    alt 倒计时内 App 恢复
+        A->>LK: reconnect
+        W->>W: 取消倒计时,上下文保留
+    else 倒计时结束仍未恢复
+        W->>W: 房间空 → 退出
+        W->>LK: disconnect
+    end
 ```
 
 ## 三层记忆系统
@@ -280,21 +287,24 @@ curl -s -X POST http://localhost:10005/connect -H "Content-Type: application/jso
 ## App 调用流程
 
 ```
-1. App → POST http://36.152.142.37:10005/connect → {room, url, token}
-2. App → WebSocket ws://36.152.142.37:10003 (LiveKit) with token
-3. Dispatcher 自动启动 Worker 加入同一房间
-4. App 说话 → VAD → ASR → LLM → TTS 播放
-5. App 断开 → Worker 退出 → 房间释放
+1. App 启动/恢复 → 读取本地存储的 room(如有)
+2. App → POST /connect {identity, room?} → Dispatcher
+3. Dispatcher:room 存在则复用/重建 Worker,否则创建新 room
+4. App → WebSocket (LiveKit) with token
+5. App 说话 → VAD → ASR → LLM → TTS 播放
+6. App 结束 → 清除本地 room,Worker 延迟退出/房间释放
 ```
 
 ### 暂停 → 继续
 
 ```
-1. App 暂停 → 保存 room 名称
-2. App 继续 → POST /connect {identity, room: "room-xxx"}
-3. Dispatcher 检测 room-xxx 是否存活
-4. 存活 → 返回相同 room,Worker 保持,上下文不丢
-5. 已释放 → 创建新 room + 新 Worker,声纹记忆恢复上下文
+1. App 暂停 → disconnect,room 名称持久化到本地
+2. Worker 进入 60s 退出倒计时,上下文保留
+3. App 继续 → 读取本地 room → POST /connect {identity, room: "room-xxx"}
+4. Dispatcher:
+   - 旧 Worker 仍在 → 返回相同 room,上下文不丢
+   - 旧 Worker 已退出 → 用 room-xxx 启动新 Worker,声纹记忆恢复上下文
+5. App 结束 → 清除本地 room
 ```
 
 ## Flutter App 编译

+ 9 - 2
asr_agent/dispatcher.py

@@ -70,9 +70,16 @@ async def handle_connect(request: web.Request) -> web.Response:
         app_identity = f"user-{int(time.time()*1000)%100000}"
         resume_room = None
 
-    if resume_room and resume_room in _workers and not _workers[resume_room].done():
+    if resume_room:
         room_name = resume_room
-        logger.info("Session resume: %s (%s)", room_name, app_identity)
+        if room_name in _workers and not _workers[room_name].done():
+            logger.info("Session resume: %s (%s)", room_name, app_identity)
+        else:
+            # Recreate worker for the same room so App can reuse the room name
+            # even after the previous worker exited.
+            task = asyncio.create_task(_spawn_worker(room_name))
+            _workers[room_name] = task
+            logger.info("Room %s recreated for %s (total workers: %d)", room_name, app_identity, len(_workers))
     else:
         room_name = f"room-{int(time.time_ns())}"
         task = asyncio.create_task(_spawn_worker(room_name))

+ 24 - 5
asr_agent/worker.py

@@ -484,6 +484,8 @@ class Worker:
         self.room = rtc.Room()
         self.tasks: dict = {}
         self.hist: list[dict] = []
+        self._exit_task = None
+        self._shutdown_event = asyncio.Event()
 
     async def run(self):
         # Pre-load Silero VAD at startup (not lazy)
@@ -499,7 +501,7 @@ class Worker:
 
         self.room.on("track_subscribed", self._on_track)
         self.room.on("participant_disconnected", self._on_off)
-        self.room.on("participant_connected", lambda p: None)
+        self.room.on("participant_connected", self._on_on)
         self.room.on("track_published", lambda pub, p: None)
         await self.room.connect(LIVEKIT_URL, _token(self.rn, self.id))
         logger.info("Worker ready (VAD+ASR+LLM+TTS)")
@@ -509,7 +511,7 @@ class Worker:
                 if pub.track and pub.kind == rtc.TrackKind.KIND_AUDIO:
                     self._on_track(pub.track, pub, p)
         try:
-            await asyncio.Future()
+            await self._shutdown_event.wait()
         finally:
             await self.room.disconnect()
 
@@ -524,10 +526,27 @@ class Worker:
         x = self.tasks.pop(p.sid, None)
         if x:
             x.cancel()
-        # Auto-exit when last participant leaves
+        # Auto-exit when room stays empty for a grace period, to allow App resume.
         if not self.room.remote_participants:
-            logger.info("Room empty, exiting")
-            asyncio.create_task(self.room.disconnect())
+            if self._exit_task is None or self._exit_task.done():
+                logger.info("Room empty, scheduling exit")
+                self._exit_task = asyncio.create_task(self._delayed_exit())
+        else:
+            if self._exit_task and not self._exit_task.done():
+                self._exit_task.cancel()
+                self._exit_task = None
+
+    def _on_on(self, p):
+        if self._exit_task and not self._exit_task.done():
+            logger.info("Participant rejoined, canceling exit")
+            self._exit_task.cancel()
+            self._exit_task = None
+
+    async def _delayed_exit(self, delay_s: float = 60.0):
+        await asyncio.sleep(delay_s)
+        if not self.room.remote_participants:
+            logger.info("Room still empty, exiting")
+            self._shutdown_event.set()
 
     async def _run(self, sid, track):
         logger.info("_run %s", sid)

+ 8 - 1
flutter_asr_client/lib/pages/recording/notifiers/recording_notifier.dart

@@ -11,6 +11,7 @@ import 'package:asr_client/models/websocket_message.dart';
 import 'package:asr_client/providers/livekit_providers.dart';
 import 'package:asr_client/services/livekit_service.dart';
 import 'package:asr_client/providers/settings_providers.dart';
+import 'package:asr_client/services/session_service.dart';
 
 enum RecordingStatus { idle, recording, paused, ending }
 
@@ -110,6 +111,7 @@ final class RecordingNotifier extends AutoDisposeAsyncNotifier<RecordingState> {
   StreamSubscription<ServerMessage>? _messageSub;
 
   LiveKitService get _liveKit => ref.read(liveKitServiceProvider);
+  SessionService get _sessionService => ref.read(sessionServiceProvider);
 
   String? _currentRoom;
 
@@ -143,6 +145,7 @@ final class RecordingNotifier extends AutoDisposeAsyncNotifier<RecordingState> {
     });
 
     final identity = 'user-${_uuid.v4().substring(0, 6)}';
+    final lastRoom = await _sessionService.getLastRoom();
 
     final initialState = const RecordingState(
       projectName: 'G15 沈海高速改扩建 · K1120+300',
@@ -150,8 +153,9 @@ final class RecordingNotifier extends AutoDisposeAsyncNotifier<RecordingState> {
       templateName: '模板 · 普通混凝土坍落度 GB/T 50080',
     );
 
-    final info = await _getRoomInfo(identity);
+    final info = await _getRoomInfo(identity, resumeRoom: lastRoom);
     _currentRoom = info.room;
+    await _sessionService.setLastRoom(info.room);
     try {
       await _liveKit.connect(url: info.url, room: info.room, identity: identity, token: info.token);
     } on Exception catch (e) {
@@ -243,6 +247,7 @@ final class RecordingNotifier extends AutoDisposeAsyncNotifier<RecordingState> {
     try {
       final identity = 'user-${_uuid.v4().substring(0, 6)}';
       final info = await _getRoomInfo(identity, resumeRoom: _currentRoom);
+      _currentRoom = info.room;
       await _liveKit.connect(url: info.url, room: info.room, identity: identity, token: info.token);
       _updateState((s) => _startedState(s));
     } on Exception catch (e) {
@@ -258,6 +263,8 @@ final class RecordingNotifier extends AutoDisposeAsyncNotifier<RecordingState> {
     _timer?.cancel();
     _messageSub?.cancel();
     await _liveKit.disconnect();
+    await _sessionService.clearLastRoom();
+    _currentRoom = null;
     _updateState((s) => s.copyWith(
       status: RecordingStatus.idle,
       elapsed: Duration.zero,

+ 5 - 0
flutter_asr_client/lib/providers/settings_providers.dart

@@ -1,10 +1,15 @@
 import 'package:flutter_riverpod/flutter_riverpod.dart';
+import 'package:asr_client/services/session_service.dart';
 import 'package:asr_client/services/settings_service.dart';
 
 final settingsServiceProvider = Provider<SettingsService>(
   (_) => SharedPreferencesSettingsService(),
 );
 
+final sessionServiceProvider = Provider<SessionService>(
+  (_) => SharedPreferencesSessionService(),
+);
+
 final serverHostProvider = FutureProvider<String>((ref) async {
   final service = ref.watch(settingsServiceProvider);
   return service.getHost();

+ 3 - 1
flutter_asr_client/lib/services/livekit_service.dart

@@ -18,7 +18,7 @@ abstract interface class LiveKitService {
 final class LiveKitServiceImpl implements LiveKitService {
   LiveKitServiceImpl({Room? room}) : _room = room ?? Room();
 
-  final Room _room;
+  Room _room;
   final _messageController = StreamController<ServerMessage>.broadcast();
   var _connected = false;
   EventsListener<RoomEvent>? _listener;
@@ -97,6 +97,8 @@ final class LiveKitServiceImpl implements LiveKitService {
     _listener?.dispose();
     _listener = null;
     await _room.disconnect();
+    // Room instances are single-use; create a fresh one for the next connect.
+    _room = Room();
   }
 
   String _makeToken(String room, String identity) {

+ 43 - 0
flutter_asr_client/lib/services/session_service.dart

@@ -0,0 +1,43 @@
+import 'package:shared_preferences/shared_preferences.dart';
+
+abstract interface class SessionService {
+  Future<String?> getLastRoom();
+  Future<void> setLastRoom(String? room);
+  Future<void> clearLastRoom();
+}
+
+final class SharedPreferencesSessionService implements SessionService {
+  SharedPreferencesSessionService({SharedPreferences? prefs}) : _prefs = prefs;
+
+  SharedPreferences? _prefs;
+
+  Future<SharedPreferences> get _preferences async {
+    if (_prefs != null) return _prefs!;
+    _prefs = await SharedPreferences.getInstance();
+    return _prefs!;
+  }
+
+  @override
+  Future<String?> getLastRoom() async {
+    final prefs = await _preferences;
+    return prefs.getString(_roomKey);
+  }
+
+  @override
+  Future<void> setLastRoom(String? room) async {
+    final prefs = await _preferences;
+    if (room == null) {
+      await prefs.remove(_roomKey);
+    } else {
+      await prefs.setString(_roomKey, room);
+    }
+  }
+
+  @override
+  Future<void> clearLastRoom() async {
+    final prefs = await _preferences;
+    await prefs.remove(_roomKey);
+  }
+
+  static const _roomKey = 'asr_last_room';
+}