diff --git a/skills/mc-prompter/SKILL.md b/skills/mc-prompter/SKILL.md index e90e7cc..d423ea4 100644 --- a/skills/mc-prompter/SKILL.md +++ b/skills/mc-prompter/SKILL.md @@ -1,33 +1,38 @@ --- name: mc-prompter -description: Browser teleprompter for the record stage and standalone shows. A service skill like mc-audio, no stage, no gate, no project.json state. Launch a local prompter server, feed it the project script.md or any text, and the creator records at their own pace with a phone remote over LAN. Classic teleprompter only in this version; voice-follow and producer mode are planned tiers. +description: Browser teleprompter for the record stage and standalone shows. A service skill like mc-audio, no stage, no gate, no project.json state. Launch a local prompter server, feed it the project script.md or any text, and the creator records at their own pace with a phone remote over LAN. Voice-follow scrolling (local streaming ASR, consent-gated model download) tracks the speaker through the script; producer mode is a planned tier. --- # mc-prompter -The record stage is creator-owned; this skill hands the creator a teleprompter for it. It is a service skill: it owns no stage, stops at no gate, and writes no project state. It launches a local web server that serves a fullscreen prompter display, a phone remote, a home page for loading and editing the script, and an OBS overlay placeholder. Everything runs offline on the creator's machine; no models, no downloads, no external requests. +The record stage is creator-owned; this skill hands the creator a teleprompter for it. It is a service skill: it owns no stage, stops at no gate, and writes no project state. It launches a local web server that serves a fullscreen prompter display, a phone remote, a home page for loading and editing the script, and an OBS overlay placeholder. The classic prompter runs offline with no models and no downloads. Voice-follow is an optional tier on top: local streaming ASR follows the speaker through the script and scrolls to match; it needs the prompter-lab workspace, whose one-time model download is consent-gated below. ## Steps -1. Load this skill's own surface (`uv run {project-root}/_bmad/scripts/resolve_customization.py --skill {skill-root}`; run `{workflow.activation_steps_prepend}` now, `{workflow.activation_steps_append}` after this step, and hold `{workflow.persistent_facts}` as standing context). Take the default port from `[prompter] port`. If a studio config exists, also load it (`uv run {project-root}/_bmad/scripts/resolve_config.py --project-root {project-root} --key modules.manticore`) and take `[prompter] port` and `[owner] wpm` when present; a missing studio config is fine here, unlike the stage skills, because standalone shows need no studio. +1. Load this skill's own surface (`uv run {project-root}/_bmad/scripts/resolve_customization.py --skill {skill-root}`; run `{workflow.activation_steps_prepend}` now, `{workflow.activation_steps_append}` after this step, and hold `{workflow.persistent_facts}` as standing context). Take the defaults from `[prompter]`: `port`, `asr-provider`, and `workspace`. If a studio config exists, also load it (`uv run {project-root}/_bmad/scripts/resolve_config.py --project-root {project-root} --key modules.manticore`) and take `[prompter]` values, `[owner] wpm`, and `engines-path` when present; a missing studio config is fine here, unlike the stage skills, because standalone shows need no studio. 2. Locate the script. Inside a pipeline project ("record with the teleprompter"), it is the project's `script.md` under the project folder. Standalone, it is any file path the creator names, markdown or plain text. No file at all is also valid: launch without `--script` and the creator pastes text on the home page. -3. Launch: `uv run {skill-root}/scripts/run_prompter.py --script ` with `--port ` when the config or the creator sets one and `--owner-wpm ` when `[owner] wpm` is known. Add `--lan` only when the creator wants the phone remote or a tablet display; it binds the LAN and on Windows triggers a firewall consent dialog. The launcher probes the port, prints the local URL, the remote URL with its session token, and the session file path, then keeps the server running until Ctrl-C. -4. Give the creator the URLs from the launcher output: the prompt page for the recording display, and the remote URL (token included) to open on a phone when `--lan` is on. Briefly explain the pages: `/` is home (load a file, paste text, edit in place, copy the remote URL), `/prompt` is the fullscreen scroller with keyboard controls and a settings drawer (press `?` there for all shortcuts), `/remote` is the phone controller, `/overlay` is an OBS browser-source placeholder for now. -5. Explain what the display does with pipeline markers: paragraphs carrying a `[TAKE ...]` marker render dimmed with a "have it already" badge because that line was already spoken well in the interview footage, and a toggle hides them entirely; sentences flagged `[INVENTED]` get a subtle badge, toggleable off; other bracketed text renders as a dimmed note and is never counted in timing. -6. If the creator edits the script from the home page and saves, the server first copies the current file to a timestamped backup under the temp session directory, then writes the edit back to the source file, so the prompted text and the pipeline artifact never silently diverge. "Session only" applies the edit without touching the file. -7. When the creator asks for voice-follow scrolling or producer mode, say plainly that those are planned tiers that have not landed yet and the classic prompter is what works today. Never pretend they work and never improvise a substitute. +3. Workspace check, only when the creator wants voice-follow or asks about it (the classic prompter needs none of this; when `asr-provider` is `none`, skip to launch). Resolve the workspace as `{engines-path}/{prompter.workspace}` when a studio config exists; otherwise ask the creator for a location or default to `~/mc-prompter-lab`. Run `uv run {skill-root}/scripts/ensure_workspace.py --workspace --check`. Ready (exit 0): proceed. Not ready (exit 4): tell the creator plainly that a bootstrap downloads about 465 MB of ASR model files plus a small Python venv, ask for their explicit go-ahead, and only then run the same command without `--check`. Declining is fine: launch the classic prompter without `--workspace`. The bootstrap is idempotent; an existing validated workspace is reused, never rebuilt, and downloads run long, so report progress rather than going silent. +4. Launch: `uv run {skill-root}/scripts/run_prompter.py --script ` with `--port ` when the config or the creator sets one and `--owner-wpm ` when `[owner] wpm` is known. Add `--workspace ` when the workspace is ready and `--asr-provider ` when the config sets one; the launcher verifies readiness itself and falls back to the classic prompter with a printed notice when the workspace is missing or incomplete. Add `--lan` only when the creator wants the phone remote or a tablet display; it binds the LAN and on Windows triggers a firewall consent dialog. The launcher probes the port, prints the local URL, the remote URL with its session token, whether voice-follow is available, and the session file path, then keeps the server running until Ctrl-C. +5. Give the creator the URLs from the launcher output: the prompt page for the recording display, and the remote URL (token included) to open on a phone when `--lan` is on. Briefly explain the pages: `/` is home (load a file, paste text, edit in place, copy the remote URL), `/prompt` is the fullscreen scroller with keyboard controls and a settings drawer (press `?` there for all shortcuts), `/remote` is the phone controller, `/overlay` is an OBS browser-source placeholder for now. +6. When voice-follow is available, explain how it works on `/prompt`: it is off by default; the toggle lives in the settings drawer and as a HUD chip, and it only appears in a browser on the server machine itself, because the microphone is captured there (a tablet or phone pointed at the page is display-only). The first enable opens a preflight panel: pick the microphone, watch the live level meter, confirm the applied audio settings, and read a few words until the tracking check passes. After that the scroll follows the voice; silence or off-script ad-libs hold the scroll (HOLD chip) and it resumes when the speaker returns to the script; clicking any word re-anchors instantly; a BEHIND chip means ASR is lagging real time on this machine. Manual controls keep working throughout, and toggling voice-follow off returns to normal wpm scrolling at the current position. +7. Explain what the display does with pipeline markers: paragraphs carrying a `[TAKE ...]` marker render dimmed with a "have it already" badge because that line was already spoken well in the interview footage, and a toggle hides them entirely; sentences flagged `[INVENTED]` get a subtle badge, toggleable off; other bracketed text renders as a dimmed note and is never counted in timing or matched by voice-follow. +8. If the creator edits the script from the home page and saves, the server first copies the current file to a timestamped backup under the temp session directory, then writes the edit back to the source file, so the prompted text and the pipeline artifact never silently diverge. "Session only" applies the edit without touching the file. +9. When the creator asks for producer mode (rundown, timing cues), say plainly that it is a planned tier that has not landed yet. The same honesty applies to `asr-provider` values: `nemotron-streaming` and `none` work today; `zipformer-small` is a planned lane and the server exits with a pointer if it is selected. Never pretend a planned lane works and never improvise a substitute. ## Rules - The prompter never advances pipeline state; when a pipeline project is being recorded, the record stage remains the creator's, and mc-pipeline stays the source of truth. +- The workspace bootstrap downloads about 465 MB; it runs only after the creator's explicit consent, and declining always leaves a working classic prompter. - The server binds localhost by default; `--lan` is opt-in and the remote URL carries a per-session token so a random LAN device cannot drive the prompter mid-show. -- Loading files by path and saving edits work only from the server machine, never from a LAN device. +- Loading files by path and saving edits work only from the server machine, never from a LAN device. Microphone capture likewise happens only on the server machine. - Stop and relay the launcher's guidance when it exits nonzero: a missing or unreadable script path, an explicit port already held by another session, or a server that failed to start are all creator-facing problems, not things to retry silently. ## Checklist -- The port and wpm came from config when a studio config exists; defaults otherwise. +- The port, wpm, provider, and workspace came from config when a studio config exists; defaults otherwise. +- Any workspace bootstrap was consented to before anything was downloaded; an existing workspace was reused, not rebuilt. - The creator got both URLs (prompt page, remote with token) and knows the pages. +- If voice-follow was requested, the creator knows it is toggled on `/prompt`, that preflight must pass before a take, and how HOLD, BEHIND, and click-to-anchor behave. - Take and invented markers were explained if the script contains them. - Any in-place save was backed up first (the server does this; confirm the backup path in its response). -- No planned tier was presented as working. +- No planned tier or provider lane was presented as working. diff --git a/skills/mc-prompter/customize.toml b/skills/mc-prompter/customize.toml index 45e4e9d..92c6ae4 100644 --- a/skills/mc-prompter/customize.toml +++ b/skills/mc-prompter/customize.toml @@ -24,3 +24,17 @@ persistent_facts = [] # Default port for the prompter server; the studio config's # [prompter] port overrides this when set. port = 8770 + +# Voice-follow ASR provider. Working values: "nemotron-streaming" +# (default, streaming English ASR via sherpa-onnx) and "none" (classic +# prompter, no ASR). "zipformer-small" is a planned lane for low-end +# hardware; selecting it makes the server exit with a pointer, it never +# pretends to run. The studio config's [prompter] asr-provider overrides +# this when set. +asr-provider = "nemotron-streaming" + +# Engine workspace folder for voice-follow, resolved by the skill as +# {engines-path}/{workspace} when a studio config exists. Holds the ASR +# venv and model files; the first bootstrap downloads ~465 MB and is +# consent-gated by the skill. The classic prompter needs no workspace. +workspace = "prompter-lab" diff --git a/skills/mc-prompter/scripts/ensure_workspace.py b/skills/mc-prompter/scripts/ensure_workspace.py new file mode 100644 index 0000000..c558e77 --- /dev/null +++ b/skills/mc-prompter/scripts/ensure_workspace.py @@ -0,0 +1,300 @@ +#!/usr/bin/env python3 +# /// script +# requires-python = ">=3.11" +# /// +"""Create or verify the prompter-lab engine workspace for mc-prompter. + +Voice-follow (tier 2) needs streaming ASR: sherpa-onnx plus its model +files. Those live in one persistent venv and a models directory in the +creator's engine workspace (default {engines-path}/prompter-lab), never in +the skill folder: + + / + .venv/ aiohttp==3.12.15 (same pin as the server's PEP 723 + header), numpy, sherpa-onnx==1.13.4, soundfile + models/ + nemotron-streaming/ encoder.int8.onnx, decoder.int8.onnx, + joiner.int8.onnx, tokens.txt (renamed from + the release tarball's dated directory) + silero_vad.onnx VAD model + out/ session artifacts + +Downloads (both from the k2-fsa/sherpa-onnx GitHub release "asr-models"): +the nemotron streaming tarball (~464 MB, deleted after extraction) and +silero_vad.onnx (~0.7 MB). Total download ~465 MB; extracted models take +~630 MB on disk. + +DEPENDENCY PIN: sherpa-onnx is pinned to 1.13.4 because model exports must +match the runtime generation; the 2026-04-25 nemotron export is validated +under exactly this version (2026-07-09). + +Idempotent: an existing, verified workspace is used as-is; only what is +missing is built. Existing venv and model files are verified (presence and +size), never rebuilt or re-downloaded. + +The CALLING SKILL asks the creator before running this (the download is +large); the script itself just does the work. + +Usage: + uv run ensure_workspace.py --workspace + [--python 3.12] [--models-only] [--check] [--dry-run] + +--check reports readiness and changes nothing (exit 0 ready, 4 not ready). +--dry-run prints the planned commands and downloads as JSON, runs nothing. +--models-only downloads and verifies the models but skips the venv build. +Exit codes: 0 ready, 1 a build step failed, 2 usage error, 4 not ready +(--check only). +""" + +import argparse +import json +import os +import shutil +import subprocess +import sys +import tarfile +import urllib.request +from pathlib import Path + +AIOHTTP_PIN = "aiohttp==3.12.15" +SHERPA_PIN = "sherpa-onnx==1.13.4" +PIN_INSTALL = [AIOHTTP_PIN, "numpy", SHERPA_PIN, "soundfile"] +ASR_RELEASE = ("https://github.com/k2-fsa/sherpa-onnx/releases/download/" + "asr-models") +NEMOTRON_TARBALL = ("sherpa-onnx-nemotron-speech-streaming-en-0.6b-560ms-" + "int8-2026-04-25.tar.bz2") +NEMOTRON_DIR = "nemotron-streaming" +VAD_FILE = "silero_vad.onnx" +TOTAL_DOWNLOAD_MB = 465 +# Size floors (bytes) for the canonical layout; real sizes measured from +# the validated 2026-04-25 release. A truncated download fails these. +MODEL_MIN_SIZES = { + f"{NEMOTRON_DIR}/encoder.int8.onnx": 500_000_000, + f"{NEMOTRON_DIR}/decoder.int8.onnx": 5_000_000, + f"{NEMOTRON_DIR}/joiner.int8.onnx": 1_000_000, + f"{NEMOTRON_DIR}/tokens.txt": 4_000, + VAD_FILE: 500_000, +} +VERIFY_SNIPPET = ( + "import importlib.metadata as md; " + "import aiohttp, numpy, sherpa_onnx, soundfile; " + "assert md.version('aiohttp') == '3.12.15', md.version('aiohttp'); " + "assert md.version('sherpa-onnx') == '1.13.4', " + "md.version('sherpa-onnx'); " + "print('ok')" +) + + +def die(msg: str, code: int = 2) -> None: + print(msg, file=sys.stderr) + sys.exit(code) + + +def venv_python(workspace: Path, platform: str = sys.platform) -> Path: + """The venv interpreter path, resolved portably.""" + if platform == "win32": + return workspace / ".venv" / "Scripts" / "python.exe" + return workspace / ".venv" / "bin" / "python" + + +def verify_layout(workspace: Path, models_only: bool = False) -> list[str]: + """File-level problems (presence and size); empty means layout is ok.""" + problems = [] + if not models_only and not venv_python(workspace).exists(): + problems.append(f"venv missing: {venv_python(workspace)}") + for rel, min_size in MODEL_MIN_SIZES.items(): + path = workspace / "models" / rel + if not path.is_file(): + problems.append(f"model missing: models/{rel}") + elif path.stat().st_size < min_size: + problems.append( + f"model truncated: models/{rel} " + f"({path.stat().st_size} < {min_size} bytes)" + ) + return problems + + +def verify_venv_imports(workspace: Path) -> list[str]: + """Run the pinned-import check inside the venv; empty means ok.""" + py = venv_python(workspace) + if not py.exists(): + return [f"venv missing: {py}"] + r = subprocess.run([str(py), "-c", VERIFY_SNIPPET], + capture_output=True, text=True) + if r.returncode != 0: + return [f"venv import check failed: {r.stderr.strip()[-300:]}"] + return [] + + +def verify(workspace: Path, models_only: bool = False) -> list[str]: + """Full readiness check; empty means the workspace is ready.""" + problems = verify_layout(workspace, models_only=models_only) + if not models_only and venv_python(workspace).exists(): + problems += verify_venv_imports(workspace) + return problems + + +def planned_commands(workspace: Path, python_version: str) -> list[list[str]]: + py = venv_python(workspace) + return [ + ["uv", "venv", "--python", python_version, + str(workspace / ".venv")], + ["uv", "pip", "install", "--python", str(py), *PIN_INSTALL], + ] + + +def planned_downloads(workspace: Path) -> list[str]: + """URLs still needed to complete the models directory.""" + urls = [] + layout = verify_layout(workspace, models_only=True) + if any(NEMOTRON_DIR in p for p in layout): + urls.append(f"{ASR_RELEASE}/{NEMOTRON_TARBALL}") + if any(VAD_FILE in p for p in layout): + urls.append(f"{ASR_RELEASE}/{VAD_FILE}") + return urls + + +def download(url: str, dest: Path) -> None: + """Download url to dest with progress lines every ~10 percent.""" + print(f"downloading {url} -> {dest}") + tmp = dest.with_suffix(dest.suffix + ".part") + last = [-1] + + def hook(blocks: int, block_size: int, total: int) -> None: + if total <= 0: + return + done = min(blocks * block_size, total) + pct = done * 100 // total + step = pct // 10 + if step > last[0]: + last[0] = step + print(f" {done // 1_000_000}/{total // 1_000_000} MB " + f"({pct}%)", flush=True) + + urllib.request.urlretrieve(url, tmp, reporthook=hook) + # os.replace overwrites an existing (e.g. truncated) file on every + # platform; Path.rename raises FileExistsError on Windows. + os.replace(tmp, dest) + + +def extract_nemotron(tar_path: Path, models_dir: Path) -> None: + """Extract the tarball, rename its dir to the stable name, clean up. + + A stale target directory (left by an interrupted earlier bootstrap) is + removed before the rename, so repairing a truncated model layout works + instead of crashing on a rename-onto-nonempty-directory. Extraction + failures die with the manual-recovery pointer rather than a traceback. + """ + print(f"extracting {tar_path.name}") + target = models_dir / NEMOTRON_DIR + try: + with tarfile.open(tar_path, "r:bz2") as tf: + top = tf.getnames()[0].split("/")[0] + try: + tf.extractall(models_dir, filter="data") + except TypeError: # Python without the filter kwarg + tf.extractall(models_dir) + extracted = models_dir / top + if extracted != target: + if target.exists(): + shutil.rmtree(target) + extracted.rename(target) + tar_path.unlink() + except (tarfile.TarError, OSError, IndexError) as e: + die(f"error: model extraction failed ({e}); grab " + f"{NEMOTRON_TARBALL} from the sherpa-onnx asr-models release " + f"page, extract it into {models_dir}, and rename the directory " + f"to {NEMOTRON_DIR}", 1) + print(f"models ready: {target}") + + +def ensure_models(workspace: Path) -> None: + models_dir = workspace / "models" + models_dir.mkdir(parents=True, exist_ok=True) + layout = verify_layout(workspace, models_only=True) + need_nemotron = any(NEMOTRON_DIR in p for p in layout) + need_vad = any(VAD_FILE in p for p in layout) + if not (need_nemotron or need_vad): + return + print(f"total download ~{TOTAL_DOWNLOAD_MB} MB (the calling skill " + "asked for consent before running this)") + if need_nemotron: + tar_path = models_dir / NEMOTRON_TARBALL + try: + download(f"{ASR_RELEASE}/{NEMOTRON_TARBALL}", tar_path) + except OSError as e: + die(f"error: model download failed ({e}); grab " + f"{NEMOTRON_TARBALL} from the sherpa-onnx asr-models " + f"release page, extract it into {models_dir}, and rename " + f"the directory to {NEMOTRON_DIR}", 1) + extract_nemotron(tar_path, models_dir) + if need_vad: + try: + download(f"{ASR_RELEASE}/{VAD_FILE}", models_dir / VAD_FILE) + except OSError as e: + die(f"error: VAD download failed ({e}); grab {VAD_FILE} from " + f"the sherpa-onnx asr-models release page into " + f"{models_dir}", 1) + + +def ensure_venv(workspace: Path, python_version: str) -> None: + commands = planned_commands(workspace, python_version) + if not venv_python(workspace).exists(): + for cmd in commands: + r = subprocess.run(cmd) + if r.returncode != 0: + die(f"error: {' '.join(cmd)} failed", 1) + else: + # Existing venv: install is idempotent and fixes a broken dep set. + r = subprocess.run(commands[1]) + if r.returncode != 0: + die("error: dependency install into existing venv failed", 1) + + +def main(argv=None) -> int: + ap = argparse.ArgumentParser( + description=__doc__, + formatter_class=argparse.RawDescriptionHelpFormatter) + ap.add_argument("--workspace", required=True, type=Path) + ap.add_argument("--python", default="3.12") + ap.add_argument("--models-only", action="store_true") + ap.add_argument("--check", action="store_true") + ap.add_argument("--dry-run", action="store_true") + args = ap.parse_args(argv) + ws = args.workspace + + if args.check: + problems = verify(ws, models_only=args.models_only) + if problems: + for p in problems: + print(p) + return 4 + print(f"ready: {ws}") + return 0 + + if args.dry_run: + print(json.dumps({ + "workspace": str(ws), + "commands": [] if args.models_only + else planned_commands(ws, args.python), + "downloads": planned_downloads(ws), + "total-download-mb": TOTAL_DOWNLOAD_MB, + }, indent=2)) + return 0 + + for sub in ("models", "out"): + (ws / sub).mkdir(parents=True, exist_ok=True) + + if not args.models_only: + ensure_venv(ws, args.python) + ensure_models(ws) + + problems = verify(ws, models_only=args.models_only) + if problems: + die("error: workspace still not ready:\n" + "\n".join(problems), 1) + print(f"ready: {ws}") + return 0 + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/skills/mc-prompter/scripts/run_prompter.py b/skills/mc-prompter/scripts/run_prompter.py index 424ca2c..01dd047 100644 --- a/skills/mc-prompter/scripts/run_prompter.py +++ b/skills/mc-prompter/scripts/run_prompter.py @@ -2,7 +2,7 @@ # /// script # requires-python = ">=3.11" # /// -"""Launch the mc-prompter teleprompter server (Phase A). +"""Launch the mc-prompter teleprompter server. Pure stdlib launcher. Probes the port, generates the session token, spawns the aiohttp server as a child process, confirms it is up, writes the session @@ -11,6 +11,7 @@ Usage: uv run {skill-root}/scripts/run_prompter.py --script [--port 8770] [--lan] [--owner-wpm 150] [--no-open] + [--workspace ] [--asr-provider nemotron-streaming|none] Behavior: port default 8770. If busy, /health on 127.0.0.1 is queried (1 s @@ -21,7 +22,18 @@ spawn `uv run --with aiohttp==3.12.15 python -m server.main ` with cwd set to this scripts directory so the server package resolves. Server args are all explicit: --port --host --script - --owner-wpm --session-file --token. + --owner-wpm --session-file --token --asr-provider, plus + --models-dir on the workspace path. + workspace when --workspace points at a ready prompter-lab (built by + ensure_workspace.py; venv present, model files present at + sane sizes; checked inline, nothing is shelled out), the + server is spawned with the workspace venv interpreter + (.venv/bin/python on POSIX, .venv\\Scripts\\python.exe on + Windows) from the same cwd, and --models-dir plus + --asr-provider are passed so voice-follow is available. + Without --workspace, or when it is not ready, the launcher + falls back to the uv-run spawn with --asr-provider none + (tier 1, classic prompter) and says so. session token = secrets.token_urlsafe(16). After /health confirms the server is up, /mc-prompter/session-.json is written: {"port": N, "pid": ..., "token": "...", @@ -32,8 +44,9 @@ it the server binds 127.0.0.1 and only localhost URLs print. --no-open skip opening the browser at the home page. -Exit codes: 0 ok, 1 server failed to start, 2 usage, 4 script path missing -or unreadable, 5 port conflict with an explicit --port. +Exit codes: 0 ok, 1 server failed to start, 2 usage, 3 planned asr provider +selected (zipformer-small; fails fast, nothing is spawned), 4 script path +missing or unreadable, 5 port conflict with an explicit --port. """ import argparse @@ -62,6 +75,21 @@ # network. wait_for_health fails immediately if the child exits. STARTUP_TIMEOUT = 300.0 PORT_SCAN_LIMIT = 20 +DEFAULT_ASR_PROVIDER = "nemotron-streaming" +ASR_PROVIDERS = ("nemotron-streaming", "zipformer-small", "none") +# Planned lanes stay valid argparse choices so the launcher (not argparse) +# owns the message, but selecting one fails fast with exit 3 before any +# spawn: nothing unvalidated pretends to run. +ASR_PLANNED_PROVIDERS = ("zipformer-small",) +# Lightweight readiness floors, mirroring ensure_workspace.py's layout +# check (presence and size only; no subprocess, no imports). +WORKSPACE_MODEL_MIN_SIZES = { + "models/nemotron-streaming/encoder.int8.onnx": 500_000_000, + "models/nemotron-streaming/decoder.int8.onnx": 5_000_000, + "models/nemotron-streaming/joiner.int8.onnx": 1_000_000, + "models/nemotron-streaming/tokens.txt": 4_000, + "models/silero_vad.onnx": 500_000, +} class PortConflict(Exception): @@ -150,15 +178,97 @@ def write_session_file(path, port, pid, token, script): return path -def build_server_cmd(port, host, script, owner_wpm, session_file, token): - """The child process command; runs with cwd = the scripts directory.""" - # Phase B seam: when the prompter-lab workspace exists, swap the - # `uv run --with` prefix for `/.venv python -m server.main` - # (POSIX .venv/bin/python, Windows .venv\Scripts\python.exe) with the - # same cwd and the same explicit args. TODO(phase-b): implement the - # workspace-present branch here; do not add config discovery. - return [ - "uv", "run", "--with", AIOHTTP_PIN, "python", "-m", "server.main", +def venv_python_path(workspace, platform=sys.platform): + """The workspace venv interpreter path, resolved portably.""" + workspace = Path(workspace) + if platform == "win32": + return workspace / ".venv" / "Scripts" / "python.exe" + return workspace / ".venv" / "bin" / "python" + + +def workspace_ready(workspace): + """True when the prompter-lab layout looks complete. + + Inline re-implementation of ensure_workspace.py --check semantics at + the file level (presence and size floors); deliberately no subprocess + and no venv import check, so launch stays instant. + """ + workspace = Path(workspace) + if not venv_python_path(workspace).exists(): + return False + for rel, min_size in WORKSPACE_MODEL_MIN_SIZES.items(): + path = workspace / rel + if not path.is_file() or path.stat().st_size < min_size: + return False + return True + + +def spawn_plan(workspace, asr_provider): + """Resolve the launch mode: (workspace-or-None, effective provider). + + A ready workspace launches through its venv with the requested + provider (default nemotron-streaming). No workspace, or a not-ready + one, falls back to the tier-1 uv-run spawn with provider "none". + """ + if workspace is not None and workspace_ready(workspace): + return workspace, (asr_provider or DEFAULT_ASR_PROVIDER) + return None, "none" + + +def downgrade_notice(requested_ws, requested_provider, workspace): + """Lines explaining a tier-1 fallback; [] when nothing was downgraded. + + An explicitly requested non-none provider that spawn_plan dropped to + "none" is always named, so the creator is never silently overridden. + """ + if workspace is not None: + return [] + explicit = requested_provider not in (None, "none") + lines = [] + if requested_ws is not None: + lines.append(f"workspace not ready: {requested_ws}") + if explicit: + lines.append( + f" requested --asr-provider {requested_provider} is " + "disabled (downgraded to none)" + ) + lines.append( + " launching the classic prompter (voice-follow off). " + "bootstrap the workspace with:" + ) + lines.append( + f" uv run {SCRIPTS_DIR / 'ensure_workspace.py'} " + f"--workspace {requested_ws}" + ) + elif explicit: + lines.append( + f"note: --asr-provider {requested_provider} needs --workspace; " + "launching the classic prompter (voice-follow off, provider " + "downgraded to none)" + ) + return lines + + +def voice_follow_line(workspace, asr_provider): + """The launch banner's voice-follow line, honest about provider none.""" + if workspace is not None and asr_provider != "none": + return ( + f" voice-follow: available ({asr_provider}, " + f"workspace {workspace})" + ) + return " voice-follow: off (classic prompter)" + + +def build_server_cmd(port, host, script, owner_wpm, session_file, token, + workspace=None, asr_provider="none"): + """The child process command; runs with cwd = the scripts directory. + + With a workspace, the command is the workspace venv interpreter running + `-m server.main` (same cwd, so the server package resolves identically + on both paths) plus --models-dir and --asr-provider. Without one, it is + the uv-run spawn with --asr-provider none. + """ + server_args = [ "--port", str(port), "--host", host, "--script", str(script) if script else "", @@ -166,6 +276,19 @@ def build_server_cmd(port, host, script, owner_wpm, session_file, token): "--session-file", str(session_file), "--token", token, ] + if workspace is not None: + workspace = Path(workspace) + return [ + str(venv_python_path(workspace)), "-m", "server.main", + *server_args, + "--models-dir", str(workspace / "models"), + "--asr-provider", asr_provider, + ] + return [ + "uv", "run", "--with", AIOHTTP_PIN, "python", "-m", "server.main", + *server_args, + "--asr-provider", asr_provider, + ] def spawn_server(cmd): @@ -251,8 +374,26 @@ def main(argv=None): help="creator wpm for read-time estimates, 0 = unset") parser.add_argument("--no-open", action="store_true", help="do not open the browser at the home page") + parser.add_argument("--workspace", default=None, + help="prompter-lab workspace path (built by " + "ensure_workspace.py); enables voice-follow " + "when ready") + parser.add_argument("--asr-provider", default=None, + choices=ASR_PROVIDERS, + help="ASR provider for voice-follow (default " + f"{DEFAULT_ASR_PROVIDER} when --workspace is " + "given; planned lanes fail fast with exit 3)") args = parser.parse_args(argv) + if args.asr_provider in ASR_PLANNED_PROVIDERS: + print( + f"error: asr-provider {args.asr_provider!r} is a planned lane " + "and is not implemented yet; use " + f"{DEFAULT_ASR_PROVIDER} or none", + file=sys.stderr, + ) + return 3 + script = Path(args.script).expanduser().resolve() if not script.is_file(): print(f"error: script not found: {script}", file=sys.stderr) @@ -284,11 +425,18 @@ def main(argv=None): print(f"error: {exc}", file=sys.stderr) return 1 + requested_ws = (Path(args.workspace).expanduser().resolve() + if args.workspace else None) + workspace, asr_provider = spawn_plan(requested_ws, args.asr_provider) + for line in downgrade_notice(requested_ws, args.asr_provider, workspace): + print(line) + host = "0.0.0.0" if args.lan else "127.0.0.1" token = secrets.token_urlsafe(16) session_file = session_file_path(port) cmd = build_server_cmd(port, host, script, args.owner_wpm, - session_file, token) + session_file, token, + workspace=workspace, asr_provider=asr_provider) child = spawn_server(cmd) try: @@ -303,6 +451,7 @@ def main(argv=None): base_url = f"http://127.0.0.1:{port}" local_url = f"{base_url}/?token={token}" print(f"mc-prompter is up (session {token[:8]})") + print(voice_follow_line(workspace, asr_provider)) print(f" home: {local_url}") print(f" prompt: {base_url}/prompt") print(f" remote: http://127.0.0.1:{port}/remote?token={token}") diff --git a/skills/mc-prompter/scripts/server/align.py b/skills/mc-prompter/scripts/server/align.py new file mode 100644 index 0000000..2db1293 --- /dev/null +++ b/skills/mc-prompter/scripts/server/align.py @@ -0,0 +1,517 @@ +#!/usr/bin/env python3 +# /// script +# requires-python = ">=3.11" +# /// +"""Voice-follow alignment engine for mc-prompter (Phase B). Pure stdlib. + +Consumes the streaming-transducer partial hypothesis stream (BPE token lists +per segment, as emitted by server/asr.py) and maintains a monotonic committed +anchor into the script's speakable words (the same global word indexing the +UI renders as data-i spans; see script_ingest.speakable_words()). + + aligner = Aligner(speakable_words(doc), + take_ranges=take_word_ranges(doc)) + result = aligner.feed(tokens, segment, final) + # result == {"anchor": int, "moved": bool, "held": bool} + aligner.set_anchor(word_index) # human override, backwards allowed + aligner.anchor # committed global word index, -1 initially + +take_ranges is the list of [start, end) global word-index pairs for the +script's take blocks (script_ingest.take_word_ranges, same indexing as +speakable_words). Take paragraphs are already-recorded footage: they are +dimmed or hidden in the UI and a presenter normally skips them aloud, so the +aligner treats their words as free to skip (see Matching below). + +Input contract (matches the observed nemotron-streaming behavior recorded in +the Phase B contract and the committed fixtures): + + feed() receives the FULL hypothesis token list for the current segment on + every partial. Hypotheses grow incrementally but the tail revises at the + word level (BPE pieces append into the previous word, e.g. "workshop" + becoming "workshopod"), so the last k_provisional merged words are treated + as provisional and never advance the committed anchor. On final=True the + whole hypothesis commits. A new segment id resets hypothesis tracking + (endpoint reset cleared the recognizer state) but never rewinds the + committed anchor. Timestamps are not used. + +Matching: + + BPE pieces merge to words (a piece starting with a space, or the first + piece, starts a new word). Both script words and hypothesis words are + normalized: casefold, punctuation stripped, hyphen/slash split, and a + small documented number/ordinal table expands digit forms to spoken words + on BOTH sides, so script "2026" matches spoken "twenty twenty six" and a + recognizer that emits digits still matches a written-out script. + + Newly committed hypothesis words (pending) are matched against a script + window starting just past the committed anchor, window = window_mult * + len(pending) + window_base expanded NON-TAKE tokens; take-block tokens + inside that span ride along for free, so a take paragraph of any length + never consumes window budget and the words after it stay reachable. A + word-level Levenshtein DP picks the window prefix end with minimum + distance. Per-word substitution cost: exact 0, prefix or difflib-fuzzy + 0.5, else 1. Insertion of an unmatched spoken word costs 1; deletion of + a script window word costs 0.25, except take-block words which delete + at cost 0 (the presenter is expected to skip them; if they DO read the + take aloud, substitution matching inside the take still works and the + anchor tracks through it). Deletions must be cheap: at cost 1, aligning + the pending words after N skipped script words (cost N) always loses to + matching nothing (cost len(pending)), so skips and recognizer word + drops would stall the anchor forever. At 0.25 a deliberate skip is + followed within the window while spurious far matches stay unprofitable + (a fuzzy match only pays for itself within about two words of + deletion). The anchor advances to the window position of the last + pending word that aligned at cost <= 0.5. The committed anchor word is + never a take-block word: if the last match lands inside a take (the + take was read aloud, or a word fuzzy-matched into it), the broadcast + anchor advances to the next non-take word so the UI always has a + visible span to scroll to (falling back to the previous non-take word + when the take runs to the end of the script). + + Hold rule (the gate): off-script speech must not move the anchor. A + match only counts toward the gate if the spoken token is informative: + length >= 4, or not in the small stopword set below. Conversational + ad-libs are stopword-dense and scripts are full of forward stopword + duplicates, so stopword coincidences ("and", "the", "we") must not be + evidence of progress. When the pending batch contains informative + tokens, moving requires (a) at least one informative token matched + EXACTLY (cost 0), and (b) informative-matches / (informative-pending + + interior-deletions) >= hold_threshold, where interior-deletions counts + unmatched non-take window words strictly between the first and last + matched position. Charging interior deletions makes scattered + coincidence matches fail (they straddle unmatched script) while a + genuine skip still passes (its matches are contiguous, the skipped + words sit BEFORE the first match). A batch of only stopwords (common in + small partial commits: "and the") falls back to the same ratio over all + pending words, so verbatim reading still advances promptly. When the + gate fails, the anchor does not move and held=True is returned. Held + pending words are consumed, so the window never fills with off-script + speech and the next on-script words re-match from the committed anchor. + + End of script: when the window is truncated by the document end and the + gate fails, the batch may be the script tail plus overflow speech + committed together in one final ("...thanks for watching" continuing + past the last script word). Overflow past the end cannot match + anything, so the gate is retried over only the pending words up to the + last matched one; if that passes, the tail anchors instead of holding + forever. + + feed() clamps each batch to the last 24 RAW words before normalization. + Pending is normally 1-5 words, but a mid-segment Aligner rebuild + (script edit while a segment is in flight) resets _consumed and the + next partial commits the whole hypothesis so far as one batch; _match + is O(len(pending) * window) with window proportional to pending, i.e. + quadratic, and a 300-word batch would block the asyncio event loop for + seconds. Older words are stale context anyway: the window matcher only + needs recent speech. The clamp is on raw words (before number expansion + multiplies tokens) so the DP input is bounded no matter what the words + expand to. + + set_anchor() jumps anywhere, backwards included (the human is the + authority), and consumes the current hypothesis so already-spoken words + are not re-matched against the new position. + +No I/O, no threads, no imports beyond stdlib. server/main.py owns the wiring +(ASR events in, anchor broadcasts out). +""" + +import difflib + +_ONES = ["zero", "one", "two", "three", "four", "five", "six", "seven", + "eight", "nine", "ten", "eleven", "twelve", "thirteen", "fourteen", + "fifteen", "sixteen", "seventeen", "eighteen", "nineteen"] +_TENS = {2: "twenty", 3: "thirty", 4: "forty", 5: "fifty", 6: "sixty", + 7: "seventy", 8: "eighty", 9: "ninety"} +# Small ordinal table: the common single-word ordinals. Larger ordinals fall +# back to cardinal expansion of the digits ("42nd" -> "forty two"), close +# enough for the fuzzy per-word cost. +_ORDINALS = {"1st": "first", "2nd": "second", "3rd": "third", "4th": "fourth", + "5th": "fifth", "6th": "sixth", "7th": "seventh", + "8th": "eighth", "9th": "ninth", "10th": "tenth", + "11th": "eleventh", "12th": "twelfth"} + +_FUZZY_RATIO = 0.75 +_FUZZY_COST = 0.5 +_DELETE_COST = 0.25 # script window word skipped; see the module docstring +_TAKE_DELETE_COST = 0.0 # take-block words are free to skip + +# Apostrophe-family characters stripped inside words so "don't", "don’t" and +# "donʼt" all normalize to "dont". U+02BC (modifier letter apostrophe) is +# Unicode category Lm and would otherwise survive isalnum() as part of the +# token; the curly quotes would split the word instead. +_APOSTROPHES = "'’ʼ‘‛" + +# Uninformative tokens for the hold gate: a match on one of these (or any +# token shorter than 4 characters that IS in this set) is a coincidence, not +# evidence the presenter is on script. Deliberately tiny and limited to +# common function words of length <= 3 (length >= 4 already counts as +# informative regardless of this set). +_STOPWORDS = frozenset( + "a an and are as at be but by do for he i in is it me my no not of on " + "or so the to us was we you".split()) + +# feed() batch clamp, in raw (pre-normalization) words; see the module +# docstring ("feed() clamps each batch..."). +_PENDING_CAP = 24 + + +def _informative(token): + """True when a match on this token counts toward the hold gate.""" + return len(token) >= 4 or token not in _STOPWORDS + + +def merge_bpe(tokens): + """Merge BPE pieces into words. + + A piece starting with a space, or the first piece, starts a new word; + every other piece (including bare punctuation pieces like "," and ".") + appends to the current word. Empty results are dropped. + """ + words = [] + for piece in tokens: + if piece.startswith(" ") or not words: + words.append(piece.strip()) + else: + words[-1] += piece.strip() + return [w for w in words if w] + + +def _expand_number(digits): + """Expand a digit string to spoken words (the documented small table). + + 0-19 and tens from the ones/tens tables; 100-999 as "N hundred [rest]"; + 1000-9999 read as digit pairs the way years are spoken ("2026" -> + "twenty twenty six", "1995" -> "nineteen ninety five"), with the special + cases "2000" -> "two thousand", "1900" -> "nineteen hundred", "2007" -> + "two thousand seven", "1907" -> "nineteen oh seven"; 10000-999999 as + "N thousand [rest]"; anything larger digit by digit. + """ + n = int(digits) + if n < 20: + return [_ONES[n]] + if n < 100: + tens, ones = divmod(n, 10) + return [_TENS[tens]] + ([_ONES[ones]] if ones else []) + if n < 1000: + hundreds, rest = divmod(n, 100) + return [_ONES[hundreds], "hundred"] + (_expand_number(str(rest)) if rest else []) + if n < 10000: + hi, lo = divmod(n, 100) + if lo == 0 and hi % 10 == 0: + return [_ONES[hi // 10], "thousand"] + if lo == 0: + return _expand_number(str(hi)) + ["hundred"] + if hi % 10 == 0 and lo < 10: + return [_ONES[hi // 10], "thousand"] + _expand_number(str(lo)) + if lo < 10: + return _expand_number(str(hi)) + ["oh"] + _expand_number(str(lo)) + return _expand_number(str(hi)) + _expand_number(str(lo)) + if n < 1_000_000: + thousands, rest = divmod(n, 1000) + return (_expand_number(str(thousands)) + ["thousand"] + + (_expand_number(str(rest)) if rest else [])) + return [_ONES[int(d)] for d in digits] + + +def normalize_word(word): + """Normalize one raw word to a list of match tokens (possibly empty). + + Casefold; "%" becomes the word "percent"; the apostrophe family + (_APOSTROPHES: ASCII, curly, and U+02BC modifier letter) is removed so + "don't", "don’t", "donʼt" and "dont" all match; commas inside digit + runs are removed; every + other non-alphanumeric character splits the word (hyphens, slashes, + trailing punctuation). Digit pieces expand through the number table; + digit+ordinal-suffix pieces go through the ordinal table; mixed pieces + split at digit/letter boundaries. + """ + w = word.casefold().replace("%", " percent ") + for ch in _APOSTROPHES: + w = w.replace(ch, "") + # Drop commas used as thousands separators before they split the number. + w = "".join(ch for i, ch in enumerate(w) + if not (ch == "," and 0 < i < len(w) - 1 + and w[i - 1].isdigit() and w[i + 1].isdigit())) + pieces = [] + buf = "" + for ch in w: + if ch.isalnum(): + buf += ch + elif buf: + pieces.append(buf) + buf = "" + if buf: + pieces.append(buf) + + out = [] + for piece in pieces: + if piece in _ORDINALS: + out.append(_ORDINALS[piece]) + continue + if piece.isdigit(): + out.extend(_expand_number(piece)) + continue + if piece.isalpha(): + out.append(piece) + continue + # Mixed digits and letters: split at boundaries ("90s", "42nd"). + if piece[:-2].isdigit() and piece[-2:] in ("st", "nd", "rd", "th"): + out.extend(_expand_number(piece[:-2])) + continue + run = "" + for ch in piece: + if run and ch.isdigit() != run[-1].isdigit(): + out.extend(_expand_number(run) if run.isdigit() else [run]) + run = "" + run += ch + if run: + out.extend(_expand_number(run) if run.isdigit() else [run]) + return out + + +def _word_cost(a, b): + """Per-word similarity cost: exact 0, prefix/fuzzy 0.5, else 1.""" + if a == b: + return 0.0 + if len(a) >= 3 and len(b) >= 3 and (a.startswith(b) or b.startswith(a)): + return _FUZZY_COST + if difflib.SequenceMatcher(None, a, b).ratio() >= _FUZZY_RATIO: + return _FUZZY_COST + return 1.0 + + +def _match(pending, window, take_flags=None): + """Word-level Levenshtein of pending against the best window prefix. + + take_flags marks window positions that belong to take blocks; those + delete at _TAKE_DELETE_COST (0) so a skipped take costs nothing to + cross. Returns (matches, best_distance) where matches is a list of + (i, j, c) triples: pending[i] aligned to window[j] at cost c <= 0.5 on + the minimum-cost path to the best prefix end (ties break to the + shortest prefix; a substitution tie beats a free take deletion, so a + take read aloud still registers its matches). + """ + m, n = len(pending), len(window) + if take_flags is None: + take_flags = [False] * n + dele = [_TAKE_DELETE_COST if t else _DELETE_COST for t in take_flags] + dist = [[0.0] * (n + 1) for _ in range(m + 1)] + for i in range(1, m + 1): + dist[i][0] = float(i) + for j in range(1, n + 1): + dist[0][j] = dist[0][j - 1] + dele[j - 1] + cost = [[0.0] * n for _ in range(m)] + for i in range(1, m + 1): + for j in range(1, n + 1): + c = _word_cost(pending[i - 1], window[j - 1]) + cost[i - 1][j - 1] = c + dist[i][j] = min(dist[i - 1][j - 1] + c, + dist[i - 1][j] + 1.0, + dist[i][j - 1] + dele[j - 1]) + best_j = min(range(n + 1), key=lambda j: (dist[m][j], j)) + matches = [] + i, j = m, best_j + while i > 0 and j > 0: + c = cost[i - 1][j - 1] + if dist[i][j] == dist[i - 1][j - 1] + c: + if c <= _FUZZY_COST: + matches.append((i - 1, j - 1, c)) + i, j = i - 1, j - 1 + elif dist[i][j] == dist[i][j - 1] + dele[j - 1]: + j -= 1 + else: + i -= 1 + matches.reverse() + return matches, dist[m][best_j] + + +class Aligner: + """Monotonic committed-anchor aligner over the script's speakable words.""" + + def __init__(self, words, k_provisional=4, window_base=10, window_mult=2, + hold_threshold=0.5, take_ranges=None): + self.k_provisional = k_provisional + self.window_base = window_base + self.window_mult = window_mult + self.hold_threshold = hold_threshold + # Take-block words (script_ingest.take_word_ranges, [start, end) + # pairs in the same global word indexing as `words`). + self._take_words = set() + for a, b in (take_ranges or []): + self._take_words.update(range(max(0, a), min(b, len(words)))) + # Expanded script: flat normalized tokens, each carrying its source + # global word index (number expansion makes this one-to-many). + self._exp = [] + self._word_last_exp = [] + last = -1 + for i, word in enumerate(words): + for tok in normalize_word(word): + self._exp.append((tok, i)) + last = len(self._exp) - 1 + self._word_last_exp.append(last) + self._exp_take = [i in self._take_words for _, i in self._exp] + self._n_words = len(words) + self._exp_anchor = -1 + self._anchor_word = -1 + self._segment = None + self._consumed = 0 + self._last_len = 0 + + @property + def anchor(self): + """Committed global word index, -1 before any match.""" + return self._anchor_word + + def set_anchor(self, word_index): + """Jump the anchor anywhere (backwards allowed) and clear tracking. + + The current hypothesis is consumed so already-spoken words are not + re-matched against the new position; the next feed matches only + newly committed words. + """ + word_index = max(-1, min(int(word_index), self._n_words - 1)) + if word_index < 0: + self._exp_anchor = -1 + self._anchor_word = -1 + else: + self._exp_anchor = self._word_last_exp[word_index] + self._anchor_word = word_index + # Consume everything already heard in the current segment so the + # next feed matches only words spoken after the jump. + self._consumed = self._last_len + + def feed(self, tokens, segment, final): + """Consume one partial or final hypothesis for a segment. + + tokens is the full BPE token list for the segment's current + hypothesis. Returns {"anchor": int, "moved": bool, "held": bool}. + """ + words = merge_bpe(tokens) + if segment != self._segment: + self._segment = segment + self._consumed = 0 + self._last_len = len(words) + if final: + commit = len(words) + else: + commit = max(self._consumed, len(words) - self.k_provisional) + raw_pending = words[self._consumed:commit] + self._consumed = max(self._consumed, commit) + # Clamp to the last _PENDING_CAP raw words: a mid-segment Aligner + # rebuild (script edit) resets _consumed and would otherwise commit + # the whole hypothesis as one batch, and _match is quadratic in the + # batch size (a 300-word batch blocks the event loop for seconds). + # Raw words, not normalized tokens, so number expansion cannot + # reinflate the bound; the dropped words are stale context the + # window matcher does not need. + raw_pending = raw_pending[-_PENDING_CAP:] + + pending = [] + for w in raw_pending: + pending.extend(normalize_word(w)) + + moved = False + held = False + if pending: + start = self._exp_anchor + 1 + width = self.window_mult * len(pending) + self.window_base + end, truncated = self._window_end(start, width) + window_toks = [tok for tok, _ in self._exp[start:end]] + take_flags = self._exp_take[start:end] + matches, _ = _match(pending, window_toks, take_flags) + ok = self._gate_passes(pending, matches, take_flags) + if not ok and truncated and matches: + # Window truncated by document end: the batch may be the + # script tail plus overflow speech. Overflow past the end + # cannot match anything, so retry the gate over only the + # pending words up to the last matched one. + last_i = max(i for i, _, _ in matches) + ok = self._gate_passes(pending[:last_i + 1], matches, + take_flags) + if not ok: + held = True + elif matches: + new_exp = start + max(j for _, j, _ in matches) + if new_exp > self._exp_anchor: + self._exp_anchor = new_exp + word = self._exp[new_exp][1] + # Never broadcast a take-block word as the anchor: the + # UI dims or hides takes, so snap to the next visible + # word (max() keeps the anchor monotonic when the + # fallback resolves backwards at end of script). + self._anchor_word = max(self._anchor_word, + self._visible_word(word)) + moved = True + return {"anchor": self._anchor_word, "moved": moved, "held": held} + + def _window_end(self, start, width): + """Window end so [start:end) holds `width` non-take tokens. + + Take-block tokens ride along for free (they delete at cost 0 in + _match), so a take span of any length never consumes window budget + and the script after it stays reachable. Returns (end, truncated) + where truncated means the document ended before the budget filled. + """ + n = len(self._exp) + end = start + budget = width + while end < n and budget > 0: + if not self._exp_take[end]: + budget -= 1 + end += 1 + return end, budget > 0 + + def _gate_passes(self, pending, matches, take_flags): + """The hold gate: is this batch evidence of on-script progress? + + Only informative matches count (see _informative); when the batch + has informative tokens, at least one must match exactly. Unmatched + non-take window words strictly between the first and last matched + position (interior deletions) are charged against the ratio, so + scattered stopword coincidences fail while a genuine contiguous + skip (deletions before the first match) passes. See the module + docstring, Hold rule. + """ + if not matches: + return False + inf_idx = {k for k, tok in enumerate(pending) if _informative(tok)} + if inf_idx: + relevant = [m for m in matches if m[0] in inf_idx] + # Require an exact informative match, UNLESS every informative + # pending token matched: a batch whose content words all match, + # merely fuzzily, is an ASR garble of an on-script read (the + # recorded streams produce e.g. "workshopod" for "workshop"), + # not an ad-lib; off-script speech carries unmatched content + # words and fails here or on the ratio below. + if (not any(c == 0.0 for _, _, c in relevant) + and {m[0] for m in relevant} != inf_idx): + return False + num = len(relevant) + denom = len(inf_idx) + else: + num = len(matches) + denom = len(pending) + matched_j = {j for _, j, _ in matches} + interior = sum(1 for j in range(min(matched_j) + 1, max(matched_j)) + if j not in matched_j and not take_flags[j]) + return num / (denom + interior) >= self.hold_threshold + + def _visible_word(self, word): + """Snap a take-block word to the nearest visible (non-take) word. + + Prefers the next non-take word (the one the presenter reads next, + and the one the UI can scroll to); falls back to the previous one + when the take runs to the end of the script. Returns -1 only when + the whole script is take blocks. + """ + if word not in self._take_words: + return word + w = word + 1 + while w < self._n_words and w in self._take_words: + w += 1 + if w < self._n_words: + return w + w = word - 1 + while w >= 0 and w in self._take_words: + w -= 1 + return w diff --git a/skills/mc-prompter/scripts/server/asr.py b/skills/mc-prompter/scripts/server/asr.py new file mode 100644 index 0000000..91f6807 --- /dev/null +++ b/skills/mc-prompter/scripts/server/asr.py @@ -0,0 +1,313 @@ +#!/usr/bin/env python3 +# /// script +# requires-python = ">=3.11" +# dependencies = ["numpy", "sherpa-onnx==1.13.4"] +# /// +"""Streaming ASR engine for mc-prompter (Phase B, voice-follow). + +This module is imported LAZILY by server/main.py, only when the server is +launched with --models-dir and an implemented --asr-provider. It normally +runs inside the prompter-lab workspace venv (which pins sherpa-onnx==1.13.4, +numpy); the PEP 723 header above documents the dependencies and allows a +direct `uv run asr.py` sanity import. CI never imports this module. + +Contract (binding, see the Phase B contract): + AsrEngine(models_dir, on_event, provider="nemotron-streaming", + num_threads=2) + models_dir the workspace models directory containing + nemotron-streaming/{encoder.int8.onnx, decoder.int8.onnx, + joiner.int8.onnx, tokens.txt} and silero_vad.onnx + on_event callable taking one dict; called from the WORKER thread. + The caller wraps it with loop.call_soon_threadsafe. + provider "nemotron-streaming" is the only implemented provider; + anything else raises ValueError (planned lanes never + pretend to run). + engine.start() spawns the worker thread; model load happens on the + worker so start() returns immediately. ready flips true + (status event) once the recognizer and VAD are built. + Restart after stop() is supported: start() begins from + fresh state (queue drained, ready/behind/dropped reset) + so a leftover stop sentinel or stale pre-stop audio can + never reach a new worker. + engine.stop() stops and joins the worker. + engine.alive True while the worker thread is running. + engine.feed(pcm16_bytes) called from the event loop with raw little + endian PCM16 mono 16 kHz bytes of any length (the + browser sends ~3840 bytes per 120 ms). Bounded queue + (maxsize 50); on overflow the OLDEST frame is dropped + and behind=True until the queue drains below half. + engine.stats {"queue": n, "behind": bool, "ready": bool} + +Events emitted through on_event: + {"kind": "partial", "segment": n, "text": str, "tokens": [...]} + on every hypothesis change. tokens are BPE pieces; a piece starting + with a space starts a new word. The tail of the hypothesis may + revise between partials; the aligner handles that. + {"kind": "final", "segment": n, "text": str, "tokens": [...]} + at an endpoint, before the recognizer resets (only when the + hypothesis is non-empty; endpoints on pure silence are not final + events). The next partial starts segment n+1 with a fresh + hypothesis. + {"kind": "vad", "speaking": bool} + on speaking/silence transitions (silero VAD, 512-sample windows). + {"kind": "status", "ready": bool, "behind": bool, "queue": int} + on ready/behind changes. If model load fails (imports included) or + the decode loop crashes mid-session, an error line goes to stderr, + ready flips False, the queue is drained (stats reflect death), a + final status event is emitted, and the worker exits; the engine + never fakes readiness. + +Pipeline per frame (all on the worker thread): PCM16 LE bytes -> float32 in +[-1, 1); 512-sample windows into the VAD (drained each window so its buffer +never grows); the full chunk into the recognizer stream; decode while ready; +endpoint -> final + reset + segment increment. +""" + +import json +import queue +import sys +import threading +from pathlib import Path + +PROVIDERS = ("nemotron-streaming",) +SAMPLE_RATE = 16000 +VAD_WINDOW = 512 +QUEUE_MAX = 50 + + +class AsrEngine: + """Dedicated-thread streaming recognizer with a bounded drop-oldest queue.""" + + def __init__(self, models_dir, on_event, provider="nemotron-streaming", + num_threads=2): + if provider not in PROVIDERS: + raise ValueError( + f"asr provider {provider!r} is not implemented; " + f"implemented: {', '.join(PROVIDERS)}" + ) + self.models_dir = Path(models_dir) + self.on_event = on_event + self.provider = provider + self.num_threads = num_threads + self._queue = queue.Queue(maxsize=QUEUE_MAX) + self._behind = False + self._ready = False + self._dropped = 0 + self._stop = threading.Event() + self._thread = None + + def start(self): + """Spawn the worker thread (idempotent). Model load happens there. + + Every start begins from fresh state: a prior stop() leaves its None + sentinel (and possibly stale pre-stop audio) in the queue, which a + restarted worker must never inherit or it would decode old audio and + then die on the sentinel while reporting ready. + """ + if self._thread is not None: + return + self._drain_queue() + self._ready = False + self._behind = False + self._dropped = 0 + self._stop.clear() + self._thread = threading.Thread( + target=self._run, name="mc-prompter-asr", daemon=True + ) + self._thread.start() + + def stop(self): + """Signal the worker and join it (safe to call more than once).""" + self._stop.set() + try: + self._queue.put_nowait(None) + except queue.Full: + pass + if self._thread is not None: + self._thread.join(timeout=10) + self._thread = None + + def feed(self, pcm16_bytes): + """Enqueue a PCM16 frame; drop the oldest frame on overflow.""" + try: + self._queue.put_nowait(pcm16_bytes) + except queue.Full: + try: + self._queue.get_nowait() + except queue.Empty: + pass + try: + self._queue.put_nowait(pcm16_bytes) + except queue.Full: + pass + self._dropped += 1 + if not self._behind: + self._behind = True + self._emit_status() + + @property + def stats(self): + return { + "queue": self._queue.qsize(), + "behind": self._behind, + "ready": self._ready, + } + + @property + def alive(self): + """True while the worker thread is running.""" + return self._thread is not None and self._thread.is_alive() + + def _drain_queue(self): + """Empty the audio queue (start reset and worker-death cleanup).""" + while True: + try: + self._queue.get_nowait() + except queue.Empty: + return + + def _emit(self, event): + try: + self.on_event(event) + except Exception as exc: # a sink bug must never kill the worker + print(f"asr: on_event raised: {exc}", file=sys.stderr) + + def _emit_status(self): + self._emit( + { + "kind": "status", + "ready": self._ready, + "behind": self._behind, + "queue": self._queue.qsize(), + } + ) + + def _build(self, sherpa_onnx): + """Construct the recognizer and VAD (worker thread only).""" + model_dir = self.models_dir / self.provider + recognizer = sherpa_onnx.OnlineRecognizer.from_transducer( + tokens=str(model_dir / "tokens.txt"), + encoder=str(model_dir / "encoder.int8.onnx"), + decoder=str(model_dir / "decoder.int8.onnx"), + joiner=str(model_dir / "joiner.int8.onnx"), + num_threads=self.num_threads, + sample_rate=SAMPLE_RATE, + feature_dim=80, + enable_endpoint_detection=True, + rule1_min_trailing_silence=2.4, + rule2_min_trailing_silence=1.2, + rule3_min_utterance_length=300, + decoding_method="greedy_search", + model_type="nemo_transducer", + ) + vad_cfg = sherpa_onnx.VadModelConfig() + vad_cfg.silero_vad.model = str(self.models_dir / "silero_vad.onnx") + vad_cfg.silero_vad.threshold = 0.5 + vad_cfg.silero_vad.min_silence_duration = 0.4 + vad_cfg.sample_rate = SAMPLE_RATE + vad = sherpa_onnx.VoiceActivityDetector( + vad_cfg, buffer_size_in_seconds=30 + ) + return recognizer, vad + + def _load_modules(self): + """Import the heavy dependencies (worker thread only; test seam).""" + import numpy + import sherpa_onnx + return numpy, sherpa_onnx + + def _run(self): + # Imports sit inside the guard: a half-installed venv (numpy or + # sherpa-onnx missing) must fail with a ready=False status, not a + # silent thread death via the default excepthook. + try: + np, sherpa_onnx = self._load_modules() + recognizer, vad = self._build(sherpa_onnx) + except Exception as exc: + print(f"asr: model load failed: {exc}", file=sys.stderr) + self._ready = False + self._drain_queue() + self._emit_status() + return + stream = recognizer.create_stream() + self._ready = True + self._emit_status() + try: + self._decode_loop(np, recognizer, vad, stream) + except Exception as exc: + # A mid-session crash (onnxruntime error, bad frame) must never + # leave ready=True on a dead worker: flag it, empty the dead + # queue so stats reflect death, tell the clients, exit. + print(f"asr: decode loop crashed: {exc}", file=sys.stderr) + self._ready = False + self._drain_queue() + self._emit_status() + + def _decode_loop(self, np, recognizer, vad, stream): + segment = 0 + prev_text = "" + speaking = False + vad_buf = np.zeros(0, dtype=np.float32) + + while not self._stop.is_set(): + try: + item = self._queue.get(timeout=0.2) + except queue.Empty: + continue + if item is None: + break + if self._behind and self._queue.qsize() < QUEUE_MAX // 2: + self._behind = False + self._emit_status() + + data = item + if len(data) % 2: + data = data[:-1] + if not data: + continue + samples = ( + np.frombuffer(data, dtype="= VAD_WINDOW: + vad.accept_waveform(vad_buf[:VAD_WINDOW]) + vad_buf = vad_buf[VAD_WINDOW:] + while not vad.empty(): + vad.pop() + now_speaking = bool(vad.is_speech_detected()) + if now_speaking != speaking: + speaking = now_speaking + self._emit({"kind": "vad", "speaking": speaking}) + + # Recognizer: any chunk length is fine. + stream.accept_waveform(SAMPLE_RATE, samples) + while recognizer.is_ready(stream): + recognizer.decode_stream(stream) + result = json.loads(recognizer.get_result_as_json_string(stream)) + text = result.get("text", "") + tokens = result.get("tokens", []) + if text != prev_text: + self._emit( + { + "kind": "partial", + "segment": segment, + "text": text, + "tokens": tokens, + } + ) + prev_text = text + if recognizer.is_endpoint(stream): + if text: + self._emit( + { + "kind": "final", + "segment": segment, + "text": text, + "tokens": tokens, + } + ) + recognizer.reset(stream) + segment += 1 + prev_text = "" diff --git a/skills/mc-prompter/scripts/server/main.py b/skills/mc-prompter/scripts/server/main.py index 5878f72..ce20568 100644 --- a/skills/mc-prompter/scripts/server/main.py +++ b/skills/mc-prompter/scripts/server/main.py @@ -3,7 +3,7 @@ # requires-python = ">=3.11" # dependencies = ["aiohttp==3.12.15"] # /// -"""mc-prompter aiohttp server (Phase A: classic teleprompter). +"""mc-prompter aiohttp server (Phase A classic teleprompter + Phase B voice-follow). Launched by run_prompter.py as `python -m server.main` with cwd set to the scripts directory (the PEP 723 header above also allows a direct @@ -14,6 +14,17 @@ python -m server.main --port 8770 --host 127.0.0.1|0.0.0.0 --script --owner-wpm --session-file --token + [--models-dir ] [--asr-provider none|nemotron-streaming] + +ASR (Phase B): when --models-dir is set AND --asr-provider is an implemented +provider (nemotron-streaming), server/asr.py is imported lazily and an +AsrEngine runs on a dedicated worker thread; server/align.py is imported +lazily too and an Aligner is rebuilt from script_ingest.speakable_words(doc) +on every script load or edit. Without --models-dir (or with provider none) +the server is pure tier 1: neither asr nor align nor numpy is ever imported, +and it runs with only aiohttp installed. Provider zipformer-small is a +documented planned lane: selecting it exits 3 at startup, it never pretends +to run. HTTP routes: GET / home page (static/home.html) @@ -28,7 +39,9 @@ "started": ""} GET /api/state {"snapshot": , "doc-version": n, "script": {"path", "title", "word-count"} or null, - "config": {"owner-wpm": N or null}} + "config": {"owner-wpm": N or null}, + "asr": {"available": bool, "provider": str, + "ready": bool}} GET /api/source {"path": , "raw": "", "doc": diff --git a/skills/mc-prompter/scripts/server/static/prompt.html b/skills/mc-prompter/scripts/server/static/prompt.html index 5d7965f..969a526 100644 --- a/skills/mc-prompter/scripts/server/static/prompt.html +++ b/skills/mc-prompter/scripts/server/static/prompt.html @@ -31,6 +31,9 @@
150 wpm manual + + + @@ -130,6 +133,13 @@

Display

+
+

Voice follow

+
+ + +
+

Script markers

@@ -158,8 +168,9 @@

Keyboard shortcuts

ccountdown, then play ssection list dsettings drawer + vtoggle voice follow (this machine only) mouse wheelspeed up / down 2 wpm - click a wordjump the eyeline to that word + click a wordjump the eyeline to that word (re-anchors in voice follow) ?this help escclose panels @@ -168,6 +179,35 @@

Keyboard shortcuts

+ +