"""Python 3.10+; pip install websocket-clientLECSYNC_API_KEY=lk_live_... python3 examples/api/python/stream.py [input.wav]No file: stream 3 seconds of silence. Only 16-bit PCM WAV, mono/stereo.Decode other formats: ffmpeg -i input.mp3 -ac 1 -ar 16000 -c:a pcm_s16le input.wav"""import jsonimport mathimport osimport queueimport structimport sysimport threadingimport timeimport urllib.errorimport urllib.requestimport wavetry: import websocketexcept ImportError as exc: raise SystemExit("pip install websocket-client") from exc# CONNECT_CONFIG_STARTCONNECT_CONFIG = { "audio": {"encoding": "pcm_s16le", "sampleRate": 16000, "channels": 1}, "transcribe": {"languages": ["en"]}, "store": False,}# CONNECT_CONFIG_ENDSAMPLE_RATE = 16000FRAME_MS = 100FRAME_BYTES = 3200MAX_QUEUE_FRAMES = 20 # 2 seconds of pending PCMdef load_audio(file): if not file: return 3 * SAMPLE_RATE, SAMPLE_RATE, lambda i: 0 with wave.open(file, "rb") as wav: channels, rate = wav.getnchannels(), wav.getframerate() if wav.getcomptype() != "NONE" or wav.getsampwidth() != 2 or channels not in (1, 2): raise ValueError("Only uncompressed 16-bit mono/stereo PCM WAV is supported") samples = wav.getnframes() data = wav.readframes(samples) if not samples or len(data) != samples * channels * 2: raise ValueError("Empty or truncated WAV") def sample(i): left = struct.unpack_from("<h", data, i * channels * 2)[0] return left if channels == 1 else (left + struct.unpack_from("<h", data, i * 4 + 2)[0]) / 2 return samples, rate, sampledef make_frame(audio, index): # Source file is in memory; resampled output is generated one frame at a time. samples, rate, sample = audio total = math.ceil(samples * SAMPLE_RATE / rate) start = index * FRAME_BYTES // 2 frame = bytearray(min(FRAME_BYTES // 2, total - start) * 2) for i in range(len(frame) // 2): position = min((start + i) * rate / SAMPLE_RATE, samples - 1) lo = math.floor(position) hi = min(lo + 1, samples - 1) value = sample(lo) + (sample(hi) - sample(lo)) * (position - lo) struct.pack_into("<h", frame, i * 2, max(-32768, min(32767, math.floor(value + 0.5)))) return bytes(frame) # Final frame can be shorter than 100 ms.def main(): key = os.environ.get("LECSYNC_API_KEY") if not key: raise ValueError("Set LECSYNC_API_KEY") audio = load_audio(sys.argv[1] if len(sys.argv) > 1 else None) origin = os.environ.get("LECSYNC_API_ORIGIN", "https://api.lecsync.com").rstrip("/") req = urllib.request.Request( f"{origin}/v1/realtime/connect", data=json.dumps(CONNECT_CONFIG).encode(), headers={"Authorization": f"Bearer {key}", "Content-Type": "application/json"}, method="POST", ) try: with urllib.request.urlopen(req, timeout=10) as res: body = json.load(res) except urllib.error.HTTPError as exc: raise RuntimeError(f"connect HTTP {exc.code}: {exc.read().decode()}") from exc url = body["url"] # Already complete. Configure only in connect; no resume. ws = websocket.create_connection(url, subprotocols=["lecsync.realtime.v1", body["token"]], timeout=10) ws.settimeout(1) pending = queue.Queue(maxsize=MAX_QUEUE_FRAMES) done, halted, cancelled = threading.Event(), threading.Event(), threading.Event() errors = [] stop_sent = threading.Event() workers = [] def produce(): try: frames = math.ceil(audio[0] * SAMPLE_RATE / audio[1] / (FRAME_BYTES // 2)) start, index = time.monotonic(), 0 while index < frames and not halted.is_set(): if halted.wait(max(0, start + (index + 1) * FRAME_MS / 1000 - time.monotonic())): break due = min(frames, math.floor((time.monotonic() - start) * 1000 / FRAME_MS)) dropped = 0 if due - index > MAX_QUEUE_FRAMES: skip = due - index - MAX_QUEUE_FRAMES index += skip dropped += skip * FRAME_BYTES while index < due: frame = make_frame(audio, index) index += 1 while True: try: pending.put_nowait(frame) break except queue.Full: try: dropped += len(pending.get_nowait()) except queue.Empty: pass if dropped: print(f"Dropped {dropped / (SAMPLE_RATE * 2):.3f} seconds of local audio", file=sys.stderr) except Exception as exc: errors.append(exc) halted.set() finally: done.set() def send(): try: while not cancelled.is_set(): if halted.is_set() or (done.is_set() and pending.empty()): ws.send(json.dumps({"type": "session.stop"})) stop_sent.set() return try: frame = pending.get(timeout=FRAME_MS / 1000) except queue.Empty: continue ws.send_binary(frame) # Blocking send has a 1 s timeout; never retry a partial frame. # Pace even after a blocked send; do not burst the queue on recovery. cancelled.wait(FRAME_MS / 1000) except Exception as exc: errors.append(exc) halted.set() cancelled.set() ready, ended, closed = False, False, False deadline = time.monotonic() + 10 try: while not cancelled.is_set() or ended: if stop_sent.is_set() and not ended: stop_sent.clear() deadline = time.monotonic() + 30 if deadline is not None and time.monotonic() > deadline: raise TimeoutError("Timed out waiting for ready, session.ended or close") try: opcode, data = ws.recv_data(control_frame=True) except websocket.WebSocketTimeoutException: continue if opcode == websocket.ABNF.OPCODE_CLOSE: code = struct.unpack("!H", data[:2])[0] if len(data) >= 2 else 1005 print("close", code, data[2:].decode("utf-8", errors="replace")) closed = True if not ended or code != 1000: raise RuntimeError("Stream closed; use a new connect") break if opcode != websocket.ABNF.OPCODE_TEXT: continue msg = json.loads(data) kind = msg.get("type") if kind in {"transcript.final", "translation.final", "digest.update", "usage.update", "session.error", "session.ended"}: print(kind, json.dumps(msg, ensure_ascii=False), flush=True) # Includes discardedSeconds. if kind == "session.ready" and not ready and not halted.is_set(): ready, deadline = True, None workers = [threading.Thread(target=produce), threading.Thread(target=send)] for worker in workers: worker.start() if kind == "session.error": # retryable means a NEW connect may succeed, never resume this stream. print("retryable:", msg.get("retryable"), "(new connect only)", file=sys.stderr) if msg.get("scope") == "session": errors.append(RuntimeError(msg.get("message") or msg.get("code"))) halted.set() deadline = time.monotonic() + 10 if kind == "session.ended": ended = True halted.set() cancelled.set() # Keep receiving the server's close frame, without sending more audio. deadline = time.monotonic() + 5 if errors: raise errors[0] finally: halted.set() cancelled.set() for worker in workers: worker.join(timeout=2) ws.close() if not closed: print("close", 1006, "No server close frame received")if __name__ == "__main__": try: main() except (Exception, KeyboardInterrupt) as exc: sys.exit(str(exc) or "Interrupted")