diff --git a/app/services/qwen3_websocket_asr.py b/app/services/qwen3_websocket_asr.py index fcd0bf4..3aba177 100644 --- a/app/services/qwen3_websocket_asr.py +++ b/app/services/qwen3_websocket_asr.py @@ -267,6 +267,8 @@ class ConnectionContext: last_partial_raw_text: str = "" last_partial_display_text: str = "" stable_partial_prefix: str = "" + pending_partial_revision_text: str = "" + pending_partial_revision_rounds: int = 0 segment_observed_text: str = "" segment_observed_language: str = "" best_partial_text: str = "" @@ -1349,6 +1351,8 @@ class Qwen3ASRService: ctx.last_partial_raw_text = "" ctx.last_partial_display_text = "" ctx.stable_partial_prefix = "" + ctx.pending_partial_revision_text = "" + ctx.pending_partial_revision_rounds = 0 ctx.segment_observed_text = "" ctx.segment_observed_language = "" ctx.best_partial_text = "" @@ -1696,6 +1700,8 @@ class Qwen3ASRService: if not text: ctx.last_partial_raw_text = "" ctx.stable_partial_prefix = "" + ctx.pending_partial_revision_text = "" + ctx.pending_partial_revision_rounds = 0 return "" previous_raw = ctx.last_partial_raw_text @@ -1719,9 +1725,30 @@ class Qwen3ASRService: 1, ) if len(stable_prefix) - stable_common >= tolerated_divergence and previous_emitted: - return previous_emitted - stable_prefix = text[:stable_common] - stable_prefix = text[: self._stable_prefix_cutoff(text, len(stable_prefix))] + # A single shorter or divergent snapshot may be a transient model + # revision. Hold the visible text until the same correction repeats. + if text == ctx.pending_partial_revision_text: + ctx.pending_partial_revision_rounds += 1 + else: + ctx.pending_partial_revision_text = text + ctx.pending_partial_revision_rounds = 1 + + if ctx.pending_partial_revision_rounds < 2: + ctx.last_partial_raw_text = text + return previous_emitted + + # The repeated revision is stable enough to replace the old prefix. + stable_prefix = text[: self._stable_prefix_cutoff(text, stable_common)] + ctx.pending_partial_revision_text = "" + ctx.pending_partial_revision_rounds = 0 + else: + ctx.pending_partial_revision_text = "" + ctx.pending_partial_revision_rounds = 0 + stable_prefix = text[:stable_common] + stable_prefix = text[: self._stable_prefix_cutoff(text, len(stable_prefix))] + else: + ctx.pending_partial_revision_text = "" + ctx.pending_partial_revision_rounds = 0 ctx.stable_partial_prefix = stable_prefix ctx.last_partial_raw_text = text @@ -2789,6 +2816,13 @@ class Qwen3ASRService: ) observed = update_observed_text(ctx, current, current_language) visible_observed = self._trim_previous_segment_overlap(ctx, observed) + if has_native_stream_snapshot: + # Qwen returns full-utterance snapshots; stabilize the + # displayed prefix while allowing repeated corrections. + visible_observed = self._stabilize_partial_text( + ctx, + visible_observed, + ) partial_display = self._clip_partial_text(visible_observed, ctx) if ( partial_display @@ -2934,6 +2968,10 @@ class Qwen3ASRService: emit_segment_start=False, ) + # Drain pending diarization before the final event; the browser closes + # the socket on `end`, so later speaker updates would otherwise be lost. + await ctx.speaker_job_queue.join() + final_updates = self._recluster_confirmed_segments(ctx) for updated in final_updates: sentence_payload = self._build_tencent_sentence(