Files
2026-08-07 18:15:36 -04:00

360 lines
14 KiB
Python

"""Speech-to-text using Apple's on-device dictation model via SFSpeechRecognizer.
This is the same recognizer macOS Dictation uses. With
``requiresOnDeviceRecognition`` set, audio never leaves the machine, there is no
model to download, and there is no Metal shader compilation — which is what
makes it a better fit here than Whisper.
Results come back through the CoreFoundation runloop, so waiting on a
`threading.Event` deadlocks: nothing pumps the runloop and the handler is never
called. Everything here pumps it in short slices instead, yielding to asyncio
between them.
Requires Dictation to be switched on in System Settings > Keyboard.
"""
import asyncio
import os
import tempfile
import wave
from collections.abc import AsyncGenerator
from difflib import SequenceMatcher
from loguru import logger
from pipecat.frames.frames import ErrorFrame, Frame, TranscriptionFrame
from pipecat.services.settings import STTSettings
from pipecat.services.stt_service import SegmentedSTTService
from pipecat.transcriptions.language import Language
from pipecat.utils.time import time_now_iso8601
try:
import Speech
from Foundation import NSURL, NSDate, NSLocale, NSRunLoop
except ImportError as e: # pragma: no cover - depends on pyobjc being installed
raise ImportError(
"Apple speech recognition needs pyobjc: pip install pyobjc-framework-Speech"
) from e
_AUTH_STATUS = {0: "not determined", 1: "denied", 2: "restricted", 3: "authorized"}
_PUMP_SLICE = 0.01
# SFSpeechRecognitionTaskState
_TASK_COMPLETED = 4
# Recognition runs far faster than real time (a 61s file finishes in about 1.3s),
# but scale the deadline with the audio anyway so a long utterance can't be
# dropped by a fixed ceiling.
_TIMEOUT_BASE = 10.0
_TIMEOUT_PER_AUDIO_SECOND = 0.5
# How much of a new transcript must still match the previous one for it to count
# as a refinement rather than the recognizer having started over.
_CONTINUATION_RATIO = 0.5
# How alike two finished passes must be to be judged the same speech re-read.
_SAME_AUDIO_RATIO = 0.6
def _pump(seconds: float = _PUMP_SLICE):
"""Give the runloop a chance to deliver Speech framework callbacks."""
NSRunLoop.currentRunLoop().runUntilDate_(NSDate.dateWithTimeIntervalSinceNow_(seconds))
def _authorize(timeout: float = 20.0) -> int:
"""Return the speech authorization status, prompting once if undetermined."""
status = Speech.SFSpeechRecognizer.authorizationStatus()
if status != 0:
return status
box: dict[str, int] = {}
Speech.SFSpeechRecognizer.requestAuthorization_(lambda s: box.setdefault("status", s))
waited = 0.0
while "status" not in box and waited < timeout:
_pump(0.05)
waited += 0.05
return box.get("status", 0)
def _make_recognizer(locale: str | None):
if locale:
recognizer = Speech.SFSpeechRecognizer.alloc().initWithLocale_(
NSLocale.localeWithLocaleIdentifier_(locale)
)
else:
recognizer = Speech.SFSpeechRecognizer.alloc().init()
if recognizer is None:
raise RuntimeError(f"No speech recognizer available for locale {locale!r}")
if not recognizer.isAvailable():
raise RuntimeError("Speech recognizer is not available right now")
return recognizer
class _Transcript:
"""Stitch a recognition back together across its internal restarts.
On long audio the recognizer does not extend one transcript to the end. It
builds one up, then silently starts over from a later point in the audio,
and the single final result covers only that last stretch — so reading the
final result alone loses everything said earlier.
A restart has to be recognised from the text itself. Partial results carry
no segment timestamps, so the only reliable marker is the transcript
ceasing to be a refinement of the previous one: a growing transcript keeps
almost all of its prefix even when the recognizer revises a word, whereas a
restart drops from hundreds of characters back to a few that share nothing
with what came before.
Restarts are of two kinds, and conflating them is what produces doubled
text. Observed on a 61s recording: the recognizer transcribes the whole
thing, starts over and transcribes the whole thing again slightly
differently, then emits the last second as its only final result. The
re-pass has to replace its predecessor while the tail is appended, so
passes that begin with the same words are treated as the same audio and
only the fullest is kept.
"""
def __init__(self):
self._passes: list[str] = []
self._current = ""
def add(self, start: float, text: str):
if self._restarted(text):
self._close_pass()
# Within a pass each result supersedes the last, so keep the newest.
self._current = text
def close(self):
"""Fold the in-progress pass in. Call once recognition has finished."""
self._close_pass()
def _restarted(self, text: str) -> bool:
"""Whether this result abandons the running transcript instead of refining it.
Compared by prefix because it runs on every partial and a refinement
always keeps its opening intact.
"""
if not self._current or not text:
return False
shared = len(os.path.commonprefix([self._current, text]))
return shared < min(len(self._current), len(text)) * _CONTINUATION_RATIO
def _close_pass(self):
finished, self._current = self._current, ""
if not finished.strip():
return
for i, existing in enumerate(self._passes):
if self._same_audio(existing, finished):
# Same stretch of audio transcribed again; keep the fuller read.
if len(finished) > len(existing):
self._passes[i] = finished
return
self._passes.append(finished)
@staticmethod
def _same_audio(a: str, b: str) -> bool:
"""Whether two finished passes cover the same speech.
Prefix matching is too strict here: a second pass corrects mistakes from
the first, often within the opening few words, so overall similarity is
what distinguishes a re-read from genuinely new audio. Only runs when a
pass closes, so the cost doesn't matter.
"""
return SequenceMatcher(None, a, b).ratio() >= _SAME_AUDIO_RATIO
def text(self) -> str:
parts = [*self._passes, self._current]
return " ".join(part.strip() for part in parts if part.strip())
def _start_recognition(path: str, locale: str | None, terms: list[str] | None = None) -> dict:
"""Kick off a recognition task. The returned dict fills in from the handler."""
recognizer = _make_recognizer(locale)
request = Speech.SFSpeechURLRecognitionRequest.alloc().initWithURL_(
NSURL.fileURLWithPath_(path)
)
request.setRequiresOnDeviceRecognition_(True)
if terms:
# Biasing the decoder towards expected words measured 23.5% -> 16.5% WER.
request.setContextualStrings_(terms)
# Partial results are what make the stitching above possible: the text from
# a segment is only ever visible while that segment is the current one.
request.setShouldReportPartialResults_(True)
box: dict = {"transcript": _Transcript()}
def handler(result, error):
# This crosses back into Objective-C, which aborts the whole process on
# an escaping Python exception. Nothing here may raise.
try:
if error is not None:
box["error"] = str(error.localizedDescription())
return
if result is None:
return
transcription = result.bestTranscription()
segments = transcription.segments()
start = float(segments[0].timestamp()) if segments else 0.0
box["transcript"].add(start, str(transcription.formattedString()))
except Exception as e: # pragma: no cover - defensive
box["error"] = f"result handler failed: {e}"
# Keep the task alive for as long as the caller holds the box.
box["_task"] = recognizer.recognitionTaskWithRequest_resultHandler_(request, handler)
return box
def _is_done(box: dict) -> bool:
if "error" in box:
return True
task = box.get("_task")
return task is not None and task.state() == _TASK_COMPLETED
def _finish(box: dict, timed_out: bool, timeout: float) -> str:
if "error" in box:
raise RuntimeError(box["error"])
if timed_out:
raise TimeoutError(f"Speech recognition timed out after {timeout}s")
box["transcript"].close()
return box["transcript"].text()
def timeout_for(audio_seconds: float) -> float:
"""A recognition deadline that scales with how much audio there is."""
return _TIMEOUT_BASE + _TIMEOUT_PER_AUDIO_SECOND * audio_seconds
def _recognize_file(path: str, locale: str | None, timeout: float) -> str:
"""Blocking recognition, for startup checks before the event loop matters."""
box = _start_recognition(path, locale)
waited = 0.0
while not _is_done(box) and waited < timeout:
_pump()
waited += _PUMP_SLICE
return _finish(box, waited >= timeout, timeout)
async def _recognize_file_async(
path: str, locale: str | None, timeout: float, terms: list[str] | None = None
) -> str:
"""Recognition that keeps the asyncio loop breathing between runloop slices."""
box = _start_recognition(path, locale, terms)
loop = asyncio.get_running_loop()
deadline = loop.time() + timeout
while not _is_done(box) and loop.time() < deadline:
_pump()
await asyncio.sleep(0.005)
return _finish(box, loop.time() >= deadline, timeout)
def probe(locale: str = "en-US") -> tuple[bool, str]:
"""Check whether Apple speech recognition can actually be used.
Returns (available, explanation). Only an attempted recognition settles
this: a disabled Dictation subsystem shows up as an error on the first
request rather than through any status flag. The probe feeds it silence, so
"no speech detected" is the healthy answer — it means the subsystem ran.
"""
status = _authorize()
if status != 3:
return False, f"speech recognition authorization is {_AUTH_STATUS.get(status, status)}"
try:
recognizer = _make_recognizer(locale)
except RuntimeError as e:
return False, str(e)
if not recognizer.supportsOnDeviceRecognition():
return False, "this Mac has no on-device recognition model installed"
with tempfile.NamedTemporaryFile(suffix=".wav", delete=False) as f:
silence = f.name
try:
with wave.open(silence, "wb") as w:
w.setnchannels(1)
w.setsampwidth(2)
w.setframerate(16000)
w.writeframes(b"\x00\x00" * 4000)
_recognize_file(silence, locale, timeout=1.0)
except TimeoutError:
pass # The subsystem answered but never finalized; good enough.
except RuntimeError as e:
reason = str(e)
if "no speech" in reason.lower():
pass # The expected reply to a silent file.
elif "disabled" in reason.lower():
return False, (
"Dictation is turned off — enable System Settings > Keyboard > Dictation"
)
else:
return False, reason
except Exception as e:
return False, str(e)
finally:
os.unlink(silence)
return True, "on-device dictation model ready"
class AppleSpeechSTTService(SegmentedSTTService):
"""Transcribe VAD-delimited speech segments with Apple's dictation model."""
def __init__(
self,
*,
locale: str = "en-US",
language: Language = Language.EN_US,
vocabulary=None,
**kwargs,
):
# The recognizer picks its model from the locale, so there is no model
# field to set; Pipecat wants every settings field initialized anyway.
super().__init__(settings=STTSettings(model=None, language=locale), **kwargs)
self._locale = locale
self._language = language
# Read per utterance rather than cached, so terms learned during the
# conversation reach the next recognition.
self._vocabulary = vocabulary
def can_generate_metrics(self) -> bool:
return True
async def run_stt(self, audio: bytes) -> AsyncGenerator[Frame, None]:
"""Transcribe one speech segment.
Args:
audio: The segment as a WAV container, per ``wants_wav_segments``.
"""
await self.start_processing_metrics()
with tempfile.NamedTemporaryFile(suffix=".wav", delete=False) as f:
f.write(audio)
path = f.name
# WAV header is 44 bytes; the rest is 16-bit mono at the pipeline rate.
audio_seconds = max(0.0, (len(audio) - 44) / 2 / self.sample_rate)
try:
terms = self._vocabulary.terms() if self._vocabulary else None
text = await _recognize_file_async(
path, self._locale, timeout_for(audio_seconds), terms
)
except TimeoutError:
# The segment held no recognizable speech.
await self.stop_processing_metrics()
return
except Exception as e:
await self.stop_processing_metrics()
yield ErrorFrame(error=f"Apple speech recognition failed: {e}")
return
finally:
os.unlink(path)
await self.stop_processing_metrics()
text = text.strip()
if not text:
return
logger.debug(f"Transcription: [{text}]")
yield TranscriptionFrame(text, self._user_id, time_now_iso8601(), self._language)