"""Local Evidence browsing and explicit consolidation, independent of DWH access.""" import os from pathlib import Path from .canonical import parse_curated_markdown from .local_archive import LocalEvidenceArchive, _content class ConsolidationError(RuntimeError): def __init__(self, message, *, saved=False): super().__init__(message) self.saved = saved def consolidate_from_config(config: Path): from tht.cli.preprocess_cmd import run_from_config from tht.config import load_config from .authoring import EvidencePreparationError, migrate_workspace_evidence cfg = load_config(config) if not cfg.evidence or not cfg.evidence.local_archive_root: raise ConsolidationError("Local Evidence is not configured for this workspace") root = cfg.evidence.local_archive_root archive = LocalEvidenceArchive(root) if not (archive.metadata / "state.yaml").exists(): try: # Explicit first consolidation performs the one-time legacy conversion. if (archive.evidence / "manifest.yaml").is_file(): migrate_workspace_evidence(root) else: archive.initialize() except (ValueError, OSError, EvidencePreparationError) as error: raise ConsolidationError(str(error)) from error result = None def activate(snapshot): nonlocal result try: result = run_from_config(config, local_snapshot=snapshot) if result.status != "succeeded": raise ConsolidationError("Evidence indexing is blocked; check unit size and review items", saved=True) except ConsolidationError: raise except Exception as error: raise ConsolidationError("Evidence files were saved, but indexing failed. Retry consolidation.", saved=True) from error try: archive.consolidate(actor=os.environ.get("THT_PRINCIPAL_SUBJECT") or "installation operator", activate=activate) except (ValueError, OSError) as error: raise ConsolidationError(str(error)) from error return result def browse(root: Path, query: dict): """Read complete working units and their active status without opening an index.""" archive = LocalEvidenceArchive(root) with archive.operation(): state = archive._state() active = archive._snapshot(state["active"]) if state["active"] else None active_units = archive._units(archive._files(active), allow_review=True) if active else {} files = archive._files(archive.evidence) items, errors = [], [] seen = set() for relative, data in files.items(): try: unit = parse_curated_markdown(data.decode(), path=Path(relative)) if unit.id in seen: raise ValueError("Duplicate Evidence identifier") seen.add(unit.id) old = active_units.get(unit.id) status = "review_required" if unit.review_items else "legacy" if unit.schema_version != 4 \ else "active" if old and _content(old[1]) == _content(unit) and old[1].provenance == unit.provenance else "modified" if old else "new" items.append({**unit.model_dump(mode="json"), "file": relative, "status": status, "revision": _content(unit)}) except (ValueError, UnicodeError) as error: errors.append({"file": relative, "message": str(error)[:1500]}) for identity, (relative, unit) in active_units.items(): if identity not in seen: items.append({**unit.model_dump(mode="json"), "file": relative, "status": "invalid" if relative in files else "removed", "revision": _content(unit)}) def matches(item): for field in ("kind", "status", "language"): if query.get(field) and item[field] != query[field]: return False if query.get("purpose") and query["purpose"] not in item["purposes"]: return False for key, field in (("concept", "concepts"), ("table", "tables"), ("column", "columns")): if query.get(key) and not any(query[key].casefold() in v.casefold() for v in item["applies_to"][field]): return False provenance = item["provenance"] if query.get("source") and query["source"].casefold() not in str(provenance).casefold(): return False return not query.get("q") or query["q"].casefold() in str(item).casefold() selected = [item for item in items if matches(item)] selected.sort(key=lambda item: (str(item.get(query.get("sort", "title"), "")).casefold(), item["id"]), reverse=query.get("direction") == "desc") page, size = int(query.get("page", 1)), int(query.get("page_size", 25)) if page < 1 or not 1 <= size <= 100: raise ValueError("Invalid Evidence page") result = {"items": selected[(page-1)*size:page*size], "total": len(selected), "page": page, "page_size": size, "errors": errors, "active_revision": state["active"], "pending_revision": state["pending"], "initialized": bool(state.get("baseline") or state["pending"] or state["active"])} if query.get("id"): result["item"] = next((item for item in items if item["id"] == query["id"]), None) from .imports import reviews result["source_reviews"] = reviews(archive) return result def source_action(config, *, action, source_id=None, revision=None, decision=None, actor="installation operator"): from tht.config import load_config from .authoring import PiEvidenceRestructurer, authoring_skill_path, migrate_workspace_evidence from .imports import acquisition_sources, decide, refresh cfg = load_config(config) if not cfg.evidence or not cfg.evidence.local_archive_root: raise ValueError("Local Evidence is not configured") archive = LocalEvidenceArchive(cfg.evidence.local_archive_root) if not (archive.metadata / "state.yaml").exists(): if (archive.evidence / "manifest.yaml").is_file(): migrate_workspace_evidence(archive.root) else: (archive.evidence / "curated").mkdir(parents=True, exist_ok=True) archive.initialize() if action == "refresh": skill = authoring_skill_path() return refresh(archive, acquisition_sources(cfg), PiEvidenceRestructurer( os.environ.get("THT_PI_EXECUTABLE", "pi"), skill)) if action != "decide": raise ValueError("Unknown source action") def activate(snapshot): from tht.cli.preprocess_cmd import run_from_config try: result = run_from_config(config, local_snapshot=snapshot) if result.status != "succeeded": raise RuntimeError("Indexing did not succeed") except Exception as error: raise ConsolidationError("Source decision saved, but indexing failed. Retry the same decision.", saved=True) from error try: return decide(archive, source_id=source_id, revision=revision, decision=decision, actor=actor, activate=activate) except (ValueError, OSError) as error: from .imports import reviews if any(r["id"] == source_id and r["status"] == "applying" for r in reviews(archive)): raise ConsolidationError(f"Source decision saved. {str(error)[:1200]}. Retry the same decision.", saved=True) from error raise