diff --git a/bridge/lan_service.py b/bridge/lan_service.py index 04395de..f643ba5 100644 --- a/bridge/lan_service.py +++ b/bridge/lan_service.py @@ -113,6 +113,8 @@ MAX_TEXT_BYTES = 65535 DEFAULT_SAMPLE_RATE = 16000 DEFAULT_MAX_AUDIO_BYTES = 512 * 1024 +DEFAULT_AUDIO_CAPTURE_ABSOLUTE_LEASE_MS = 14_500 +DEFAULT_AUDIO_CAPTURE_INACTIVITY_LEASE_MS = 4_000 DEFAULT_DOWNLINK_AUDIO_CHUNK_BYTES = 4096 DEFAULT_DOWNLINK_BINARY_FRAME_DELAY_MS = 180 DEFAULT_DOWNLINK_TEXT_FRAME_DELAY_MS = 40 @@ -791,6 +793,8 @@ class LanBridgeConfig: client_idle_timeout_s: float = DEFAULT_CLIENT_IDLE_TIMEOUT_S disable_audio_downlink: bool = False max_audio_bytes: int = DEFAULT_MAX_AUDIO_BYTES + audio_capture_absolute_lease_ms: int = DEFAULT_AUDIO_CAPTURE_ABSOLUTE_LEASE_MS + audio_capture_inactivity_lease_ms: int = DEFAULT_AUDIO_CAPTURE_INACTIVITY_LEASE_MS audio_evidence_dir: Path | None = None memory_file: Path | None = None turn_log_file: Path | None = None @@ -851,6 +855,14 @@ def __post_init__(self) -> None: ) if not 0.0 <= float(self.stt_min_confidence) <= 1.0: raise ValueError("stt_min_confidence must be between zero and one") + if not 12_500 <= int(self.audio_capture_absolute_lease_ms) <= 15_000: + raise ValueError("audio_capture_absolute_lease_ms must be between 12500 and 15000") + if not 1_000 <= int(self.audio_capture_inactivity_lease_ms) < int( + self.audio_capture_absolute_lease_ms + ): + raise ValueError( + "audio_capture_inactivity_lease_ms must be between 1000 and the absolute lease" + ) if self.stt_diagnostic_expected_text: validate_expected_text(self.stt_diagnostic_expected_text) if self.turn_log_file is None: @@ -883,26 +895,56 @@ class AudioUpload: chunks: int = 0 truncated: bool = False buffer: bytearray = field(default_factory=bytearray) - started_at_monotonic: float = 0.0 + seq: int = 0 + absolute_lease_ms: int = DEFAULT_AUDIO_CAPTURE_ABSOLUTE_LEASE_MS + inactivity_lease_ms: int = DEFAULT_AUDIO_CAPTURE_INACTIVITY_LEASE_MS + started_at_monotonic: float | None = None + last_activity_at_monotonic: float | None = None - def start(self, sample_rate: object = DEFAULT_SAMPLE_RATE) -> None: + def start( + self, + sample_rate: object = DEFAULT_SAMPLE_RATE, + *, + seq: int = 0, + absolute_lease_ms: int = DEFAULT_AUDIO_CAPTURE_ABSOLUTE_LEASE_MS, + inactivity_lease_ms: int = DEFAULT_AUDIO_CAPTURE_INACTIVITY_LEASE_MS, + ) -> None: self.clear() try: parsed_rate = int(sample_rate) except (TypeError, ValueError): parsed_rate = DEFAULT_SAMPLE_RATE self.sample_rate = max(8000, min(48000, parsed_rate)) + self.seq = max(0, int(seq)) + self.absolute_lease_ms = max(1, int(absolute_lease_ms)) + self.inactivity_lease_ms = max(1, int(inactivity_lease_ms)) self.active = True self.started_at_monotonic = time.perf_counter() + self.last_activity_at_monotonic = self.started_at_monotonic def clear(self) -> None: self.active = False self.bytes_received = 0 self.chunks = 0 self.truncated = False - self.started_at_monotonic = 0.0 + self.seq = 0 + self.started_at_monotonic = None + self.last_activity_at_monotonic = None self.buffer.clear() + def expiry_code(self, at_monotonic: float | None = None) -> str: + if not self.active or self.started_at_monotonic is None: + return "" + current = time.perf_counter() if at_monotonic is None else float(at_monotonic) + absolute_elapsed_ms = (current - self.started_at_monotonic) * 1000.0 + if absolute_elapsed_ms >= self.absolute_lease_ms: + return "audio_capture_absolute_lease_expired" + if self.last_activity_at_monotonic is not None: + inactivity_elapsed_ms = (current - self.last_activity_at_monotonic) * 1000.0 + if inactivity_elapsed_ms >= self.inactivity_lease_ms: + return "audio_capture_inactivity_expired" + return "" + def append(self, payload: bytes, max_bytes: int) -> None: if not self.active: raise WebSocketProtocolError("audio received before utterance_start") @@ -913,6 +955,7 @@ def append(self, payload: bytes, max_bytes: int) -> None: self.truncated = True if allowed > 0: self.buffer.extend(payload[:allowed]) + self.last_activity_at_monotonic = time.perf_counter() @property def stored_bytes(self) -> int: @@ -932,12 +975,21 @@ def summary(self) -> dict[str, object]: "audio_sample_rate": self.sample_rate, "audio_duration_ms": self.duration_ms, "audio_truncated": self.truncated, + "audio_capture_seq": self.seq, + "audio_capture_absolute_lease_ms": self.absolute_lease_ms, + "audio_capture_inactivity_lease_ms": self.inactivity_lease_ms, } - if self.started_at_monotonic > 0.0: + if self.started_at_monotonic is not None: + current = time.perf_counter() summary["audio_capture_elapsed_ms"] = round( - (time.perf_counter() - self.started_at_monotonic) * 1000.0, + (current - self.started_at_monotonic) * 1000.0, 2, ) + if self.last_activity_at_monotonic is not None: + summary["audio_capture_inactivity_ms"] = round( + (current - self.last_activity_at_monotonic) * 1000.0, + 2, + ) return summary def finish_and_clear(self) -> dict[str, object]: @@ -1520,6 +1572,12 @@ def __init__( self.playback_response_seq = 0 self.conversation_playback_complete_seq = 0 self.audio_protocol_errors = 0 + self.audio_uploads_started = 0 + self.audio_uploads_cancelled = 0 + self.audio_uploads_expired = 0 + self.audio_stale_ends_rejected = 0 + self.audio_last_reject_code = "" + self._rejected_audio_turns: dict[int, str] = {} self._stt_diagnostic_pending = bool(config.stt_diagnostic_expected_text) saved_initiative = self.memory.fact_value("user.initiative_enabled") if self.initiative_policy is not None and saved_initiative == "false": @@ -1776,6 +1834,10 @@ def _observe_conversation_transition(self, transition) -> None: self._finalize_memory_session() def connection_closed(self) -> None: + self._cancel_audio_capture( + reason="bridge connection closed before a terminal marker", + code="audio_capture_connection_closed", + ) if self.conversation is not None and self.conversation.phase != ConversationPhase.IDLE: self._conversation_payload(self.conversation.bridge_lost()) @@ -2009,18 +2071,165 @@ def _append_audio_error_log( record[key] = audio_summary[key] self._append_turn_log(record) - def _append_audio_protocol_event(self, *, code: str, payload_bytes: int) -> None: + def _append_audio_protocol_event( + self, + *, + code: str, + payload_bytes: int, + seq: int = 0, + detail: str = "", + ) -> None: self.audio_protocol_errors += 1 - self._append_turn_log( - { - "schema": "stackchan.audio-protocol-event.v1", - "generated_at": utc_timestamp(), - "session": self.session, - "code": code, - "payload_bytes": max(0, int(payload_bytes)), - "audio_protocol_errors": self.audio_protocol_errors, - } + record: dict[str, object] = { + "schema": "stackchan.audio-protocol-event.v1", + "generated_at": utc_timestamp(), + "session": self.session, + "code": code, + "payload_bytes": max(0, int(payload_bytes)), + "audio_protocol_errors": self.audio_protocol_errors, + } + if seq > 0: + record["seq"] = seq + if detail: + record["detail"] = detail[:160] + self._append_turn_log(record) + + def _append_audio_capture_event( + self, + *, + code: str, + payload_bytes: int, + seq: int = 0, + detail: str = "", + ) -> None: + record: dict[str, object] = { + "schema": "stackchan.audio-capture-event.v1", + "generated_at": utc_timestamp(), + "session": self.session, + "code": code, + "payload_bytes": max(0, int(payload_bytes)), + } + if seq > 0: + record["seq"] = seq + if detail: + record["detail"] = detail[:160] + self._append_turn_log(record) + + @staticmethod + def _audio_message_seq(message: dict[str, Any]) -> int: + try: + return max(0, int(message.get("seq") or 0)) + except (TypeError, ValueError): + return 0 + + def _audio_upload_telemetry(self) -> dict[str, object]: + return { + "audio_upload_active": self.audio.active, + "audio_uploads_started": self.audio_uploads_started, + "audio_uploads_cancelled": self.audio_uploads_cancelled, + "audio_uploads_expired": self.audio_uploads_expired, + "audio_stale_ends_rejected": self.audio_stale_ends_rejected, + "audio_protocol_errors": self.audio_protocol_errors, + "audio_last_reject_code": self.audio_last_reject_code, + "audio_capture_absolute_lease_ms": self.config.audio_capture_absolute_lease_ms, + "audio_capture_inactivity_lease_ms": self.config.audio_capture_inactivity_lease_ms, + } + + def _remember_audio_rejection(self, seq: int, code: str) -> None: + self.audio_last_reject_code = code + if seq <= 0: + return + self._rejected_audio_turns[seq] = code + while len(self._rejected_audio_turns) > 16: + self._rejected_audio_turns.pop(next(iter(self._rejected_audio_turns))) + + def _cancel_audio_capture(self, *, reason: str, code: str) -> dict[str, object] | None: + if not self.audio.active: + return None + summary = self.audio.summary() + seq = self.audio.seq + payload_bytes = self.audio.bytes_received + self.audio.clear() + self.audio_uploads_cancelled += 1 + self._remember_audio_rejection(seq, code) + self._append_audio_capture_event( + code=code, + payload_bytes=payload_bytes, + seq=seq, + detail=reason, + ) + summary.update(self._audio_upload_telemetry()) + return summary + + def _expire_audio_capture(self) -> dict[str, object] | None: + code = self.audio.expiry_code() + if not code: + return None + summary = self.audio.summary() + seq = self.audio.seq + payload_bytes = self.audio.bytes_received + self.audio.clear() + self.audio_uploads_expired += 1 + self._remember_audio_rejection(seq, code) + self._append_audio_capture_event( + code=code, + payload_bytes=payload_bytes, + seq=seq, + detail="partial PCM discarded before STT", ) + transition = self.conversation.cancel(now_ms(), code) if self.conversation is not None else None + frame = error_frame(code, "partial PCM discarded before STT") + frame.update(summary) + frame.update(self._audio_upload_telemetry()) + frame.update(self._conversation_payload(transition)) + return frame + + def prepare_utterance_end( + self, + message: dict[str, Any], + *, + finalized_audio: FinalizedAudioUpload | None = None, + ) -> dict[str, object] | None: + seq = self._audio_message_seq(message) + if finalized_audio is not None: + finalized_reject = str( + finalized_audio.summary.get("audio_capture_reject_code", "") + ) + if finalized_reject: + frame = error_frame(finalized_reject, "partial PCM discarded before STT") + frame.update(finalized_audio.summary) + frame.update(self._audio_upload_telemetry()) + return frame + expired = self._expire_audio_capture() + if expired is not None: + return expired + if self.audio.active and self.audio.seq > 0 and seq > 0 and seq != self.audio.seq: + self.audio_stale_ends_rejected += 1 + self.audio_last_reject_code = "utterance_end_seq_mismatch" + self._append_audio_protocol_event( + code="utterance_end_seq_mismatch", + payload_bytes=0, + seq=seq, + detail=f"active capture seq is {self.audio.seq}", + ) + frame = error_frame("utterance_end_seq_mismatch", f"active capture seq is {self.audio.seq}") + frame.update(self._audio_upload_telemetry()) + return frame + rejected_code = self._rejected_audio_turns.get(seq, "") if seq > 0 else "" + if not self.audio.active and rejected_code: + self.audio_stale_ends_rejected += 1 + self.audio_last_reject_code = "utterance_end_stale" + self._append_audio_protocol_event( + code="utterance_end_stale", + payload_bytes=0, + seq=seq, + detail=rejected_code, + ) + frame = error_frame("utterance_end_stale", rejected_code) + frame["audio_capture_reject_code"] = rejected_code + frame.update(self._audio_upload_telemetry()) + return frame + return None @staticmethod def _validate_audio_end_declaration( @@ -2209,6 +2418,7 @@ def handle_text( self.endpoint_id = str(frame.get("endpoint_id", self.endpoint_id)) if frame.get("type") != "error" else self.endpoint_id return [frame] if message_type == "heartbeat": + audio_expired = self._expire_audio_capture() self.robot_embodiment.update(message) self._last_robot_heartbeat = dict(message) if ( @@ -2233,7 +2443,7 @@ def handle_text( if endpoint_id: frame["endpoint_id"] = endpoint_id frame.update(self._conversation_payload(conversation_transition)) - return [frame] + return ([audio_expired] if audio_expired is not None else []) + [frame] if message_type == "claim_brain": return [self.control_state.claim_brain(message)] if message_type == "release_brain": @@ -2250,6 +2460,9 @@ def handle_text( return [self._handle_settings_set(message)] if message_type == "diagnostics_request": frame = self.control_state.diagnostics_snapshot(self.config) + audio_diagnostics = frame.get("audio") + if isinstance(audio_diagnostics, dict): + audio_diagnostics.update(self._audio_upload_telemetry()) frame.update(self._conversation_payload()) return [frame] if message_type == "capability_update": @@ -2263,11 +2476,68 @@ def handle_text( ) if conversation_error is not None: return [conversation_error] - self.audio.start(message.get("sample_rate", DEFAULT_SAMPLE_RATE)) - return [{"type": "listening", **self.audio.summary(), **self._conversation_payload()}] + capture_seq = self._audio_message_seq(message) + self._cancel_audio_capture( + reason="a newer utterance_start superseded the capture", + code="audio_capture_superseded", + ) + if capture_seq > 0: + self._rejected_audio_turns.pop(capture_seq, None) + self.audio.start( + message.get("sample_rate", DEFAULT_SAMPLE_RATE), + seq=capture_seq, + absolute_lease_ms=self.config.audio_capture_absolute_lease_ms, + inactivity_lease_ms=self.config.audio_capture_inactivity_lease_ms, + ) + self.audio_uploads_started += 1 + return [ + { + "type": "listening", + **self.audio.summary(), + **self._audio_upload_telemetry(), + **self._conversation_payload(), + } + ] + if message_type == "utterance_cancel": + owner_error = self._owner_gate(message) + if owner_error is not None: + return [owner_error] + reason = str(message.get("reason") or "capture_cancelled") + seq = self._audio_message_seq(message) + if self.audio.active and self.audio.seq > 0 and seq > 0 and seq != self.audio.seq: + return [ + error_frame("utterance_cancel_seq_mismatch", f"active capture seq is {self.audio.seq}") + | self._audio_upload_telemetry() + ] + known_duplicate = seq > 0 and seq in self._rejected_audio_turns + summary = self._cancel_audio_capture( + reason=reason, + code="audio_capture_cancelled", + ) + transition = None + if summary is not None or not known_duplicate: + if summary is None: + self._remember_audio_rejection(seq, "audio_capture_cancelled") + self.cancel_active_turn(reason) + transition = ( + self.conversation.cancel(now_ms(), reason) + if self.conversation is not None + else None + ) + frame: dict[str, object] = { + "type": "heartbeat", + "utterance_cancelled": True, + "utterance_cancel_seq": seq, + "utterance_cancel_duplicate": summary is None, + **self._audio_upload_telemetry(), + } + if summary is not None: + frame.update(summary) + frame.update(self._conversation_payload(transition)) + return [frame] if message_type == "cancel": - self.audio.clear() reason = str(message.get("reason") or "cancelled") + self._cancel_audio_capture(reason=reason, code="audio_capture_cancelled") self.cancel_active_turn(reason) if self.conversation is not None: transition = self.conversation.cancel(now_ms(), reason) @@ -2282,6 +2552,12 @@ def handle_text( owner_error = self._owner_gate(message) if owner_error is not None: return [owner_error] + audio_end_error = self.prepare_utterance_end( + message, + finalized_audio=finalized_audio, + ) + if audio_end_error is not None: + return [audio_end_error] return self._handle_utterance_end( message, suppress_thinking=suppress_thinking, @@ -2579,6 +2855,15 @@ def early_thinking_frame(self, text: str) -> dict[str, object] | None: return frame def finalize_audio_upload(self) -> FinalizedAudioUpload: + expired = self._expire_audio_capture() + if expired is not None: + summary = { + key: value + for key, value in expired.items() + if key not in {"type", "code", "detail"} + } + summary["audio_capture_reject_code"] = str(expired["code"]) + return FinalizedAudioUpload(pcm=b"", summary=summary) return self.audio.finalize() def handle_binary(self, payload: bytes) -> list[dict[str, object]]: @@ -2588,6 +2873,9 @@ def handle_binary(self, payload: bytes) -> list[dict[str, object]]: self.control_state.reconcile_owner() if self.endpoint_id and self.control_state.active_brain_owner and self.endpoint_id != self.control_state.active_brain_owner: return [error_frame("brain_owner_mismatch", self.endpoint_id)] + expired = self._expire_audio_capture() + if expired is not None: + return [expired] try: self.audio.append(payload, self.config.max_audio_bytes) except WebSocketProtocolError as exc: @@ -2598,7 +2886,13 @@ def handle_binary(self, payload: bytes) -> list[dict[str, object]]: frame = error_frame("audio_without_utterance", str(exc)) frame["audio_protocol_errors"] = self.audio_protocol_errors return [frame] - return [{"type": "heartbeat", **self.audio.summary()}] + return [ + { + "type": "heartbeat", + **self.audio.summary(), + **self._audio_upload_telemetry(), + } + ] def _handle_text_audio(self, message: dict[str, Any]) -> list[dict[str, object]]: encoded = str(message.get("pcm_b64") or message.get("audio_b64") or "").strip() @@ -2899,6 +3193,12 @@ def _run_utterance_end( finalized = finalized_audio if finalized_audio is not None else self.finalize_audio_upload() pcm = finalized.pcm audio_summary = dict(finalized.summary) + finalized_reject = str(audio_summary.get("audio_capture_reject_code", "")) + if finalized_reject: + frame = error_frame(finalized_reject, "partial PCM discarded before STT") + frame.update(audio_summary) + frame.update(self._audio_upload_telemetry()) + return [frame] has_audio = int(audio_summary["audio_bytes"]) > 0 declaration_error = self._validate_audio_end_declaration(message, audio_summary) if declaration_error: @@ -4101,6 +4401,11 @@ def run_initiative_turn(decision: InitiativeDecision) -> None: close_interrupted_response("playback_interrupted") if text_message_type == "utterance_end": + if isinstance(parsed_text, dict): + audio_end_error = session.prepare_utterance_end(parsed_text) + if audio_end_error is not None: + send_live(audio_end_error) + continue if turn_thread is not None and turn_thread.is_alive(): turn_thread.join(timeout=1.5) if turn_thread is not None and turn_thread.is_alive(): @@ -4108,6 +4413,14 @@ def run_initiative_turn(decision: InitiativeDecision) -> None: continue early_frame = session.early_thinking_frame(text) finalized_audio = session.finalize_audio_upload() + if isinstance(parsed_text, dict): + audio_end_error = session.prepare_utterance_end( + parsed_text, + finalized_audio=finalized_audio, + ) + if audio_end_error is not None: + send_live(audio_end_error) + continue if early_frame is not None: sent_at = send_live(early_frame) if sent_at is not None and isinstance(parsed_text, dict): @@ -4380,6 +4693,16 @@ def build_arg_parser() -> argparse.ArgumentParser: parser.add_argument("--client-idle-timeout-s", type=float, default=DEFAULT_CLIENT_IDLE_TIMEOUT_S) parser.add_argument("--disable-audio-downlink", action="store_true") parser.add_argument("--max-audio-bytes", type=int, default=DEFAULT_MAX_AUDIO_BYTES) + parser.add_argument( + "--audio-capture-absolute-lease-ms", + type=int, + default=DEFAULT_AUDIO_CAPTURE_ABSOLUTE_LEASE_MS, + ) + parser.add_argument( + "--audio-capture-inactivity-lease-ms", + type=int, + default=DEFAULT_AUDIO_CAPTURE_INACTIVITY_LEASE_MS, + ) parser.add_argument("--audio-evidence-dir", type=Path) parser.add_argument("--memory-file", type=Path) parser.add_argument("--turn-log-file", type=Path) @@ -4436,6 +4759,12 @@ def main() -> int: parser.error("--conversation-reply-window-step-ms cannot be negative") if not 0.0 <= args.stt_min_confidence <= 1.0: parser.error("--stt-min-confidence must be between 0 and 1") + if not 12_500 <= args.audio_capture_absolute_lease_ms <= 15_000: + parser.error("--audio-capture-absolute-lease-ms must be between 12500 and 15000") + if not 1_000 <= args.audio_capture_inactivity_lease_ms < args.audio_capture_absolute_lease_ms: + parser.error( + "--audio-capture-inactivity-lease-ms must be between 1000 and the absolute lease" + ) if not 0 <= args.conversation_acoustic_tail_ms <= 2_000: parser.error("--conversation-acoustic-tail-ms must be between 0 and 2000") if args.initiative_min_interval_seconds < MIN_UNPROMPTED_INTERVAL_MS // 1000: @@ -4487,6 +4816,8 @@ def main() -> int: client_idle_timeout_s=max(1.0, args.client_idle_timeout_s), disable_audio_downlink=args.disable_audio_downlink, max_audio_bytes=args.max_audio_bytes, + audio_capture_absolute_lease_ms=args.audio_capture_absolute_lease_ms, + audio_capture_inactivity_lease_ms=args.audio_capture_inactivity_lease_ms, audio_evidence_dir=args.audio_evidence_dir, memory_file=args.memory_file, turn_log_file=args.turn_log_file, diff --git a/bridge/test_lan_service.py b/bridge/test_lan_service.py index 85c76d3..b2c7e59 100644 --- a/bridge/test_lan_service.py +++ b/bridge/test_lan_service.py @@ -2200,6 +2200,178 @@ def test_binary_audio_upload_tracks_telemetry_and_requires_stt_or_transcript(sel self.assertFalse(session.audio.active) self.assertEqual(0, session.audio.bytes_received) + def test_audio_capture_absolute_lease_exceeds_firmware_max_but_cannot_be_refreshed(self): + session = LanBridgeSession(LanBridgeConfig()) + self.assertGreater(session.config.audio_capture_absolute_lease_ms, 12_000) + self.assertLessEqual(session.config.audio_capture_absolute_lease_ms, 15_000) + session.handle_text( + json.dumps({"type": "utterance_start", "seq": 41, "sample_rate": 16000}) + ) + session.handle_binary(b"\x01\x00\x02\x00") + current = time.perf_counter() + session.audio.started_at_monotonic = current - ( + session.config.audio_capture_absolute_lease_ms + 1 + ) / 1000.0 + session.audio.last_activity_at_monotonic = current + + expired = session.handle_binary(b"\x03\x00\x04\x00") + + self.assertEqual("audio_capture_absolute_lease_expired", expired[0]["code"]) + self.assertEqual(4, expired[0]["audio_bytes"]) + self.assertFalse(session.audio.active) + self.assertEqual(0, session.audio.bytes_received) + self.assertEqual(1, expired[0]["audio_uploads_expired"]) + + def test_audio_capture_valid_after_twelve_seconds_with_recent_pcm(self): + session = LanBridgeSession(LanBridgeConfig()) + session.handle_text( + json.dumps({"type": "utterance_start", "seq": 42, "sample_rate": 16000}) + ) + session.handle_binary(b"\x01\x00\x02\x00") + current = time.perf_counter() + session.audio.started_at_monotonic = current - 12.25 + session.audio.last_activity_at_monotonic = current + runner_result = SimpleNamespace( + raw_response=json.dumps( + { + "spoken_text": "I heard the complete sentence.", + "mode": "speak", + "earcon": "none", + "emotion": {"arousal": 0.0, "valence": 0.0}, + "memory_write": {}, + "memory_forget": [], + } + ), + command_source="test", + elapsed_ms=1.0, + approx_tokens_per_sec=10.0, + ) + + with patch("lan_service.run_runner_profile", return_value=runner_result) as runner: + frames = session.handle_text( + json.dumps( + {"type": "utterance_end", "seq": 42, "text": "complete sentence"} + ) + ) + + runner.assert_called_once() + thinking = next(frame for frame in frames if frame.get("type") == "thinking") + self.assertGreaterEqual(thinking["audio_capture_elapsed_ms"], 12_000) + self.assertNotIn("audio_capture_reject_code", thinking) + self.assertFalse(session.audio.active) + + def test_audio_capture_inactivity_discards_partial_pcm_and_rejects_late_end(self): + session = LanBridgeSession(LanBridgeConfig()) + session.handle_text( + json.dumps({"type": "utterance_start", "seq": 43, "sample_rate": 16000}) + ) + session.handle_binary(b"\x01\x00\x02\x00") + current = time.perf_counter() + session.audio.last_activity_at_monotonic = current - ( + session.config.audio_capture_inactivity_lease_ms + 1 + ) / 1000.0 + + with ( + patch("lan_service.transcribe_pcm") as stt, + patch("lan_service.run_runner_profile") as runner, + patch("lan_service.synthesize_speech") as tts, + ): + expired = session.handle_binary(b"\x03\x00\x04\x00") + late_end = session.handle_text( + json.dumps({"type": "utterance_end", "seq": 43, "text": "stale"}) + ) + + self.assertEqual("audio_capture_inactivity_expired", expired[0]["code"]) + self.assertEqual("utterance_end_stale", late_end[0]["code"]) + self.assertEqual( + "audio_capture_inactivity_expired", + late_end[0]["audio_capture_reject_code"], + ) + stt.assert_not_called() + runner.assert_not_called() + tts.assert_not_called() + + def test_heartbeat_expires_inactive_capture_without_dropping_heartbeat(self): + session = LanBridgeSession(LanBridgeConfig()) + session.handle_text(json.dumps({"type": "utterance_start", "seq": 48})) + session.handle_binary(b"\x01\x00\x02\x00") + current = time.perf_counter() + session.audio.last_activity_at_monotonic = current - ( + session.config.audio_capture_inactivity_lease_ms + 1 + ) / 1000.0 + + frames = session.handle_text( + json.dumps({"type": "heartbeat", "power_source": "wall"}) + ) + + self.assertEqual("audio_capture_inactivity_expired", frames[0]["code"]) + self.assertEqual("heartbeat", frames[1]["type"]) + self.assertFalse(session.audio.active) + + def test_utterance_cancel_is_idempotent_and_late_end_runs_no_pipeline(self): + session = LanBridgeSession(LanBridgeConfig()) + session.handle_text( + json.dumps({"type": "utterance_start", "seq": 44, "sample_rate": 16000}) + ) + session.handle_binary(b"\x01\x00\x02\x00") + + with ( + patch("lan_service.transcribe_pcm") as stt, + patch("lan_service.run_runner_profile") as runner, + patch("lan_service.synthesize_speech") as tts, + ): + cancelled = session.handle_text( + json.dumps( + {"type": "utterance_cancel", "seq": 44, "reason": "capture_discontinuity"} + ) + ) + duplicate = session.handle_text( + json.dumps( + {"type": "utterance_cancel", "seq": 44, "reason": "capture_discontinuity"} + ) + ) + late_end = session.handle_text( + json.dumps({"type": "utterance_end", "seq": 44, "text": "unfinished"}) + ) + + self.assertTrue(cancelled[0]["utterance_cancelled"]) + self.assertFalse(cancelled[0]["utterance_cancel_duplicate"]) + self.assertEqual(4, cancelled[0]["audio_bytes"]) + self.assertTrue(duplicate[0]["utterance_cancel_duplicate"]) + self.assertEqual("utterance_end_stale", late_end[0]["code"]) + self.assertEqual("audio_capture_cancelled", late_end[0]["audio_capture_reject_code"]) + self.assertFalse(session.audio.active) + self.assertEqual(0, session.audio.bytes_received) + stt.assert_not_called() + runner.assert_not_called() + tts.assert_not_called() + + def test_wrong_terminal_sequence_does_not_finalize_active_capture(self): + session = LanBridgeSession(LanBridgeConfig()) + session.handle_text( + json.dumps({"type": "utterance_start", "seq": 45, "sample_rate": 16000}) + ) + session.handle_binary(b"\x01\x00\x02\x00") + + mismatch = session.handle_text( + json.dumps({"type": "utterance_end", "seq": 46, "text": "wrong turn"}) + ) + + self.assertEqual("utterance_end_seq_mismatch", mismatch[0]["code"]) + self.assertTrue(session.audio.active) + self.assertEqual(4, session.audio.bytes_received) + + def test_audio_lease_telemetry_is_exposed_in_diagnostics(self): + session = LanBridgeSession(LanBridgeConfig()) + session.handle_text(json.dumps({"type": "utterance_start", "seq": 47})) + + diagnostics = session.handle_text(json.dumps({"type": "diagnostics_request"}))[0] + + self.assertTrue(diagnostics["audio"]["audio_upload_active"]) + self.assertEqual(1, diagnostics["audio"]["audio_uploads_started"]) + self.assertEqual(14_500, diagnostics["audio"]["audio_capture_absolute_lease_ms"]) + self.assertEqual(4_000, diagnostics["audio"]["audio_capture_inactivity_lease_ms"]) + def test_empty_utterance_end_does_not_run_runner(self): session = LanBridgeSession(LanBridgeConfig(runner_case="greeting")) diff --git a/src/io/BridgeAudioUplink.cpp b/src/io/BridgeAudioUplink.cpp index aa526d8..ed58c13 100644 --- a/src/io/BridgeAudioUplink.cpp +++ b/src/io/BridgeAudioUplink.cpp @@ -23,13 +23,15 @@ bool BridgeAudioUplink::begin(const BridgeAudioUplinkConfig& config, return fail("audio_uplink_session_missing"); } if (config_.sampleRate == 0 || config_.maxAudioBytes == 0 || - config_.maxChunkBytes == 0 || + config_.maxChunkBytes == 0 || config_.terminalRetryMs == 0 || config_.maxChunkBytes > kBridgeAudioStreamChunkPayloadMax) { telemetry_.ready = false; return fail("audio_uplink_bad_config"); } telemetry_.lastError[0] = '\0'; + pendingTerminalSeq_ = 0; + pendingTerminalReason_[0] = '\0'; return true; } @@ -39,6 +41,8 @@ void BridgeAudioUplink::reset() { telemetry_.ready = ready; telemetry_.enabled = config_.enabled; telemetry_.wakeGateRequired = config_.wakeGateRequired; + pendingTerminalSeq_ = 0; + pendingTerminalReason_[0] = '\0'; if (!config_.enabled) { copyError("audio_uplink_disabled"); } @@ -55,7 +59,7 @@ bool BridgeAudioUplink::beginTurn(uint32_t seq, uint32_t nowMs, bool wakeGateOpe if (session_ == nullptr) { return fail("audio_uplink_session_missing"); } - if (telemetry_.active) { + if (telemetry_.active || telemetry_.terminalPending) { return fail("audio_uplink_already_active"); } if (config_.wakeGateRequired && !wakeGateOpen) { @@ -74,6 +78,7 @@ bool BridgeAudioUplink::beginTurn(uint32_t seq, uint32_t nowMs, bool wakeGateOpe telemetry_.lastSeq = seq; telemetry_.activeBytes = 0; telemetry_.activeChunks = 0; + telemetry_.lastTerminal = BridgeAudioTerminalKind::None; telemetry_.lastError[0] = '\0'; return true; } @@ -131,7 +136,6 @@ bool BridgeAudioUplink::submitPcmBytes(uint32_t seq, } bool BridgeAudioUplink::endTurn(uint32_t seq, uint32_t nowMs) { - (void)nowMs; if (!configured()) { return fail("audio_uplink_not_ready"); } @@ -145,26 +149,67 @@ bool BridgeAudioUplink::endTurn(uint32_t seq, uint32_t nowMs) { return fail("audio_uplink_seq_mismatch"); } + return requestTerminal(BridgeAudioTerminalKind::End, seq, nowMs, nullptr); +} + +void BridgeAudioUplink::abort(uint32_t nowMs, const char* reason) { + if (!telemetry_.active && !telemetry_.terminalPending) { + copyError(reason != nullptr ? reason : "audio_uplink_aborted"); + return; + } + requestTerminal(BridgeAudioTerminalKind::Cancel, + telemetry_.lastSeq, + nowMs, + reason != nullptr ? reason : "audio_uplink_aborted"); +} + +BridgeAudioTerminalServiceResult BridgeAudioUplink::servicePendingTerminal(uint32_t nowMs) { + if (!telemetry_.terminalPending) { + return BridgeAudioTerminalServiceResult::Idle; + } + char frame[kBridgeEndpointControlResponseMax] = {}; - if (!writeEndFrame(seq, frame, sizeof(frame)) || !queueText(frame)) { - telemetry_.queueFailures++; - return fail("utterance_end_queue_failed"); + const bool encoded = telemetry_.pendingTerminal == BridgeAudioTerminalKind::End + ? writeEndFrame(pendingTerminalSeq_, frame, sizeof(frame)) + : writeCancelFrame(pendingTerminalSeq_, + pendingTerminalReason_, + frame, + sizeof(frame)); + telemetry_.terminalAttempts++; + if (encoded && queueText(frame)) { + const BridgeAudioTerminalKind delivered = telemetry_.pendingTerminal; + telemetry_.terminalPending = false; + telemetry_.pendingTerminal = BridgeAudioTerminalKind::None; + telemetry_.lastTerminal = delivered; + telemetry_.lastSeq = pendingTerminalSeq_; + if (delivered == BridgeAudioTerminalKind::End) { + telemetry_.turnsCompleted++; + telemetry_.lastError[0] = '\0'; + } else { + telemetry_.turnsAborted++; + telemetry_.cancelFramesQueued++; + copyError(pendingTerminalReason_); + } + return BridgeAudioTerminalServiceResult::Queued; } - telemetry_.active = false; - telemetry_.turnsCompleted++; - telemetry_.lastSeq = seq; - telemetry_.lastError[0] = '\0'; - return true; -} + telemetry_.queueFailures++; + telemetry_.terminalRetries++; + if (nowMs - telemetry_.terminalRequestedAtMs < config_.terminalRetryMs) { + copyError("utterance_terminal_queue_pending"); + return BridgeAudioTerminalServiceResult::Pending; + } -void BridgeAudioUplink::abort(uint32_t nowMs, const char* reason) { - (void)nowMs; - if (telemetry_.active) { - telemetry_.turnsAborted++; + if (session_ != nullptr) { + session_->stop(nowMs); } - telemetry_.active = false; - copyError(reason != nullptr ? reason : "audio_uplink_aborted"); + telemetry_.terminalTimeouts++; + telemetry_.turnsAborted++; + telemetry_.terminalPending = false; + telemetry_.pendingTerminal = BridgeAudioTerminalKind::None; + telemetry_.lastTerminal = BridgeAudioTerminalKind::Cancel; + fail("utterance_terminal_delivery_timeout"); + return BridgeAudioTerminalServiceResult::FailedClosed; } bool BridgeAudioUplink::configured() const { @@ -206,6 +251,68 @@ bool BridgeAudioUplink::writeEndFrame(uint32_t seq, char* out, size_t outSize) c return written > 0 && static_cast(written) < outSize; } +bool BridgeAudioUplink::writeCancelFrame(uint32_t seq, + const char* reason, + char* out, + size_t outSize) const { + if (out == nullptr || outSize == 0 || reason == nullptr || reason[0] == '\0') { + return false; + } + const int written = snprintf(out, + outSize, + "{\"type\":\"utterance_cancel\",\"seq\":%lu," + "\"reason\":\"%s\",\"audio_bytes\":%lu,\"chunks\":%lu}", + static_cast(seq), + reason, + static_cast(telemetry_.activeBytes), + static_cast(telemetry_.activeChunks)); + return written > 0 && static_cast(written) < outSize; +} + +bool BridgeAudioUplink::requestTerminal(BridgeAudioTerminalKind kind, + uint32_t seq, + uint32_t nowMs, + const char* reason) { + if (kind == BridgeAudioTerminalKind::None) { + return fail("audio_uplink_bad_terminal"); + } + + const uint32_t resolvedSeq = seq != 0 ? seq : telemetry_.lastSeq; + const bool alreadyPending = telemetry_.terminalPending; + telemetry_.active = false; + telemetry_.terminalPending = true; + telemetry_.pendingTerminal = kind; + if (!alreadyPending) { + telemetry_.terminalRequestedAtMs = nowMs; + } + pendingTerminalSeq_ = resolvedSeq; + copyTerminalReason(reason != nullptr ? reason : "audio_uplink_aborted"); + const BridgeAudioTerminalServiceResult result = servicePendingTerminal(nowMs); + return result == BridgeAudioTerminalServiceResult::Queued || + result == BridgeAudioTerminalServiceResult::Pending; +} + +void BridgeAudioUplink::copyTerminalReason(const char* reason) { + if (reason == nullptr) { + pendingTerminalReason_[0] = '\0'; + return; + } + size_t out = 0; + while (reason[out] != '\0' && out < sizeof(pendingTerminalReason_) - 1u) { + const char value = reason[out]; + const bool safe = (value >= 'a' && value <= 'z') || + (value >= 'A' && value <= 'Z') || + (value >= '0' && value <= '9') || + value == '_' || value == '-' || value == '.' || value == ':'; + pendingTerminalReason_[out] = safe ? value : '_'; + ++out; + } + pendingTerminalReason_[out] = '\0'; + if (out == 0) { + copyTerminalReason("audio_uplink_aborted"); + } +} + bool BridgeAudioUplink::fail(const char* reason) { telemetry_.errors++; copyError(reason); diff --git a/src/io/BridgeAudioUplink.hpp b/src/io/BridgeAudioUplink.hpp index e76aec8..6de544c 100644 --- a/src/io/BridgeAudioUplink.hpp +++ b/src/io/BridgeAudioUplink.hpp @@ -16,12 +16,29 @@ constexpr uint32_t kBridgeAudioUplinkSampleRate = 16000; constexpr uint32_t kBridgeAudioUplinkMaxBytes = 512u * 1024u; constexpr size_t kBridgeAudioUplinkErrorMax = kBridgeErrorMax; +enum class BridgeAudioTerminalKind : uint8_t { + None = 0, + End, + Cancel, +}; + +enum class BridgeAudioTerminalServiceResult : uint8_t { + Idle = 0, + Queued, + Pending, + FailedClosed, +}; + struct BridgeAudioUplinkConfig { bool enabled = STACKCHAN_ENABLE_BRIDGE_AUDIO_UPLINK != 0; bool wakeGateRequired = true; // privacy gate: audio may leave only after wake/explicit activation. uint32_t sampleRate = kBridgeAudioUplinkSampleRate; uint32_t maxAudioBytes = kBridgeAudioUplinkMaxBytes; uint16_t maxChunkBytes = kBridgeAudioStreamChunkPayloadMax; + // A terminal frame may wait for the single-slot socket writer to drain, but + // it must never wait forever. Expiry closes the socket so the host discards + // the unterminated upload. + uint32_t terminalRetryMs = 1000; }; struct BridgeAudioUplinkTelemetry { @@ -29,6 +46,7 @@ struct BridgeAudioUplinkTelemetry { bool enabled = false; bool active = false; bool wakeGateRequired = true; + bool terminalPending = false; uint32_t turnsStarted = 0; uint32_t turnsCompleted = 0; uint32_t turnsAborted = 0; @@ -40,6 +58,13 @@ struct BridgeAudioUplinkTelemetry { uint32_t lastSeq = 0; uint32_t activeBytes = 0; uint32_t activeChunks = 0; + uint32_t terminalAttempts = 0; + uint32_t terminalRetries = 0; + uint32_t terminalTimeouts = 0; + uint32_t cancelFramesQueued = 0; + uint32_t terminalRequestedAtMs = 0; + BridgeAudioTerminalKind pendingTerminal = BridgeAudioTerminalKind::None; + BridgeAudioTerminalKind lastTerminal = BridgeAudioTerminalKind::None; char lastError[kBridgeAudioUplinkErrorMax] = {}; }; @@ -54,6 +79,7 @@ class BridgeAudioUplink { bool submitPcmBytes(uint32_t seq, const uint8_t* payload, size_t length, uint32_t nowMs); bool endTurn(uint32_t seq, uint32_t nowMs); void abort(uint32_t nowMs, const char* reason = nullptr); + BridgeAudioTerminalServiceResult servicePendingTerminal(uint32_t nowMs); const BridgeAudioUplinkTelemetry& telemetry() const { return telemetry_; @@ -65,12 +91,20 @@ class BridgeAudioUplink { bool queueBinary(const uint8_t* payload, size_t length); bool writeStartFrame(uint32_t seq, char* out, size_t outSize) const; bool writeEndFrame(uint32_t seq, char* out, size_t outSize) const; + bool writeCancelFrame(uint32_t seq, const char* reason, char* out, size_t outSize) const; + bool requestTerminal(BridgeAudioTerminalKind kind, + uint32_t seq, + uint32_t nowMs, + const char* reason); + void copyTerminalReason(const char* reason); bool fail(const char* reason); void copyError(const char* reason); BridgeAudioUplinkConfig config_; BridgeAudioUplinkTelemetry telemetry_; BridgeNetworkSession* session_ = nullptr; + uint32_t pendingTerminalSeq_ = 0; + char pendingTerminalReason_[kBridgeAudioUplinkErrorMax] = {}; }; } // namespace stackchan diff --git a/src/io/BridgeSocketWriter.cpp b/src/io/BridgeSocketWriter.cpp index 14a940f..c69e617 100644 --- a/src/io/BridgeSocketWriter.cpp +++ b/src/io/BridgeSocketWriter.cpp @@ -162,6 +162,12 @@ BridgeSocketWriterDrainResult BridgeSocketWriter::drainPending(uint32_t nowMs, const size_t remaining = frameBytes_ - frameOffset_; const size_t written = sink_->write(frame_ + frameOffset_, remaining); if (written == 0) { + if (sink_->lastWriteWouldBlock()) { + telemetry_.writeDeferrals++; + telemetry_.frameBuffered = true; + telemetry_.lastError[0] = '\0'; + return BridgeSocketWriterDrainResult::Partial; + } telemetry_.writeFailures++; copyError("socket_write_failed"); return BridgeSocketWriterDrainResult::WriteFailed; diff --git a/src/io/BridgeSocketWriter.hpp b/src/io/BridgeSocketWriter.hpp index aca473a..e050f69 100644 --- a/src/io/BridgeSocketWriter.hpp +++ b/src/io/BridgeSocketWriter.hpp @@ -41,6 +41,7 @@ struct BridgeSocketWriterTelemetry { uint32_t binaryFramesEncoded = 0; uint32_t binaryFramesWritten = 0; uint32_t partialWrites = 0; + uint32_t writeDeferrals = 0; uint32_t writeFailures = 0; uint32_t textBytesQueued = 0; uint32_t textBytesWritten = 0; @@ -57,6 +58,10 @@ class BridgeSocketWriterSink { virtual bool isConnected() const = 0; virtual size_t write(const uint8_t* data, size_t length) = 0; + // A zero-byte write normally means the connection failed. Network sinks may + // instead report transient backpressure so the retained frame is retried by + // a later service pass without blocking the real-time audio capture path. + virtual bool lastWriteWouldBlock() const { return false; } }; class BridgeSocketWriter { diff --git a/src/io/BridgeWakeGate.cpp b/src/io/BridgeWakeGate.cpp index 3983258..1200a18 100644 --- a/src/io/BridgeWakeGate.cpp +++ b/src/io/BridgeWakeGate.cpp @@ -66,10 +66,13 @@ void BridgeWakeGate::applyEvent(const RobotEvent& event, uint32_t nowMs) { case EventType::SpeechEnded: case EventType::ResponseStarted: case EventType::ResponseEnded: - case EventType::Error: completeTurn(nowMs, "bridge_wake_gate_event_end"); expireGate(nowMs); break; + case EventType::Error: + abortTurn(nowMs, "bridge_wake_gate_event_error"); + expireGate(nowMs); + break; default: break; } @@ -79,17 +82,31 @@ void BridgeWakeGate::update(uint32_t nowMs) { if (!telemetry_.ready || !config_.enabled) { return; } + serviceTerminal(nowMs); + if (telemetry_.turnActive && uplink_ != nullptr && + !uplink_->telemetry().active && !uplink_->telemetry().terminalPending && + uplink_->telemetry().lastTerminal != BridgeAudioTerminalKind::None) { + settleTerminal(BridgeAudioTerminalServiceResult::Queued); + } if (telemetry_.gateOpen && nowMs >= telemetry_.closeAtMs) { - completeTurn(nowMs, "bridge_wake_gate_timeout"); + abortTurn(nowMs, "bridge_wake_gate_timeout"); expireGate(nowMs); } if (telemetry_.turnActive && nowMs - telemetry_.turnStartedAtMs >= config_.maxTurnMs) { - completeTurn(nowMs, "bridge_wake_gate_max_turn"); + abortTurn(nowMs, "bridge_wake_gate_max_turn"); expireGate(nowMs); } } +void BridgeWakeGate::cancelActiveTurn(uint32_t nowMs, const char* reason) { + if (!telemetry_.ready || !config_.enabled) { + return; + } + abortTurn(nowMs, reason != nullptr ? reason : "bridge_wake_gate_cancelled"); + expireGate(nowMs); +} + bool BridgeWakeGate::isGateOpen(uint32_t nowMs) const { return telemetry_.ready && config_.enabled && telemetry_.gateOpen && nowMs < telemetry_.closeAtMs; @@ -141,6 +158,7 @@ void BridgeWakeGate::startTurnIfPossible(uint32_t nowMs) { } void BridgeWakeGate::completeTurn(uint32_t nowMs, const char* reason) { + (void)reason; if (!telemetry_.turnActive || uplink_ == nullptr) { return; } @@ -148,17 +166,66 @@ void BridgeWakeGate::completeTurn(uint32_t nowMs, const char* reason) { if (uplink_->telemetry().active) { if (!uplink_->endTurn(telemetry_.lastSeq, nowMs)) { telemetry_.endFailures++; - uplink_->abort(nowMs, reason); - telemetry_.turnsAborted++; - telemetry_.turnActive = false; copyError(uplink_->telemetry().lastError); return; } } - telemetry_.turnsCompleted++; + if (uplink_->telemetry().terminalPending) { + copyError("bridge_wake_gate_terminal_pending"); + return; + } + + settleTerminal(BridgeAudioTerminalServiceResult::Queued); +} + +void BridgeWakeGate::abortTurn(uint32_t nowMs, const char* reason) { + if (!telemetry_.turnActive || uplink_ == nullptr) { + return; + } + uplink_->abort(nowMs, reason); + if (uplink_->telemetry().terminalPending) { + copyError("bridge_wake_gate_terminal_pending"); + return; + } + settleTerminal(BridgeAudioTerminalServiceResult::Queued); +} + +void BridgeWakeGate::serviceTerminal(uint32_t nowMs) { + if (!telemetry_.turnActive || uplink_ == nullptr || + !uplink_->telemetry().terminalPending) { + return; + } + const BridgeAudioTerminalServiceResult result = uplink_->servicePendingTerminal(nowMs); + if (result == BridgeAudioTerminalServiceResult::Pending) { + telemetry_.terminalRetries++; + copyError("bridge_wake_gate_terminal_pending"); + return; + } + if (result == BridgeAudioTerminalServiceResult::FailedClosed) { + telemetry_.terminalTimeouts++; + telemetry_.endFailures++; + } + if (result == BridgeAudioTerminalServiceResult::Queued || + result == BridgeAudioTerminalServiceResult::FailedClosed) { + settleTerminal(result); + } +} + +void BridgeWakeGate::settleTerminal(BridgeAudioTerminalServiceResult result) { + const bool completed = result == BridgeAudioTerminalServiceResult::Queued && + uplink_ != nullptr && + uplink_->telemetry().lastTerminal == BridgeAudioTerminalKind::End; + if (completed) { + telemetry_.turnsCompleted++; + telemetry_.lastError[0] = '\0'; + } else { + telemetry_.turnsAborted++; + copyError(uplink_ != nullptr ? uplink_->telemetry().lastError + : "bridge_wake_gate_terminal_failed"); + } + telemetry_.turnActive = false; - telemetry_.lastError[0] = '\0'; } bool BridgeWakeGate::uplinkReadyForTurn() const { @@ -166,7 +233,7 @@ bool BridgeWakeGate::uplinkReadyForTurn() const { return false; } const BridgeAudioUplinkTelemetry& uplink = uplink_->telemetry(); - return uplink.ready && uplink.enabled && !uplink.active; + return uplink.ready && uplink.enabled && !uplink.active && !uplink.terminalPending; } void BridgeWakeGate::copyError(const char* reason) { diff --git a/src/io/BridgeWakeGate.hpp b/src/io/BridgeWakeGate.hpp index 63409d1..a287e63 100644 --- a/src/io/BridgeWakeGate.hpp +++ b/src/io/BridgeWakeGate.hpp @@ -3,6 +3,7 @@ #include #include "io/BridgeAudioUplink.hpp" +#include "io/VoiceActivityEndpoint.hpp" #include "persona/EventBus.hpp" namespace stackchan { @@ -31,6 +32,26 @@ constexpr bool dedicatedWakeCaptureMaySubmit(bool gateOpen, return gateOpen && gateTurnActive && uplinkActive; } +constexpr bool dedicatedWakeCaptureChunkContinuous(uint32_t micCaptureUs, + uint32_t micDeadlineUs, + uint32_t previousChunkAtMs, + uint32_t currentChunkAtMs, + uint32_t gapDeadlineMs) { + return micCaptureUs <= micDeadlineUs && + (previousChunkAtMs == 0 || + currentChunkAtMs - previousChunkAtMs <= gapDeadlineMs); +} + +constexpr bool dedicatedWakeCaptureMayCommit(VoiceActivityEndpointReason endpointReason, + bool endpointEnabled, + bool chunkLimitReached, + bool submitFailed, + bool captureDiscontinuity) { + return !submitFailed && !captureDiscontinuity && + (endpointReason == VoiceActivityEndpointReason::TrailingSilence || + (!endpointEnabled && chunkLimitReached)); +} + struct BridgeWakeGateConfig { bool enabled = true; bool speechStartsTurn = STACKCHAN_BRIDGE_WAKE_ON_SPEECH != 0; @@ -52,6 +73,8 @@ struct BridgeWakeGateTelemetry { uint32_t turnsAborted = 0; uint32_t beginFailures = 0; uint32_t endFailures = 0; + uint32_t terminalRetries = 0; + uint32_t terminalTimeouts = 0; uint32_t suppressedStarts = 0; uint32_t lastSeq = 0; uint32_t openedAtMs = 0; @@ -68,6 +91,7 @@ class BridgeWakeGate { void applyEvent(const RobotEvent& event, uint32_t nowMs); void update(uint32_t nowMs); + void cancelActiveTurn(uint32_t nowMs, const char* reason); bool isGateOpen(uint32_t nowMs) const; const BridgeWakeGateTelemetry& telemetry() const { @@ -80,6 +104,9 @@ class BridgeWakeGate { void expireGate(uint32_t nowMs); void startTurnIfPossible(uint32_t nowMs); void completeTurn(uint32_t nowMs, const char* reason); + void abortTurn(uint32_t nowMs, const char* reason); + void serviceTerminal(uint32_t nowMs); + void settleTerminal(BridgeAudioTerminalServiceResult result); bool uplinkReadyForTurn() const; void copyError(const char* reason); diff --git a/src/io/BridgeWiFiClientSocket.cpp b/src/io/BridgeWiFiClientSocket.cpp index 9ac709e..e8361b5 100644 --- a/src/io/BridgeWiFiClientSocket.cpp +++ b/src/io/BridgeWiFiClientSocket.cpp @@ -2,6 +2,10 @@ #include +#if defined(ARDUINO_ARCH_ESP32) +#include +#endif + namespace stackchan { namespace { @@ -102,10 +106,23 @@ int BridgeWiFiClientSocket::read(uint8_t* out, size_t outSize) { size_t BridgeWiFiClientSocket::write(const uint8_t* data, size_t length) { #if defined(ARDUINO_ARCH_ESP32) + lastWriteWouldBlock_ = false; if (data == nullptr || length == 0) { return 0; } - return client_.write(data, length); + const int socketFd = client_.fd(); + if (socketFd < 0 || !client_.connected()) { + return 0; + } + errno = 0; + const int result = ::send(socketFd, data, length, MSG_DONTWAIT); + if (result > 0) { + return static_cast(result); + } + if (result < 0 && (errno == EAGAIN || errno == EWOULDBLOCK)) { + lastWriteWouldBlock_ = true; + } + return 0; #else (void)data; (void)length; @@ -114,6 +131,7 @@ size_t BridgeWiFiClientSocket::write(const uint8_t* data, size_t length) { } void BridgeWiFiClientSocket::stop() { + lastWriteWouldBlock_ = false; #if defined(ARDUINO_ARCH_ESP32) client_.stop(); client_ = WiFiClient {}; diff --git a/src/io/BridgeWiFiClientSocket.hpp b/src/io/BridgeWiFiClientSocket.hpp index f5c6179..7a51b43 100644 --- a/src/io/BridgeWiFiClientSocket.hpp +++ b/src/io/BridgeWiFiClientSocket.hpp @@ -19,6 +19,7 @@ class BridgeWiFiClientSocket final : public BridgeNetworkSocket { int available() override; int read(uint8_t* out, size_t outSize) override; size_t write(const uint8_t* data, size_t length) override; + bool lastWriteWouldBlock() const override { return lastWriteWouldBlock_; } void stop() override; uint32_t connectAttempts() const { return connectAttempts_; } @@ -38,6 +39,7 @@ class BridgeWiFiClientSocket final : public BridgeNetworkSocket { uint32_t maxConnectDurationMs_ = 0; int lastConnectErrno_ = 0; int lastConnectResult_ = 0; + bool lastWriteWouldBlock_ = false; }; } // namespace stackchan diff --git a/src/io/VoiceActivityEndpoint.cpp b/src/io/VoiceActivityEndpoint.cpp index 5b4e161..0b685c4 100644 --- a/src/io/VoiceActivityEndpoint.cpp +++ b/src/io/VoiceActivityEndpoint.cpp @@ -23,11 +23,15 @@ bool VoiceActivityEndpoint::begin(const VoiceActivityEndpointConfig& config, uin telemetry_.speechSeen = false; telemetry_.captureStartedAtMs = nowMs; telemetry_.lastSpeechAtMs = 0; + telemetry_.capturedAudioMs = 0; + telemetry_.lastSpeechAudioMs = 0; telemetry_.lastLevel = 0.0f; telemetry_.lastZeroCrossingRate = 0.0f; telemetry_.noiseFloor = clamp01(config_.initialNoiseFloor); telemetry_.lastReason = VoiceActivityEndpointReason::None; - consecutiveSpeechMs_ = 0; + capturedSamples_ = 0; + consecutiveSpeechSamples_ = 0; + lastSpeechEndSample_ = 0; if (telemetry_.active) { ++telemetry_.capturesStarted; } @@ -44,8 +48,9 @@ VoiceActivityEndpointReason VoiceActivityEndpoint::process(const int16_t* sample ++telemetry_.chunksProcessed; telemetry_.lastLevel = level(samples, sampleCount); telemetry_.lastZeroCrossingRate = zeroCrossingRate(samples, sampleCount); - const uint32_t chunkMs = static_cast( - (sampleCount * 1000u + config_.sampleRate - 1u) / config_.sampleRate); + capturedSamples_ += static_cast(sampleCount); + telemetry_.capturedAudioMs = static_cast( + (capturedSamples_ * 1000u) / config_.sampleRate); const float dynamicThreshold = telemetry_.noiseFloor * config_.speechNoiseMultiplier; const float speechThreshold = dynamicThreshold > config_.minimumSpeechLevel ? dynamicThreshold @@ -56,13 +61,16 @@ VoiceActivityEndpointReason VoiceActivityEndpoint::process(const int16_t* sample if (speech) { ++telemetry_.speechChunks; - consecutiveSpeechMs_ += chunkMs > 0 ? chunkMs : 1u; + consecutiveSpeechSamples_ += static_cast(sampleCount); + lastSpeechEndSample_ = capturedSamples_; telemetry_.lastSpeechAtMs = nowMs; - if (consecutiveSpeechMs_ >= config_.minimumSpeechMs) { + telemetry_.lastSpeechAudioMs = telemetry_.capturedAudioMs; + if (consecutiveSpeechSamples_ * 1000u >= + static_cast(config_.minimumSpeechMs) * config_.sampleRate) { telemetry_.speechSeen = true; } } else { - consecutiveSpeechMs_ = 0; + consecutiveSpeechSamples_ = 0; if (!telemetry_.speechSeen) { const float adapt = telemetry_.lastLevel < telemetry_.noiseFloor ? 0.04f : 0.01f; telemetry_.noiseFloor += (telemetry_.lastLevel - telemetry_.noiseFloor) * adapt; @@ -72,12 +80,17 @@ VoiceActivityEndpointReason VoiceActivityEndpoint::process(const int16_t* sample } } - const uint32_t elapsedMs = nowMs - telemetry_.captureStartedAtMs; - if (telemetry_.speechSeen && !speech && elapsedMs >= config_.minimumCaptureMs && - nowMs - telemetry_.lastSpeechAtMs >= config_.trailingSilenceMs) { + const bool minimumCaptured = + capturedSamples_ * 1000u >= + static_cast(config_.minimumCaptureMs) * config_.sampleRate; + const bool trailingSilenceCaptured = + (capturedSamples_ - lastSpeechEndSample_) * 1000u >= + static_cast(config_.trailingSilenceMs) * config_.sampleRate; + if (telemetry_.speechSeen && !speech && minimumCaptured && trailingSilenceCaptured) { return finish(VoiceActivityEndpointReason::TrailingSilence, nowMs); } - if (elapsedMs >= config_.maximumCaptureMs) { + if (capturedSamples_ * 1000u >= + static_cast(config_.maximumCaptureMs) * config_.sampleRate) { return finish(VoiceActivityEndpointReason::MaxDuration, nowMs); } return VoiceActivityEndpointReason::None; @@ -85,7 +98,7 @@ VoiceActivityEndpointReason VoiceActivityEndpoint::process(const int16_t* sample void VoiceActivityEndpoint::cancel() { telemetry_.active = false; - consecutiveSpeechMs_ = 0; + consecutiveSpeechSamples_ = 0; } VoiceActivityEndpointReason VoiceActivityEndpoint::forceMaximum(uint32_t nowMs) { @@ -127,7 +140,7 @@ VoiceActivityEndpointReason VoiceActivityEndpoint::finish(VoiceActivityEndpointR if (reason == VoiceActivityEndpointReason::MaxDuration) { ++telemetry_.maxDurationFallbacks; } - consecutiveSpeechMs_ = 0; + consecutiveSpeechSamples_ = 0; return reason; } diff --git a/src/io/VoiceActivityEndpoint.hpp b/src/io/VoiceActivityEndpoint.hpp index ca1827b..73f2c4f 100644 --- a/src/io/VoiceActivityEndpoint.hpp +++ b/src/io/VoiceActivityEndpoint.hpp @@ -51,6 +51,11 @@ struct VoiceActivityEndpointTelemetry { uint32_t captureStartedAtMs = 0; uint32_t lastSpeechAtMs = 0; uint32_t lastEndpointAtMs = 0; + // Semantic endpoint timing is based only on PCM that was actually captured. + // Wall time above remains diagnostic and must not turn a transport/microphone + // stall into apparent silence. + uint32_t capturedAudioMs = 0; + uint32_t lastSpeechAudioMs = 0; float lastLevel = 0.0f; float lastZeroCrossingRate = 0.0f; float noiseFloor = 0.015f; @@ -73,7 +78,9 @@ class VoiceActivityEndpoint { VoiceActivityEndpointConfig config_; VoiceActivityEndpointTelemetry telemetry_; - uint32_t consecutiveSpeechMs_ = 0; + uint64_t capturedSamples_ = 0; + uint64_t consecutiveSpeechSamples_ = 0; + uint64_t lastSpeechEndSample_ = 0; }; const char* voiceActivityEndpointReasonName(VoiceActivityEndpointReason reason); diff --git a/src/main.cpp b/src/main.cpp index a25dc00..35abfef 100644 --- a/src/main.cpp +++ b/src/main.cpp @@ -1021,6 +1021,13 @@ struct WakeMwwDedicatedCaptureRuntime { uint16_t chunksSubmitted = 0; uint32_t serviceCalls = 0; uint32_t maxServiceUs = 0; + uint32_t maxMicCaptureUs = 0; + uint32_t maxVadUs = 0; + uint32_t maxSubmitUs = 0; + uint32_t captureTimeouts = 0; + uint32_t discontinuities = 0; + uint32_t captureStartedAtMs = 0; + uint32_t lastChunkCompletedAtMs = 0; VoiceActivityEndpoint endpoint; }; WakeMwwDedicatedCaptureRuntime gWakeMwwDedicatedCapture; @@ -3571,13 +3578,33 @@ void highPassWakeMwwAudio(int16_t* audioBuf, size_t sampleCount, int32_t& previo #endif } -bool recordWakeMwwAudioBlocking(int16_t* audioBuf, size_t sampleCount, uint32_t sampleRate, bool stereo) { +bool recordWakeMwwAudioBlocking(int16_t* audioBuf, + size_t sampleCount, + uint32_t sampleRate, + bool stereo, + uint32_t timeoutMs = 0, + uint32_t* elapsedUs = nullptr) { + const uint32_t startedAtUs = micros(); + const uint32_t startedAtMs = millis(); if (!M5.Mic.record(audioBuf, sampleCount, sampleRate, stereo)) { + if (elapsedUs != nullptr) { + *elapsedUs = micros() - startedAtUs; + } return false; } while (M5.Mic.isRecording() != 0) { + if (timeoutMs > 0 && millis() - startedAtMs >= timeoutMs) { + M5.Mic.end(); + if (elapsedUs != nullptr) { + *elapsedUs = micros() - startedAtUs; + } + return false; + } vTaskDelay(pdMS_TO_TICKS(1)); } + if (elapsedUs != nullptr) { + *elapsedUs = micros() - startedAtUs; + } return true; } @@ -4655,37 +4682,27 @@ DedicatedWakeCaptureSubmitResult submitDedicatedWakeCaptureChunk(uint32_t seq, uint16_t sampleCount, uint32_t nowMs) { (void)nowMs; - constexpr uint16_t kSubmitAttempts = - STACKCHAN_MWW_WAKE_UPLINK_SUBMIT_RETRY_ATTEMPTS > 0 - ? STACKCHAN_MWW_WAKE_UPLINK_SUBMIT_RETRY_ATTEMPTS - : 1; - constexpr uint16_t kSubmitDelayMs = STACKCHAN_MWW_WAKE_UPLINK_SUBMIT_RETRY_DELAY_MS; - - for (uint16_t attempt = 0; attempt < kSubmitAttempts; ++attempt) { - const uint32_t attemptMs = millis(); - gBridgeNetworkSession.update(attemptMs); - const uint32_t authorityNowMs = millis(); - const BridgeWakeGateTelemetry& gateBeforeAttempt = gBridgeWakeGate.telemetry(); - if (!dedicatedWakeCaptureMaySubmit( - gBridgeWakeGate.isGateOpen(authorityNowMs), - gateBeforeAttempt.turnActive, - gBridgeAudioUplink.telemetry().active)) { - return DedicatedWakeCaptureSubmitResult::AuthorityExpired; - } - const BridgeSocketWriterTelemetry& writer = gBridgeNetworkSession.writer().telemetry(); - if (!writer.frameBuffered && !writer.binaryFrameQueued && - gBridgeAudioUplink.submitPcmChunk(seq, samples, sampleCount, authorityNowMs)) { - gBridgeNetworkSession.update(millis()); - return DedicatedWakeCaptureSubmitResult::Submitted; - } - gBridgeNetworkSession.update(millis()); - if (kSubmitDelayMs > 0) { - vTaskDelay(pdMS_TO_TICKS(kSubmitDelayMs)); - } else { - taskYIELD(); - } + // Socket writes are nonblocking. Give the retained writer one service pass, + // then fail this utterance closed if the binary slot is still occupied. The + // old 40-attempt loop could stop microphone acquisition for multiple seconds + // and silently remove the middle of a sentence. + const uint32_t attemptMs = millis(); + gBridgeNetworkSession.update(attemptMs); + const uint32_t authorityNowMs = millis(); + const BridgeWakeGateTelemetry& gateBeforeAttempt = gBridgeWakeGate.telemetry(); + if (!dedicatedWakeCaptureMaySubmit( + gBridgeWakeGate.isGateOpen(authorityNowMs), + gateBeforeAttempt.turnActive, + gBridgeAudioUplink.telemetry().active)) { + return DedicatedWakeCaptureSubmitResult::AuthorityExpired; + } + const BridgeSocketWriterTelemetry& writer = gBridgeNetworkSession.writer().telemetry(); + if (writer.frameBuffered || writer.binaryFrameQueued) { + return DedicatedWakeCaptureSubmitResult::Failed; } - return DedicatedWakeCaptureSubmitResult::Failed; + return gBridgeAudioUplink.submitPcmChunk(seq, samples, sampleCount, authorityNowMs) + ? DedicatedWakeCaptureSubmitResult::Submitted + : DedicatedWakeCaptureSubmitResult::Failed; } void finishDedicatedWakeCaptureTurn(uint32_t seq, uint32_t nowMs) { @@ -4752,6 +4769,8 @@ bool beginDedicatedWakeCaptureAfterCue(const RobotEvent& wakeEvent) { gWakeMwwDedicatedCapture.seq = started.lastSeq; gWakeMwwDedicatedCapture.chunksAttempted = 0; gWakeMwwDedicatedCapture.chunksSubmitted = 0; + gWakeMwwDedicatedCapture.captureStartedAtMs = captureStartMs; + gWakeMwwDedicatedCapture.lastChunkCompletedAtMs = 0; VoiceActivityEndpointConfig endpointConfig; // Endpoint the first wake-gated utterance too, not just conversation replies. // Previously the initial capture ran a fixed length regardless of when the @@ -4763,9 +4782,22 @@ bool beginDedicatedWakeCaptureAfterCue(const RobotEvent& wakeEvent) { return true; } -void finishDedicatedWakeCaptureSession(bool captured) { +void finishDedicatedWakeCaptureSession(bool commit, const char* cancelReason = nullptr) { if (gWakeMwwDedicatedCapture.active) { - finishDedicatedWakeCaptureTurn(gWakeMwwDedicatedCapture.seq, millis()); + const uint32_t terminalAtMs = millis(); + if (commit) { + finishDedicatedWakeCaptureTurn(gWakeMwwDedicatedCapture.seq, terminalAtMs); + } else { + gBridgeWakeGate.cancelActiveTurn( + terminalAtMs, + cancelReason != nullptr ? cancelReason : "dedicated_capture_cancelled"); + if (gBridgeAudioUplink.telemetry().active) { + gBridgeAudioUplink.abort( + terminalAtMs, + cancelReason != nullptr ? cancelReason : "dedicated_capture_cancelled"); + } + gBridgeNetworkSession.update(terminalAtMs); + } } M5.Mic.end(); releaseWakeMwwAudioPause(millis()); @@ -4773,7 +4805,7 @@ void finishDedicatedWakeCaptureSession(bool captured) { const uint32_t captureEndMs = millis(); if (gWakeCueSequence.phase() == WakeCueSequencePhase::Capturing) { - gWakeCueSequence.finishCapture(captureEndMs, captured); + gWakeCueSequence.finishCapture(captureEndMs, commit); } else { gWakeCueSequence.abort(captureEndMs); } @@ -4782,7 +4814,7 @@ void finishDedicatedWakeCaptureSession(bool captured) { gWakeMwwDedicatedCapture.endpoint.cancel(); gWakeMwwDedicatedCapture.active = false; gWakeMwwDedicatedCapture.conversationReplyCapture = false; - if (captured) { + if (commit) { ++gWakeSrProbe.wakeEventsApplied; } } @@ -4798,7 +4830,7 @@ void serviceDedicatedWakeCaptureChunk() { gBridgeWakeGate.isGateOpen(serviceNowMs), gateBeforeCapture.turnActive, gBridgeAudioUplink.telemetry().active)) { - finishDedicatedWakeCaptureSession(gWakeMwwDedicatedCapture.chunksSubmitted > 0); + finishDedicatedWakeCaptureSession(false, "dedicated_capture_authority_expired"); return; } @@ -4809,16 +4841,43 @@ void serviceDedicatedWakeCaptureChunk() { constexpr size_t kMonoSamples = STACKCHAN_MWW_WAKE_UPLINK_CHUNK_SAMPLES; constexpr size_t kRecordChannels = kRecordStereo ? 2u : 1u; constexpr size_t kRecordSamples = kMonoSamples * kRecordChannels; + constexpr uint32_t kChunkAudioMs = + static_cast((kMonoSamples * 1000u) / kSampleRate); + constexpr uint32_t kChunkPhaseDeadlineMs = 250; + constexpr uint32_t kChunkGapDeadlineMs = kChunkAudioMs + kChunkPhaseDeadlineMs; static int16_t recordBuf[kRecordSamples]; static int16_t monoBuf[kMonoSamples]; ++gWakeMwwDedicatedCapture.serviceCalls; ++gWakeMwwDedicatedCapture.chunksAttempted; bool submitFailed = false; + bool captureDiscontinuity = false; VoiceActivityEndpointReason endpointReason = VoiceActivityEndpointReason::None; - if (!recordWakeMwwAudioBlocking(recordBuf, kRecordSamples, kSampleRate, kRecordStereo)) { + uint32_t micCaptureUs = 0; + if (!recordWakeMwwAudioBlocking(recordBuf, + kRecordSamples, + kSampleRate, + kRecordStereo, + kChunkPhaseDeadlineMs, + &micCaptureUs)) { gWakeSrProbe.recordDrops++; + ++gWakeMwwDedicatedCapture.captureTimeouts; + captureDiscontinuity = true; } else { + if (micCaptureUs > gWakeMwwDedicatedCapture.maxMicCaptureUs) { + gWakeMwwDedicatedCapture.maxMicCaptureUs = micCaptureUs; + } + const uint32_t chunkCompletedAtMs = millis(); + if (!dedicatedWakeCaptureChunkContinuous( + micCaptureUs, + kChunkPhaseDeadlineMs * 1000u, + gWakeMwwDedicatedCapture.lastChunkCompletedAtMs, + chunkCompletedAtMs, + kChunkGapDeadlineMs)) { + ++gWakeMwwDedicatedCapture.discontinuities; + captureDiscontinuity = true; + } + gWakeMwwDedicatedCapture.lastChunkCompletedAtMs = chunkCompletedAtMs; if constexpr (kRecordStereo) { constexpr size_t selectedChannel = STACKCHAN_MWW_WAKE_STEREO_MONO_CHANNEL ? 1u : 0u; @@ -4834,7 +4893,9 @@ void serviceDedicatedWakeCaptureChunk() { gWakeSrProbe.samplesFed += kMonoSamples; gWakeSrProbe.lastRecordMs = millis(); - if (gWakeMwwDedicatedCapture.endpoint.telemetry().enabled) { + const uint32_t vadStartedAtUs = micros(); + if (!captureDiscontinuity && + gWakeMwwDedicatedCapture.endpoint.telemetry().enabled) { const uint32_t speechAtBefore = gWakeMwwDedicatedCapture.endpoint.telemetry().lastSpeechAtMs; endpointReason = gWakeMwwDedicatedCapture.endpoint.process( @@ -4846,8 +4907,8 @@ void serviceDedicatedWakeCaptureChunk() { // from bridge UserSpeaking events, which arrive only after host STT. During // our own dedicated capture nothing told it, so the gate closed 6 s after // the wake word and sent utterance_end while the microphone was still - // recording. Every chunk after that hit an uplink with no active turn, - // which is where the audio_uplink_not_active errors came from. + // recording. The gate now cancels rather than commits on expiry, but + // renewing it on observed speech still prevents avoidable turn loss. // // Renewing only on real speech preserves the privacy semantics: the gate // still closes gateOpenMs after the speaker actually stops. @@ -4859,22 +4920,35 @@ void serviceDedicatedWakeCaptureChunk() { gBridgeWakeGate.applyEvent(speakingEvent, gWakeSrProbe.lastRecordMs); } } + const uint32_t vadUs = micros() - vadStartedAtUs; + if (vadUs > gWakeMwwDedicatedCapture.maxVadUs) { + gWakeMwwDedicatedCapture.maxVadUs = vadUs; + } + if (captureDiscontinuity) { + finishDedicatedWakeCaptureSession(false, "dedicated_capture_discontinuity"); + return; + } const BridgeWakeGateTelemetry& gateBeforeSubmit = gBridgeWakeGate.telemetry(); if (!dedicatedWakeCaptureMaySubmit( gBridgeWakeGate.isGateOpen(gWakeSrProbe.lastRecordMs), gateBeforeSubmit.turnActive, gBridgeAudioUplink.telemetry().active)) { - finishDedicatedWakeCaptureSession(gWakeMwwDedicatedCapture.chunksSubmitted > 0); + finishDedicatedWakeCaptureSession(false, "dedicated_capture_authority_expired"); return; } + const uint32_t submitStartedAtUs = micros(); const DedicatedWakeCaptureSubmitResult submitResult = submitDedicatedWakeCaptureChunk( gWakeMwwDedicatedCapture.seq, monoBuf, kMonoSamples, gWakeSrProbe.lastRecordMs); + const uint32_t submitUs = micros() - submitStartedAtUs; + if (submitUs > gWakeMwwDedicatedCapture.maxSubmitUs) { + gWakeMwwDedicatedCapture.maxSubmitUs = submitUs; + } if (submitResult == DedicatedWakeCaptureSubmitResult::Submitted) { ++gWakeMwwDedicatedCapture.chunksSubmitted; } else if (submitResult == DedicatedWakeCaptureSubmitResult::AuthorityExpired) { - finishDedicatedWakeCaptureSession(gWakeMwwDedicatedCapture.chunksSubmitted > 0); + finishDedicatedWakeCaptureSession(false, "dedicated_capture_authority_expired"); return; } else { gWakeMwwUplinkSubmitFailed = gWakeMwwUplinkSubmitFailed + 1u; @@ -4892,10 +4966,21 @@ void serviceDedicatedWakeCaptureChunk() { endpointReason == VoiceActivityEndpointReason::None) { endpointReason = gWakeMwwDedicatedCapture.endpoint.forceMaximum(millis()); } - const bool complete = submitFailed || endpointReason != VoiceActivityEndpointReason::None || - chunkLimitReached; + const bool complete = captureDiscontinuity || submitFailed || + endpointReason != VoiceActivityEndpointReason::None || chunkLimitReached; if (complete) { - finishDedicatedWakeCaptureSession(gWakeMwwDedicatedCapture.chunksSubmitted > 0); + const bool semanticEnd = dedicatedWakeCaptureMayCommit( + endpointReason, + gWakeMwwDedicatedCapture.endpoint.telemetry().enabled, + chunkLimitReached, + submitFailed, + captureDiscontinuity); + const char* cancelReason = captureDiscontinuity + ? "dedicated_capture_discontinuity" + : (submitFailed + ? "dedicated_capture_backpressure" + : "dedicated_capture_maximum"); + finishDedicatedWakeCaptureSession(semanticEnd, semanticEnd ? nullptr : cancelReason); } } @@ -4967,7 +5052,7 @@ void serviceDedicatedWakeCapture(uint32_t nowMs) { const bool started = gWakeMwwPendingCaptureEventReady && beginDedicatedWakeCaptureAfterCue(gWakeMwwPendingCaptureEvent); if (!started) { - finishDedicatedWakeCaptureSession(false); + finishDedicatedWakeCaptureSession(false, "dedicated_capture_start_failed"); } } @@ -6086,6 +6171,10 @@ void printRuntimeStatus() { Serial.print(gBridgeNetworkSession.writer().telemetry().binaryFramesQueued); Serial.print(F(" bridge_network_binary_dropped=")); Serial.print(gBridgeNetworkSession.writer().telemetry().binaryFramesDropped); + Serial.print(F(" bridge_network_write_deferrals=")); + Serial.print(gBridgeNetworkSession.writer().telemetry().writeDeferrals); + Serial.print(F(" bridge_network_write_failures=")); + Serial.print(gBridgeNetworkSession.writer().telemetry().writeFailures); const BridgeEndpointRegistryTelemetry& endpoints = gBridgeEndpointRegistry.telemetry(); Serial.print(F(" bridge_endpoint_registry_ready=")); Serial.print(endpoints.ready ? 1 : 0); @@ -6184,6 +6273,14 @@ void printRuntimeStatus() { Serial.print(uplink.queueFailures); Serial.print(F(" bridge_uplink_last_seq=")); Serial.print(uplink.lastSeq); + Serial.print(F(" bridge_uplink_terminal_pending=")); + Serial.print(uplink.terminalPending ? 1 : 0); + Serial.print(F(" bridge_uplink_terminal_retries=")); + Serial.print(uplink.terminalRetries); + Serial.print(F(" bridge_uplink_terminal_timeouts=")); + Serial.print(uplink.terminalTimeouts); + Serial.print(F(" bridge_uplink_cancel_frames=")); + Serial.print(uplink.cancelFramesQueued); #if STACKCHAN_HAS_MWW_WAKE_PROBE && STACKCHAN_ENABLE_BRIDGE_AUDIO_UPLINK && STACKCHAN_MWW_WAKE_DRIVES_AUDIO_UPLINK Serial.print(F(" mww_uplink_pending=")); Serial.print(gWakeMwwUplinkPendingReady ? 1 : 0); @@ -8411,6 +8508,17 @@ void serveBridgeLeanStatusJson(WiFiClient& client, append(",\"bridge_uplink_bytes\":%lu", static_cast(uplink.bytesQueued)); append(",\"bridge_uplink_errors\":%lu", static_cast(uplink.errors)); append(",\"bridge_uplink_queue_failures\":%lu", static_cast(uplink.queueFailures)); + append(",\"bridge_uplink_terminal_pending\":%s", uplink.terminalPending ? "true" : "false"); + append(",\"bridge_uplink_terminal_retries\":%lu", + static_cast(uplink.terminalRetries)); + append(",\"bridge_uplink_terminal_timeouts\":%lu", + static_cast(uplink.terminalTimeouts)); + append(",\"bridge_uplink_cancel_frames\":%lu", + static_cast(uplink.cancelFramesQueued)); + append(",\"bridge_network_write_deferrals\":%lu", + static_cast(gBridgeNetworkSession.writer().telemetry().writeDeferrals)); + append(",\"bridge_network_write_failures\":%lu", + static_cast(gBridgeNetworkSession.writer().telemetry().writeFailures)); append(",\"bridge_uplink_last_error\":\"%s\"", uplink.lastError); #if STACKCHAN_HAS_MWW_WAKE_PROBE && STACKCHAN_ENABLE_BRIDGE_AUDIO_UPLINK && STACKCHAN_MWW_WAKE_DRIVES_AUDIO_UPLINK append(",\"mww_uplink_pending\":%s", gWakeMwwUplinkPendingReady ? "true" : "false"); @@ -8580,6 +8688,16 @@ void serveBridgeLeanStatusJson(WiFiClient& client, static_cast(gWakeMwwDedicatedCapture.serviceCalls)); append(",\"wake_capture_max_service_us\":%lu", static_cast(gWakeMwwDedicatedCapture.maxServiceUs)); + append(",\"wake_capture_max_mic_us\":%lu", + static_cast(gWakeMwwDedicatedCapture.maxMicCaptureUs)); + append(",\"wake_capture_max_vad_us\":%lu", + static_cast(gWakeMwwDedicatedCapture.maxVadUs)); + append(",\"wake_capture_max_submit_us\":%lu", + static_cast(gWakeMwwDedicatedCapture.maxSubmitUs)); + append(",\"wake_capture_timeouts\":%lu", + static_cast(gWakeMwwDedicatedCapture.captureTimeouts)); + append(",\"wake_capture_discontinuities\":%lu", + static_cast(gWakeMwwDedicatedCapture.discontinuities)); const VoiceActivityEndpointTelemetry& replyVad = gWakeMwwDedicatedCapture.endpoint.telemetry(); append(",\"conversation_reply_vad_enabled\":%s", replyVad.enabled ? "true" : "false"); @@ -8598,6 +8716,10 @@ void serveBridgeLeanStatusJson(WiFiClient& client, static_cast(replyVad.endpointsDetected)); append(",\"conversation_reply_vad_max_fallbacks\":%lu", static_cast(replyVad.maxDurationFallbacks)); + append(",\"conversation_reply_vad_captured_audio_ms\":%lu", + static_cast(replyVad.capturedAudioMs)); + append(",\"conversation_reply_vad_last_speech_audio_ms\":%lu", + static_cast(replyVad.lastSpeechAudioMs)); append(",\"conversation_reply_vad_last_reason\":\"%s\"", voiceActivityEndpointReasonName(replyVad.lastReason)); append(",\"conversation_reply_vad_last_level\":%.4f", diff --git a/test/test_native_logic/test_main.cpp b/test/test_native_logic/test_main.cpp index e01f779..68dacb5 100644 --- a/test/test_native_logic/test_main.cpp +++ b/test/test_native_logic/test_main.cpp @@ -1640,8 +1640,12 @@ void test_voice_activity_endpoint_ends_after_sustained_speech_and_trailing_silen TEST_ASSERT_TRUE(endpoint.telemetry().speechSeen); TEST_ASSERT_EQUAL(static_cast(VoiceActivityEndpointReason::None), static_cast(endpoint.process(silence, 800, 200))); - TEST_ASSERT_EQUAL(static_cast(VoiceActivityEndpointReason::TrailingSilence), + TEST_ASSERT_EQUAL(static_cast(VoiceActivityEndpointReason::None), static_cast(endpoint.process(silence, 800, 300))); + TEST_ASSERT_EQUAL(static_cast(VoiceActivityEndpointReason::None), + static_cast(endpoint.process(silence, 800, 350))); + TEST_ASSERT_EQUAL(static_cast(VoiceActivityEndpointReason::TrailingSilence), + static_cast(endpoint.process(silence, 800, 400))); TEST_ASSERT_FALSE(endpoint.telemetry().active); TEST_ASSERT_EQUAL_UINT32(1, endpoint.telemetry().endpointsDetected); TEST_ASSERT_EQUAL_STRING( @@ -1823,6 +1827,65 @@ void test_voice_activity_endpoint_disabled_path_preserves_fixed_capture() { TEST_ASSERT_EQUAL_UINT32(0, endpoint.telemetry().capturesStarted); } +void test_voice_activity_endpoint_ignores_missing_wall_time_for_semantic_silence() { + VoiceActivityEndpointConfig config; + config.enabled = true; + config.minimumCaptureMs = 100; + config.minimumSpeechMs = 100; + config.trailingSilenceMs = 200; + config.maximumCaptureMs = 500; + VoiceActivityEndpoint endpoint; + TEST_ASSERT_TRUE(endpoint.begin(config, 1000)); + + int16_t speech[800] = {}; + int16_t silence[800] = {}; + fillVoiceEndpointSpeech(speech, 800); + TEST_ASSERT_EQUAL(static_cast(VoiceActivityEndpointReason::None), + static_cast(endpoint.process(speech, 800, 1050))); + TEST_ASSERT_EQUAL(static_cast(VoiceActivityEndpointReason::None), + static_cast(endpoint.process(speech, 800, 1100))); + TEST_ASSERT_TRUE(endpoint.telemetry().speechSeen); + + // An IntentTask/microphone/transport stall is not captured silence. Only one + // 50 ms silence chunk follows, despite the 7.5 second wall-clock jump. + TEST_ASSERT_EQUAL(static_cast(VoiceActivityEndpointReason::None), + static_cast(endpoint.process(silence, 800, 8600))); + TEST_ASSERT_TRUE(endpoint.telemetry().active); + TEST_ASSERT_EQUAL_UINT32(150, endpoint.telemetry().capturedAudioMs); + TEST_ASSERT_EQUAL_UINT32(100, endpoint.telemetry().lastSpeechAudioMs); + + TEST_ASSERT_EQUAL(static_cast(VoiceActivityEndpointReason::None), + static_cast(endpoint.process(silence, 800, 8650))); + TEST_ASSERT_EQUAL(static_cast(VoiceActivityEndpointReason::None), + static_cast(endpoint.process(silence, 800, 8700))); + TEST_ASSERT_EQUAL(static_cast(VoiceActivityEndpointReason::TrailingSilence), + static_cast(endpoint.process(silence, 800, 8750))); + TEST_ASSERT_EQUAL_UINT32(300, endpoint.telemetry().capturedAudioMs); +} + +void test_voice_activity_endpoint_ignores_missing_wall_time_for_maximum() { + VoiceActivityEndpointConfig config; + config.enabled = true; + config.minimumCaptureMs = 100; + config.minimumSpeechMs = 100; + config.trailingSilenceMs = 200; + config.maximumCaptureMs = 500; + VoiceActivityEndpoint endpoint; + TEST_ASSERT_TRUE(endpoint.begin(config, 1000)); + + int16_t silence[800] = {}; + for (uint32_t chunk = 1; chunk < 10; ++chunk) { + const uint32_t sparseWallMs = 1000u + chunk * 5000u; + TEST_ASSERT_EQUAL(static_cast(VoiceActivityEndpointReason::None), + static_cast(endpoint.process(silence, 800, sparseWallMs))); + } + TEST_ASSERT_TRUE(endpoint.telemetry().active); + TEST_ASSERT_EQUAL_UINT32(450, endpoint.telemetry().capturedAudioMs); + TEST_ASSERT_EQUAL(static_cast(VoiceActivityEndpointReason::MaxDuration), + static_cast(endpoint.process(silence, 800, 51000))); + TEST_ASSERT_EQUAL_UINT32(500, endpoint.telemetry().capturedAudioMs); +} + void test_embodied_energy_classifies_with_hysteresis_and_charging_priority() { EmbodiedEnergy energy; energy.reset(0); @@ -5875,16 +5938,29 @@ class CapturingBridgeSocketSink final : public BridgeSocketWriterSink { size_t write(const uint8_t* data, size_t length) override { if (!connected || data == nullptr || length == 0) { + writeWouldBlock = false; + return 0; + } + if (deferredWrites > 0) { + --deferredWrites; + writeWouldBlock = true; return 0; } + writeWouldBlock = false; const size_t allowed = maxWriteBytes == 0 || maxWriteBytes > length ? length : maxWriteBytes; bytes.insert(bytes.end(), data, data + allowed); writes++; return allowed; } + bool lastWriteWouldBlock() const override { + return writeWouldBlock; + } + bool connected = true; size_t maxWriteBytes = 0; + uint32_t deferredWrites = 0; + bool writeWouldBlock = false; uint32_t writes = 0; std::vector bytes; }; @@ -7013,6 +7089,204 @@ void test_bridge_audio_uplink_rejects_bad_sequence_and_limits() { TEST_ASSERT_TRUE(socket.outgoing.empty()); } +void test_bridge_socket_writer_retains_frame_during_transient_backpressure() { + BridgeClient bridge; + TEST_ASSERT_TRUE(bridge.begin()); + BridgeWebSocketTransport transport; + TEST_ASSERT_TRUE(transport.begin(bridge, 701)); + TEST_ASSERT_TRUE(transport.acceptHandshakeResponse( + "HTTP/1.1 101 Switching Protocols\r\n" + "Upgrade: websocket\r\n" + "Connection: Upgrade\r\n" + "Sec-WebSocket-Accept: ok\r\n" + "\r\n", + 702)); + + CapturingBridgeSocketSink sink; + sink.deferredWrites = 1; + BridgeSocketWriter writer; + TEST_ASSERT_TRUE(writer.begin(transport, sink, 0x44556679)); + TEST_ASSERT_TRUE(writer.queueTextFrame("{\"type\":\"utterance_cancel\",\"seq\":7}")); + + TEST_ASSERT_EQUAL(static_cast(BridgeSocketWriterDrainResult::Partial), + static_cast(writer.drainPendingFrame(703))); + TEST_ASSERT_TRUE(writer.telemetry().frameBuffered); + TEST_ASSERT_EQUAL_UINT32(1, writer.telemetry().writeDeferrals); + TEST_ASSERT_EQUAL_UINT32(0, writer.telemetry().writeFailures); + TEST_ASSERT_EQUAL_UINT32(0, writer.telemetry().framesWritten); + + TEST_ASSERT_EQUAL(static_cast(BridgeSocketWriterDrainResult::WroteFrame), + static_cast(writer.drainPendingFrame(704))); + TEST_ASSERT_FALSE(writer.telemetry().frameBuffered); + TEST_ASSERT_EQUAL_UINT32(1, writer.telemetry().framesWritten); +} + +void test_bridge_audio_uplink_retries_retained_end_after_writer_drains() { + BridgeClient bridge; + FakeBridgeNetworkSocket socket; + BridgeNetworkSession session; + connectBridgeNetworkSession(bridge, socket, session, 960); + + BridgeAudioUplinkConfig config; + config.enabled = true; + BridgeAudioUplink uplink; + TEST_ASSERT_TRUE(uplink.begin(config, &session)); + TEST_ASSERT_TRUE(uplink.beginTurn(21, 1000, true)); + session.update(1001); + socket.clearOutgoing(); + + const int16_t samples[] = {100, -200, 300, -400}; + TEST_ASSERT_TRUE(uplink.submitPcmChunk(21, samples, 4, 1002)); + session.update(1003); + std::vector decodedAudio; + TEST_ASSERT_TRUE(decodeMaskedClientBinaryFrame(socket.outgoing, decodedAudio)); + socket.clearOutgoing(); + TEST_ASSERT_TRUE(session.queueTextFrame("{\"type\":\"busy\"}")); + // The text slot is occupied. Ending must stop capture but retain the terminal + // rather than discard the only utterance_end attempt. + TEST_ASSERT_TRUE(uplink.endTurn(21, 1004)); + TEST_ASSERT_FALSE(uplink.telemetry().active); + TEST_ASSERT_TRUE(uplink.telemetry().terminalPending); + TEST_ASSERT_EQUAL_UINT32(0, uplink.telemetry().turnsCompleted); + + session.update(1005); + socket.clearOutgoing(); + TEST_ASSERT_EQUAL(static_cast(BridgeAudioTerminalServiceResult::Queued), + static_cast(uplink.servicePendingTerminal(1006))); + TEST_ASSERT_FALSE(uplink.telemetry().terminalPending); + TEST_ASSERT_EQUAL_UINT32(1, uplink.telemetry().turnsCompleted); + TEST_ASSERT_EQUAL(static_cast(BridgeAudioTerminalKind::End), + static_cast(uplink.telemetry().lastTerminal)); + + session.update(1007); + char decodedEnd[kBridgeEndpointControlResponseMax] = {}; + TEST_ASSERT_TRUE(decodeMaskedClientTextFrame(socket.outgoing, decodedEnd, sizeof(decodedEnd))); + TEST_ASSERT_NOT_NULL(std::strstr(decodedEnd, "\"type\":\"utterance_end\"")); + TEST_ASSERT_NOT_NULL(std::strstr(decodedEnd, "\"audio_bytes\":8")); +} + +void test_bridge_audio_uplink_abort_queues_explicit_cancel() { + BridgeClient bridge; + FakeBridgeNetworkSocket socket; + BridgeNetworkSession session; + connectBridgeNetworkSession(bridge, socket, session, 970); + + BridgeAudioUplinkConfig config; + config.enabled = true; + BridgeAudioUplink uplink; + TEST_ASSERT_TRUE(uplink.begin(config, &session)); + TEST_ASSERT_TRUE(uplink.beginTurn(22, 1100, true)); + session.update(1101); + socket.clearOutgoing(); + + uplink.abort(1102, "capture stalled"); + TEST_ASSERT_FALSE(uplink.telemetry().active); + TEST_ASSERT_FALSE(uplink.telemetry().terminalPending); + TEST_ASSERT_EQUAL_UINT32(1, uplink.telemetry().turnsAborted); + TEST_ASSERT_EQUAL_UINT32(1, uplink.telemetry().cancelFramesQueued); + session.update(1103); + char decodedCancel[kBridgeEndpointControlResponseMax] = {}; + TEST_ASSERT_TRUE( + decodeMaskedClientTextFrame(socket.outgoing, decodedCancel, sizeof(decodedCancel))); + TEST_ASSERT_NOT_NULL(std::strstr(decodedCancel, "\"type\":\"utterance_cancel\"")); + TEST_ASSERT_NOT_NULL(std::strstr(decodedCancel, "\"seq\":22")); + TEST_ASSERT_NOT_NULL(std::strstr(decodedCancel, "\"reason\":\"capture_stalled\"")); +} + +void test_bridge_audio_uplink_terminal_timeout_closes_socket_fail_closed() { + BridgeClient bridge; + FakeBridgeNetworkSocket socket; + BridgeNetworkSession session; + connectBridgeNetworkSession(bridge, socket, session, 980); + + BridgeAudioUplinkConfig config; + config.enabled = true; + config.terminalRetryMs = 100; + BridgeAudioUplink uplink; + TEST_ASSERT_TRUE(uplink.begin(config, &session)); + TEST_ASSERT_TRUE(uplink.beginTurn(23, 1200, true)); + session.update(1201); + socket.clearOutgoing(); + + TEST_ASSERT_TRUE(session.queueTextFrame("{\"type\":\"busy\"}")); + TEST_ASSERT_TRUE(uplink.endTurn(23, 1202)); + TEST_ASSERT_TRUE(uplink.telemetry().terminalPending); + TEST_ASSERT_EQUAL(static_cast(BridgeAudioTerminalServiceResult::Pending), + static_cast(uplink.servicePendingTerminal(1301))); + TEST_ASSERT_TRUE(socket.connected); + TEST_ASSERT_EQUAL(static_cast(BridgeAudioTerminalServiceResult::FailedClosed), + static_cast(uplink.servicePendingTerminal(1302))); + TEST_ASSERT_FALSE(socket.connected); + TEST_ASSERT_EQUAL_UINT32(1, socket.stops); + TEST_ASSERT_EQUAL_UINT32(1, uplink.telemetry().terminalTimeouts); + TEST_ASSERT_EQUAL_UINT32(1, uplink.telemetry().turnsAborted); + TEST_ASSERT_EQUAL_STRING("utterance_terminal_delivery_timeout", + uplink.telemetry().lastError); +} + +void test_bridge_wake_gate_services_retained_terminal_after_writer_drains() { + BridgeClient bridge; + FakeBridgeNetworkSocket socket; + BridgeNetworkSession session; + connectBridgeNetworkSession(bridge, socket, session, 990); + + BridgeAudioUplinkConfig uplinkConfig; + uplinkConfig.enabled = true; + BridgeAudioUplink uplink; + TEST_ASSERT_TRUE(uplink.begin(uplinkConfig, &session)); + BridgeWakeGate gate; + TEST_ASSERT_TRUE(gate.begin(BridgeWakeGateConfig {}, &uplink)); + + RobotEvent wake; + wake.type = EventType::WakeWord; + gate.applyEvent(wake, 1400); + // Do not drain utterance_start before speech end: this reproduces the full + // writer slot that formerly lost utterance_end. + RobotEvent ended; + ended.type = EventType::SpeechEnded; + gate.applyEvent(ended, 1401); + TEST_ASSERT_TRUE(gate.telemetry().turnActive); + TEST_ASSERT_FALSE(uplink.telemetry().active); + TEST_ASSERT_TRUE(uplink.telemetry().terminalPending); + + session.update(1402); + socket.clearOutgoing(); + gate.update(1403); + TEST_ASSERT_FALSE(gate.telemetry().turnActive); + TEST_ASSERT_FALSE(uplink.telemetry().terminalPending); + TEST_ASSERT_EQUAL_UINT32(1, gate.telemetry().turnsCompleted); + session.update(1404); + char decodedEnd[kBridgeEndpointControlResponseMax] = {}; + TEST_ASSERT_TRUE(decodeMaskedClientTextFrame(socket.outgoing, decodedEnd, sizeof(decodedEnd))); + TEST_ASSERT_NOT_NULL(std::strstr(decodedEnd, "\"type\":\"utterance_end\"")); +} + +void test_bridge_wake_gate_observes_direct_uplink_cancel() { + BridgeClient bridge; + FakeBridgeNetworkSocket socket; + BridgeNetworkSession session; + connectBridgeNetworkSession(bridge, socket, session, 995); + + BridgeAudioUplinkConfig uplinkConfig; + uplinkConfig.enabled = true; + BridgeAudioUplink uplink; + TEST_ASSERT_TRUE(uplink.begin(uplinkConfig, &session)); + BridgeWakeGate gate; + TEST_ASSERT_TRUE(gate.begin(BridgeWakeGateConfig {}, &uplink)); + + RobotEvent wake; + wake.type = EventType::WakeWord; + gate.applyEvent(wake, 1500); + session.update(1501); + socket.clearOutgoing(); + uplink.abort(1502, "capture_phase_timeout"); + TEST_ASSERT_TRUE(gate.telemetry().turnActive); + gate.update(1503); + TEST_ASSERT_FALSE(gate.telemetry().turnActive); + TEST_ASSERT_EQUAL_UINT32(1, gate.telemetry().turnsAborted); + TEST_ASSERT_EQUAL_UINT32(0, gate.telemetry().turnsCompleted); +} + void test_bridge_wake_gate_suppresses_turn_when_uplink_disabled() { BridgeAudioUplink uplink; TEST_ASSERT_TRUE(uplink.begin()); @@ -7360,7 +7634,15 @@ void test_bridge_wake_gate_survives_a_long_utterance() { gate.update(turnStartedMs + kBridgeWakeGateMaxTurnMs); TEST_ASSERT_FALSE(gate.telemetry().gateOpen); TEST_ASSERT_FALSE(gate.telemetry().turnActive); - TEST_ASSERT_EQUAL_UINT32(1, gate.telemetry().turnsCompleted); + TEST_ASSERT_EQUAL_UINT32(0, gate.telemetry().turnsCompleted); + TEST_ASSERT_EQUAL_UINT32(1, gate.telemetry().turnsAborted); + session.update(turnStartedMs + kBridgeWakeGateMaxTurnMs + 1u); + char decodedCancel[kBridgeEndpointControlResponseMax] = {}; + TEST_ASSERT_TRUE( + decodeMaskedClientTextFrame(socket.outgoing, decodedCancel, sizeof(decodedCancel))); + TEST_ASSERT_NOT_NULL(std::strstr(decodedCancel, "\"type\":\"utterance_cancel\"")); + TEST_ASSERT_NOT_NULL(std::strstr(decodedCancel, + "\"reason\":\"bridge_wake_gate_max_turn\"")); } void test_bridge_wake_gate_max_turn_outlasts_the_capture_ceiling() { @@ -7418,7 +7700,15 @@ void test_bridge_wake_gate_renews_on_speech_and_expires() { TEST_ASSERT_FALSE(gate.telemetry().gateOpen); TEST_ASSERT_FALSE(gate.telemetry().turnActive); TEST_ASSERT_EQUAL_UINT32(1, gate.telemetry().gatesExpired); - TEST_ASSERT_EQUAL_UINT32(1, gate.telemetry().turnsCompleted); + TEST_ASSERT_EQUAL_UINT32(0, gate.telemetry().turnsCompleted); + TEST_ASSERT_EQUAL_UINT32(1, gate.telemetry().turnsAborted); + session.update(3601); + char decodedCancel[kBridgeEndpointControlResponseMax] = {}; + TEST_ASSERT_TRUE( + decodeMaskedClientTextFrame(socket.outgoing, decodedCancel, sizeof(decodedCancel))); + TEST_ASSERT_NOT_NULL(std::strstr(decodedCancel, "\"type\":\"utterance_cancel\"")); + TEST_ASSERT_NOT_NULL(std::strstr(decodedCancel, + "\"reason\":\"bridge_wake_gate_timeout\"")); } void test_bridge_network_session_reconnects_after_socket_disconnect() { @@ -8277,6 +8567,24 @@ void test_dedicated_wake_capture_submission_requires_all_owners() { TEST_ASSERT_FALSE(dedicatedWakeCaptureMaySubmit(false, false, false)); } +void test_dedicated_wake_capture_continuity_and_commit_fail_closed() { + TEST_ASSERT_TRUE(dedicatedWakeCaptureChunkContinuous(50000, 250000, 0, 100, 300)); + TEST_ASSERT_TRUE(dedicatedWakeCaptureChunkContinuous(50000, 250000, 100, 350, 300)); + TEST_ASSERT_FALSE(dedicatedWakeCaptureChunkContinuous(7503802, 250000, 100, 150, 300)); + TEST_ASSERT_FALSE(dedicatedWakeCaptureChunkContinuous(50000, 250000, 100, 401, 300)); + + TEST_ASSERT_TRUE(dedicatedWakeCaptureMayCommit( + VoiceActivityEndpointReason::TrailingSilence, true, false, false, false)); + TEST_ASSERT_FALSE(dedicatedWakeCaptureMayCommit( + VoiceActivityEndpointReason::MaxDuration, true, true, false, false)); + TEST_ASSERT_FALSE(dedicatedWakeCaptureMayCommit( + VoiceActivityEndpointReason::TrailingSilence, true, false, true, false)); + TEST_ASSERT_FALSE(dedicatedWakeCaptureMayCommit( + VoiceActivityEndpointReason::TrailingSilence, true, false, false, true)); + TEST_ASSERT_TRUE(dedicatedWakeCaptureMayCommit( + VoiceActivityEndpointReason::None, false, true, false, false)); +} + void test_dedicated_wake_capture_keeps_queue_failures_visible_while_authorized() { BridgeClient bridge; FakeBridgeNetworkSocket socket; @@ -8362,15 +8670,16 @@ void test_dedicated_wake_capture_retry_stops_when_backpressure_reaches_gate_edge TEST_ASSERT_EQUAL_UINT32(0, uplink.telemetry().queueFailures); TEST_ASSERT_EQUAL_UINT32(1, uplink.telemetry().chunksQueued); - RobotEvent ended; - ended.type = EventType::SpeechEnded; - gate.applyEvent(ended, kGateEdgeMs); + // Expiry is an integrity cancel, not a semantic utterance end. The host must + // discard the partial PCM rather than transcribe it as if the user stopped. + gate.update(kGateEdgeMs); session.update(kGateEdgeMs + 1u); - char decodedEnd[kBridgeEndpointControlResponseMax] = {}; - TEST_ASSERT_TRUE(decodeMaskedClientTextFrame(socket.outgoing, decodedEnd, sizeof(decodedEnd))); - TEST_ASSERT_NOT_NULL(std::strstr(decodedEnd, "\"type\":\"utterance_end\"")); - TEST_ASSERT_NOT_NULL(std::strstr(decodedEnd, "\"audio_bytes\":8")); - TEST_ASSERT_NOT_NULL(std::strstr(decodedEnd, "\"chunks\":1")); + char decodedCancel[kBridgeEndpointControlResponseMax] = {}; + TEST_ASSERT_TRUE( + decodeMaskedClientTextFrame(socket.outgoing, decodedCancel, sizeof(decodedCancel))); + TEST_ASSERT_NOT_NULL(std::strstr(decodedCancel, "\"type\":\"utterance_cancel\"")); + TEST_ASSERT_NOT_NULL(std::strstr(decodedCancel, "\"audio_bytes\":8")); + TEST_ASSERT_NOT_NULL(std::strstr(decodedCancel, "\"chunks\":1")); } void test_dedicated_wake_capture_stops_cleanly_at_release_gate_boundary() { @@ -8457,9 +8766,7 @@ void test_dedicated_wake_capture_stops_cleanly_at_release_gate_boundary() { gate.isGateOpen(capturedAtMs), gate.telemetry().turnActive, uplink.telemetry().active)) { - RobotEvent ended; - ended.type = EventType::SpeechEnded; - gate.applyEvent(ended, capturedAtMs); + gate.update(capturedAtMs); endpoint.cancel(); captureEndedAtMs = capturedAtMs; expiredDuringCapture = true; @@ -8481,15 +8788,16 @@ void test_dedicated_wake_capture_stops_cleanly_at_release_gate_boundary() { socket.clearOutgoing(); } - // Clean expiry still emits the one closing control frame, after all accepted - // binary frames. Its declaration must match exactly what crossed the wire. + // Privacy-gate expiry emits an explicit cancel after all accepted binary + // frames. Its declaration still matches exactly what crossed the wire. session.update(captureEndedAtMs + 1u); - char decodedEnd[kBridgeEndpointControlResponseMax] = {}; - TEST_ASSERT_TRUE(decodeMaskedClientTextFrame(socket.outgoing, decodedEnd, sizeof(decodedEnd))); - TEST_ASSERT_NOT_NULL(std::strstr(decodedEnd, "\"type\":\"utterance_end\"")); - TEST_ASSERT_NOT_NULL(std::strstr(decodedEnd, "\"seq\":1")); - TEST_ASSERT_NOT_NULL(std::strstr(decodedEnd, "\"audio_bytes\":190400")); - TEST_ASSERT_NOT_NULL(std::strstr(decodedEnd, "\"chunks\":119")); + char decodedCancel[kBridgeEndpointControlResponseMax] = {}; + TEST_ASSERT_TRUE( + decodeMaskedClientTextFrame(socket.outgoing, decodedCancel, sizeof(decodedCancel))); + TEST_ASSERT_NOT_NULL(std::strstr(decodedCancel, "\"type\":\"utterance_cancel\"")); + TEST_ASSERT_NOT_NULL(std::strstr(decodedCancel, "\"seq\":1")); + TEST_ASSERT_NOT_NULL(std::strstr(decodedCancel, "\"audio_bytes\":190400")); + TEST_ASSERT_NOT_NULL(std::strstr(decodedCancel, "\"chunks\":119")); ++endFrames; socket.clearOutgoing(); session.update(captureEndedAtMs + 2u); @@ -8515,7 +8823,8 @@ void test_dedicated_wake_capture_stops_cleanly_at_release_gate_boundary() { TEST_ASSERT_FALSE(gate.telemetry().turnActive); TEST_ASSERT_FALSE(uplink.telemetry().active); TEST_ASSERT_EQUAL_UINT32(1, gate.telemetry().turnsStarted); - TEST_ASSERT_EQUAL_UINT32(1, gate.telemetry().turnsCompleted); + TEST_ASSERT_EQUAL_UINT32(0, gate.telemetry().turnsCompleted); + TEST_ASSERT_EQUAL_UINT32(1, gate.telemetry().turnsAborted); TEST_ASSERT_FALSE(endpoint.telemetry().speechSeen); TEST_ASSERT_EQUAL_UINT32(0, endpoint.telemetry().speechChunks); } @@ -8961,6 +9270,8 @@ int main() { RUN_TEST(test_voice_activity_endpoint_default_preserves_one_and_a_half_second_natural_pause); RUN_TEST(test_voice_activity_endpoint_default_continuous_speech_still_hits_maximum); RUN_TEST(test_voice_activity_endpoint_disabled_path_preserves_fixed_capture); + RUN_TEST(test_voice_activity_endpoint_ignores_missing_wall_time_for_semantic_silence); + RUN_TEST(test_voice_activity_endpoint_ignores_missing_wall_time_for_maximum); RUN_TEST(test_audio_capture_adapter_disabled_default_is_ready_without_source); RUN_TEST(test_audio_capture_adapter_rejects_oversized_window); RUN_TEST(test_audio_capture_adapter_records_pcm_and_emits_reflex_events); @@ -9172,6 +9483,7 @@ int main() { RUN_TEST(test_bridge_socket_writer_writes_pending_endpoint_response_frame); RUN_TEST(test_bridge_socket_writer_retains_partial_frame_until_complete); RUN_TEST(test_bridge_socket_writer_disconnected_keeps_pending_response); + RUN_TEST(test_bridge_socket_writer_retains_frame_during_transient_backpressure); RUN_TEST(test_bridge_socket_writer_writes_queued_text_frame); RUN_TEST(test_bridge_socket_writer_bounds_queued_text_frame); RUN_TEST(test_bridge_socket_writer_writes_queued_binary_frame); @@ -9193,12 +9505,18 @@ int main() { RUN_TEST(test_bridge_audio_uplink_requires_wake_gate_before_start); RUN_TEST(test_bridge_audio_uplink_queues_start_chunk_and_end_frames); RUN_TEST(test_bridge_audio_uplink_rejects_bad_sequence_and_limits); + RUN_TEST(test_bridge_audio_uplink_retries_retained_end_after_writer_drains); + RUN_TEST(test_bridge_audio_uplink_abort_queues_explicit_cancel); + RUN_TEST(test_bridge_audio_uplink_terminal_timeout_closes_socket_fail_closed); RUN_TEST(test_bridge_wake_gate_suppresses_turn_when_uplink_disabled); RUN_TEST(test_bridge_wake_gate_starts_and_completes_uplink_turn); + RUN_TEST(test_bridge_wake_gate_services_retained_terminal_after_writer_drains); + RUN_TEST(test_bridge_wake_gate_observes_direct_uplink_cancel); RUN_TEST(test_bridge_wake_gate_can_start_uplink_turn_from_speech_when_enabled); RUN_TEST(test_bridge_wake_gate_renews_on_speech_and_expires); RUN_TEST(test_bridge_wake_gate_survives_a_long_utterance); RUN_TEST(test_dedicated_wake_capture_submission_requires_all_owners); + RUN_TEST(test_dedicated_wake_capture_continuity_and_commit_fail_closed); RUN_TEST(test_dedicated_wake_capture_keeps_queue_failures_visible_while_authorized); RUN_TEST(test_dedicated_wake_capture_retry_stops_when_backpressure_reaches_gate_edge); RUN_TEST(test_dedicated_wake_capture_stops_cleanly_at_release_gate_boundary);