diff options
| author | Craig Jennings <c@cjennings.net> | 2026-09-25 13:05:42 -0400 |
|---|---|---|
| committer | Craig Jennings <c@cjennings.net> | 2026-09-25 13:05:42 -0400 |
| commit | d447db5ac9063ec12afb8f5bc7b3a0b05c75bad4 (patch) | |
| tree | 1c0231a2528f779a83741937e0da6fea0c409da9 /working/meeting-transcription-service/src/transcribe-worker | |
| parent | adab9abae11792eaa17711a63a74efa68885a0be (diff) | |
| download | archsetup-d447db5ac9063ec12afb8f5bc7b3a0b05c75bad4.tar.gz archsetup-d447db5ac9063ec12afb8f5bc7b3a0b05c75bad4.zip | |
This is a self-hosted transcription service: whisper.cpp for the words, pyannote for the speaker labels. It has been running on ratio since 2026-09-17, with velox as the offline fallback. A systemd path unit watches a filesystem queue and starts a oneshot worker per job. There is no network listener. ssh is the transport, systemd is the daemon, and the filesystem is the queue.
It lands in working/ rather than its final home because two decisions come first. I haven't picked where the code lives in this repo. The Hugging Face token the diarization model needs on its first download also has to be handled, since anyone can read this repo.
Neither blocks the service, which already runs. Both block the install path this repo owes it.
The accompanying note lists what each machine needs. The torch venv is 1.3 GB and the whisper model is a separate download, so the note describes both rather than carrying them here.
Diffstat (limited to 'working/meeting-transcription-service/src/transcribe-worker')
| -rwxr-xr-x | working/meeting-transcription-service/src/transcribe-worker | 255 |
1 files changed, 255 insertions, 0 deletions
diff --git a/working/meeting-transcription-service/src/transcribe-worker b/working/meeting-transcription-service/src/transcribe-worker new file mode 100755 index 0000000..12e64b0 --- /dev/null +++ b/working/meeting-transcription-service/src/transcribe-worker @@ -0,0 +1,255 @@ +#!/usr/bin/env python3 +"""Drain the meeting-transcription queue, one job at a time. + +Layout under the state directory (default ~/.local/state/meeting-transcribe): + + incoming/<id>/ audio.<ext> and an optional job.json, dropped by the client + work/<id>/ the job being processed + done/<id>.txt the transcript, plus done/<id>.json with run metadata + failed/<id>.log what went wrong, by stage + +The client uploads into uploading/<id>/ and renames the folder into incoming/ when +the copy is complete, so the queue only ever lists whole jobs. A folder whose name +starts with a dot is skipped as well, as a second line of defence. Nothing is +written into done/ except by rename, so a reader never sees half a transcript. + +Standard library only. whisper-cli and ffmpeg come from PATH; the diarizer runs in +its own virtualenv and is called as a subprocess. +""" + +from __future__ import annotations + +import argparse +import fcntl +import json +import os +import re +import shutil +import subprocess +import sys +import time +from dataclasses import dataclass, field +from pathlib import Path + +sys.path.insert(0, str(Path(__file__).resolve().parent)) + +import merge_transcript # noqa: E402 + +MAX_ATTEMPTS = 2 +LANGUAGE_RE = re.compile(r"^[A-Za-z]{2,8}(-[A-Za-z0-9]{1,8})*$") +SHARE = Path.home() / ".local/share" + + +@dataclass +class Config: + """Where things live and how the tools are called.""" + + home: Path = Path.home() / ".local/state/meeting-transcribe" + whisper_model: Path = SHARE / "whisper-models/ggml-large-v3-turbo-q5_0.bin" + diarize_cmd: list[str] = field( + default_factory=lambda: [ + str(SHARE / "pyannote-diarize/.venv/bin/python"), + str(Path(__file__).resolve().parent / "diarize.py"), + ] + ) + threads: int = max(1, (os.cpu_count() or 4) // 2) + + +class JobError(Exception): + """A job failed at a named stage.""" + + def __init__(self, stage: str, detail: str) -> None: + super().__init__(f"{stage}: {detail}") + self.stage = stage + self.detail = detail + + +def _run(stage: str, argv: list[str]) -> None: + """Run one tool; raise JobError carrying the tail of its stderr on failure.""" + try: + result = subprocess.run(argv, capture_output=True, text=True, check=False, stdin=subprocess.DEVNULL) + except OSError as err: + raise JobError(stage, f"cannot run {argv[0]}: {err.strerror}") from err + if result.returncode != 0: + tail = "\n".join(result.stderr.strip().splitlines()[-15:]) + raise JobError(stage, f"exit {result.returncode}\n{tail}") + + +def _read_job(folder: Path) -> dict: + """Load and validate job.json. A missing file means all defaults.""" + path = folder / "job.json" + if not path.exists(): + return {} + try: + job = json.loads(path.read_text(encoding="utf-8")) + if not isinstance(job, dict): + raise ValueError("expected an object") + except (OSError, ValueError) as err: + raise JobError("job", f"job.json is unreadable: {err}") from err + language = job.get("language", "en") + if not isinstance(language, str) or not LANGUAGE_RE.match(language): + raise JobError("job", f"job.json: bad language {language!r}") + for key in ("speakers", "min_speakers", "max_speakers"): + value = job.get(key) + if value is not None and (isinstance(value, bool) or not isinstance(value, int) or value < 1): + raise JobError("job", f"job.json: {key} must be a positive whole number, got {value!r}") + return job + + +def _diarize_options(job: dict) -> list[str]: + options: list[str] = [] + if job.get("speakers") is not None: + return ["--speakers", str(job["speakers"])] + if job.get("min_speakers") is not None: + options += ["--min-speakers", str(job["min_speakers"])] + if job.get("max_speakers") is not None: + options += ["--max-speakers", str(job["max_speakers"])] + return options + + +def process(folder: Path, config: Config) -> tuple[list[str], dict]: + """Run one job folder through the pipeline. Returns transcript lines and metadata.""" + job = _read_job(folder) + audio = next((p for p in sorted(folder.iterdir()) if p.name.startswith("audio.")), None) + if audio is None: + raise JobError("job", "no audio file in the job folder") + + wav = folder / "speech.wav" + _run("ffmpeg", ["ffmpeg", "-v", "error", "-y", "-i", str(audio), "-ar", "16000", "-ac", "1", str(wav)]) + + # -ojf keeps normal segments (better text) and adds token offsets for the merge. + # -mc 0 stops whisper feeding its own output back in, which is what sends it + # into repetition loops on long meetings. + prefix = folder / "words" + _run("whisper", [ + "whisper-cli", "-m", str(config.whisper_model), "-f", str(wav), + "-l", job.get("language", "en"), "-ojf", "-mc", "0", + "-t", str(config.threads), "-of", str(prefix), + ]) + + turns = folder / "turns.json" + _run("diarize", [*config.diarize_cmd, str(wav), str(turns), *_diarize_options(job)]) + + try: + loaded_turns = merge_transcript.load_turns_json(turns) + # Whisper invents words in silence. Drop them before the loop guard, so a + # quiet meeting isn't mistaken for a repetition loop. + units, dropped = merge_transcript.drop_outside_speech( + merge_transcript.load_whisper_json(prefix.with_suffix(".json")), loaded_turns + ) + # A short stutter is collapsed and recorded; a long loop still fails the job. + units, collapsed = merge_transcript.collapse_repetitions(units) + lines = merge_transcript.merge(units, loaded_turns) + except ValueError as err: + raise JobError("merge", str(err)) from err + + meta = { + "original_name": job.get("original_name"), + "language": job.get("language", "en"), + "speakers_requested": job.get("speakers"), + "speakers_found": len({t.speaker for t in loaded_turns}), + "lines": len(lines), + "words": sum(len(line.split()) - 3 for line in lines), + "dropped_outside_speech": dropped, + "loops_collapsed": [[phrase, count] for phrase, count in collapsed], + } + return lines, meta + + +def _write_atomic(path: Path, text: str) -> None: + partial = path.with_name(path.name + ".partial") + partial.write_text(text, encoding="utf-8") + partial.replace(path) + + +def _fail(failed_dir: Path, job_id: str, stage: str, detail: str) -> None: + """Record a failed job. Falls back to stderr if even the log cannot be written.""" + try: + _write_atomic(failed_dir / f"{job_id}.log", f"stage: {stage}\n{detail}\n") + except OSError as err: + print(f"{job_id}: {stage}: {detail} (and the failure log could not be written: {err})", file=sys.stderr) + + +def _bump_attempts(folder: Path) -> int: + """Count this run against the job, tolerating a missing or broken job.json.""" + path = folder / "job.json" + try: + job = json.loads(path.read_text(encoding="utf-8")) if path.exists() else {} + if not isinstance(job, dict): + return 1 + except (OSError, ValueError): + return 1 # _read_job reports the real problem + try: + attempts = int(job.get("attempts") or 0) + 1 + except (TypeError, ValueError): + attempts = 1 # a malformed counter counts as a first try + job["attempts"] = attempts + path.write_text(json.dumps(job), encoding="utf-8") + return attempts + + +def drain(config: Config) -> dict[str, int]: + """Process every waiting job. Returns counts of done and failed jobs.""" + dirs = {name: config.home / name for name in ("incoming", "work", "done", "failed")} + for folder in dirs.values(): + folder.mkdir(parents=True, exist_ok=True) + + counts = {"done": 0, "failed": 0} + with open(config.home / "lock", "w", encoding="utf-8") as lock: + try: + fcntl.flock(lock, fcntl.LOCK_EX | fcntl.LOCK_NB) + except OSError: + return counts # another worker is draining; it will reach these jobs + + while True: + # Jobs stranded in work/ by a crash go first, then arrivals, oldest first. + waiting = sorted( + (p for d in (dirs["work"], dirs["incoming"]) for p in d.iterdir() + if p.is_dir() and not p.name.startswith(".")), + key=lambda p: (p.parent.name != "work", p.stat().st_mtime), + ) + if not waiting: + return counts + source = waiting[0] + job_id = source.name + done_txt = dirs["done"] / f"{job_id}.txt" + if done_txt.exists(): + shutil.rmtree(source) + continue + + folder = dirs["work"] / job_id + if source != folder: + source.rename(folder) + started = time.monotonic() + try: + if _bump_attempts(folder) > MAX_ATTEMPTS: + raise JobError("worker", f"gave up after {MAX_ATTEMPTS} attempts; the job kept crashing") + lines, meta = process(folder, config) + meta["seconds"] = round(time.monotonic() - started, 1) + _write_atomic(dirs["done"] / f"{job_id}.json", json.dumps(meta)) + _write_atomic(done_txt, "\n".join(lines) + "\n") + (dirs["failed"] / f"{job_id}.log").unlink(missing_ok=True) + counts["done"] += 1 + except JobError as err: + _fail(dirs["failed"], job_id, err.stage, err.detail) + counts["failed"] += 1 + except Exception as err: # noqa: BLE001 - one bad job must never stop the queue + _fail(dirs["failed"], job_id, "worker", f"{type(err).__name__}: {err}") + counts["failed"] += 1 + finally: + shutil.rmtree(folder, ignore_errors=True) + + +def main(argv: list[str]) -> int: + """CLI: drain the queue once and exit. systemd's path unit calls this.""" + parser = argparse.ArgumentParser(description="Drain the meeting-transcription queue.") + parser.add_argument("--home", type=Path, help="state directory (default ~/.local/state/meeting-transcribe)") + args = parser.parse_args(argv) + config = Config(home=args.home) if args.home else Config() + counts = drain(config) + print(f"done={counts['done']} failed={counts['failed']}", file=sys.stderr) + return 0 + + +if __name__ == "__main__": + sys.exit(main(sys.argv[1:])) |
