360 lines
14 KiB
Python
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)
|