Voice / vlib /reports.py
Wiself's picture
small tiny fixes
271bd15
Raw History Blame Contribute Delete
20.7 kB
"""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).")