mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-03 20:27:56 +00:00
v7next F1: domain D16 - the first real leaf transplant, proof-green
usage_accounting.py (sitting exactly at the 1600 hard cap upstream) gives up
its L-C2 extraction again: ouroboros/usage_legacy_import.py (270 lines, 5
symbols, 3 handle rewrites matching the oracle's declared set exactly) emitted
by scripts/v7next_transplant.py with the triple proof green on every symbol
(ast/tokens/bytes), and the facade shrinks 1600 -> 1384 with the oracle's EOF
re-export block and noqa discipline. The one residual diff vs the reference
leaf is a genuine upstream comment rewrite (e9bf6f14) - upstream bytes win;
the reference ledger row is byte-falsified as a copy source and recorded in
docs/v7next/LEDGER_CORRECTIONS.md (transform stays valid).
Pins transplanted with the campaign's reverse-mapping discipline: LEAVES /
LC2 owner tables reduced to the rows whose leaves exist on this tree (32
v7-only modules absent, noted with oracle-SHA keys); the settings-writer
exempt table is re-keyed to the leaf exactly as the oracle's own adaptation
did. size-ratchet manifest regenerated with the official tool: the new 1384
band entry carries a FLAGGED rationale (not self-approved), and the
regeneration also heals the stale giant entry the D15 pilot left behind
(the base was red on the size_ratchet lane - caught by this lane).
512/110/206/33 tests green in isolation at the lane base; ratchet lane 5
passed; independently re-verified.
This commit is contained in:
parent
b044bac582
commit
1a17218dc4
7 changed files with 488 additions and 231 deletions
|
|
@ -20,3 +20,22 @@ with evidence, found lane by lane. Applied to the campaign's carried ledger at F
|
|||
4. MIGRATION row 166 (retirement of 4 CLAUDE_CODE markers, id "none") — needs an
|
||||
explicit ADOPTION disposition (umbrella under D02 or its own row): zero
|
||||
production emitters of those markers exist at this tip (claim re-proven).
|
||||
|
||||
## From the D16 split pilot (base 5d3398c1, 2026-08-30)
|
||||
5. MIGRATION row 3911 (`usage_accounting.py::_legacy_snapshot` ->
|
||||
`usage_legacy_import.py::_legacy_snapshot`, "verbatim extraction") —
|
||||
BYTE-FALSIFIED as a copy source, transform still valid: upstream e9bf6f14
|
||||
rewrote the settings-hash comment inside the span (two lines "... prove
|
||||
non-mutation by hash, but never copy / their contents into the usage
|
||||
archive." became one line "... never copy contents."). The tool's --check of
|
||||
the reference leaf against tip bytes fails token-lockstep on exactly this
|
||||
span (ast=True, tokens=False); re-emitting from tip bytes is proof-green on
|
||||
the first round with the reference declared set {_legacy_snapshot, _locked,
|
||||
_read_records_locked} unchanged. Copying the reference leaf verbatim would
|
||||
have silently reverted an upstream comment edit.
|
||||
6. MIGRATION rows 3910-3914 status "pending upstream transfer" — RE-CONFIRMED
|
||||
at this tip (contrast with the D15 project_facts case, entry 1 above):
|
||||
upstream still carries the unsplit legacy import inside
|
||||
ouroboros/usage_accounting.py (1600 lines, exactly at the hard cap;
|
||||
IMPORT_REL at :60, the four defs at :1374-:1600). The extraction was
|
||||
performed by this lane from tip bytes.
|
||||
|
|
|
|||
|
|
@ -26,7 +26,6 @@ GIANT_PATHS = (
|
|||
"tests/test_delegated_subagent_transport.py",
|
||||
"tests/test_delivery_forced_finalization.py",
|
||||
"tests/test_devtools_benchmarks.py",
|
||||
"tests/test_evolution_state_integrity_v3.py",
|
||||
"tests/test_extension_loader.py",
|
||||
"tests/test_extensions_api.py",
|
||||
"tests/test_git_ops_recovery.py",
|
||||
|
|
@ -166,6 +165,7 @@ BAND_PATHS = {
|
|||
"ouroboros/tools/review_context_atlas.py": "Grew INTO the band by the #284 pack-arithmetic fixes: measured render charged at admission, exact per-row costs, target capped at the hard rail, honest eviction diagnostics \u2014 all in the module that owns the arithmetic.",
|
||||
"ouroboros/tools/skill_exec.py": None,
|
||||
"ouroboros/tools/skill_publish.py": "Entered the band from 952 lines: publish now writes the OuroborosHub publication receipt at pr_opened through the shared locked-update seam and maps the receipt from the validated serialized form (hubflow sprint, receipt-as-only-stored-fact design).",
|
||||
"ouroboros/usage_accounting.py": "Entered the band from the 1501-1600 zone (1600 lines) by the v7 L-C2 extraction of the one-time legacy usage import into ouroboros/usage_legacy_import.py; shrink-only residue of the split, not new growth.",
|
||||
"ouroboros/utils.py": None,
|
||||
"ouroboros/workspace_executor.py": None,
|
||||
"scripts/run_external_review.py": None,
|
||||
|
|
|
|||
|
|
@ -46,7 +46,7 @@ from ouroboros.usage_ledger import ( # noqa: F401 — re-exported substrate
|
|||
_validate_records,
|
||||
_write_bytes_atomic_fsync,
|
||||
)
|
||||
from ouroboros.utils import append_jsonl, atomic_write_json, utc_now_iso
|
||||
from ouroboros.utils import append_jsonl, atomic_write_json, utc_now_iso # noqa: F401 -- the accounting module keeps its historical import surface for the L-C2 leaf
|
||||
from ouroboros._usage_rows import ( # noqa: F401 (re-exported substrate vocabulary)
|
||||
REVIEW_ATTRIBUTION_KEYS,
|
||||
_breakdown_bucket,
|
||||
|
|
@ -57,7 +57,6 @@ from ouroboros._usage_rows import ( # noqa: F401 (re-exported substrate vocabu
|
|||
)
|
||||
from ouroboros.skill_review_usage import skill_review_usage
|
||||
log = logging.getLogger(__name__)
|
||||
IMPORT_REL = pathlib.Path("state/usage_import_watermark.json")
|
||||
__all__ = (
|
||||
"AttemptRequest", "AttemptReservation", "BudgetExceeded", "PhysicalAttemptCapture",
|
||||
"PhysicalAttemptContext", "PhysicalAttemptLimitExceeded", "PhysicalAttemptPreconditionFailed",
|
||||
|
|
@ -1371,230 +1370,15 @@ async def execute_physical_attempt_async(
|
|||
return response
|
||||
|
||||
|
||||
def _legacy_snapshot(root: pathlib.Path) -> Tuple[list[Dict[str, Any]], Dict[str, Any], Dict[str, str]]:
|
||||
events_path = root / "logs" / "events.jsonl"
|
||||
state_path = root / "state" / "state.json"
|
||||
settings_path = pathlib.Path(os.environ.get("OUROBOROS_SETTINGS_PATH") or root / "settings.json")
|
||||
sources = {"events.jsonl": events_path, "state.json": state_path}
|
||||
snapshots: Dict[str, bytes] = {}
|
||||
for name, path in sources.items():
|
||||
try:
|
||||
snapshots[name] = path.read_bytes()
|
||||
except FileNotFoundError:
|
||||
continue
|
||||
except OSError as exc:
|
||||
raise UsageAccountingError(f"cannot snapshot legacy usage source {path}: {exc}") from exc
|
||||
hashes = {name: hashlib.sha256(snapshots[name]).hexdigest() if name in snapshots else "" for name in sources}
|
||||
# Settings are owner-secret state: prove non-mutation by hash, never copy contents.
|
||||
try:
|
||||
hashes["settings.json"] = hashlib.sha256(settings_path.read_bytes()).hexdigest()
|
||||
except FileNotFoundError:
|
||||
hashes["settings.json"] = ""
|
||||
except OSError as exc:
|
||||
raise UsageAccountingError(f"cannot hash settings file {settings_path}: {exc}") from exc
|
||||
rows: list[Dict[str, Any]] = []
|
||||
try:
|
||||
event_text = snapshots.get("events.jsonl", b"").decode("utf-8")
|
||||
for line_no, line in enumerate(event_text.splitlines(), 1):
|
||||
try:
|
||||
value = json.loads(line)
|
||||
except json.JSONDecodeError:
|
||||
continue
|
||||
if isinstance(value, dict) and value.get("type") == "llm_usage":
|
||||
rows.append({**value, "_legacy_line": line_no})
|
||||
except UnicodeDecodeError:
|
||||
pass
|
||||
try:
|
||||
state = json.loads(snapshots.get("state.json", b"{}").decode("utf-8"))
|
||||
if not isinstance(state, dict):
|
||||
state = {}
|
||||
except (UnicodeDecodeError, json.JSONDecodeError):
|
||||
state = {}
|
||||
|
||||
combined = hashlib.sha256(json.dumps(hashes, sort_keys=True).encode("utf-8")).hexdigest()[:16]
|
||||
archive = root / "archive" / "usage_import" / combined
|
||||
archive.mkdir(parents=True, exist_ok=True)
|
||||
for name, payload in snapshots.items():
|
||||
target = archive / name
|
||||
if target.exists():
|
||||
if target.read_bytes() != payload:
|
||||
raise UsageAccountingError(f"legacy usage archive mismatch: {target}")
|
||||
else:
|
||||
_write_bytes_atomic_fsync(target, payload)
|
||||
try:
|
||||
target.chmod(0o400)
|
||||
except OSError:
|
||||
pass
|
||||
atomic_write_json(archive / "sha256.json", hashes, trailing_newline=True, fsync=True)
|
||||
return rows, state, hashes
|
||||
|
||||
|
||||
def ensure_legacy_imported(
|
||||
drive_root: Optional[pathlib.Path] = None,
|
||||
) -> Dict[str, Any]:
|
||||
"""One resumable import of legacy usage telemetry and the state cost delta."""
|
||||
root = _drive_root(drive_root)
|
||||
completed = _completed_import_watermark(root)
|
||||
if completed is not None:
|
||||
return completed
|
||||
# Separate from the hot budget lock: source snapshot/archive may do I/O,
|
||||
# while concurrent startup importers still serialize on one generation.
|
||||
with _named_lock(root, "usage_import.lock", timeout_sec=60.0, stale_sec=600.0):
|
||||
return _ensure_legacy_imported_locked(root)
|
||||
|
||||
|
||||
def _completed_import_watermark(root: pathlib.Path) -> Optional[Dict[str, Any]]:
|
||||
try:
|
||||
value = json.loads((root / IMPORT_REL).read_text(encoding="utf-8"))
|
||||
except (FileNotFoundError, OSError, UnicodeDecodeError, json.JSONDecodeError):
|
||||
return None
|
||||
return value if isinstance(value, dict) and value.get("completed") else None
|
||||
|
||||
|
||||
def _ensure_legacy_imported_locked(
|
||||
root: pathlib.Path,
|
||||
) -> Dict[str, Any]:
|
||||
watermark = root / IMPORT_REL
|
||||
existing = _completed_import_watermark(root)
|
||||
if existing is not None:
|
||||
return existing
|
||||
|
||||
legacy_rows, state, hashes = _legacy_snapshot(root)
|
||||
baseline_source = "state.json"
|
||||
candidates: list[Dict[str, Any]] = []
|
||||
seen_fingerprints: set[str] = set()
|
||||
imported_cost = 0.0
|
||||
usage_count = 0
|
||||
for event in legacy_rows:
|
||||
line_no = int(event.pop("_legacy_line", 0) or 0)
|
||||
legacy_usage = event.get("usage") if isinstance(event.get("usage"), dict) else {}
|
||||
fingerprint = hashlib.sha256(
|
||||
json.dumps(event, ensure_ascii=False, sort_keys=True, separators=(",", ":"), default=str).encode("utf-8")
|
||||
).hexdigest()
|
||||
if fingerprint in seen_fingerprints:
|
||||
continue
|
||||
seen_fingerprints.add(fingerprint)
|
||||
task_id = str(event.get("task_id") or "")
|
||||
root_task_id = str(event.get("root_task_id") or task_id)
|
||||
raw_cost = event.get("cost")
|
||||
if raw_cost is None:
|
||||
raw_cost = legacy_usage.get("cost", legacy_usage.get("total_cost"))
|
||||
cost = _number(raw_cost)
|
||||
|
||||
def legacy_int(field: str, *aliases: str) -> int:
|
||||
for candidate in (field, *aliases):
|
||||
value = event.get(candidate)
|
||||
if value in (None, ""):
|
||||
value = legacy_usage.get(candidate)
|
||||
try:
|
||||
return max(0, int(float(value or 0)))
|
||||
except (TypeError, ValueError):
|
||||
continue
|
||||
return 0
|
||||
|
||||
prompt = legacy_int("prompt_tokens", "input_tokens")
|
||||
completion = legacy_int("completion_tokens", "output_tokens")
|
||||
provider = str(event.get("provider") or event.get("api_key_type") or "unknown")
|
||||
if cost == 0 and (prompt or completion) and provider != "local":
|
||||
cost = None # legacy zero may mean unknown pricing, never "free"
|
||||
usage_count += 1
|
||||
if cost is not None:
|
||||
imported_cost += cost
|
||||
candidates.append(
|
||||
{
|
||||
"kind": "legacy_usage",
|
||||
"attempt_id": f"legacy-{fingerprint[:24]}",
|
||||
"state": "settled",
|
||||
"model": str(event.get("model") or ""),
|
||||
"provider": provider,
|
||||
"cost_usd": cost,
|
||||
"cost_final": bool(cost is not None and not event.get("cost_estimated")),
|
||||
"reservation_upper_bound_usd": None,
|
||||
"prompt_tokens": prompt,
|
||||
"completion_tokens": completion,
|
||||
"cached_tokens": legacy_int("cached_tokens", "cache_read_input_tokens"),
|
||||
"cache_write_tokens": legacy_int("cache_write_tokens", "cache_creation_input_tokens"),
|
||||
"prompt_cache_ttl": str(
|
||||
event.get("prompt_cache_ttl") or legacy_usage.get("prompt_cache_ttl") or ""
|
||||
),
|
||||
"task_id": task_id,
|
||||
"root_task_id": root_task_id,
|
||||
"parent_task_id": str(event.get("parent_task_id") or ""),
|
||||
"category": str(event.get("category") or "legacy"),
|
||||
"source": "legacy_llm_usage",
|
||||
"legacy_line": line_no,
|
||||
}
|
||||
)
|
||||
legacy_calls = max(0, int(state.get("spent_calls") or state.get("calls") or 0))
|
||||
metadata_count = max(0, legacy_calls - usage_count)
|
||||
if metadata_count:
|
||||
identity = hashlib.sha256(
|
||||
f"legacy-metadata:{metadata_count}:{hashes.get('state.json', '')}".encode()
|
||||
).hexdigest()
|
||||
candidates.append(
|
||||
{
|
||||
"kind": "legacy_metadata",
|
||||
"attempt_id": f"legacy-{identity[:24]}",
|
||||
"state": "unresolved",
|
||||
"model": "",
|
||||
"provider": "legacy",
|
||||
"reservation_upper_bound_usd": None,
|
||||
"ambiguous_call_count": metadata_count,
|
||||
"task_id": "",
|
||||
"root_task_id": "",
|
||||
"parent_task_id": "",
|
||||
"category": "legacy",
|
||||
"source": "legacy_state_call_delta",
|
||||
}
|
||||
)
|
||||
state_spent = _number(state.get("spent_usd")) or 0.0
|
||||
delta = round(max(0.0, state_spent - imported_cost), 6)
|
||||
if delta:
|
||||
identity = hashlib.sha256(f"legacy-delta:{delta:.6f}:{hashes.get('state.json', '')}".encode()).hexdigest()
|
||||
candidates.append(
|
||||
{
|
||||
"kind": "legacy_delta",
|
||||
"attempt_id": f"legacy-{identity[:24]}",
|
||||
"state": "settled",
|
||||
"model": "",
|
||||
"provider": "legacy",
|
||||
"cost_usd": delta,
|
||||
"cost_final": False,
|
||||
"reservation_upper_bound_usd": None,
|
||||
"task_id": "",
|
||||
"root_task_id": "",
|
||||
"parent_task_id": "",
|
||||
"category": "legacy",
|
||||
"source": "legacy_state_delta",
|
||||
}
|
||||
)
|
||||
|
||||
with _locked(root):
|
||||
current_watermark = _completed_import_watermark(root)
|
||||
if current_watermark is not None:
|
||||
return current_watermark
|
||||
records = _read_records_locked(root)
|
||||
existing_ids = {str(row.get("attempt_id") or "") for row in records}
|
||||
missing = [row for row in candidates if row["attempt_id"] not in existing_ids]
|
||||
_append_rows_locked(root, records, missing)
|
||||
result = {
|
||||
"completed": True,
|
||||
"completed_at": utc_now_iso(),
|
||||
"source_sha256": hashes,
|
||||
"legacy_baseline_source": baseline_source,
|
||||
"legacy_baseline_spent_usd": state_spent,
|
||||
"legacy_baseline_spent_calls": legacy_calls,
|
||||
"legacy_usage_count": usage_count,
|
||||
"legacy_metadata_count": metadata_count,
|
||||
"legacy_delta_usd": delta,
|
||||
# The legacy schema has no trustworthy typed test/operator bit.
|
||||
# Never invent exclusions from names, task ids, or source strings.
|
||||
"quarantined_test_operator_rows": 0,
|
||||
"test_operator_quarantine_policy": "typed_evidence_only_no_inference",
|
||||
"events_exceed_state_calls": max(0, usage_count - legacy_calls),
|
||||
"events_exceed_state_usd": round(max(0.0, imported_cost - state_spent), 6),
|
||||
"rows_appended": len(missing),
|
||||
}
|
||||
atomic_write_json(watermark, result, trailing_newline=True, fsync=True)
|
||||
append_jsonl(root / "logs" / "events.jsonl", {"type": "usage_import_completed", **result})
|
||||
return result
|
||||
# v7 L-C2 split: the one-time legacy usage-telemetry import (source snapshot and
|
||||
# archive, candidate rows, state-baseline reconciliation, completed watermark)
|
||||
# lives in ouroboros/usage_legacy_import.py. Re-exported under the historical
|
||||
# names so callers and monkeypatching tests keep working unchanged (facade
|
||||
# identity pinned in tests/test_lc2_owner_facades.py).
|
||||
from ouroboros.usage_legacy_import import ( # noqa: E402, F401 -- intentional public re-exports
|
||||
IMPORT_REL,
|
||||
_completed_import_watermark,
|
||||
_ensure_legacy_imported_locked,
|
||||
_legacy_snapshot,
|
||||
ensure_legacy_imported,
|
||||
)
|
||||
|
|
|
|||
270
ouroboros/usage_legacy_import.py
Normal file
270
ouroboros/usage_legacy_import.py
Normal file
|
|
@ -0,0 +1,270 @@
|
|||
"""The one-time legacy usage-telemetry import (v7 L-C2 split).
|
||||
|
||||
The resumable import of pre-ledger usage telemetry (``llm_usage`` events and
|
||||
the ``state.json`` cost baseline) into the append-only attempt ledger: source
|
||||
snapshot and read-only archive, deduplicated candidate rows, the metadata/delta
|
||||
reconciliation against the state baseline, and the completed-import watermark.
|
||||
Extracted from usage_accounting.py; usage_accounting re-exports every name, so
|
||||
historical imports and monkeypatch targets keep working."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import hashlib
|
||||
import json
|
||||
import os
|
||||
import pathlib
|
||||
from typing import Any, Dict, Optional, Tuple
|
||||
|
||||
from ouroboros.usage_ledger import (
|
||||
UsageAccountingError,
|
||||
_append_rows_locked,
|
||||
_drive_root,
|
||||
_named_lock,
|
||||
_number,
|
||||
_write_bytes_atomic_fsync,
|
||||
)
|
||||
from ouroboros.utils import append_jsonl, atomic_write_json, utc_now_iso
|
||||
|
||||
IMPORT_REL = pathlib.Path("state/usage_import_watermark.json")
|
||||
|
||||
|
||||
def _usage():
|
||||
"""The parent usage_accounting module, read at call time.
|
||||
|
||||
The accounting members stay monkeypatch-addressable at their historical
|
||||
``ouroboros.usage_accounting`` bindings (tests rebind them there), so this
|
||||
leaf resolves every cross-reference through the module at each call instead
|
||||
of freezing whatever object a from-import saw at import time.
|
||||
"""
|
||||
from ouroboros import usage_accounting
|
||||
|
||||
return usage_accounting
|
||||
|
||||
|
||||
def _legacy_snapshot(root: pathlib.Path) -> Tuple[list[Dict[str, Any]], Dict[str, Any], Dict[str, str]]:
|
||||
events_path = root / "logs" / "events.jsonl"
|
||||
state_path = root / "state" / "state.json"
|
||||
settings_path = pathlib.Path(os.environ.get("OUROBOROS_SETTINGS_PATH") or root / "settings.json")
|
||||
sources = {"events.jsonl": events_path, "state.json": state_path}
|
||||
snapshots: Dict[str, bytes] = {}
|
||||
for name, path in sources.items():
|
||||
try:
|
||||
snapshots[name] = path.read_bytes()
|
||||
except FileNotFoundError:
|
||||
continue
|
||||
except OSError as exc:
|
||||
raise UsageAccountingError(f"cannot snapshot legacy usage source {path}: {exc}") from exc
|
||||
hashes = {name: hashlib.sha256(snapshots[name]).hexdigest() if name in snapshots else "" for name in sources}
|
||||
# Settings are owner-secret state: prove non-mutation by hash, never copy contents.
|
||||
try:
|
||||
hashes["settings.json"] = hashlib.sha256(settings_path.read_bytes()).hexdigest()
|
||||
except FileNotFoundError:
|
||||
hashes["settings.json"] = ""
|
||||
except OSError as exc:
|
||||
raise UsageAccountingError(f"cannot hash settings file {settings_path}: {exc}") from exc
|
||||
rows: list[Dict[str, Any]] = []
|
||||
try:
|
||||
event_text = snapshots.get("events.jsonl", b"").decode("utf-8")
|
||||
for line_no, line in enumerate(event_text.splitlines(), 1):
|
||||
try:
|
||||
value = json.loads(line)
|
||||
except json.JSONDecodeError:
|
||||
continue
|
||||
if isinstance(value, dict) and value.get("type") == "llm_usage":
|
||||
rows.append({**value, "_legacy_line": line_no})
|
||||
except UnicodeDecodeError:
|
||||
pass
|
||||
try:
|
||||
state = json.loads(snapshots.get("state.json", b"{}").decode("utf-8"))
|
||||
if not isinstance(state, dict):
|
||||
state = {}
|
||||
except (UnicodeDecodeError, json.JSONDecodeError):
|
||||
state = {}
|
||||
|
||||
combined = hashlib.sha256(json.dumps(hashes, sort_keys=True).encode("utf-8")).hexdigest()[:16]
|
||||
archive = root / "archive" / "usage_import" / combined
|
||||
archive.mkdir(parents=True, exist_ok=True)
|
||||
for name, payload in snapshots.items():
|
||||
target = archive / name
|
||||
if target.exists():
|
||||
if target.read_bytes() != payload:
|
||||
raise UsageAccountingError(f"legacy usage archive mismatch: {target}")
|
||||
else:
|
||||
_write_bytes_atomic_fsync(target, payload)
|
||||
try:
|
||||
target.chmod(0o400)
|
||||
except OSError:
|
||||
pass
|
||||
atomic_write_json(archive / "sha256.json", hashes, trailing_newline=True, fsync=True)
|
||||
return rows, state, hashes
|
||||
|
||||
|
||||
def ensure_legacy_imported(
|
||||
drive_root: Optional[pathlib.Path] = None,
|
||||
) -> Dict[str, Any]:
|
||||
"""One resumable import of legacy usage telemetry and the state cost delta."""
|
||||
root = _drive_root(drive_root)
|
||||
completed = _completed_import_watermark(root)
|
||||
if completed is not None:
|
||||
return completed
|
||||
# Separate from the hot budget lock: source snapshot/archive may do I/O,
|
||||
# while concurrent startup importers still serialize on one generation.
|
||||
with _named_lock(root, "usage_import.lock", timeout_sec=60.0, stale_sec=600.0):
|
||||
return _ensure_legacy_imported_locked(root)
|
||||
|
||||
|
||||
def _completed_import_watermark(root: pathlib.Path) -> Optional[Dict[str, Any]]:
|
||||
try:
|
||||
value = json.loads((root / IMPORT_REL).read_text(encoding="utf-8"))
|
||||
except (FileNotFoundError, OSError, UnicodeDecodeError, json.JSONDecodeError):
|
||||
return None
|
||||
return value if isinstance(value, dict) and value.get("completed") else None
|
||||
|
||||
|
||||
def _ensure_legacy_imported_locked(
|
||||
root: pathlib.Path,
|
||||
) -> Dict[str, Any]:
|
||||
watermark = root / IMPORT_REL
|
||||
existing = _completed_import_watermark(root)
|
||||
if existing is not None:
|
||||
return existing
|
||||
|
||||
legacy_rows, state, hashes = _usage()._legacy_snapshot(root)
|
||||
baseline_source = "state.json"
|
||||
candidates: list[Dict[str, Any]] = []
|
||||
seen_fingerprints: set[str] = set()
|
||||
imported_cost = 0.0
|
||||
usage_count = 0
|
||||
for event in legacy_rows:
|
||||
line_no = int(event.pop("_legacy_line", 0) or 0)
|
||||
legacy_usage = event.get("usage") if isinstance(event.get("usage"), dict) else {}
|
||||
fingerprint = hashlib.sha256(
|
||||
json.dumps(event, ensure_ascii=False, sort_keys=True, separators=(",", ":"), default=str).encode("utf-8")
|
||||
).hexdigest()
|
||||
if fingerprint in seen_fingerprints:
|
||||
continue
|
||||
seen_fingerprints.add(fingerprint)
|
||||
task_id = str(event.get("task_id") or "")
|
||||
root_task_id = str(event.get("root_task_id") or task_id)
|
||||
raw_cost = event.get("cost")
|
||||
if raw_cost is None:
|
||||
raw_cost = legacy_usage.get("cost", legacy_usage.get("total_cost"))
|
||||
cost = _number(raw_cost)
|
||||
|
||||
def legacy_int(field: str, *aliases: str) -> int:
|
||||
for candidate in (field, *aliases):
|
||||
value = event.get(candidate)
|
||||
if value in (None, ""):
|
||||
value = legacy_usage.get(candidate)
|
||||
try:
|
||||
return max(0, int(float(value or 0)))
|
||||
except (TypeError, ValueError):
|
||||
continue
|
||||
return 0
|
||||
|
||||
prompt = legacy_int("prompt_tokens", "input_tokens")
|
||||
completion = legacy_int("completion_tokens", "output_tokens")
|
||||
provider = str(event.get("provider") or event.get("api_key_type") or "unknown")
|
||||
if cost == 0 and (prompt or completion) and provider != "local":
|
||||
cost = None # legacy zero may mean unknown pricing, never "free"
|
||||
usage_count += 1
|
||||
if cost is not None:
|
||||
imported_cost += cost
|
||||
candidates.append(
|
||||
{
|
||||
"kind": "legacy_usage",
|
||||
"attempt_id": f"legacy-{fingerprint[:24]}",
|
||||
"state": "settled",
|
||||
"model": str(event.get("model") or ""),
|
||||
"provider": provider,
|
||||
"cost_usd": cost,
|
||||
"cost_final": bool(cost is not None and not event.get("cost_estimated")),
|
||||
"reservation_upper_bound_usd": None,
|
||||
"prompt_tokens": prompt,
|
||||
"completion_tokens": completion,
|
||||
"cached_tokens": legacy_int("cached_tokens", "cache_read_input_tokens"),
|
||||
"cache_write_tokens": legacy_int("cache_write_tokens", "cache_creation_input_tokens"),
|
||||
"prompt_cache_ttl": str(
|
||||
event.get("prompt_cache_ttl") or legacy_usage.get("prompt_cache_ttl") or ""
|
||||
),
|
||||
"task_id": task_id,
|
||||
"root_task_id": root_task_id,
|
||||
"parent_task_id": str(event.get("parent_task_id") or ""),
|
||||
"category": str(event.get("category") or "legacy"),
|
||||
"source": "legacy_llm_usage",
|
||||
"legacy_line": line_no,
|
||||
}
|
||||
)
|
||||
legacy_calls = max(0, int(state.get("spent_calls") or state.get("calls") or 0))
|
||||
metadata_count = max(0, legacy_calls - usage_count)
|
||||
if metadata_count:
|
||||
identity = hashlib.sha256(
|
||||
f"legacy-metadata:{metadata_count}:{hashes.get('state.json', '')}".encode()
|
||||
).hexdigest()
|
||||
candidates.append(
|
||||
{
|
||||
"kind": "legacy_metadata",
|
||||
"attempt_id": f"legacy-{identity[:24]}",
|
||||
"state": "unresolved",
|
||||
"model": "",
|
||||
"provider": "legacy",
|
||||
"reservation_upper_bound_usd": None,
|
||||
"ambiguous_call_count": metadata_count,
|
||||
"task_id": "",
|
||||
"root_task_id": "",
|
||||
"parent_task_id": "",
|
||||
"category": "legacy",
|
||||
"source": "legacy_state_call_delta",
|
||||
}
|
||||
)
|
||||
state_spent = _number(state.get("spent_usd")) or 0.0
|
||||
delta = round(max(0.0, state_spent - imported_cost), 6)
|
||||
if delta:
|
||||
identity = hashlib.sha256(f"legacy-delta:{delta:.6f}:{hashes.get('state.json', '')}".encode()).hexdigest()
|
||||
candidates.append(
|
||||
{
|
||||
"kind": "legacy_delta",
|
||||
"attempt_id": f"legacy-{identity[:24]}",
|
||||
"state": "settled",
|
||||
"model": "",
|
||||
"provider": "legacy",
|
||||
"cost_usd": delta,
|
||||
"cost_final": False,
|
||||
"reservation_upper_bound_usd": None,
|
||||
"task_id": "",
|
||||
"root_task_id": "",
|
||||
"parent_task_id": "",
|
||||
"category": "legacy",
|
||||
"source": "legacy_state_delta",
|
||||
}
|
||||
)
|
||||
|
||||
with _usage()._locked(root):
|
||||
current_watermark = _completed_import_watermark(root)
|
||||
if current_watermark is not None:
|
||||
return current_watermark
|
||||
records = _usage()._read_records_locked(root)
|
||||
existing_ids = {str(row.get("attempt_id") or "") for row in records}
|
||||
missing = [row for row in candidates if row["attempt_id"] not in existing_ids]
|
||||
_append_rows_locked(root, records, missing)
|
||||
result = {
|
||||
"completed": True,
|
||||
"completed_at": utc_now_iso(),
|
||||
"source_sha256": hashes,
|
||||
"legacy_baseline_source": baseline_source,
|
||||
"legacy_baseline_spent_usd": state_spent,
|
||||
"legacy_baseline_spent_calls": legacy_calls,
|
||||
"legacy_usage_count": usage_count,
|
||||
"legacy_metadata_count": metadata_count,
|
||||
"legacy_delta_usd": delta,
|
||||
# The legacy schema has no trustworthy typed test/operator bit.
|
||||
# Never invent exclusions from names, task ids, or source strings.
|
||||
"quarantined_test_operator_rows": 0,
|
||||
"test_operator_quarantine_policy": "typed_evidence_only_no_inference",
|
||||
"events_exceed_state_calls": max(0, usage_count - legacy_calls),
|
||||
"events_exceed_state_usd": round(max(0.0, imported_cost - state_spent), 6),
|
||||
"rows_appended": len(missing),
|
||||
}
|
||||
atomic_write_json(watermark, result, trailing_newline=True, fsync=True)
|
||||
append_jsonl(root / "logs" / "events.jsonl", {"type": "usage_import_completed", **result})
|
||||
return result
|
||||
49
tests/test_lc2_owner_facades.py
Normal file
49
tests/test_lc2_owner_facades.py
Normal file
|
|
@ -0,0 +1,49 @@
|
|||
"""Facade-identity contract for the v7 L-C2 leaf owners.
|
||||
|
||||
Every member the L-C2 split moved out of ``ouroboros/agent.py``,
|
||||
``ouroboros/agent_task_pipeline.py`` and ``ouroboros/usage_accounting.py``
|
||||
keeps a parent re-export under its historical name, so existing callers and
|
||||
monkeypatching tests keep working unchanged — the parent binding IS the leaf's
|
||||
object, the same way the queue, loop and update_merge splits pin their leaves.
|
||||
The hot-code parity clause pins the update_merge direction of the rule: none of
|
||||
these parents is a HOT_CODE_PATHS member, so a leaf that merely moved code out
|
||||
of one must not silently acquire the label either.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import importlib
|
||||
|
||||
# leaf module -> (parent module, every member the leaf owns; the parent
|
||||
# re-exports each name).
|
||||
# v7next transplant note: the reference table (ouroboros_v7_wip @ 9f691656)
|
||||
# also carries the ouroboros.agent_dispatch and ouroboros.post_task_synthesis
|
||||
# rows; those extractions belong to other domains and land with their lanes —
|
||||
# only the D16 usage split exists on this integration branch so far.
|
||||
LC2_LEAF_OWNERS: dict[str, tuple[str, str]] = {
|
||||
"ouroboros.usage_legacy_import": (
|
||||
"ouroboros.usage_accounting",
|
||||
"IMPORT_REL _legacy_snapshot ensure_legacy_imported _completed_import_watermark "
|
||||
"_ensure_legacy_imported_locked"
|
||||
),
|
||||
}
|
||||
|
||||
|
||||
def test_lc2_owner_facades_preserve_identity():
|
||||
for leaf, (parent, names) in LC2_LEAF_OWNERS.items():
|
||||
parent_module = importlib.import_module(parent)
|
||||
leaf_module = importlib.import_module(leaf)
|
||||
for name in names.split():
|
||||
assert getattr(parent_module, name) is getattr(leaf_module, name), f"{leaf}.{name}"
|
||||
|
||||
|
||||
def test_lc2_leaves_keep_hot_code_label_parity():
|
||||
"""Managed-update conflict labelling names none of the L-C2 parents; the
|
||||
split must not silently upgrade or downgrade the label for code that merely
|
||||
moved — parent and leaves carry the SAME membership."""
|
||||
from supervisor.update_merge_policy import HOT_CODE_PATHS
|
||||
|
||||
for leaf, (parent, _names) in LC2_LEAF_OWNERS.items():
|
||||
parent_path = parent.replace(".", "/") + ".py"
|
||||
leaf_path = leaf.replace(".", "/") + ".py"
|
||||
assert (leaf_path in HOT_CODE_PATHS) == (parent_path in HOT_CODE_PATHS), leaf
|
||||
135
tests/test_module_handle_extraction.py
Normal file
135
tests/test_module_handle_extraction.py
Normal file
|
|
@ -0,0 +1,135 @@
|
|||
"""The module-handle extraction (spec §1.9 batch №8, delta D18) and its invariants.
|
||||
|
||||
`supervisor/queue.py` and `supervisor/workers.py` could not be split the way every
|
||||
other v7 module was. Their bodies read module globals that ``init`` /
|
||||
``init_queue_refs`` REBIND — PENDING, RUNNING, DRIVE_ROOT, WORKERS and the rest —
|
||||
so a leaf holding `from supervisor.queue import PENDING` would freeze the object it
|
||||
saw at import time, and a leaf keeping its own copy would be a second answer to the
|
||||
same question (67 test sites rebind these names on the parent and must keep
|
||||
working). The owner approved ONE mechanical exception: a declared parent name X is
|
||||
read as ``_queue().X`` / ``_pool().X`` — a function-local import of the parent — so
|
||||
the binding is resolved at call time.
|
||||
|
||||
The one-time proof that each moved body is otherwise unchanged (AST-equal modulo
|
||||
exactly that substitution, over the declared set, with zero other differences) is
|
||||
recorded in the extraction commits. What is pinned HERE is the property that has to
|
||||
survive every later edit:
|
||||
|
||||
* the parent is reached only through a call-time handle, never a top-level import;
|
||||
* every declared name is really bound by the parent (a typo would silently match
|
||||
nothing and make the proof vacuous);
|
||||
* the declared set is exactly the set the leaf actually reads through the handle —
|
||||
neither a stale name nor an undeclared one;
|
||||
* and, the load-bearing one, NO leaf reads a parent-owned name directly. That is
|
||||
the bug class the handle exists to prevent, and it is the one a later "tidy-up"
|
||||
would reintroduce by adding an innocent-looking from-import.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import ast
|
||||
import pathlib
|
||||
|
||||
import pytest
|
||||
|
||||
REPO = pathlib.Path(__file__).resolve().parents[1]
|
||||
|
||||
# leaf -> (parent, handle, declared substitution set)
|
||||
# v7next transplant note: the reference table (ouroboros_v7_wip @ 9f691656)
|
||||
# carries one row per extracted leaf across every domain, plus the
|
||||
# queue/worker/loop/git_ops-specific invariant tests below the parametrized
|
||||
# trio. On this integration branch each row lands with the lane that
|
||||
# transplants its domain; only the D16 L-C2 usage split exists so far. The
|
||||
# domain-specific standalone tests travel with their own rows.
|
||||
LEAVES: dict[str, tuple[str, str, frozenset[str]]] = {
|
||||
"ouroboros/usage_legacy_import.py": ("ouroboros/usage_accounting.py", "_usage", frozenset({
|
||||
"_legacy_snapshot", "_locked", "_read_records_locked",
|
||||
})),
|
||||
}
|
||||
|
||||
|
||||
def _tree(rel: str) -> ast.Module:
|
||||
return ast.parse((REPO / rel).read_text(encoding="utf-8"))
|
||||
|
||||
|
||||
def _module_bindings(tree: ast.Module) -> set[str]:
|
||||
bound: set[str] = set()
|
||||
for node in tree.body:
|
||||
if isinstance(node, (ast.FunctionDef, ast.AsyncFunctionDef, ast.ClassDef)):
|
||||
bound.add(node.name)
|
||||
elif isinstance(node, ast.Assign):
|
||||
bound.update(t.id for t in node.targets if isinstance(t, ast.Name))
|
||||
elif isinstance(node, ast.AnnAssign) and isinstance(node.target, ast.Name):
|
||||
bound.add(node.target.id)
|
||||
elif isinstance(node, (ast.Import, ast.ImportFrom)):
|
||||
bound.update(a.asname or a.name.split(".")[0] for a in node.names)
|
||||
elif (isinstance(node, ast.If) and isinstance(node.test, ast.Name)
|
||||
and node.test.id == "TYPE_CHECKING"):
|
||||
# Annotation-only bindings: lazy under future annotations, never
|
||||
# imported at runtime, so nothing is frozen at import time.
|
||||
for sub in node.body:
|
||||
if isinstance(sub, (ast.Import, ast.ImportFrom)):
|
||||
bound.update(a.asname or a.name.split(".")[0] for a in sub.names)
|
||||
return bound
|
||||
|
||||
|
||||
def _handle_reads(tree: ast.AST, handle: str) -> set[str]:
|
||||
reads: set[str] = set()
|
||||
for node in ast.walk(tree):
|
||||
if (isinstance(node, ast.Attribute) and isinstance(node.value, ast.Call)
|
||||
and isinstance(node.value.func, ast.Name) and node.value.func.id == handle):
|
||||
reads.add(node.attr)
|
||||
return reads
|
||||
|
||||
|
||||
@pytest.mark.parametrize("leaf", sorted(LEAVES))
|
||||
def test_each_leaf_reaches_its_parent_only_through_a_call_time_handle(leaf: str) -> None:
|
||||
parent, handle, _declared = LEAVES[leaf]
|
||||
parent_module = parent[:-3].replace("/", ".")
|
||||
tree = _tree(leaf)
|
||||
for node in tree.body: # module scope only: a lazy import inside the handle is the point
|
||||
if isinstance(node, ast.ImportFrom):
|
||||
assert node.module != parent_module, f"{leaf} imports its parent at module scope"
|
||||
if isinstance(node, ast.Import):
|
||||
assert all(a.name != parent_module for a in node.names), leaf
|
||||
handles = [n for n in tree.body if isinstance(n, ast.FunctionDef) and n.name == handle]
|
||||
assert len(handles) == 1, f"{leaf}: expected exactly one {handle}() definition"
|
||||
assert [n for n in ast.walk(handles[0]) if isinstance(n, (ast.Import, ast.ImportFrom))], (
|
||||
f"{leaf}: {handle}() must import the parent at call time"
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.parametrize("leaf", sorted(LEAVES))
|
||||
def test_the_declared_set_is_exactly_what_the_leaf_reads_through_the_handle(leaf: str) -> None:
|
||||
parent, handle, declared = LEAVES[leaf]
|
||||
actual = _handle_reads(_tree(leaf), handle)
|
||||
assert actual == set(declared), (
|
||||
f"{leaf}: declared {sorted(declared)} but reads {sorted(actual)}"
|
||||
)
|
||||
bound = _module_bindings(_tree(parent))
|
||||
missing = sorted(set(declared) - bound)
|
||||
assert missing == [], f"{leaf}: declared names absent from {parent}: {missing}"
|
||||
|
||||
|
||||
@pytest.mark.parametrize("leaf", sorted(LEAVES))
|
||||
def test_no_leaf_reads_a_parent_owned_name_directly(leaf: str) -> None:
|
||||
"""The bug the handle exists to prevent: a direct read freezes the binding the
|
||||
leaf saw at import time, so `init` rebinding the parent's name — or a test doing
|
||||
the same — would leave this module looking at the old object forever."""
|
||||
parent, _handle, _declared = LEAVES[leaf]
|
||||
leaf_tree = _tree(leaf)
|
||||
parent_defs: set[str] = set()
|
||||
for node in _tree(parent).body:
|
||||
if isinstance(node, (ast.FunctionDef, ast.AsyncFunctionDef, ast.ClassDef)):
|
||||
parent_defs.add(node.name)
|
||||
elif isinstance(node, ast.Assign):
|
||||
parent_defs.update(t.id for t in node.targets if isinstance(t, ast.Name))
|
||||
elif isinstance(node, ast.AnnAssign) and isinstance(node.target, ast.Name):
|
||||
parent_defs.add(node.target.id)
|
||||
own = _module_bindings(leaf_tree)
|
||||
direct = {
|
||||
node.id for node in ast.walk(leaf_tree)
|
||||
if isinstance(node, ast.Name) and isinstance(node.ctx, ast.Load)
|
||||
and node.id in parent_defs and node.id not in own
|
||||
}
|
||||
assert direct == set(), f"{leaf} reads {sorted(direct)} directly instead of through the handle"
|
||||
|
|
@ -890,7 +890,7 @@ def test_every_settings_writer_routes_through_the_shared_prologue():
|
|||
"immune-system ROLLBACK: rewrites the exact bytes snapshotted before an agent shell "
|
||||
"command. It authors no value, and filtering a restore would corrupt it — an "
|
||||
"owner-authored default would be dropped instead of restored.",
|
||||
("ouroboros/usage_accounting.py", "_legacy_snapshot"):
|
||||
("ouroboros/usage_legacy_import.py", "_legacy_snapshot"):
|
||||
"reads/hashes the settings file for the usage archive; its writes target the archive.",
|
||||
("ouroboros/tools/core.py", "_data_write"):
|
||||
"names SETTINGS_PATH only to REFUSE agent writes to it.",
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue