#!/usr/bin/env python3 """Replay an explicit local audio manifest into Oruk, one attempt per item. Default: offline preflight only. Python 3.10+, macOS/Linux, standard library. No Hume access, automatic HTTP retries, redirects, resampling or directory crawl. See https://oruk.ai/examples/hume-batch-manifest.example.json and the repository operator guide docs/hume-batch-replay.md. Code: MIT; input rights are unchanged. """ from __future__ import annotations import argparse from concurrent.futures import ThreadPoolExecutor from contextlib import contextmanager from decimal import Decimal, InvalidOperation, ROUND_CEILING import datetime import hashlib import http.client import io import json import math import os from pathlib import Path import re import stat import sys import threading from urllib.parse import urlsplit import uuid import wave MAX_FILES = 500 MAX_AUDIO_BYTES = 4 * 1024 * 1024 MAX_TOTAL_BYTES = 256 * 1024 * 1024 MAX_JSON_BYTES = 2 * 1024 * 1024 MAX_STATE_BYTES = 8 * 1024 * 1024 MAX_RESPONSE_BYTES = 4 * 1024 * 1024 MAX_JSON_DEPTH = 24 MAX_JSON_NODES = 50_000 MAX_CONCURRENCY = 4 API_BASE = "https://speech-api.oruk.ai" CONTRACTS = { ("oruk-spectra-2", "analysis"), ("oruk-spectra-2", "transcriptions"), ("oruk-resonance", "analysis"), ("oruk-resonance", "affect"), ("oruk-resonance", "emotions"), ("oruk-resonance", "styles"), } RESPONSE_TASKS = {"analysis": "analysis", "affect": "affect", "transcriptions": "transcription", "emotions": "emotion", "styles": "style"} # Spectra's public contract returns every label exactly once on both routes. # Keep this standalone snapshot aligned with lib/api-product.ts; no Hume mapping. SPECTRA_SCORE_LABELS = { "emotions": frozenset("happy excited hopeful sad worried angry frustrated disappointed scared disgusted surprised embarrassed proud relieved neutral".split()), "styles": frozenset("energetic passionate irritated warm playful sarcastic deadpan hesitant confident sincere skeptical tired formal casual impatient distracted".split()), } HEX = re.compile(r"[0-9a-f]{64}\Z") IDENTIFIER = re.compile(r"[A-Za-z0-9_-]{1,64}\Z") STATE_KEYS = {"schema_version", "manifest_sha256", "batch_id", "items"} def utc_now(): return datetime.datetime.now(datetime.timezone.utc).isoformat() def digest(data: bytes) -> str: return hashlib.sha256(data).hexdigest() def canonical(value) -> bytes: return json.dumps(value, sort_keys=True, separators=(",", ":"), ensure_ascii=True, allow_nan=False).encode() def strict_json(data: bytes): def pairs(items): result = {} for key, value in items: if key in result: raise ValueError("duplicate JSON key") result[key] = value return result def constant(_): raise ValueError("non-finite JSON number") try: value = json.loads(data.decode("utf-8"), object_pairs_hook=pairs, parse_constant=constant) except (UnicodeError, RecursionError, OverflowError) as exc: raise ValueError("invalid or excessive JSON") from exc pending, count = [(value, 0)], 0 while pending: item, depth = pending.pop() count += 1 if count > MAX_JSON_NODES or depth > MAX_JSON_DEPTH: raise ValueError("JSON structure exceeds its limit") if isinstance(item, float) and not math.isfinite(item): raise ValueError("non-finite JSON number") if isinstance(item, dict): pending.extend((child, depth + 1) for child in item.values()) elif isinstance(item, list): pending.extend((child, depth + 1) for child in item) return value def fields(value, expected, name): if not isinstance(value, dict) or set(value) != set(expected): raise ValueError(f"{name} must contain exactly: {', '.join(sorted(expected))}") def positive_int(value, maximum, name): if type(value) is not int or not 1 <= value <= maximum: raise ValueError(f"{name} must be an integer from 1 through {maximum}") return value def relative_parts(value): if not isinstance(value, str) or not value or "\\" in value or "\x00" in value: raise ValueError("audio_path must be a plain relative path") parts = value.split("/") if any(p in {"", ".", ".."} for p in parts) or any(ord(c) < 32 for c in value): raise ValueError("absolute paths, traversal and control characters are forbidden") return parts @contextmanager def directory(path): """Open every directory component without following symlinks (including parents).""" if not hasattr(os, "O_NOFOLLOW") or not hasattr(os, "O_DIRECTORY"): raise ValueError("this example requires macOS/Linux no-follow directory support") target = Path(os.path.abspath(path)) fd = os.open("/", os.O_RDONLY | os.O_DIRECTORY) try: for component in target.parts[1:]: next_fd = os.open(component, os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW, dir_fd=fd) os.close(fd) fd = next_fd yield fd finally: os.close(fd) def read_at(root_fd, relative, maximum): parts = relative_parts(relative) fd = os.dup(root_fd) try: for component in parts[:-1]: next_fd = os.open(component, os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW, dir_fd=fd) os.close(fd) fd = next_fd file_fd = os.open(parts[-1], os.O_RDONLY | os.O_NOFOLLOW | os.O_NONBLOCK, dir_fd=fd) with os.fdopen(file_fd, "rb") as stream: before = os.fstat(stream.fileno()) if not stat.S_ISREG(before.st_mode) or before.st_size > maximum: raise ValueError("input must be a bounded regular file") value = stream.read(maximum + 1) after = os.fstat(stream.fileno()) if len(value) > maximum or len(value) != before.st_size or (before.st_size, before.st_mtime_ns) != (after.st_size, after.st_mtime_ns): raise ValueError("input changed while being read or exceeded its limit") return value finally: os.close(fd) def read_path(path, maximum=MAX_JSON_BYTES): path = Path(os.path.abspath(path)) with directory(path.parent) as root_fd: return read_at(root_fd, path.name, maximum) def wav_frames(data): if len(data) < 44 or data[:4] != b"RIFF" or data[8:12] != b"WAVE" or int.from_bytes(data[4:8], "little") != len(data) - 8: raise ValueError("expected an intact RIFF WAV, with no trailing or missing bytes") try: with wave.open(io.BytesIO(data), "rb") as audio: if (audio.getnchannels(), audio.getsampwidth(), audio.getframerate(), audio.getcomptype()) != (1, 2, 16000, "NONE"): raise ValueError("this runner requires original mono 16 kHz PCM16 WAV; it never converts audio") count = audio.getnframes() if not 720 <= count <= 960000 or len(audio.readframes(count + 1)) != count * 2: raise ValueError("audio must contain 45 ms–60 s of intact PCM samples") return count except (wave.Error, EOFError) as exc: raise ValueError("invalid WAV structure") from exc def money(value, nullable=False): if nullable and value is None: return None if not isinstance(value, str) or not re.fullmatch(r"(?:0|[1-9][0-9]{0,5})(?:\.[0-9]{1,9})?", value): raise ValueError("declared rate must be a nonnegative decimal string or null") try: return Decimal(value) except InvalidOperation as exc: raise ValueError("invalid declared rate") from exc def multipart(item, data, model): boundary = "oruk-batch-" + item["request_id"] body = (f"--{boundary}\r\nContent-Disposition: form-data; name=\"model\"\r\n\r\n{model}\r\n" f"--{boundary}\r\nContent-Disposition: form-data; name=\"file\"; filename=\"audio.wav\"\r\n" "Content-Type: audio/wav\r\n\r\n").encode() + data + f"\r\n--{boundary}--\r\n".encode() return body, f"multipart/form-data; boundary={boundary}" def preflight(manifest_path, audio_root): raw_manifest = read_path(manifest_path) configured_key = os.environ.get("ORUK_API_KEY", "") if configured_key and configured_key.encode() in raw_manifest: raise ValueError("credential detected in manifest") manifest = strict_json(raw_manifest) if configured_key and configured_key.encode() in canonical(manifest): raise ValueError("credential detected in manifest") fields(manifest, {"schema_version", "batch_id", "account_reference", "model", "task", "pricing", "items"}, "manifest") if manifest["schema_version"] != 1 or type(manifest["schema_version"]) is not int: raise ValueError("unsupported manifest schema_version") if not isinstance(manifest["batch_id"], str) or str(uuid.UUID(manifest["batch_id"])) != manifest["batch_id"]: raise ValueError("batch_id must be a canonical UUID, unique to this batch") if not isinstance(manifest["account_reference"], str) or not IDENTIFIER.fullmatch(manifest["account_reference"]): raise ValueError("account_reference must be a non-secret account alias") if not isinstance(manifest["model"], str) or not isinstance(manifest["task"], str) or (manifest["model"], manifest["task"]) not in CONTRACTS: raise ValueError("unsupported model/task; EVI, TTS, Resonance-2 and streaming are outside this runner") fields(manifest["pricing"], {"declared_usd_per_minute", "source"}, "pricing") rate = money(manifest["pricing"]["declared_usd_per_minute"], nullable=True) source = manifest["pricing"]["source"] if not isinstance(source, str) or not 1 <= len(source) <= 300 or any(ord(c) < 32 for c in source): raise ValueError("pricing.source must state the account offer/source or why pricing is unknown") entries = manifest["items"] if not isinstance(entries, list) or not 1 <= len(entries) <= MAX_FILES: raise ValueError(f"items must contain 1–{MAX_FILES} explicit recordings") seen_ids, seen_audio = set(), set() checked, errors, total_bytes, total_frames, costs = [], [], 0, 0, [] with directory(audio_root) as root_fd: for entry in entries: fields(entry, {"id", "audio_path", "sha256", "bytes", "frames"}, "item") item_id = entry["id"] if not isinstance(item_id, str) or not IDENTIFIER.fullmatch(item_id) or item_id in seen_ids: raise ValueError("item IDs must be unique, 1–64 ASCII letters/numbers/underscores/hyphens") seen_ids.add(item_id) relative_parts(entry["audio_path"]) if not isinstance(entry["sha256"], str) or not HEX.fullmatch(entry["sha256"]) or entry["sha256"] in seen_audio: raise ValueError("audio SHA-256 values must be lowercase, valid and unique within the batch") seen_audio.add(entry["sha256"]) positive_int(entry["bytes"], MAX_AUDIO_BYTES, "item.bytes") positive_int(entry["frames"], 960000, "item.frames") total_bytes += entry["bytes"] if total_bytes > MAX_TOTAL_BYTES: raise ValueError("batch exceeds 256 MiB") try: data = read_at(root_fd, entry["audio_path"], MAX_AUDIO_BYTES) frames = wav_frames(data) if len(data) != entry["bytes"] or digest(data) != entry["sha256"] or frames != entry["frames"]: raise ValueError("audio size, SHA-256 or frame count does not match manifest") except (ValueError, OSError) as exc: errors.append({"id": item_id, "error": "audio_preflight_failed", "detail": str(exc) if isinstance(exc, ValueError) else "file inaccessible, nonregular or symlinked"}) continue identity = {"endpoint": "/v1/audio/" + manifest["task"], "model": manifest["model"], "task": manifest["task"], "audio_sha256": entry["sha256"], "filename": "audio.wav", "content_type": "audio/wav"} fingerprint = digest(canonical(identity)) request_id = "hume-batch-" + digest(canonical([manifest["batch_id"], item_id, fingerprint]))[:48] item = {**entry, "payload_fingerprint": fingerprint, "request_id": request_id} item["multipart_sha256"] = digest(multipart(item, data, manifest["model"])[0]) checked.append(item) total_frames += frames if rate is not None: cost = (max(Decimal(frames) / 16000, Decimal(1)) * rate / 60).quantize(Decimal("0.000001"), rounding=ROUND_CEILING) costs.append(cost) if errors: return {"valid": False, "errors": errors, "network_requests": 0} fingerprint = digest(canonical(manifest)) return {"valid": True, "manifest": manifest, "manifest_sha256": fingerprint, "items": checked, "estimate": {"requests": len(checked), "audio_bytes": total_bytes, "audio_seconds": total_frames / 16000, "billable_seconds_with_one_second_minimum": sum(max(x["frames"] / 16000, 1) for x in checked), "declared_usd_per_minute": str(rate) if rate is not None else None, "estimated_usd": str(sum(costs)) if rate is not None else None, "pricing_source": source, "invoice_or_spend_cap": False, "note": "Operator-declared rate only; allowances, tax, credits and account billing may differ."}} class Ledger: def __init__(self, state_dir, plan): self.state_dir, self.plan = Path(os.path.abspath(state_dir)), plan self.mutex = threading.Lock() def __enter__(self): import fcntl with directory(self.state_dir.parent) as parent: try: os.mkdir(self.state_dir.name, mode=0o700, dir_fd=parent) except FileExistsError: pass self.context = directory(self.state_dir) self.fd = self.context.__enter__() if stat.S_IMODE(os.fstat(self.fd).st_mode) & 0o077: self.close() raise ValueError("state directory must be private (mode 0700)") self.lock = os.open("lock", os.O_RDWR | os.O_CREAT | os.O_NOFOLLOW | os.O_NONBLOCK, 0o600, dir_fd=self.fd) if not stat.S_ISREG(os.fstat(self.lock).st_mode): self.close() raise ValueError("ledger lock must be a regular file") try: fcntl.flock(self.lock, fcntl.LOCK_EX | fcntl.LOCK_NB) try: value = strict_json(read_at(self.fd, "ledger.json", MAX_STATE_BYTES)) except FileNotFoundError: # A pre-existing result is never silently overwritten by a fresh ledger. for item in self.plan["items"]: try: os.stat(item["request_id"] + ".json", dir_fd=self.fd, follow_symlinks=False) except FileNotFoundError: continue raise ValueError("result exists without ledger; reconcile before continuing") value = {"schema_version": 1, "manifest_sha256": self.plan["manifest_sha256"], "batch_id": self.plan["manifest"]["batch_id"], "items": {item["id"]: {"request_id": item["request_id"], "payload_fingerprint": item["payload_fingerprint"], "audio_sha256": item["sha256"], "multipart_sha256": item["multipart_sha256"], "status": "pending", "attempts": 0} for item in self.plan["items"]}} self.validate(value) self.value = value for row in value["items"].values(): if row["status"] == "dispatching": row.update(status="uncertain", error="interrupted_after_durable_dispatch_intent") self.save() return self except BaseException: self.close() raise def validate(self, value): configured_key = os.environ.get("ORUK_API_KEY", "") if configured_key and configured_key.encode() in canonical(value): raise ValueError("credential detected in ledger") fields(value, STATE_KEYS, "ledger") if value["schema_version"] != 1 or value["manifest_sha256"] != self.plan["manifest_sha256"] or value["batch_id"] != self.plan["manifest"]["batch_id"]: raise ValueError("ledger belongs to another immutable manifest; do not replace or reset it") if not isinstance(value["items"], dict) or set(value["items"]) != {x["id"] for x in self.plan["items"]}: raise ValueError("ledger item set differs from manifest") statuses = {"pending", "dispatching", "completed", "uncertain", "rejected", "local_error", "reconciled_completed", "abandoned"} for item in self.plan["items"]: row = value["items"][item["id"]] for key, expected in [("request_id", item["request_id"]), ("payload_fingerprint", item["payload_fingerprint"]), ("audio_sha256", item["sha256"]), ("multipart_sha256", item["multipart_sha256"])]: if not isinstance(row, dict) or row.get(key) != expected: raise ValueError("ledger request/payload identity changed") if row.get("status") not in statuses or type(row.get("attempts")) is not int or row["attempts"] not in (0, 1): raise ValueError("invalid ledger state") if row["status"] == "pending" and row["attempts"] != 0: raise ValueError("an attempted item cannot become pending") if row["status"] in {"dispatching", "completed", "uncertain", "rejected", "reconciled_completed"} and row["attempts"] != 1: raise ValueError("attempted state must retain its dispatch count") if row["status"] == "completed": result = read_at(self.fd, item["request_id"] + ".json", MAX_RESPONSE_BYTES) if digest(result) != row.get("response_sha256"): raise ValueError("completed result is missing or modified; no inference will be repeated") def atomic(self, filename, data): temporary = ".write-" + uuid.uuid4().hex fd = os.open(temporary, os.O_WRONLY | os.O_CREAT | os.O_EXCL | os.O_NOFOLLOW, 0o600, dir_fd=self.fd) try: with os.fdopen(fd, "wb") as stream: stream.write(data) stream.flush() os.fsync(stream.fileno()) os.replace(temporary, filename, src_dir_fd=self.fd, dst_dir_fd=self.fd) os.fsync(self.fd) finally: try: os.unlink(temporary, dir_fd=self.fd) except FileNotFoundError: pass def save(self): self.atomic("ledger.json", canonical(self.value) + b"\n") def close(self): if getattr(self, "lock", None) is not None: os.close(self.lock) self.lock = None if getattr(self, "context", None) is not None: self.context.__exit__(None, None, None) self.context = None def __exit__(self, *_): self.close() def endpoint(base_url, allow_local_http=False): parsed = urlsplit(base_url) if parsed.username or parsed.password or parsed.query or parsed.fragment or parsed.path not in {"", "/"}: raise ValueError("API base URL must have no credentials, path, query or fragment") if base_url.rstrip("/") == API_BASE: return parsed if allow_local_http and parsed.scheme == "http" and parsed.hostname in {"127.0.0.1", "::1"}: return parsed raise ValueError("only the canonical HTTPS API is allowed; loopback HTTP is solely for explicit local fixture tests") def send_once(base_url, task, body, content_type, key, request_id, timeout, allow_local_http=False): parsed = endpoint(base_url, allow_local_http) connection_type = http.client.HTTPSConnection if parsed.scheme == "https" else http.client.HTTPConnection connection = connection_type(parsed.hostname, port=parsed.port, timeout=timeout) try: connection.request("POST", "/v1/audio/" + task, body=body, headers={"Authorization": "Bearer " + key, "X-Request-ID": request_id, "Content-Type": content_type}) response = connection.getresponse() length = response.getheader("Content-Length") transfer = response.getheader("Transfer-Encoding") if transfer is not None and (transfer.lower() != "chunked" or length is not None): raise ValueError("ambiguous response framing") if length is not None and (not re.fullmatch(r"[0-9]{1,10}", length) or int(length) > MAX_RESPONSE_BYTES): raise ValueError("invalid or excessive response length") raw = response.read(MAX_RESPONSE_BYTES + 1) if len(raw) > MAX_RESPONSE_BYTES: raise ValueError("response limit exceeded") # read(amt) may return a short body without raising IncompleteRead. if length is not None and len(raw) != int(length): raise ValueError("response ended before its declared length") return response.status, response.getheader("X-Request-ID"), raw finally: connection.close() def validate_response(value, manifest): """Check the selected native contract without translating Hume labels.""" def positive_number(number): return type(number) in (int, float) and 0 < number <= 10**12 and math.isfinite(number) if (value.get("object") != "speech.result" or value.get("model") != manifest["model"] or value.get("task") != RESPONSE_TASKS[manifest["task"]] or not isinstance(value.get("id"), str) or not 1 <= len(value["id"]) <= 128 or not positive_number(value.get("duration"))): raise ValueError("response contract mismatch") usage = value.get("usage") if not isinstance(usage, dict) or not all(positive_number(usage.get(key)) for key in ("audio_seconds", "billable_seconds")): raise ValueError("response usage is incomplete") if manifest["task"] in {"analysis", "transcriptions"} and not isinstance(value.get("text"), str): raise ValueError("response transcript is missing") score_fields = set() if manifest["model"] == "oruk-spectra-2" or manifest["task"] in {"analysis", "affect", "emotions"}: score_fields.add("emotions") if manifest["model"] == "oruk-spectra-2" or manifest["task"] in {"analysis", "affect", "styles"}: score_fields.add("styles") for field in score_fields: scores = value.get(field) if not isinstance(scores, list) or any( not isinstance(row, dict) or not isinstance(row.get("label"), str) or not row["label"] or type(row.get("score")) not in (int, float) or not 0 <= row["score"] <= 1 for row in scores): raise ValueError("response scores are missing or invalid") if manifest["model"] == "oruk-spectra-2": expected = SPECTRA_SCORE_LABELS[field] if len(scores) != len(expected) or {row["label"] for row in scores} != expected: raise ValueError("response scores do not contain each contracted label exactly once") def execute(plan, audio_root, ledger, key, concurrency=1, timeout=120, base_url=API_BASE, allow_local_http=False): positive_int(concurrency, MAX_CONCURRENCY, "concurrency") endpoint(base_url, allow_local_http) if not isinstance(key, str) or not key or any(ord(c) < 33 or ord(c) > 126 for c in key): raise ValueError("ORUK_API_KEY must contain a nonempty printable token; no key is saved") if type(timeout) not in (int, float) or not math.isfinite(timeout) or not 1 <= timeout <= 300: raise ValueError("timeout must be 1–300 seconds") # An accidental credential in a free-text manifest or restored ledger is refused. if key.encode() in canonical(plan["manifest"]) or key.encode() in canonical(ledger.value): raise ValueError("credential detected in manifest or ledger") manifest = plan["manifest"] stop = threading.Event() with directory(audio_root) as root_fd: def run(item): with ledger.mutex: row = ledger.value["items"][item["id"]] if stop.is_set() or row["status"] != "pending": return try: data = read_at(root_fd, item["audio_path"], MAX_AUDIO_BYTES) if digest(data) != item["sha256"] or len(data) != item["bytes"] or wav_frames(data) != item["frames"]: raise ValueError("audio changed after preflight") body, content_type = multipart(item, data, manifest["model"]) if digest(body) != item["multipart_sha256"]: raise ValueError("multipart identity changed") except (OSError, ValueError): with ledger.mutex: row.update(status="local_error", error="audio_changed_or_unreadable_before_dispatch") ledger.save() return with ledger.mutex: if stop.is_set(): return row.update(status="dispatching", attempts=1, dispatched_at=utc_now()) ledger.save() # Durable intent MUST precede even the TCP connection. outcome = {"status": "uncertain", "error": "network_or_response_outcome_unknown"} result_data = None try: status, response_id, raw = send_once(base_url, manifest["task"], body, content_type, key, item["request_id"], timeout, allow_local_http) if key.encode() in raw: raise ValueError("response contains a credential") outcome["http_status"] = status value = strict_json(raw) if key.encode() in canonical(value): raise ValueError("response contains a credential") if not isinstance(value, dict): raise ValueError("response is not an object") if response_id is not None and response_id != item["request_id"]: raise ValueError("response request ID mismatch") if status == 200: validate_response(value, manifest) result_data = raw outcome = {"status": "completed", "http_status": status, "response_sha256": digest(raw)} elif 400 <= status < 500 and status != 409: outcome.update(status="rejected", error="http_rejected_no_automatic_retry") else: outcome["error"] = "http_outcome_requires_reconciliation" code = value.get("error", {}).get("code") if isinstance(value.get("error"), dict) else None if isinstance(code, str) and re.fullmatch(r"[a-z0-9_]{1,100}", code): outcome["api_error_code"] = code except Exception: pass # Never log raw exceptions, keys, payloads or response content. with ledger.mutex: if result_data is not None: ledger.atomic(item["request_id"] + ".json", result_data) row.update(outcome, observed_at=utc_now()) ledger.save() pool = ThreadPoolExecutor(max_workers=concurrency) try: list(pool.map(run, plan["items"])) except BaseException: stop.set() pool.shutdown(wait=True, cancel_futures=True) raise else: pool.shutdown(wait=True) return summary(plan, ledger.value) def summary(plan, value): return {"manifest_sha256": plan["manifest_sha256"], "estimate": plan["estimate"], "items": [{"id": name, **row} for name, row in value["items"].items()], "automatic_retries": 0, "note": "Resume sends only never-attempted pending entries. Resolve uncertain outcomes against Usage/support."} def resolve(ledger, item_id, resolution, evidence, acknowledged): if not acknowledged or resolution not in {"completed", "abandoned"} or not isinstance(evidence, str) or not 1 <= len(evidence) <= 300 or any(ord(c) < 32 for c in evidence): raise ValueError("manual resolution requires explicit acknowledgement and a non-secret evidence reference") key = os.environ.get("ORUK_API_KEY", "") if key and key in evidence: raise ValueError("evidence reference must not contain a key") row = ledger.value["items"].get(item_id) if not row or row["status"] not in {"uncertain", "rejected", "local_error"}: raise ValueError("only unresolved items may be manually reconciled") if resolution == "completed" and row["attempts"] != 1: raise ValueError("an item never dispatched cannot be reconciled as completed") row.update(status="reconciled_completed" if resolution == "completed" else "abandoned", resolution={"recorded_at": utc_now(), "operator_attestation": True, "outcome": resolution, "evidence_reference": evidence}) ledger.save() class SafeParser(argparse.ArgumentParser): def error(self, message): self.exit(2, "Invalid arguments. Use --help; argument values are not echoed.\n") def main(argv=None): parser = SafeParser(description=__doc__) parser.add_argument("manifest", type=Path) parser.add_argument("--audio-root", required=True, type=Path) parser.add_argument("--state-dir", type=Path, help="private local ledger/results; required for execution or reconciliation") parser.add_argument("--execute", action="store_true", help="send pending items once; requires ORUK_API_KEY") parser.add_argument("--confirm-manifest", help="SHA-256 shown by the offline preflight") parser.add_argument("--acknowledge-unpriced", action="store_true", help="explicitly accept an unknown monetary estimate") parser.add_argument("--concurrency", type=int, default=1) parser.add_argument("--timeout", type=float, default=120) parser.add_argument("--api-base-url", default=API_BASE) parser.add_argument("--allow-local-http", action="store_true", help="loopback fixture server only") parser.add_argument("--resolve-item") parser.add_argument("--resolution", choices=["completed", "abandoned"]) parser.add_argument("--evidence-reference") parser.add_argument("--acknowledge-manual-reconciliation", action="store_true") args = parser.parse_args(argv) try: plan = preflight(args.manifest, args.audio_root) if not plan["valid"]: print(json.dumps(plan, indent=2)) return 2 if not args.resolve_item and any((args.resolution, args.evidence_reference, args.acknowledge_manual_reconciliation)): raise ValueError("reconciliation options require resolve-item") if args.execute and args.resolve_item: raise ValueError("execution and manual reconciliation are separate actions") if not args.execute and not args.resolve_item: print(json.dumps({k: v for k, v in plan.items() if k != "manifest"}, indent=2)) return 0 if args.state_dir is None or args.confirm_manifest != plan["manifest_sha256"]: raise ValueError("state-dir and exact confirm-manifest digest are required") if args.execute and plan["estimate"]["estimated_usd"] is None and not args.acknowledge_unpriced: raise ValueError("unknown price: supply a declared account rate or explicitly acknowledge-unpriced") with Ledger(args.state_dir, plan) as ledger: if args.resolve_item: resolve(ledger, args.resolve_item, args.resolution, args.evidence_reference, args.acknowledge_manual_reconciliation) result = summary(plan, ledger.value) else: result = execute(plan, args.audio_root, ledger, os.environ.get("ORUK_API_KEY", ""), args.concurrency, args.timeout, args.api_base_url, args.allow_local_http) print(json.dumps(result, indent=2)) return 0 if all(x["status"] in {"completed", "reconciled_completed", "abandoned"} for x in result["items"]) else 3 except (ValueError, OSError, KeyError, TypeError) as exc: print(json.dumps({"error": str(exc) if isinstance(exc, ValueError) else "local state/input error; no automatic retry", "exception_type": type(exc).__name__}), file=sys.stderr) return 2 if __name__ == "__main__": raise SystemExit(main())