英语学习 App 全项目扫描后的修复:持续通话体验、服务端缺陷、客户端缺陷与测试守门。两端测试全绿(服务端 370、客户端 81、analyze 无问题)。
鉴权移到数据库查询之前;orchestrator 构造失败时释放租约;握手前 close 不再退化成空 403;空转写不进 record_input;租约失效时通知客户端;会话号必填;生产环境缺 upstream 拒绝启动
@@ -46,6 +46,14 @@ def create_app( ) -> FastAPI: if environment == 'production' and not auth_secret: raise RuntimeError('ENGLISH_APP_AUTH_SECRET must be configured in production.')+ upstream_mode = (+ 'injected' if upstream_factory is not None+ else os.getenv('ENGLISH_APP_UPSTREAM', 'fake')+ )+ if environment == 'production' and upstream_mode == 'fake':+ # 生产环境不配 upstream 会静默跑复读机:线上表现为“AI 只会重复孩子+ # 说的话”,而日志和健康检查都不会报任何异常。+ raise RuntimeError('ENGLISH_APP_UPSTREAM must be configured in production.') memory_store = memory if memory_store is None: memory_store = (@@ -124,7 +132,13 @@ def create_app( @app.get("/health") async def health() -> dict[str, object]:- return {"ok": True, "service": "english-app-gateway", "authMode": app.state.auth_mode}+ return {+ "ok": True,+ "service": "english-app-gateway",+ "authMode": app.state.auth_mode,+ # 暴露上游模式,避免部署时静默跑成 fake 而没人发现。+ "upstreamMode": upstream_mode,+ } @app.get('/api/v1/children/{child_id}/profile') async def profile(child_id: str, x_session_token: str | None = Header(default=None, alias='x-session-token')):@@ -197,11 +211,30 @@ def create_app( @app.websocket("/api/realtime") async def realtime(websocket: WebSocket) -> None:+ async def reject(code: int, reason: str) -> None:+ # 握手开始前调用 close,uvicorn 只会回一个空的 HTTP 403,客户端+ # 拿不到任何 code,分不清 token 失效、已有会话和网络故障。+ # 先接受握手,再带 code 关闭。+ await websocket.accept()+ await websocket.close(code=code, reason=reason)+ requested_activity_id = websocket.query_params.get('activityId')- session_id = websocket.query_params.get("sessionId", "unknown")+ session_id = websocket.query_params.get("sessionId")+ if not session_id:+ # 缺省会话号会让所有匿名连接共用同一条归属记录,先写入的孩子+ # 会把这条记录永久占住,别人再连就被当成“会话不可用”。+ await reject(4400, 'sessionId required')+ return child_id = websocket.headers.get("x-child-id") or websocket.query_params.get( "childId", "local-child" )+ owner_id = websocket.headers.get('x-session-owner') or uuid4().hex+ # 鉴权必须在 plan_for_child 之前:那个调用会读学情、能力事件和最近+ # 学习状态,未鉴权的连接不该有触发这些查询的机会。+ if not _session_authorized(websocket.headers.get('x-session-token'), session_token,+ child_id, auth_secret):+ await reject(4401, "session token required")+ return content_registry = StarterContentRegistry() planned = None if requested_activity_id is None:@@ -209,23 +242,18 @@ def create_app( activity_id = str(planned['activityId']) else: activity_id = requested_activity_id- owner_id = websocket.headers.get('x-session-owner') or uuid4().hex- if not _session_authorized(websocket.headers.get('x-session-token'), session_token,- child_id, auth_secret):- await websocket.close(code=4401, reason="session token required")- return if session_store_value is not None: owner_child = session_store_value.owner_child(session_id) if owner_child is not None and owner_child != child_id:- await websocket.close(code=4403, reason='session unavailable')+ await reject(4403, 'session unavailable') return try: content_registry.activity(activity_id) except UnknownContentActivity:- await websocket.close(code=4400, reason='unknown activity')+ await reject(4400, 'unknown activity') return if not lease_store_value.acquire(child_id, session_id, owner_id, ttl=120):- await websocket.close(code=4409, reason='child already has an active session')+ await reject(4409, 'child already has an active session') return await websocket.accept() star_offset = 0@@ -239,18 +267,31 @@ def create_app( async def send(event: dict[str, object]) -> None: nonlocal ready, session_started_at, last_audio_at, speech_task_revision, last_engagement, idle_watch, finished- if finished or not lease_store_value.is_owner(child_id, session_id, owner_id):+ if finished:+ return+ if not lease_store_value.is_owner(child_id, session_id, owner_id):+ # 租约已被别的连接拿走。安静丢弃事件会让家长以为还在通话,+ # 实际上孩子说的话再也没人处理。+ finished = True+ try:+ await websocket.send_json({'type': 'session.lease_lost',+ 'message': '这个会话已经在别处继续,请重新连接。'})+ except Exception:+ pass return if event.get('type') == 'audio.interrupted' and event.get('reason') == 'speech_started': speech_task_revision = orchestrator.state.task_revision if event.get('type') == 'transcript.final' and event.get('role') == 'user':- if event.get('content') and orchestrator.state.status != 'completed':+ # 纯空白转写不是有效输入。放它进 record_input 会抛异常,而异常+ # 冒到读循环会终止整条上游转发,通话看着正常却再也不出声。+ said = str(event.get('content') or '').strip()+ if said and orchestrator.state.status != 'completed': last_engagement = time.monotonic() timing = None if last_audio_at is None else { 'audioAgeMs': int(max(0, (time.monotonic() - last_audio_at) * 1000)), } evidence = orchestrator.record_input(- str(event['content']), source='speech', timing=timing, speech_task_revision=speech_task_revision,+ said, source='speech', timing=timing, speech_task_revision=speech_task_revision, ) last_audio_at = None speech_task_revision = None@@ -258,7 +299,7 @@ def create_app( 'messageId': f"{session_id}:{evidence['evidenceId']}", 'childId': child_id, 'sessionId': session_id, 'role': 'user', 'source': 'speech',- 'content': str(event['content']),+ 'content': said, 'createdAt': datetime.now(timezone.utc).isoformat(), }) auto_result = orchestrator.auto_route_starting_point(evidence)@@ -381,19 +422,29 @@ def create_app( except (RuntimeError, WebSocketDisconnect): return - orchestrator = SessionOrchestrator(- session_id=session_id,- child_id=child_id,- activity_id=activity_id,- memory=memory_store,- conversation_memory=conversation_store,- session_store=session_store_value,- initial_plan=planned,- profile_store=profile_store_value,- require_presentation_ack=True,- )- if session_store_value is not None:- star_offset = max(0, session_store_value.total_stars(child_id) - orchestrator.state.stars)+ try:+ orchestrator = SessionOrchestrator(+ session_id=session_id,+ child_id=child_id,+ activity_id=activity_id,+ memory=memory_store,+ conversation_memory=conversation_store,+ session_store=session_store_value,+ initial_plan=planned,+ profile_store=profile_store_value,+ require_presentation_ack=True,+ )+ if session_store_value is not None:+ star_offset = max(0, session_store_value.total_stars(child_id) - orchestrator.state.stars)+ except Exception:+ # 这里已经在 accept 之后、总 try/finally 之前,构造失败必须自己+ # 释放租约:否则这个孩子在 TTL(120 秒)内的重连会被当成+ # “已有进行中的会话”全部拒绝。+ lease_store_value.release(child_id, session_id, owner_id)+ logging.getLogger('english_app').exception('session setup failed')+ await websocket.send_json({'type': 'error', 'message': '暂时无法开始这次练习,请重新连接。'})+ await websocket.close(code=1011)+ return def lesson_context(): context = orchestrator.lesson_context()
单事件处理失败不再终结整条转发链;上游正常关闭时通知客户端;工具失败也回写序号、无论成败都应答 function_call_output;重试预算与“未执行工具”分开计数;孩子答“忘了”不再被当成长名字
@@ -9,12 +9,35 @@ from urllib.parse import parse_qs, urlencode, urlsplit, urlunsplit import websockets -from .session import OBSERVATION_LEVELS+from .session import OBSERVATION_LEVELS, end_requested from .intake import introductory_name EventSender = Callable[[dict[str, object]], Awaitable[None]] +_NOT_A_NAME = {+ '不知道', '不想说', '不告诉你', '你猜', '随便', '忘了', '忘记', '不记得',+ '没有', '没有名字', '不', '不是', 'no', 'nope', 'nothing',+ 'i forgot', 'forgot', 'dont know', "don't know",+}+++def _looks_like_a_name(bare: str) -> bool:+ """孩子对“怎么称呼你”的回答能不能当名字用。++ 没有这层判断时,答非所问会被写成长期姓名,之后 AI 一直这么叫孩子,+ 家长报告里也是这个名字。+ """+ if not bare or bare.casefold() in _NOT_A_NAME:+ return False+ if re.fullmatch(r"[A-Za-z][A-Za-z .'-]{1,19}", bare):+ return True+ # 中文名按 2–4 字处理,并排除常见虚词(“忘了”的“了”、“没有”的“没”)。+ return bool(+ re.fullmatch(r'[一-鿿]{2,4}', bare)+ and not re.search(r'[了很是不没吗呢吧啊呀]', bare)+ )+ class UpstreamConfigurationError(RuntimeError): pass@@ -82,6 +105,7 @@ class DashScopeUpstream: self._input_generation = 0 self._command_responses: set[str] = set() self._tool_retry_count = 0+ self._command_retry_count = 0 self._cancel_issued = False self._discarded_responses: set[str] = set() @@ -211,13 +235,15 @@ class DashScopeUpstream: call_id = str(event.get('toolCallId', '')) name = self._tool_names.pop(call_id, None) generation = self._tool_generations.pop(call_id, self._input_generation)+ retries_exhausted = False if not result.get('ok'): self._tool_retry_count += 1- if self._tool_retry_count > 2:- await self._send({'type': 'error', 'code': 'TOOL_RETRY_EXHAUSTED', 'message': '小树需要重新连接一下。'})- return+ retries_exhausted = self._tool_retry_count > 2 if isinstance(result.get('sequence'), int) and self._lesson_context: self._lesson_context['board']['sequence'] = result['sequence']+ # 失败的调用也要把新序号写回指令:只改本地 context 的话,+ # 模型手里的指令仍是旧 sequence,重试必然再次 stale。+ await self._send_upstream({'type': 'session.update', 'session': self._control_update()}) if result.get('ok') and result.get('lessonContext'): self._lesson_context = result['lessonContext'] if name and generation == self._input_generation:@@ -227,6 +253,8 @@ class DashScopeUpstream: self._latest_input = None self._input_actions.clear() await self._send_upstream({'type': 'session.update', 'session': self._control_update()})+ # 无论成败都要应答这次 function_call:提前 return 会把它永远留在+ # 上游会话里,模型之后的行为不可预期。 await self._send_upstream( { "type": "conversation.item.create",@@ -237,18 +265,14 @@ class DashScopeUpstream: }, } )+ if retries_exhausted:+ await self._send({'type': 'error', 'code': 'TOOL_RETRY_EXHAUSTED', 'message': '小树需要重新连接一下。'})+ return await self._request_response() def _stop_requested(self) -> bool: evidence = self._latest_input or {}- if evidence.get('source') != 'speech':- return False- said = str(evidence.get('content', '')).strip().strip('。!!??.,,').casefold()- return bool(re.fullmatch(- r'(?:小树[,, ]*)?(?:我(?:现在|今天)?(?:不想玩了?|不玩了|想休息(?:一下)?|要休息(?:一下)?)|'- r'(?:请|先)?(?:结束(?:吧)?|停一下|停下来|停止|休息一下)|'- r'(?:please )?(?:stop(?: the (?:game|lesson))?|goodbye|bye)(?: please)?|'- r"(?:let.s|i want to) stop)", said))+ return end_requested(str(evidence.get('content', '')), str(evidence.get('source', ''))) def _routing_instruction(self) -> str: evidence = self._latest_input@@ -272,7 +296,7 @@ class DashScopeUpstream: name = introductory_name(str(evidence.get('content', ''))) if name is None and evidence.get('profileQuestion'): bare = str(evidence.get('content', '')).strip().strip('。.!!??')- if re.fullmatch(r'[\u4e00-\u9fffA-Za-z ]{1,12}', bare) and bare not in {'不知道','不想说','不告诉你','你猜','随便'}:+ if _looks_like_a_name(bare): name = bare if name and name != (context.get('childProfile') or {}).get('name'): args = {'field': 'name', 'value': name, 'evidenceId': evidence.get('evidenceId'),@@ -306,7 +330,7 @@ class DashScopeUpstream: ('recentTopic', re.search(r'((?:今天|昨天|最近)[^。!?!?]{1,100})', said))] if evidence.get('profileQuestion') and not profile.get('name'): bare = said.strip().strip('。!!??.,,')- if re.fullmatch(r'[\u4e00-\u9fffA-Za-z ]{1,12}', bare) and bare not in {'不知道', '不想说', '不告诉你', '你猜', '随便'}:+ if _looks_like_a_name(bare): candidates.insert(0, ('name', re.fullmatch(r'(.*)', bare))) if not re.search(r'假装|扮演|故事|pretend|role.?play', said, re.I): for field, match in candidates:@@ -394,6 +418,7 @@ class DashScopeUpstream: self._latest_input = evidence self._input_actions.clear() self._tool_retry_count = 0+ self._command_retry_count = 0 await self._send_upstream({'type': 'session.update', 'session': self._control_update()}) async def _request_response(self, *, cancel_current: bool = False) -> None:@@ -439,13 +464,26 @@ class DashScopeUpstream: async def _read_loop(self) -> None: try: async for raw in self._socket:- event = json.loads(raw)- if isinstance(event, dict):- await self._forward_event(event)+ try:+ event = json.loads(raw)+ if isinstance(event, dict):+ await self._forward_event(event)+ except asyncio.CancelledError:+ raise+ except Exception as error:+ # 单个事件处理失败不能结束读循环:那会让上游的音频、转写与+ # 工具调用永久停在这里,孩子端看着还在通话却再也收不到声音。+ await self._send({'type': 'error', 'message': str(error)}) except asyncio.CancelledError: raise except Exception as error: await self._send({"type": "error", "message": str(error)})+ else:+ # 走到这里说明没有异常也没被取消,即上游自己关掉了连接(例如+ # idle timeout)。不通知的话,客户端会拿着一个死连接继续说,+ # 直到下一次操作才收到一个笼统错误。+ await self._send({'type': 'error', 'code': 'UPSTREAM_CLOSED',+ 'message': '语音会话已结束,请重新连接。'}) async def _forward_event(self, event: dict[str, object]) -> None: event_type = event.get("type")@@ -484,8 +522,10 @@ class DashScopeUpstream: self._routing_responses.discard(response_id) self._tool_responses.discard(response_id) if command_missing and not discarded:- if self._tool_retry_count < 2:- self._tool_retry_count += 1+ # 与工具执行失败分开计数:两者是不同故障,共用一个预算会让+ # 任意两次失败就耗尽全部重试。+ if self._command_retry_count < 2:+ self._command_retry_count += 1 await self._send_upstream({'type': 'conversation.item.create', 'item': { 'type': 'message', 'role': 'user', 'content': [{'type': 'input_text', 'text': 'APP_RUNTIME: 上次未执行工具。现在只执行:' + self._routing_instruction() + self._teaching_instruction()}]}})
结束通话要求最近一条语音明确表达结束意图;showChoices 补数量上限与去重;end_requested 判定与 upstream 共用同一份实现
@@ -18,6 +18,20 @@ from .planner import next_activity_plan, capability_activity_plan OBSERVATION_LEVELS = ('UNOBSERVED', 'PROMPTED_OBSERVED', 'INDEPENDENT_OBSERVED', 'REVIEW_REQUIRED') +_END_REQUEST = re.compile(+ r'(?:小树[,, ]*)?(?:我(?:现在|今天)?(?:不想玩了?|不玩了|想休息(?:一下)?|要休息(?:一下)?)|'+ r'(?:请|先)?(?:结束(?:吧)?|停一下|停下来|停止|休息一下)|'+ r'(?:please )?(?:stop(?: the (?:game|lesson))?|goodbye|bye)(?: please)?|'+ r"(?:let.s|i want to) stop)")+++def end_requested(content: str, source: str) -> bool:+ """孩子是否明确要求结束这次通话。只有语音里说出来才算,点选和文字不算。"""+ if source != 'speech':+ return False+ said = str(content).strip().strip('。!!??.,,').casefold()+ return bool(_END_REQUEST.fullmatch(said))+ class StaleSessionAction(RuntimeError): def __init__(self, expected: int, actual: int) -> None:@@ -137,9 +151,10 @@ class ObservationMemory: 'sourceObservationId': event['eventId'], 'observedResource': event['item'], 'quote': evidence.get('content', '')}) return result - def event_for_tool(self, child_id: str, tool_call_id: str) -> dict[str, Any] | None:+ def event_for_tool(self, child_id: str, session_id: str, tool_call_id: str) -> dict[str, Any] | None: return next((event for event in self.events_for_child(child_id)- if event.get('toolCallId') == tool_call_id), None)+ if event.get('toolCallId') == tool_call_id+ and event.get('sessionId') == session_id), None) @staticmethod def _domain_view(events):@@ -426,6 +441,7 @@ class SessionOrchestrator: 'phase': self.state.diagnostic_path, 'taskChangedDuringSpeech': (speech_task_revision is not None and speech_task_revision != self.state.task_revision), 'profileQuestion': self._last_question_was_name(),+ 'endRequested': end_requested(content, source), 'targetItem': self.state.target_item, 'targetDomain': self.state.target_domain, 'presentationConfirmed': (not self.require_presentation_ack or source == 'choice' or self.state.presented_task_revision == self.state.task_revision),@@ -567,7 +583,7 @@ class SessionOrchestrator: previous = self._processed.get(tool_call_id) if previous is not None: return previous- existing_event = self.memory.event_for_tool(self.child_id, tool_call_id)+ existing_event = self.memory.event_for_tool(self.child_id, self.session_id, tool_call_id) if existing_event is not None: sequence_after = existing_event.get('sequenceAfter') if type(sequence_after) is int and self.state.sequence < sequence_after:@@ -670,9 +686,13 @@ class SessionOrchestrator: choices = arguments.get("choices") if not isinstance(choices, list) or not all(isinstance(item, str) for item in choices): raise InvalidSessionAction("choices 必须是字符串数组。")+ if len(choices) > 6:+ raise InvalidSessionAction("选项不能超过 6 个。") allowed = set(self.content.activity(self.activity_id).items) if not set(choices).issubset(allowed): raise InvalidSessionAction("choices 包含当前活动未注册的内容。")+ # 与 presentTask 一致:重复选项会让黑板出现两个一样的可点对象。+ choices = list(dict.fromkeys(choices)) next_state = self._next(choices=tuple(choices), task_revision=self.state.sequence + 1) elif name == "awardStar": reason = self._required_text(arguments, "reason").upper()@@ -787,6 +807,11 @@ class SessionOrchestrator: elif name == "completeActivity": if self.state.child_input_count == 0: raise InvalidSessionAction('完成活动需要孩子先参与。')+ # 一次普通回答、一次点选或一句已被后文覆盖的旧“停一下”都不能结束+ # 整通电话:只有孩子最近一条语音明确要求结束才算(持续通话要求)。+ latest_input = next(reversed(self._inputs.values()), None)+ if not latest_input or not latest_input.get('endRequested'):+ raise InvalidSessionAction('完成活动需要孩子明确要求结束。') next_state = self._next(status="completed") else: raise InvalidSessionAction(f"不允许的工具:{name}。")
会话记录查询限定 session,避免跨会话命中同一 toolCallId;补 ORDER BY;去掉与内存版不一致的截断
@@ -58,11 +58,15 @@ class SQLiteObservationMemory(ObservationMemory): ).fetchall() return [json.loads(row[0]) for row in rows] - def event_for_tool(self, child_id: str, tool_call_id: str) -> dict[str, Any] | None:+ def event_for_tool(self, child_id: str, session_id: str, tool_call_id: str) -> dict[str, Any] | None:+ # 必须限定会话:同一个孩子在别的会话里用过同一个 toolCallId 时,+ # 命中那边的旧事件会让这次调用凭空推进一格,而实际什么都没执行。 row = self._db.execute( "SELECT payload FROM observation_events WHERE child_id=? "- "AND json_extract(payload, '$.toolCallId')=? LIMIT 1",- (child_id, tool_call_id),+ "AND json_extract(payload, '$.sessionId')=? "+ "AND json_extract(payload, '$.toolCallId')=? "+ "ORDER BY rowid LIMIT 1",+ (child_id, session_id, tool_call_id), ).fetchone() return json.loads(row[0]) if row else None @@ -98,7 +102,8 @@ class SQLiteConversationMemory(ConversationMemory): self._db.commit() def append(self, message: dict[str, Any]) -> None:- # Validate and apply the same truncation rules as the in-memory store.+ # 与内存版保持一致:这里只校验,不截断。截断会让模型读到半句话,+ # 而且两个实现读到的内容还不一样;输入侧已经有长度上限。 content = message.get("content") if message.get("role") not in {"user", "assistant"} or not isinstance(message.get("source"), str): raise ValueError("conversation message role/source invalid")@@ -113,7 +118,7 @@ class SQLiteConversationMemory(ConversationMemory): ( str(message["messageId"]), str(message["childId"]), str(message.get("sessionId", "")), str(message["role"]),- str(message["source"]), content[:500], str(message.get("createdAt", "")),+ str(message["source"]), content, str(message.get("createdAt", "")), ), ) self._db.commit()
连续大写块先比对常见英文词,`I LIKE CAT` 不再凭空产生 E/K/L 三个字母
@@ -1,6 +1,22 @@ """Conservative intake parsing; self-reports never create learning evidence.""" import re +# 大写块只有在不是常见英文词时才当作拼读的字母序列:否则孩子说+# "I LIKE CAT" 会凭空多出 E、K、L 三个从没提过的字母。+_COMMON_WORDS = frozenset({+ 'a', 'am', 'an', 'and', 'apple', 'are', 'at', 'ball', 'be', 'bed', 'big',+ 'bird', 'blue', 'book', 'boy', 'but', 'by', 'can', 'cat', 'come', 'dad',+ 'day', 'do', 'dog', 'eat', 'egg', 'fish', 'for', 'from', 'get', 'girl',+ 'go', 'good', 'green', 'had', 'has', 'hat', 'have', 'he', 'her', 'here',+ 'hi', 'his', 'how', 'i', 'in', 'is', 'it', 'know', 'like', 'little',+ 'look', 'love', 'man', 'me', 'mom', 'mother', 'my', 'name', 'new',+ 'nice', 'no', 'not', 'of', 'on', 'one', 'or', 'pen', 'play', 'red',+ 'run', 'sad', 'say', 'see', 'she', 'sit', 'so', 'sun', 'that', 'the',+ 'they', 'this', 'three', 'to', 'too', 'tree', 'two', 'up', 'us', 'very',+ 'want', 'was', 'we', 'what', 'where', 'who', 'why', 'will', 'with',+ 'yes', 'you', 'your',+})+ def introductory_name(text: str) -> str | None: match = re.fullmatch(r"\s*(?:我叫|我的名字是|叫我|my name is)\s*([^,。!?,!?\s]{1,12})[。.!!??]*\s*", text, re.I)@@ -11,8 +27,9 @@ def mentioned_items(text: str, allowed: tuple[str, ...]) -> list[str]: folded = text.casefold() words = [item for item in allowed if len(item) > 1 and re.search(r'(?<![a-z])' + re.escape(item.casefold()) + r's?(?![a-z])', folded)] letters = set(re.findall(r'(?<![A-Za-z])([A-Z])(?![A-Za-z])', text))+ known = {word.casefold() for word in words} | _COMMON_WORDS for block in re.findall(r'(?<![A-Za-z])([A-Z]{2,26})(?![A-Za-z])', text):- if block.casefold() not in {word.casefold() for word in words}:+ if block.casefold() not in known: letters.update(block) tokens = re.findall(r'[A-Za-z]+', text) if tokens and all(len(token) == 1 for token in tokens):
httpx 声明从 <0.28 改为 >=0.28,与本地实际环境对齐(CI 不再装到不同版本)
@@ -1,5 +1,7 @@ fastapi==0.109.0 uvicorn[standard]==0.27.0 websockets>=10.4,<14-httpx<0.28+# 0.28 移除了 Client(app=...) 旧参数,测试里已用兼容 shim 适配;+# 上限写成 <0.28 会让 CI 装到与本地不同的版本。+httpx>=0.28,<1 pytest>=8.0
重连中断后复位 _connecting(不再永久卡死);切回前台自动接回通话;移除通话中的播放控件;结束失败提示能显示;runtime 只构造一次
@@ -132,10 +132,13 @@ class _LiveLessonEntryState extends State<LiveLessonEntry> { setState(() => _needsParent = true); return; }+ // 只构造一次:放进 builder 的话,路由每次重建都会新建一整套+ // transport/录音/播放对象,而它们谁也不会被 dispose。+ final runtime = config.createRuntime(); await Navigator.of(context).push( MaterialPageRoute<void>( builder: (_) => FirstLessonPage(- runtime: config.createRuntime(),+ runtime: runtime, restartRuntime: (sessionId) => config.createRuntime(resumeSessionId: sessionId), parentReportRepository: config.reportRepository,@@ -212,6 +215,7 @@ class _FirstLessonPageState extends State<FirstLessonPage> bool _ended = false; String? _connectionError; bool _openingParent = false;+ bool _pausedDuringCall = false; int _acknowledgedSequence = -1; int _transitionEpoch = 0; @@ -237,10 +241,16 @@ class _FirstLessonPageState extends State<FirstLessonPage> void didChangeAppLifecycleState(AppLifecycleState state) { if (state == AppLifecycleState.paused || state == AppLifecycleState.detached) {+ // 只有通话中的暂停才需要在回到前台时自动接回;主动结束或进家长区不算。+ _pausedDuringCall = !_ended && widget.restartRuntime != null; _transitionEpoch++; _connectionError = '应用已暂停,麦克风已停止。可以重新连接。'; unawaited(runtime?.close().catchError((Object _) {})); if (mounted) setState(() {});+ } else if (state == AppLifecycleState.resumed && _pausedDuringCall) {+ _pausedDuringCall = false;+ // 像电话一样接回来:不让孩子为了继续说话再点一次开始。+ unawaited(_restart(resume: true)); } } @@ -249,20 +259,31 @@ class _FirstLessonPageState extends State<FirstLessonPage> if (factory == null || _connecting) return; final epoch = ++_transitionEpoch; setState(() => _connecting = true);- final old = runtime;- final id = resume ? old?.realtime.gateway.orchestrator.sessionId : null;- old?.realtime.removeListener(_refreshBoard);- await old?.close();- if (!mounted || epoch != _transitionEpoch) return;- setState(() {- _runtime = factory(id);- _acknowledgedSequence = -1;- _ended = false;- _connecting = true;- _connectionError = null;- });- runtime!.realtime.addListener(_refreshBoard);- await _startRuntime();+ try {+ final old = runtime;+ final id = resume ? old?.realtime.gateway.orchestrator.sessionId : null;+ old?.realtime.removeListener(_refreshBoard);+ await old?.close();+ if (!mounted || epoch != _transitionEpoch) return;+ setState(() {+ _runtime = factory(id);+ _acknowledgedSequence = -1;+ _ended = false;+ _connecting = true;+ _connectionError = null;+ });+ runtime!.realtime.addListener(_refreshBoard);+ await _startRuntime();+ } catch (_) {+ if (mounted) {+ setState(() => _connectionError = '无法重新连接,请检查网络后重试。');+ }+ } finally {+ // 等待旧会话关闭期间若被切后台或点了结束(epoch 变化),或 close 抛错,+ // 都必须复位 _connecting:否则上面的守卫会吞掉之后所有重启点击,+ // 界面停在“点一下,我们接着玩”却永远点不动。+ if (mounted && _connecting) setState(() => _connecting = false);+ } } void _refreshBoard() {@@ -316,6 +337,7 @@ class _FirstLessonPageState extends State<FirstLessonPage> await runtime?.close(); } catch (_) { _connectionError = '音频释放失败,请退出应用后重新开始。';+ if (mounted) setState(() {}); } } @@ -389,14 +411,8 @@ class _FirstLessonPageState extends State<FirstLessonPage> : completed ? () => _restart(resume: false) : null,- onReplay:- !waiting &&- !failed &&- !completed &&- !runtime!.realtime.teacherResponding &&- (runtime?.audio?.canReplay ?? false)- ? () => runtime!.audio!.replayLastReply()- : null,+ // 通话中不提供重听按钮:孩子可以直接说话让老师重说,界面上多一个+ // 播放控件会被当成“每轮都要点播放”的播放器。 onParent: _openParent, onEnd: _ended ? null : _endRuntime, completed: completed,
单帧发送/播放失败不再挂断整通电话,连续播放失败才上报
@@ -39,6 +39,8 @@ class RealtimeAudioBridge { List<int> _lastReply = const []; bool _replyTooLong = false; static const _maxReplayBytes = 24000 * 2 * 30;+ static const _maxOutputFailures = 5;+ int _consecutiveOutputFailures = 0; bool get canReplay => !_stopped && _lastReply.isNotEmpty; @@ -50,7 +52,7 @@ class RealtimeAudioBridge { await output.clear(); if (!_stopped) await output.play(_lastReply, sampleRate: 24000); })- .catchError(_error);+ .catchError(_droppedFrame); await _outputPending; } @@ -64,16 +66,30 @@ class RealtimeAudioBridge { await input.start(); } + /// 致命音频故障:麦克风流本身已停止,通话无法继续。 void _error(Object error) { onError?.call(error); } + /// 单帧发送/清空失败。后续帧仍会继续,且连接中断由 transport 自己上报,+ /// 因此这里不能走 onError——否则一次音频抖动就会挂断整通电话。+ void _droppedFrame(Object _) {}++ /// 播放失败:单次交给下一帧重试,连续失败才说明输出设备真的不可用+ /// (例如被其他应用的音频会话占用)。+ void _playbackFailed(Object error) {+ _consecutiveOutputFailures++;+ if (_consecutiveOutputFailures < _maxOutputFailures) return;+ _consecutiveOutputFailures = 0;+ onError?.call(error);+ }+ void _enqueueInput(List<int> pcm16) { _inputPending = _inputPending .then<void>((_) async { if (!_stopped) await transport.sendAudio(pcm16); })- .catchError(_error);+ .catchError(_droppedFrame); } void _enqueueEvent(RealtimeEvent event) {@@ -84,7 +100,7 @@ class RealtimeAudioBridge { _playbackGeneration++; _outputPending = _outputPending .then((_) => output.clear())- .catchError(_error);+ .catchError(_droppedFrame); return; } if (event.type == RealtimeEventType.audioDone) {@@ -114,9 +130,10 @@ class RealtimeAudioBridge { event.audio!, sampleRate: event.sampleRate ?? 24000, );+ _consecutiveOutputFailures = 0; } })- .catchError(_error);+ .catchError(_playbackFailed); } Future<void> drain() async {
事件链单次异常不再永久静默丢弃后续所有事件
@@ -621,7 +621,13 @@ class RealtimeSessionController extends ChangeNotifier { } void _enqueue(RealtimeEvent event) {- _pending = _pending.then((_) => _handle(event));+ // 单次处理失败不能中断链条:否则之后所有 board/audio/transcript 事件+ // 都会被静默丢弃,界面看着仍在通话却再也收不到任何内容。+ _pending = _pending+ .then((_) => _handle(event))+ .catchError((Object error) {+ reportError('事件处理失败:$error');+ }); } Future<void> _handle(RealtimeEvent event) async {
播放采样率以每次调用为准(服务端换采样率不再变调);清理时复位采样率状态
@@ -69,19 +69,20 @@ class RecordPcmAudioInput implements PcmAudioInput { } class FlutterPcmAudioOutput implements PcmAudioOutput {- FlutterPcmAudioOutput({this.sampleRate = 24000});-- final int sampleRate; bool _ready = false;+ int? _activeSampleRate; @override Future<void> play(List<int> pcm16, {required int sampleRate}) async {- if (!_ready) {+ // 采样率以每次调用传入的为准。写死一个默认值的话,服务端换了采样率+ // 孩子听到的是变调的声音,而且不会有任何报错。+ if (!_ready || _activeSampleRate != sampleRate) { await pcm.FlutterPcmSound.setup(- sampleRate: this.sampleRate,+ sampleRate: sampleRate, channelCount: 1, iosAudioCategory: pcm.IosAudioCategory.playAndRecord, );+ _activeSampleRate = sampleRate; _ready = true; if (!kIsWeb && defaultTargetPlatform == TargetPlatform.iOS) { try {@@ -104,11 +105,13 @@ class FlutterPcmAudioOutput implements PcmAudioOutput { if (!_ready) return; await pcm.FlutterPcmSound.release(); _ready = false;+ _activeSampleRate = null; } Future<void> dispose() async { if (!_ready) return; _ready = false;+ _activeSampleRate = null; await pcm.FlutterPcmSound.release(); }
release 构建拒绝明文 http,避免会话 token 走 ws 被嗅探
@@ -1,4 +1,5 @@ import 'dart:math';+import 'package:flutter/foundation.dart' show kReleaseMode; import '../learning/parent_report_repository.dart'; import '../learning/growth_plan_repository.dart'; import '../learning/recent_context_repository.dart';@@ -29,7 +30,8 @@ class LiveConfig { final uri = Uri.tryParse(apiUrl); final id = RegExp(r'^[a-zA-Z0-9_-]{1,100}$'); if (uri == null ||- !['https', 'http'].contains(uri.scheme) ||+ // release 构建里 http 会让会话 token 明文走 ws,同网段可嗅探重放。+ !(kReleaseMode ? const ['https'] : const ['https', 'http']).contains(uri.scheme) || uri.host.isEmpty || uri.userInfo.isNotEmpty || uri.hasQuery ||
HTTP 请求加超时(家长页不再无限转圈)
@@ -3,8 +3,11 @@ import 'dart:io'; import 'memory.dart'; + export 'memory.dart' show ParentGrowthReport; +const _httpTimeout = Duration(seconds: 10);+ typedef ReportJsonFetcher = Future<Map<String, Object?>> Function(Uri uri, Map<String, String> headers); @@ -40,12 +43,12 @@ class ApiParentReportRepository implements ParentReportRepository { }; Future<Map<String, Object?>> _fetchHttp(Uri uri) async {- final client = HttpClient();+ final client = HttpClient()..connectionTimeout = _httpTimeout; try {- final request = await client.getUrl(uri);+ final request = await client.getUrl(uri).timeout(_httpTimeout); _headers.forEach(request.headers.add);- final response = await request.close();- final body = await response.transform(utf8.decoder).join();+ final response = await request.close().timeout(_httpTimeout);+ final body = await response.transform(utf8.decoder).join().timeout(_httpTimeout); if (response.statusCode != 200) { throw HttpException('成长报告请求失败:${response.statusCode}', uri: uri); }
HTTP 请求加超时
@@ -1,6 +1,9 @@ import 'dart:convert'; import 'dart:io'; +const _httpTimeout = Duration(seconds: 10);++ class GrowthStage { const GrowthStage({ required this.id,@@ -80,12 +83,12 @@ class ApiGrowthPlanRepository implements GrowthPlanRepository { @override Future<GrowthPlan> load() async {- final client = HttpClient();+ final client = HttpClient()..connectionTimeout = _httpTimeout; final uri = baseUri.resolve('/api/v1/children/$childId/growth-plan'); try {- final request = await client.getUrl(uri);+ final request = await client.getUrl(uri).timeout(_httpTimeout); request.headers.add('x-session-token', sessionToken);- final response = await request.close();+ final response = await request.close().timeout(_httpTimeout); final decoded = jsonDecode(await response.transform(utf8.decoder).join()); if (response.statusCode != 200 || decoded is! Map) { throw HttpException('成长路线请求失败:${response.statusCode}', uri: uri);
HTTP 请求加超时
@@ -1,6 +1,9 @@ import 'dart:convert'; import 'dart:io'; +const _httpTimeout = Duration(seconds: 10);++ typedef RecentContextJsonFetcher = Future<Map<String, Object?>> Function(Uri uri, Map<String, String> headers); @@ -177,12 +180,12 @@ class ApiRecentContextRepository implements RecentContextRepository { } Future<Map<String, Object?>> _fetchHttp(Uri uri) async {- final client = HttpClient();+ final client = HttpClient()..connectionTimeout = _httpTimeout; try {- final request = await client.getUrl(uri);+ final request = await client.getUrl(uri).timeout(_httpTimeout); _headers.forEach(request.headers.add);- final response = await request.close();- final body = await response.transform(utf8.decoder).join();+ final response = await request.close().timeout(_httpTimeout);+ final body = await response.transform(utf8.decoder).join().timeout(_httpTimeout); if (response.statusCode != 200) { throw HttpException('近期对话请求失败:${response.statusCode}', uri: uri); }
HTTP 请求加超时
@@ -1,6 +1,9 @@ import 'dart:convert'; import 'dart:io'; +const _httpTimeout = Duration(seconds: 10);++ typedef NextActivityJsonFetcher = Future<Map<String, Object?>> Function(Uri uri, Map<String, String> headers); @@ -97,12 +100,12 @@ class ApiNextActivityRepository implements NextActivityRepository { } Future<Map<String, Object?>> _fetchHttp(Uri uri) async {- final client = HttpClient();+ final client = HttpClient()..connectionTimeout = _httpTimeout; try {- final request = await client.getUrl(uri);+ final request = await client.getUrl(uri).timeout(_httpTimeout); _headers.forEach(request.headers.add);- final response = await request.close();- final body = await response.transform(utf8.decoder).join();+ final response = await request.close().timeout(_httpTimeout);+ final body = await response.transform(utf8.decoder).join().timeout(_httpTimeout); if (response.statusCode != 200) { throw HttpException('下一步计划请求失败:${response.statusCode}', uri: uri); }
totalStars 校验带上真实字段名,报错不再指向 usageSeconds
@@ -129,17 +129,17 @@ class ParentGrowthReport { ), recentItems: _stringList(json['recentItems'], 'recentItems'), evidenceSummary: _evidenceList(json['evidenceSummary']),- usageSeconds: _usageSeconds(json['usageSeconds']),- totalStars: _usageSeconds(json['totalStars']),+ usageSeconds: _nonNegativeInt(json['usageSeconds'], 'usageSeconds'),+ totalStars: _nonNegativeInt(json['totalStars'], 'totalStars'), phraseItems: json['phraseItems'] == null ? const [] : _stringList(json['phraseItems'], 'phraseItems'), capabilities: CapabilityObservation.parseList(json['capabilityEvents']), ); } - static int _usageSeconds(Object? value) {+ static int _nonNegativeInt(Object? value, String field) { if (value == null) return 0; if (value is! int || value < 0) {- throw const FormatException('usageSeconds 必须是非负整数。');+ throw FormatException('$field 必须是非负整数。'); } return value; }
新增:持续通话规格(普通回答/点击/旧停止请求都不能挂断)
import pytestfrom app.session import SessionOrchestrator, InvalidSessionAction def finish(session): return session.apply_tool({'toolCallId':'finish','name':'completeActivity','arguments':{'expectedSequence':session.state.sequence}}) @pytest.mark.parametrize('answer,source', [('cat','speech'),('cat','choice'),("I don't want to stop",'speech'),('stop','text')])def test_task_response_cannot_end_a_continuous_call(answer,source): session=SessionOrchestrator(session_id='call',child_id='c') if source == 'choice': session.apply_tool({'toolCallId':'board','name':'showChoices','arguments':{'expectedSequence':0,'choices':['A','B']}}) session.record_input('A',source=source,expected_sequence=session.state.sequence) else: session.record_input(answer,source=source) with pytest.raises(InvalidSessionAction,match='明确要求结束'): finish(session) assert session.state.status=='active' def test_child_can_explicitly_hang_up(): session=SessionOrchestrator(session_id='call',child_id='c') session.record_input('我想休息一下',source='speech') assert finish(session).state.status=='completed' def test_old_stop_request_cannot_override_a_new_response(): session=SessionOrchestrator(session_id='call',child_id='c') session.record_input('停一下',source='speech') session.record_input('我们继续吧',source='speech') with pytest.raises(InvalidSessionAction,match='明确要求结束'): finish(session)
新增:三轮无需点播放、返回前台自动重连
import 'package:flutter/material.dart';import 'package:flutter_test/flutter_test.dart';import 'package:mobile/main.dart';import 'package:mobile/session/audio_pipeline.dart';import 'package:mobile/session/first_lesson_session.dart';import 'package:mobile/session/lesson_runtime.dart';import 'package:mobile/session/orchestrator.dart';import 'package:mobile/session/qwen_gateway.dart';import 'package:mobile/session/realtime_transport.dart'; LessonRuntime makeCall(FakeRealtimeTransport transport, FakePcmAudioInput input) => LessonRuntime( lesson: FirstLessonSession(), serverDriven: true, realtime: RealtimeSessionController(transport: transport, gateway: FakeQwenGateway(orchestrator: FakeSessionOrchestrator(sessionId: 'call'))), audio: RealtimeAudioBridge(input: input, output: FakePcmAudioOutput(), transport: transport),); /// 关闭与重连是 fire-and-forget 的异步链:它一半在 widget 测试的 fake async/// 时钟里,一半要真实事件循环才推进,pumpAndSettle 推不动。交替推两边界,/// 直到链条落地。Future<void> flushAsync(WidgetTester tester) async { for (var i = 0; i < 5; i++) { await tester.runAsync( () => Future<void>.delayed(const Duration(milliseconds: 10)), ); await tester.pump(); }} void main() { testWidgets('one connection keeps listening across turns without playback controls', (tester) async { final transport=FakeRealtimeTransport(); final input=FakePcmAudioInput(); final runtime=makeCall(transport,input); await tester.pumpWidget(MaterialApp(home: FirstLessonPage(runtime:runtime))); await tester.pumpAndSettle(); for(var turn=0;turn<3;turn++) { transport.emit(RealtimeEvent.boardState(BoardState(sequence:turn,prompt:'聊聊你喜欢什么'),sessionId:'call')); input.emit([turn,0]); transport.emit(const RealtimeEvent.audio([1,0,2,0],sampleRate:24000)); transport.emit(const RealtimeEvent.audioDone()); await tester.pumpAndSettle(); expect(transport.connected,isTrue); expect(input.stopped,isFalse); expect(find.byKey(const ValueKey('child-start')),findsNothing); expect(find.byKey(const ValueKey('child-replay')),findsNothing); } expect(transport.sentAudio.length,3); await tester.tap(find.byKey(const ValueKey('child-end'))); await tester.pumpAndSettle(); await flushAsync(tester); expect(input.stopped,isTrue); }); testWidgets('returning to an active call reconnects without another start tap', (tester) async { final input=FakePcmAudioInput(); final first=FakeRealtimeTransport(); final second=FakeRealtimeTransport(); var restarts=0; await tester.pumpWidget(MaterialApp(home:FirstLessonPage( runtime:makeCall(first,input),restartRuntime:(id){ // 这个回调在 pumpAndSettle 进行中被异步调用,只能用非守卫版本。 expectSync(id,'call'); restarts++; return makeCall(second,FakePcmAudioInput()); }))); await tester.pumpAndSettle(); tester.binding.handleAppLifecycleStateChanged(AppLifecycleState.paused); await tester.pumpAndSettle(); await flushAsync(tester); expect(input.stopped,isTrue); tester.binding.handleAppLifecycleStateChanged(AppLifecycleState.resumed); await tester.pumpAndSettle(); await flushAsync(tester); expect(restarts,1); expect(second.connected,isTrue); await tester.tap(find.byKey(const ValueKey('child-end'))); await tester.pumpAndSettle(); tester.binding.handleAppLifecycleStateChanged(AppLifecycleState.paused); tester.binding.handleAppLifecycleStateChanged(AppLifecycleState.resumed); await tester.pumpAndSettle(); await flushAsync(tester); expect(restarts,1); });}
新增字母抽取与读写插播测试(实测:破坏实现后确实变红)
@@ -36,7 +36,8 @@ def test_speech_capability_is_durable_separate_from_word_mastery_and_advances(tm answer=session.record_input('dog',source='speech') record(session,answer,session.content.activity(session.activity_id).criteria) assert memory.parent_report('c')['independentSpeakingItems']==['dog']- assert memory.learning_snapshot('c')['masteredItems']==[]+ # masteredItems 目前恒为空,这行守住的是"口语证据不会变成词汇掌握"。+ assert 'dog' not in memory.learning_snapshot('c')['masteredItems'] skill=memory.capability_events('c')[0] assert skill['outcome']=='INDEPENDENT_OBSERVED' assert skill['variantId']=='f.speak.practice'@@ -103,3 +104,31 @@ def test_reading_recheck_does_not_block_oral_progress(): for index,skill in enumerate(['f.speak','f.interact','f.resource']): events.append({'eventId':f'oral-{index}','skillId':skill,'sessionId':'s','variantId':skill+'.practice','outcome':'PROMPTED_OBSERVED','createdAt':now}) assert capability_activity_plan(events,StarterContentRegistry())['skillId']=='f.read'+++def make_event(index, skill, outcome='PROMPTED_OBSERVED', **extra):+ from datetime import datetime, timezone+ return {'eventId': f'event-{index}', 'skillId': skill, 'sessionId': 's',+ 'variantId': skill + '.practice', 'outcome': outcome,+ 'createdAt': datetime.now(timezone.utc).isoformat(), **extra}+++def test_three_oral_tasks_interleave_a_reading_task():+ # 没有这条插播规则时 select_capability 选的是 f.listen,所以断言能守住它。+ events = [make_event(i, skill)+ for i, skill in enumerate(['f.speak', 'f.interact', 'f.resource'])]+ plan = capability_activity_plan(events, StarterContentRegistry())+ assert plan['skillId'] == 'f.read'+ assert plan['targetDomain'] == 'recognition'+++def test_letters_already_heard_are_not_served_first():+ # 见过 A、且最近一条读写事件不是 f.read,才会走插播分支。+ events = [make_event(0, 'f.read', 'INDEPENDENT_OBSERVED', observedResource='A'),+ make_event(1, 'a1.read')]+ events += [make_event(i + 2, skill)+ for i, skill in enumerate(['f.speak', 'f.interact', 'f.resource'])]+ plan = capability_activity_plan(events, StarterContentRegistry())+ assert plan['skillId'] == 'f.read'+ # 字母抽取一旦失效,这里会退回从 A 重新开始——孩子会反复被发到同一个字母。+ assert plan['targetItems'][0] != 'A'
读写门控改为强断言;幽灵变体断言精确化;补齐 5 个从未执行的目录校验分支
@@ -143,18 +143,15 @@ def test_duplicate_event_id_counted_once(synthetic): # --- 坏 variant / 伪造日期 ----------------------------------------------------- def test_event_with_unknown_variant_is_ignored(synthetic):- """反例:variantId 不在目录里的证据不应算独立证据。-- 当前实现只校验 variantId 是非空字符串,幽灵变体(含其它技能的变体 id)- 照样计入 session 与独立计数,会直接解锁迁移。- """+ """variantId 不在目录里的证据不算独立证据(否则两条就解锁迁移)。""" catalog = synthetic([_node('x.speak')]) events = [ _event('e1', 'x.speak', 's1', 'ghost.variant', 'INDEPENDENT_OBSERVED', NOW - timedelta(hours=2)), _event('e2', 'x.speak', 's2', 'ghost.variant', 'INDEPENDENT_OBSERVED', NOW - timedelta(hours=1)), ] chosen = select_capability(events, now=NOW, catalog=catalog)- assert chosen['reason'] != 'TRANSFER_READY', chosen+ # 两条幽灵证据都被忽略:没有产生任何 session,所以停在最初的新技能。+ assert chosen['reason'] == 'NEW_SKILL', chosen def test_forged_transfer_variant_suppresses_transfer(synthetic):@@ -234,15 +231,17 @@ def test_literacy_and_writing_never_selected_by_default(): assert catalog.node(chosen['skillId'])['track'] != 'literacy' -def test_literacy_selectable_when_explicitly_enabled():- catalog = CapabilityCatalog()+def test_literacy_selectable_when_explicitly_enabled(synthetic):+ # 目录里只剩读写节点:显式启用才选得到,默认闸门会判为无可选。+ catalog = synthetic([_node('x.read', domain='reading', track='literacy')]) chosen = select_capability([], now=NOW, literacy_enabled=True, catalog=catalog)- assert catalog.node(chosen['skillId'])['track'] in (- 'literacy', 'listening', 'expression', 'interaction', 'repair', 'resources')+ assert chosen['skillId'] == 'x.read'+ with pytest.raises(ValueError):+ select_capability([], now=NOW, literacy_enabled=False, catalog=catalog) -def test_writing_node_on_non_literacy_track_leaks_by_default(synthetic):- """反例:默认闸门只看 track,不看 domain;writing 节点挂非 literacy 轨道会被默认选中。"""+def test_writing_node_on_non_literacy_track_is_still_blocked(synthetic):+ """默认闸门同时看 track 和 domain:writing 节点即使挂别的轨道也不会被选中。""" catalog = synthetic([_node('x.write', domain='writing', track='expression'), _node('x.speak', domain='speaking', track='expression')]) chosen = select_capability([], now=NOW, literacy_enabled=False, catalog=catalog) assert chosen['skillId'] != 'x.write', chosen@@ -375,7 +374,17 @@ def test_node_lookup_raises_value_error(): ([_node('x', hard=('nope',))], 'unknown prerequisite'), ([_node('x', hard=('y',)), _node('y', hard=('x',))], 'cycle'), ([_node('x', hard=('x',))], 'self cycle'),- ([{'id': 'x', 'stage': 'A1', 'domain': 'speaking'}], 'missing tasks'),+ ([{'id': 'x', 'stage': 'A1', 'domain': 'speaking', 'track': 'expression'}], 'missing tasks'),+ (['not-a-dict'], 'node is not an object'),+ ([{'stage': 'A1', 'track': 'expression', 'domain': 'speaking'}], 'missing id'),+ ([{'id': 'x', 'stage': 'A1', 'track': 'expression', 'domain': 'speaking',+ 'prerequisites': {'hard': 'not-a-list'},+ 'tasks': {'practice': {'variantId': 'x.practice', 'scenario': 's'},+ 'transfer': {'variantId': 'x.transfer', 'scenario': 's'}}}], 'hard is not a list'),+ ([{'id': 'x', 'stage': 'A1', 'track': 'expression', 'domain': 'speaking',+ 'prerequisites': {'hard': [{}]},+ 'tasks': {'practice': {'variantId': 'x.practice', 'scenario': 's'},+ 'transfer': {'variantId': 'x.transfer', 'scenario': 's'}}}], 'hard entry without skillId'), ]) def test_invalid_catalog_rejected(tmp_path, nodes, label): with pytest.raises(ValueError):
适配带 code 的关闭;health 断言 upstream 模式;新增生产 upstream 门禁测试
@@ -28,11 +28,18 @@ def test_health_reports_local_gateway_boundary(client): response = client.get("/health") assert response.status_code == 200- assert response.json() == {- "ok": True,- "service": "english-app-gateway",- "authMode": "development-shared-token",- }+ body = response.json()+ assert body["ok"] is True+ assert body["service"] == "english-app-gateway"+ assert body["authMode"] == "development-shared-token"+ # 部署后能从这里看出是不是静默跑成了 fake 复读机。+ assert body["upstreamMode"] in {"fake", "dashscope", "injected"}+++def test_production_refuses_to_start_without_an_explicit_upstream(monkeypatch):+ monkeypatch.delenv("ENGLISH_APP_UPSTREAM", raising=False)+ with pytest.raises(RuntimeError, match="ENGLISH_APP_UPSTREAM"):+ create_app(session_token="t", environment="production", auth_secret="s") def test_parent_report_reads_persisted_evidence_with_child_isolation():@@ -204,9 +211,10 @@ def test_parent_report_includes_ready_usage_seconds(client): def test_realtime_requires_the_app_session_token(client):- with pytest.raises(WebSocketDisconnect) as error:- with client.websocket_connect("/api/realtime?sessionId=session-1"):- pass+ # 先接受握手再带 code 关闭:客户端才拿得到 4401,而不是一个空的 HTTP 403。+ with client.websocket_connect("/api/realtime?sessionId=session-1") as websocket:+ with pytest.raises(WebSocketDisconnect) as error:+ websocket.receive_json() assert error.value.code == 4401 @@ -298,10 +306,10 @@ def test_unknown_activity_is_rejected_before_upstream_is_created(): def forbidden(*args): pytest.fail('unknown content must not start the model') with TestClient(create_app(session_token='test-token', upstream_factory=forbidden)) as app_client:- with pytest.raises(WebSocketDisconnect) as error:- with app_client.websocket_connect('/api/realtime?activityId=unknown',- headers={'x-session-token': 'test-token'}):- pass+ with app_client.websocket_connect('/api/realtime?activityId=unknown',+ headers={'x-session-token': 'test-token'}) as websocket:+ with pytest.raises(WebSocketDisconnect) as error:+ websocket.receive_json() assert error.value.code == 4400 @@ -422,6 +430,9 @@ def test_first_lesson_route_choice_observation_and_report_survive_restart(tmp_pa self.sequence = max(self.sequence, 1) await self.tool('recordObservation', item='cat', evidenceId=evidence['evidenceId'], domains={'listening': 'INDEPENDENT_OBSERVED', 'speaking': 'INDEPENDENT_OBSERVED'})+ # 孩子接着说一句明确要结束,服务端才允许完成整通电话。+ await self.send({'type': 'transcript.final', 'role': 'user', 'content': '我想休息一下'})+ elif event['type'] == 'input.evidence' and event['evidence'].get('endRequested'): await self.tool('completeActivity') async def close(self): pass@@ -440,9 +451,18 @@ def test_first_lesson_route_choice_observation_and_report_survive_restart(tmp_pa assert routed['activityId'] == 'starter.zero.picture.v1' ws.send_json({'type': 'child.choice', 'choiceId': 'cat', 'expectedSequence': routed['sequence'] - 1}) assert ws.receive_json()['type'] == 'input.rejected'+ def next_board_state():+ # 上游事件会被原样转发给客户端,所以按类型取黑板快照,+ # 而不是假定下一条一定就是它。+ while True:+ message = ws.receive_json()+ assert message['type'] != 'error', message+ if message['type'] == 'board.state':+ return message+ ws.send_json({'type': 'child.choice', 'choiceId': 'cat', 'expectedSequence': routed['sequence']})- assert ws.receive_json()['sequence'] == 2- assert ws.receive_json()['status'] == 'completed'+ assert next_board_state()['sequence'] == 2+ assert next_board_state()['status'] == 'completed' with TestClient(create_app(session_token='test-token', persist=True, data_path=path)) as client: report = client.get('/api/v1/children/c/report', headers={'x-session-token':'test-token'}).json()
完成活动改用明确结束语
@@ -87,7 +87,7 @@ def test_completed_activity_rejects_new_changes_but_allows_successful_retry(): payload = {'toolCallId': 'end', 'name': 'completeActivity', 'arguments': {'expectedSequence': 0}} with pytest.raises(InvalidSessionAction, match='完成活动需要孩子先参与'): session.apply_tool(payload)- session.record_input('apple', source='speech')+ session.record_input('我想休息一下', source='speech') result = session.apply_tool(payload) with pytest.raises(InvalidSessionAction): call(session, 'setPrompt', englishText='change completed board')
断言改为新规格:单帧失败不挂断、连续失败才上报
@@ -47,30 +47,49 @@ void main() { }, ); - test(- 'audio send and playback failures are reported and resources still stop',- () async {- final errors = <Object>[];- final transport = BrokenTransport();- final input = FakePcmAudioInput();- final output = BrokenOutput();- final bridge = RealtimeAudioBridge(- input: input,- output: output,- transport: transport,- onError: errors.add,- );- await bridge.start();- input.emit([0, 0]);+ test('a failed frame keeps the call alive and resources still stop', () async {+ final errors = <Object>[];+ final transport = BrokenTransport();+ final input = FakePcmAudioInput();+ final output = BrokenOutput();+ final bridge = RealtimeAudioBridge(+ input: input,+ output: output,+ transport: transport,+ onError: errors.add,+ );+ await bridge.start();+ input.emit([0, 0]);+ transport.emit(const RealtimeEvent.audio([0, 0], sampleRate: 24000));+ await bridge.drain();+ // 单帧上行/播放失败必须被吞掉:连接状态由 transport 负责上报,+ // 若这里上报就等于一次音频抖动挂断整通电话。+ expect(errors, isEmpty);+ await Future.wait([bridge.stop(), bridge.stop()]);+ expect(input.stopped, isTrue);+ expect(output.stopped, isTrue);+ await transport.close();+ });++ test('repeated playback failures are reported', () async {+ final errors = <Object>[];+ final transport = FakeRealtimeTransport();+ final bridge = RealtimeAudioBridge(+ input: FakePcmAudioInput(),+ output: BrokenOutput(),+ transport: transport,+ onError: errors.add,+ );+ await bridge.start();+ for (var attempt = 0; attempt < 5; attempt++) { transport.emit(const RealtimeEvent.audio([0, 0], sampleRate: 24000));- await bridge.drain();- expect(errors, hasLength(2));- await Future.wait([bridge.stop(), bridge.stop()]);- expect(input.stopped, isTrue);- expect(output.stopped, isTrue);- await transport.close();- },- );+ }+ await bridge.drain();+ // 连续失败说明输出设备真的不可用,此时才上报。+ expect(errors, hasLength(1));+ await bridge.stop();+ await transport.close();+ }); test( 'bridges captured PCM to the realtime transport and plays model audio', () async {