"""Local reports and community sharing.""" from pathlib import Path import json import os import re import sys import urllib.error import urllib.request from vlib import ctx from vlib.ui import _fail, _noninteractive, _note, _now_iso, _ok, _say, _step, _warn from vlib.net import _get_hf_token, _parse_retry_after from vlib.tensors import _is_gguf_file, read_safetensors from vlib.sources import _family, _public_source, detect_output_tensor, largest_2d from vlib.registry import _append_report, _iter_voices, _read_meta, _registry_resolve def _hash_file(path): import hashlib h = hashlib.sha256() with open(path, "rb") as f: while True: b = f.read(8 << 20) if not b: break h.update(b) return h.hexdigest() def _sanitize_note(s): if s is None: return None s = str(s).replace("\r", " ").replace("\n", " ").replace("\x00", " ") s = re.sub(r"\s+", " ", s).strip() if not s: return None # strip CSV injection prefixes if s and s[0] in "=+@-": s = "'" + s if len(s) > 240: s = s[:239] + "…" return s def _sanitize_reporter(s): if not s: return None s = re.sub(r"[^A-Za-z0-9._-]", "_", str(s).strip())[:64].strip("_") return s or None def _prompt_share(args, preview=None): if getattr(args, "no_push", False): return False if getattr(args, "yes", False): return True if _noninteractive(args): return False if preview: _say(f" Will share: {preview}") try: ans = input(" Leave a breadcrumb to help the next person? [y/N] (N): ").strip().lower() except (EOFError, KeyboardInterrupt): _say(" ⚠ Cancelled.") sys.exit(1) return ans in ("y", "yes") def _sync_pending(): """Sync outbox without creating anything. Shared by sync-only mode.""" pending = sorted(ctx.REPORTS_OUTBOX.glob("*.json")) if not pending: _say(" No pending reports to sync.") return _step(f"Syncing {len(pending)} pending report(s)…") recs = [] for p in pending: try: recs.append(json.loads(p.read_text())) except Exception as e: _warn(f" Could not read {p.name}: {e}") if recs: n = _push_reports_batch(recs) if n < len(recs): _warn(f" Synced {n}/{len(recs)} report(s) — rest stay queued.") else: _ok(f"Synced {n}/{len(recs)} report(s).") return def _push_reports_batch(records): """Batch push: group by shard, single fetch+upload per shard. Returns pushed count. Non-collaborators lack write access to Wiself/voice-reports → fallback to PR (create_pr=True). Fixes tempfile.mktemp racy API → mkstemp 0600; 3 retries on 429/409/412. Note: download-modify-upload still racy for direct pushes (two owners racing may lose); dedup on report_id softens it, and PR path eliminates race for crowd (each PR is isolated).""" import time, tempfile if not records: return 0 tok = _get_hf_token() if not tok: _warn(" No HF token — saved to outbox only. Run: voice auth") return 0 try: from huggingface_hub import HfApi from huggingface_hub import CommitOperationAdd api = HfApi(token=tok) except Exception as e: _warn(f" huggingface_hub not available ({e}) — kept in outbox.") return 0 try: api.create_repo(repo_id=ctx.REPORTS_DATASET_ID, repo_type="dataset", exist_ok=True) except Exception: pass from collections import defaultdict by_shard = defaultdict(list) for r in records: shard = f"{ctx.REPORTS_REMOTE_PREFIX}{r['ts'][:7]}.jsonl" by_shard[shard].append(r) pushed = 0 for shard, recs in by_shard.items(): for attempt in range(3): try: existing = "" try: from huggingface_hub import hf_hub_download tmp = hf_hub_download(repo_id=ctx.REPORTS_DATASET_ID, repo_type="dataset", filename=shard, token=tok) existing = Path(tmp).read_text(encoding="utf-8", errors="replace") except Exception: existing = "" # parse existing_ids properly (not substring) for correct dedup existing_ids = set() if existing: for line in existing.splitlines(): if not line.strip(): continue try: existing_ids.add(json.loads(line).get("report_id")) except Exception: continue # headroom check per Tim Cook: warn if shard near 10k if len(existing_ids) > 10000: _warn(f" Shard {shard} has {len(existing_ids)} reports — consider rotation (reports-YYYY-MM-2).") todo = [r for r in recs if r["report_id"] not in existing_ids] if not todo: for r in recs: try: (ctx.REPORTS_OUTBOX / f"{r['report_id']}.json").unlink(missing_ok=True) except Exception: pass pushed += len(recs) break if not existing: new_content = "\n".join(json.dumps(r, ensure_ascii=False, separators=(",", ":")) for r in todo) + "\n" else: if not existing.endswith("\n"): existing += "\n" new_content = existing + "\n".join(json.dumps(r, ensure_ascii=False, separators=(",", ":")) for r in todo) + "\n" # mkstemp 0600, not deprecated mktemp fd_tmp, tmp_path = tempfile.mkstemp(prefix="voice-report-") tmpf = Path(tmp_path) try: os.write(fd_tmp, new_content.encode("utf-8")) try: os.fsync(fd_tmp) except Exception: pass finally: try: os.close(fd_tmp) except Exception: pass try: os.chmod(tmpf, 0o600) except Exception: pass # try direct push first; on 403 (no write access) fallback to PR try: api.upload_file(path_or_fileobj=str(tmpf), path_in_repo=shard, repo_id=ctx.REPORTS_DATASET_ID, repo_type="dataset", commit_message=f"report batch {len(todo)} -> {shard}") except Exception as e_direct: is_403 = False try: from huggingface_hub.utils import HfHubHTTPError as _HE if isinstance(e_direct, _HE) and getattr(e_direct.response, "status_code", None) == 403: is_403 = True elif "403" in str(e_direct) or "Forbidden" in str(e_direct): is_403 = True except Exception: pass if is_403: _step(" No write access — creating pull request…") api.create_commit( repo_id=ctx.REPORTS_DATASET_ID, repo_type="dataset", operations=[CommitOperationAdd(path_in_repo=shard, path_or_fileobj=str(tmpf))], commit_message=f"report batch {len(todo)} -> {shard} ({todo[0]['report_id'][:8]}…)", create_pr=True, ) _ok(" PR opened — thank you! Maintainer will merge shortly.") else: raise try: tmpf.unlink(missing_ok=True) except Exception: pass for r in todo: try: (ctx.REPORTS_OUTBOX / f"{r['report_id']}.json").unlink(missing_ok=True) except Exception: pass pushed += len(todo) break except Exception as e: try: from huggingface_hub.utils import HfHubHTTPError as _HE code = getattr(getattr(e, "response", None), "status_code", None) if isinstance(e, _HE) and code in (429, 409, 412): ra = _parse_retry_after(getattr(e.response, "headers", {}) or {}) time.sleep(ra or (0.7 * (2 ** attempt))) continue # 403 already handled via PR above; other 403 should not retry as direct if isinstance(e, _HE) and code == 403: # if PR also failed, keep in outbox _warn(f" Push failed for {shard} (403 Forbidden) — kept in outbox. PR fallback also failed.") break except Exception: pass if attempt < 2 and isinstance(e, (urllib.error.URLError, TimeoutError, OSError)): time.sleep(0.7 * (2 ** attempt)) continue if attempt >= 2: if ctx.VERBOSE: import traceback _say(traceback.format_exc()) _warn(f" Push failed for {shard} ({type(e).__name__}: {e}) — kept in outbox.") break time.sleep(0.7 * (2 ** attempt)) return pushed def _push_report(record): return _push_reports_batch([record]) == 1 def _last_operation(): try: if not ctx.OPERATIONS_LOG.exists(): return None # read last non-empty line with open(ctx.OPERATIONS_LOG, "rb") as f: f.seek(0, 2) size = f.tell() if size == 0: return None # read tail 8KB read = min(size, 8192) f.seek(size - read) tail = f.read().decode("utf-8", "replace").strip().splitlines() for line in reversed(tail): if not line.strip(): continue try: rec = json.loads(line) op = rec.get("op") if op: return op except Exception: continue except Exception: pass return None def cmd_report(args): msg = _sanitize_note(getattr(args, "message", None)) voice_arg = getattr(args, "voice", None) target_raw = getattr(args, "target", None) if getattr(args, "sync", False) and not (msg or voice_arg or target_raw): # Sync-only: never mint a report just to sync. _sync_pending() return if not msg: _fail(' ✗ Empty note — pass -m "what you tried, what happened".') if isinstance(getattr(args, "message", None), str) and len(getattr(args, "message")) > 240: _note(" Note truncated to 240 chars.") import uuid, platform report_id = uuid.uuid4().hex ts = _now_iso() # auto-pull voice provenance _voice_missed = None voice_meta = None voice_name = _sanitize_note(voice_arg.strip()) if voice_arg and voice_arg.strip() else None try: if voice_arg: vp, vm = _registry_resolve(voice_arg) if vp and vp.exists(): if vm is None: vm = _read_meta(Path(vp).parent / "voice.json") voice_meta = dict(vm or {}) if "tensor_sha256" not in voice_meta and vp.exists(): try: voice_meta["tensor_sha256"] = _hash_file(str(vp)) except Exception: pass else: _voice_missed = voice_arg voices = sorted([d for d in _iter_voices() if (d / "voice.json").exists()], key=lambda p: p.stat().st_mtime, reverse=True) if voices: vm = _read_meta(voices[0] / "voice.json") if vm: voice_meta = dict(vm) if not voice_name: voice_name = _sanitize_note(vm.get("name") or voices[0].name) else: # zero flags: auto-pull most recent voice for full traceability voices = sorted([d for d in _iter_voices() if (d / "voice.json").exists()], key=lambda p: p.stat().st_mtime, reverse=True) if voices: vm = _read_meta(voices[0] / "voice.json") if vm: voice_meta = dict(vm) if not voice_name: voice_name = _sanitize_note(vm.get("name") or voices[0].name) except Exception: pass if _voice_missed: _warn(f" Voice '{_voice_missed}' not in registry — note saved without its provenance.") target_basename = _sanitize_note(Path(target_raw).name) if target_raw else None target_info = None if target_raw and Path(target_raw).expanduser().exists(): try: p = Path(target_raw).expanduser() if _is_gguf_file(str(p)): from gguf import GGUFReader r = GGUFReader(str(p)) arch = None for fld in r.fields.values(): if fld.name == "general.architecture": try: v = fld.parts[fld.data[0]] if isinstance(v, bytes): arch = v.decode() elif hasattr(v, "tobytes"): try: arch = v.tobytes().decode() except Exception: arch = str(v) else: arch = str(v) arch = arch.strip("\x00").strip() except Exception: pass break # find output tensor tname = None for pref in ("output.weight", "token_embd.weight"): for t in r.tensors: if t.name == pref: tname = t.name break if tname: break shape = None dtype = None if tname: for t in r.tensors: if t.name == tname: shape = list(reversed([int(x) for x in t.shape])) dtype = t.tensor_type.name break target_info = {"basename": target_basename, "format": "gguf", "tensor_name": tname, "architecture": arch, "family": _family(arch), "shape": shape, "dtype": dtype} else: # try safetensors probe (header only) try: hdr, tens = read_safetensors(str(p)) # pick output tensor tname = detect_output_tensor(list(tens.keys()), None) or largest_2d(list(tens.keys()), {k: v["shape"] for k, v in tens.items()}) info = tens.get(tname) if tname else None target_info = {"basename": target_basename, "format": "safetensors", "tensor_name": tname, "shape": list(info["shape"]) if info else None, "dtype": info["dtype"] if info else None} except Exception: target_info = {"basename": target_basename, "format": "safetensors"} except Exception: pass if target_info is None and target_basename: target_info = {"basename": target_basename} # env enrichment try: import numpy as _np np_ver = getattr(_np, "__version__", None) except Exception: np_ver = None try: import gguf as _gg gg_ver = getattr(_gg, "__version__", "installed") except Exception: gg_ver = None env_info = {"voice_ver": ctx.VERSION, "python": platform.python_version(), "platform": platform.platform(), "numpy": np_ver, "gguf": gg_ver} operation = getattr(args, "operation", None) or _last_operation() or "report" record = { "v": 2, "report_id": report_id, "ts": ts, "voice": voice_name, "voice_meta": voice_meta, "target": target_basename, "target_meta": target_info, "operation": operation, "status": "ok", "message": msg, "env": env_info, } # keep backward compat flat fields for old dataset readers if voice_meta: record["voice_source"] = _public_source(voice_meta.get("source_hf_model") or voice_meta.get("source")) record["voice_tensor"] = voice_meta.get("tensor_name") record["voice_dtype"] = voice_meta.get("dtype") record["voice_shape"] = voice_meta.get("shape") reporter = _sanitize_reporter(getattr(args, "reporter", None)) if reporter: record["reporter"] = reporter # auto-reporter from HF token if available (opt-in via --reporter not set) if "reporter" not in record: # do not auto-fill PII; keep null unless user opted pass # validate size try: rec_json = json.dumps(record, ensure_ascii=False) if len(rec_json) > 4096: _warn(" Note too large — truncating.") record["message"] = _sanitize_note((msg or "")[:240]) except Exception: pass _append_report(record) _ok(f"Thanks for noting — saved locally ({report_id[:8]}…) → {ctx.REPORTS_JSONL}") _say(f" Outbox: {ctx.REPORTS_OUTBOX / (report_id + '.json')} (private, stays until you choose to share)") # also handle --sync of pending (record already saved above) if getattr(args, "sync", False): _sync_pending() return # warm preview — what will be shared preview = None try: pv = voice_name or (voice_meta.get("name") if voice_meta else None) or "—" pt = (voice_meta.get("tensor_name") if voice_meta else None) or (target_info.get("tensor_name") if target_info else None) or "—" ps = voice_meta.get("shape") if voice_meta and voice_meta.get("shape") else (target_info.get("shape") if target_info else None) pd = voice_meta.get("dtype") if voice_meta and voice_meta.get("dtype") else (target_info.get("dtype") if target_info else None) pv_str = f"voice: {pv}, tensor: {pt}" if ps: pv_str += f" {ps}" if pd: pv_str += f" {pd}" if target_basename: pv_str += f", target: {target_basename}" if target_info and target_info.get("family"): pv_str += f" {target_info.get('family')}" preview = pv_str except Exception: preview = None should_push = _prompt_share(args, preview=preview) if should_push: _step("Sharing to community to help others…") ok = _push_report(record) if ok: _ok("Shared to community — thank you! Helps the next person.") pending = sorted(ctx.REPORTS_OUTBOX.glob("*.json")) if pending: recs = [] for p in pending: try: recs.append(json.loads(p.read_text())) except Exception: continue if recs: _step(f"Syncing {len(recs)} pending report(s)…") _push_reports_batch(recs) # headroom warning try: sz = ctx.REPORTS_JSONL.stat().st_size if ctx.REPORTS_JSONL.exists() else 0 if sz > 50 * 1024 * 1024: _warn(f" reports.jsonl {sz/1e6:.1f} MB — consider archiving (rotation to reports-YYYY-MM.jsonl).") if len(list(ctx.REPORTS_OUTBOX.glob("*.json"))) > 5000: _warn(" Outbox >5000 pending — run voice report --sync --yes to flush.") except Exception: pass else: _warn(" Share pending — will retry on next report or voice doctor.") else: _say(" Saved locally. Outbox retained for later push (rerun with --yes to share).")