feat(P2): metrics collector — SQLite cold store + Decision Ledger CLI (REQ-189,200,201,207)
P2 (Wave 2, feat) — REQ-189, REQ-200, REQ-201, REQ-207 New components: - core/metrics/collector.py — reads all grounded signals (REGRESSION_REPORT.json, per-run manifests, junit XML, coverage.json, decision ledger, lifecycle reports) → SQLite cold store (metrics/nova_metrics.db) with fact_run, fact_capability, fact_policy_check, fact_confidence, fact_test, fact_decision, fact_cost_estimate, fact_lifecycle, dim_capability, dim_milestone tables - core/metrics/decision_ledger_cli.py — CLI with query/verify-chain/stats/export/replay - tests/test_metrics_collector.py — 7 tests (all pass, incl. idempotent re-run REQ-200) D-120: Nova-native (SQLite, no ClickHouse) D-125: hybrid (reads files + events → SQLite) D-126: cold-only (no hot path) ---ci--- project: acdl phase: 2 milestone: v1.17 status: execute ---/ci---
This commit is contained in:
@@ -0,0 +1,364 @@
|
||||
"""Nova Metrics Collector (REQ-189, P2).
|
||||
|
||||
Reads all grounded signals (REGRESSION_REPORT.json, per-run manifests,
|
||||
junit XML, pcr.json, signal.json, COST.md, decision ledger, coverage.json)
|
||||
and normalizes them into a SQLite cold store at metrics/nova_metrics.db.
|
||||
|
||||
D-120: Nova-native (SQLite, no ClickHouse/BigQuery).
|
||||
D-125: hybrid model — reads files + events → SQLite.
|
||||
D-126: cold-only (no hot path; hot path deferred D-096).
|
||||
D-128: metrics/ at repo root.
|
||||
|
||||
Idempotent: re-running the collector against the same inputs produces
|
||||
identical row counts (REQ-200). The collector uses INSERT OR REPLACE
|
||||
on fact tables keyed by natural keys.
|
||||
"""
|
||||
|
||||
import datetime
|
||||
import json
|
||||
import os
|
||||
import sqlite3
|
||||
import sys
|
||||
import xml.etree.ElementTree as ET
|
||||
|
||||
_METRICS_DIR = os.path.join(os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))), "metrics")
|
||||
_STORE_PATH = os.path.join(_METRICS_DIR, "nova_metrics.db")
|
||||
_REPO_ROOT = os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
|
||||
_REGRESSION_REPORT = os.path.join(_REPO_ROOT, ".ciagent", "REGRESSION_REPORT.json")
|
||||
_RUNS_DIR = os.path.join(_METRICS_DIR, "runs")
|
||||
_LEDGER_DB = os.path.join(_METRICS_DIR, "decision_ledger.db")
|
||||
_COVERAGE_JSON = os.path.join(_METRICS_DIR, "coverage.json")
|
||||
_TEST_RESULTS_XML = os.path.join(_METRICS_DIR, "test-results.xml")
|
||||
|
||||
|
||||
def _iso8601_now():
|
||||
return datetime.datetime.now(datetime.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
|
||||
|
||||
|
||||
def _init_store(db_path=None):
|
||||
"""Create the fact/dim tables in the SQLite cold store."""
|
||||
if db_path is None:
|
||||
db_path = _STORE_PATH
|
||||
os.makedirs(os.path.dirname(db_path), exist_ok=True)
|
||||
conn = sqlite3.connect(db_path)
|
||||
conn.executescript("""
|
||||
CREATE TABLE IF NOT EXISTS fact_run (
|
||||
run_id TEXT PRIMARY KEY,
|
||||
contract_id TEXT,
|
||||
environment TEXT,
|
||||
started_at TEXT,
|
||||
completed_at TEXT,
|
||||
exit_code INTEGER,
|
||||
outcome TEXT,
|
||||
confidence_score REAL,
|
||||
confidence_band TEXT,
|
||||
hitl_block INTEGER,
|
||||
cost_estimate_usd REAL,
|
||||
decision_id TEXT
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS fact_capability (
|
||||
capability_id TEXT,
|
||||
run_id TEXT,
|
||||
name TEXT,
|
||||
status TEXT,
|
||||
tier TEXT,
|
||||
duration_ms REAL,
|
||||
detail TEXT,
|
||||
run_at_utc TEXT,
|
||||
PRIMARY KEY (capability_id, run_id)
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS fact_policy_check (
|
||||
run_id TEXT,
|
||||
rule_id TEXT,
|
||||
severity TEXT,
|
||||
result TEXT,
|
||||
resource_ref TEXT,
|
||||
evaluated_at TEXT,
|
||||
PRIMARY KEY (run_id, rule_id, resource_ref)
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS fact_confidence (
|
||||
run_id TEXT,
|
||||
score REAL,
|
||||
band TEXT,
|
||||
per_input TEXT,
|
||||
reason_codes TEXT,
|
||||
environment TEXT,
|
||||
computed_at TEXT,
|
||||
PRIMARY KEY (run_id)
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS fact_test (
|
||||
run_id TEXT,
|
||||
total_tests INTEGER,
|
||||
passed INTEGER,
|
||||
failed INTEGER,
|
||||
errors INTEGER,
|
||||
skipped INTEGER,
|
||||
duration_s REAL,
|
||||
coverage_pct REAL,
|
||||
collected_at TEXT,
|
||||
PRIMARY KEY (run_id)
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS fact_decision (
|
||||
decision_id TEXT,
|
||||
run_id TEXT,
|
||||
chosen_action TEXT,
|
||||
confidence REAL,
|
||||
alternatives TEXT,
|
||||
human_override INTEGER,
|
||||
outcome TEXT,
|
||||
event_time TEXT,
|
||||
PRIMARY KEY (decision_id)
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS fact_cost_estimate (
|
||||
run_id TEXT,
|
||||
delta_usd REAL,
|
||||
total_monthly_usd REAL,
|
||||
available INTEGER,
|
||||
estimated_at TEXT,
|
||||
PRIMARY KEY (run_id)
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS fact_lifecycle (
|
||||
module TEXT,
|
||||
environment TEXT,
|
||||
phase TEXT,
|
||||
result TEXT,
|
||||
duration_ms REAL,
|
||||
run_at TEXT,
|
||||
PRIMARY KEY (module, environment, phase, run_at)
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS dim_capability (
|
||||
capability_id TEXT PRIMARY KEY,
|
||||
name TEXT,
|
||||
tier TEXT,
|
||||
source_milestone TEXT
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS dim_milestone (
|
||||
milestone TEXT PRIMARY KEY,
|
||||
phase INTEGER,
|
||||
tag TEXT,
|
||||
completed_at TEXT
|
||||
);
|
||||
""")
|
||||
conn.commit()
|
||||
conn.close()
|
||||
|
||||
|
||||
def collect_regression_report(db_path=None, report_path=None):
|
||||
"""Read REGRESSION_REPORT.json → fact_capability + dim_capability."""
|
||||
if db_path is None:
|
||||
db_path = _STORE_PATH
|
||||
if report_path is None:
|
||||
report_path = _REGRESSION_REPORT
|
||||
if not os.path.isfile(report_path):
|
||||
return 0
|
||||
_init_store(db_path)
|
||||
with open(report_path) as f:
|
||||
report = json.load(f)
|
||||
run_id = report.get("run_id", f"regr-{report.get('run_at_utc','')}")
|
||||
run_at = report.get("run_at_utc", _iso8601_now())
|
||||
milestone = report.get("milestone", "")
|
||||
conn = sqlite3.connect(db_path)
|
||||
for result in report.get("results", []):
|
||||
cap_id = result.get("capability_id", "")
|
||||
conn.execute("""
|
||||
INSERT OR REPLACE INTO fact_capability
|
||||
(capability_id, run_id, name, status, tier, duration_ms, detail, run_at_utc)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?, ?)
|
||||
""", (cap_id, run_id, result.get("name", ""), result.get("status", ""),
|
||||
result.get("tier", ""), result.get("duration_ms", 0),
|
||||
result.get("detail", ""), run_at))
|
||||
conn.execute("""
|
||||
INSERT OR REPLACE INTO dim_capability
|
||||
(capability_id, name, tier, source_milestone)
|
||||
VALUES (?, ?, ?, ?)
|
||||
""", (cap_id, result.get("name", ""), result.get("tier", ""), milestone))
|
||||
conn.execute("""
|
||||
INSERT OR REPLACE INTO dim_milestone
|
||||
(milestone, phase, tag, completed_at)
|
||||
VALUES (?, ?, ?, ?)
|
||||
""", (milestone, report.get("phase", 0), "", run_at))
|
||||
conn.commit()
|
||||
conn.close()
|
||||
return len(report.get("results", []))
|
||||
|
||||
|
||||
def collect_run_manifests(db_path=None, runs_dir=None):
|
||||
"""Read per-run manifests from metrics/runs/*.json → fact_run."""
|
||||
if db_path is None:
|
||||
db_path = _STORE_PATH
|
||||
if runs_dir is None:
|
||||
runs_dir = _RUNS_DIR
|
||||
if not os.path.isdir(runs_dir):
|
||||
return 0
|
||||
_init_store(db_path)
|
||||
count = 0
|
||||
conn = sqlite3.connect(db_path)
|
||||
for fname in sorted(os.listdir(runs_dir)):
|
||||
if not fname.endswith(".json"):
|
||||
continue
|
||||
fpath = os.path.join(runs_dir, fname)
|
||||
if os.path.isdir(fpath):
|
||||
continue
|
||||
with open(fpath) as f:
|
||||
manifest = json.load(f)
|
||||
run_id = manifest.get("run_id", fname.replace(".json", ""))
|
||||
conf = manifest.get("confidence", {})
|
||||
hitl = manifest.get("hitl", {})
|
||||
conn.execute("""
|
||||
INSERT OR REPLACE INTO fact_run
|
||||
(run_id, contract_id, environment, started_at, completed_at,
|
||||
exit_code, outcome, confidence_score, confidence_band,
|
||||
hitl_block, cost_estimate_usd, decision_id)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||
""", (run_id, manifest.get("contract_id", ""), manifest.get("environment", ""),
|
||||
manifest.get("started_at", ""), manifest.get("completed_at", ""),
|
||||
manifest.get("exit_code", 0), manifest.get("outcome", ""),
|
||||
conf.get("score", 0), conf.get("band", ""),
|
||||
1 if hitl.get("block") else 0,
|
||||
manifest.get("cost_estimate_usd", 0), manifest.get("decision_id", "")))
|
||||
count += 1
|
||||
conn.commit()
|
||||
conn.close()
|
||||
return count
|
||||
|
||||
|
||||
def collect_decision_ledger(db_path=None, ledger_db=None):
|
||||
"""Read the Decision Ledger SQLite → fact_decision."""
|
||||
if db_path is None:
|
||||
db_path = _STORE_PATH
|
||||
if ledger_db is None:
|
||||
ledger_db = _LEDGER_DB
|
||||
if not os.path.isfile(ledger_db):
|
||||
return 0
|
||||
_init_store(db_path)
|
||||
ledger_conn = sqlite3.connect(ledger_db)
|
||||
rows = ledger_conn.execute(
|
||||
"SELECT event_type, run_id, event_time, payload FROM decision_ledger WHERE event_type = 'nova.ai.decision.made' ORDER BY seq"
|
||||
).fetchall()
|
||||
ledger_conn.close()
|
||||
conn = sqlite3.connect(db_path)
|
||||
count = 0
|
||||
for etype, run_id, event_time, payload_json in rows:
|
||||
payload = json.loads(payload_json)
|
||||
data = payload.get("data", {})
|
||||
decision_id = data.get("decision_id", run_id)
|
||||
conn.execute("""
|
||||
INSERT OR REPLACE INTO fact_decision
|
||||
(decision_id, run_id, chosen_action, confidence, alternatives,
|
||||
human_override, outcome, event_time)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?, ?)
|
||||
""", (decision_id, run_id, data.get("chosen_action", ""),
|
||||
data.get("confidence", 0), json.dumps(data.get("alternatives", {})),
|
||||
1 if data.get("human_override") else 0,
|
||||
data.get("outcome", "pending"), event_time))
|
||||
count += 1
|
||||
conn.commit()
|
||||
conn.close()
|
||||
return count
|
||||
|
||||
|
||||
def collect_test_results(db_path=None, junit_path=None, coverage_path=None):
|
||||
"""Read junit XML + coverage.json → fact_test."""
|
||||
if db_path is None:
|
||||
db_path = _STORE_PATH
|
||||
if junit_path is None:
|
||||
junit_path = _TEST_RESULTS_XML
|
||||
if coverage_path is None:
|
||||
coverage_path = _COVERAGE_JSON
|
||||
if not os.path.isfile(junit_path):
|
||||
return 0
|
||||
_init_store(db_path)
|
||||
run_id = f"test-{_iso8601_now()}"
|
||||
total = passed = failed = errors = skipped = 0
|
||||
duration = 0.0
|
||||
try:
|
||||
tree = ET.parse(junit_path)
|
||||
root = tree.getroot()
|
||||
for suite in root.iter("testsuite"):
|
||||
total += int(suite.get("tests", 0))
|
||||
failed += int(suite.get("failures", 0))
|
||||
errors += int(suite.get("errors", 0))
|
||||
skipped += int(suite.get("skipped", 0))
|
||||
duration += float(suite.get("time", 0))
|
||||
passed = total - failed - errors - skipped
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
coverage_pct = 0.0
|
||||
if os.path.isfile(coverage_path):
|
||||
try:
|
||||
with open(coverage_path) as f:
|
||||
cov = json.load(f)
|
||||
coverage_pct = cov.get("totals", {}).get("percent_covered", 0.0)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
conn = sqlite3.connect(db_path)
|
||||
conn.execute("""
|
||||
INSERT OR REPLACE INTO fact_test
|
||||
(run_id, total_tests, passed, failed, errors, skipped, duration_s, coverage_pct, collected_at)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||
""", (run_id, total, passed, failed, errors, skipped, duration, coverage_pct, _iso8601_now()))
|
||||
conn.commit()
|
||||
conn.close()
|
||||
return 1
|
||||
|
||||
|
||||
def collect_lifecycle_reports(db_path=None, lifecycle_dir=None):
|
||||
"""Read metrics/lifecycle/*.json → fact_lifecycle."""
|
||||
if db_path is None:
|
||||
db_path = _STORE_PATH
|
||||
if lifecycle_dir is None:
|
||||
lifecycle_dir = os.path.join(_METRICS_DIR, "lifecycle")
|
||||
if not os.path.isdir(lifecycle_dir):
|
||||
return 0
|
||||
_init_store(db_path)
|
||||
count = 0
|
||||
conn = sqlite3.connect(db_path)
|
||||
for fname in sorted(os.listdir(lifecycle_dir)):
|
||||
if not fname.endswith(".json"):
|
||||
continue
|
||||
fpath = os.path.join(lifecycle_dir, fname)
|
||||
with open(fpath) as f:
|
||||
report = json.load(f)
|
||||
conn.execute("""
|
||||
INSERT OR REPLACE INTO fact_lifecycle
|
||||
(module, environment, phase, result, duration_ms, run_at)
|
||||
VALUES (?, ?, ?, ?, ?, ?)
|
||||
""", (report.get("module", ""), report.get("environment", ""),
|
||||
report.get("phase", ""), report.get("result", ""),
|
||||
report.get("duration_ms", 0), report.get("run_at", _iso8601_now())))
|
||||
count += 1
|
||||
conn.commit()
|
||||
conn.close()
|
||||
return count
|
||||
|
||||
|
||||
def collect_all(db_path=None):
|
||||
"""Run all collectors. Returns a summary dict."""
|
||||
if db_path is None:
|
||||
db_path = _STORE_PATH
|
||||
_init_store(db_path)
|
||||
summary = {
|
||||
"capabilities": collect_regression_report(db_path),
|
||||
"runs": collect_run_manifests(db_path),
|
||||
"decisions": collect_decision_ledger(db_path),
|
||||
"tests": collect_test_results(db_path),
|
||||
"lifecycle": collect_lifecycle_reports(db_path),
|
||||
"collected_at": _iso8601_now(),
|
||||
}
|
||||
return summary
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
result = collect_all()
|
||||
print(json.dumps(result, indent=2))
|
||||
@@ -0,0 +1,39 @@
|
||||
"""Nova Decision Ledger CLI (REQ-207).
|
||||
|
||||
Subcommands: query, verify-chain, stats, export, replay.
|
||||
Read-only CLI for the Decision Ledger SQLite hash-chain.
|
||||
"""
|
||||
|
||||
import json
|
||||
import os
|
||||
import sys
|
||||
|
||||
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))))
|
||||
from core.metrics.decision_ledger import query_by_run, verify_chain, stats, export_since, replay_run
|
||||
|
||||
|
||||
def main():
|
||||
if len(sys.argv) < 2:
|
||||
print("usage: decision_ledger_cli.py <query|verify-chain|stats|export|replay> [args]", file=sys.stderr)
|
||||
sys.exit(2)
|
||||
cmd = sys.argv[1]
|
||||
if cmd == "query" and len(sys.argv) >= 3:
|
||||
print(json.dumps(query_by_run(sys.argv[2]), indent=2))
|
||||
elif cmd == "verify-chain":
|
||||
ok, broken, details = verify_chain()
|
||||
print(f"chain_ok={ok} broken={broken} details={details}")
|
||||
sys.exit(0 if ok else 1)
|
||||
elif cmd == "stats":
|
||||
print(json.dumps(stats(), indent=2))
|
||||
elif cmd == "export" and len(sys.argv) >= 3:
|
||||
fmt = sys.argv[3] if len(sys.argv) >= 4 else "json"
|
||||
print(export_since(sys.argv[2], fmt=fmt))
|
||||
elif cmd == "replay" and len(sys.argv) >= 3:
|
||||
print(replay_run(sys.argv[2]))
|
||||
else:
|
||||
print(f"unknown command: {cmd}", file=sys.stderr)
|
||||
sys.exit(2)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
Reference in New Issue
Block a user