feat(spectral): Cosmos Gremlin loader for the codebook graph (dry-run by default)

Loads codebook/codeword/fingerprint/equation vertices with in_cluster,
shares_fingerprint and within_cartan_floor(fisher_distance) edges into
the mathblob/concepts graph, mirroring the Research Stack
load_module_graph.py conventions (pk partition key, coalesce upserts,
RU-throttle retry). Applied 2026-07-02: 273 vertices + 339 edges ok;
250 pre-existing equation vertices refreshed in place with exact values.
Docs note the Spark-side JDBC verification of ene.rrc_predictions.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
allaun 2026-07-01 21:32:48 -05:00
parent 6eae21513e
commit 3ba610bdf3
2 changed files with 286 additions and 0 deletions

View file

@ -288,6 +288,28 @@ DB with pgvector, `concept_citations`). It has **no published port** —
access is podman-exec only. It is the citation/grounding layer for access is podman-exec only. It is the citation/grounding layer for
equation provenance; this sync does not (and cannot) connect to it. equation provenance; this sync does not (and cannot) connect to it.
### Gremlin graph layer (Cosmos DB mathblob/concepts)
`python/spectral_codebook_gremlin.py` loads the codebook into the Azure
Cosmos DB Gremlin graph used by the Research Stack module/concept loaders
(credentials in `Research Stack/.env.gremlin`, same conventions as
`load_module_graph.py`). Dry-run by default; `--apply` upserts.
Model: 1 `codebook` provenance vertex, 9 `codeword` cluster vertices, 13
`fingerprint` vertices (cospectral collision classes only), 250 `equation`
vertices (ids verbatim; spiral_index stored as string — it exceeds int64),
with `has_codeword` / `in_cluster` / `shares_fingerprint` /
`within_cartan_floor(fisher_distance)` edges. First applied 2026-07-02:
273 vertices + 339 edges, all ok; the 250 equation vertices already
existed in the graph and were refreshed in place with exact values.
Spark-side availability of the synced rows was verified the same day by a
client-mode PySpark 3.5.5 job (JDK 17 via nix, driver on the tailnet,
`spark.driver.host` set) reading `ene.rrc_predictions` through JDBC from
the spark://100.92.88.64:7077 cluster: 250 rows, C0C8 and shape counts
matching the sync exactly. CouchDB (:5984) remains unresponsive — nothing
loaded there.
## Output schema (`data/spectral_codebook.json`) ## Output schema (`data/spectral_codebook.json`)
- header: `schema` (`spectral_codebook_v2`), `phi_corkscrew` (packing - header: `schema` (`spectral_codebook_v2`), `phi_corkscrew` (packing

View file

@ -0,0 +1,264 @@
#!/usr/bin/env -S uv run
# /// script
# requires-python = ">=3.11"
# dependencies = ["gremlinpython", "python-dotenv"]
# ///
"""
spectral_codebook_gremlin.py Load the spectral codebook into the Cosmos DB
Gremlin graph (mathblob / concepts), alongside the Research Stack module graph.
Graph model (labels chosen to not collide with the existing 'module' layer):
codebook 1 vertex provenance root: schema, packing base/offset, date
codeword 9 vertices gap-rule clusters C0..C8 with λ ranges and counts
fingerprint one per charpoly COLLISION class only (same char
poly, >1 equation): the cospectral classes
equation 250 vertices one per corpus id (identifiers preserved verbatim),
with codeword, exact spectral radius, charpoly csv,
spiral index (string exceeds int64), corkscrew
angle_frac, sparse-tail flag
codebook -has_codeword-> codeword
equation -in_cluster-> codeword
equation -shares_fingerprint-> fingerprint (collision classes only)
equation -within_cartan_floor-> equation (Fisher d < Δ = 17/1792;
property fisher_distance)
Credentials: Research Stack/.env.gremlin (same file as load_module_graph.py).
Reads data/spectral_codebook.json regenerate with
python3 python/spectral_codebook.py --check-manifold --cartan-floor
if it predates the phi_corkscrew/cartan_floor blocks.
DRY RUN by default: prints the plan, writes nothing. --apply upserts
(idempotent coalesce pattern, small batches for Cosmos RU limits).
"""
from __future__ import annotations
import argparse
import json
import os
import sys
import time
from pathlib import Path
REPO_ROOT = Path(__file__).resolve().parent.parent
CODEBOOK_JSON = REPO_ROOT / "data" / "spectral_codebook.json"
ENV_FILE = Path.home() / "Research Stack" / ".env.gremlin"
BATCH_PAUSE_EVERY = 25 # Cosmos free tier: pace requests to avoid RU 429s
BATCH_PAUSE_SECS = 0.5
MAX_RETRIES = 5
# ── Graph plan from the codebook JSON ─────────────────────────────────────
def build_plan(doc: dict) -> dict:
"""Vertices/edges to upsert, derived purely from the codebook JSON."""
pc = doc.get("phi_corkscrew", {})
packing = pc.get("packing", {})
root_id = "spectral_codebook_v2"
vertices = [{
"label": "codebook", "id": root_id,
"props": {
"schema": doc["schema"],
"generated": doc["generated"],
"source": doc["source"],
"matrix_count": doc["matrix_count"],
"pack_base": packing.get("base", 0),
"pack_offset": packing.get("offset", 0),
},
}]
edges = []
for row in doc["clusters"]:
cw = row["codeword"]
vertices.append({
"label": "codeword", "id": f"codeword:{cw}",
"props": {
"name": cw,
"count": row["count"],
"distinct_matrices": row["distinct_matrices"],
"lambda_min": row["lambda_min"],
"lambda_max": row["lambda_max"],
},
})
edges.append((root_id, f"codeword:{cw}", "has_codeword", {}))
collision_classes = doc["collisions"]["charpoly_collision_classes"]
for fp_csv in collision_classes:
vertices.append({
"label": "fingerprint", "id": f"fp:{fp_csv}",
"props": {"charpoly": fp_csv,
"members": len(collision_classes[fp_csv])},
})
for e in doc["entries"]:
eid = e["equation_id"] # preserved verbatim
fp_csv = ",".join(str(c) for c in e["charpoly"])
vertices.append({
"label": "equation", "id": eid,
"props": {
"codeword": e["codeword"],
"cluster_index": e["index"],
"spectral_radius": e["spectral_radius"],
"charpoly": fp_csv,
# spiral indices reach ~1e44 — far beyond int64, so string
"spiral_index": str(e["spiral_index"]),
"angle_frac": e["layout"]["angle_frac"],
"sparse_tail": bool(e.get("sparse_tail", False)),
},
})
edges.append((eid, f"codeword:{e['codeword']}", "in_cluster", {}))
if fp_csv in collision_classes:
edges.append((eid, f"fp:{fp_csv}", "shares_fingerprint", {}))
for pair in doc.get("cartan_floor", {}).get("sub_delta_pairs", []):
a, b = pair["representatives"]
edges.append((a, b, "within_cartan_floor",
{"fisher_distance": pair["fisher_distance"]}))
return {"vertices": vertices, "edges": edges}
# ── Gremlin plumbing (mirrors Research Stack load_module_graph.py) ────────
def make_client():
from dotenv import load_dotenv
from gremlin_python.driver import client as gremlin_client, serializer
if not ENV_FILE.exists():
sys.exit(f"{ENV_FILE} not found — run setup_mathblob.py first")
load_dotenv(ENV_FILE)
return gremlin_client.Client(
os.environ["GREMLIN_ENDPOINT"], "g",
username=os.environ["GREMLIN_USERNAME"],
password=os.environ["GREMLIN_PASSWORD"],
message_serializer=serializer.GraphSONSerializersV2d0(),
)
def submit(c, query: str, bindings: dict | None = None):
"""Submit with retry on Cosmos RU throttling (429)."""
for attempt in range(MAX_RETRIES):
try:
cb = c.submitAsync(query, bindings or {})
return cb.result().all().result()
except Exception as e: # noqa: BLE001 — driver raises plain Exception
msg = str(e)
if "429" in msg or "TooManyRequests" in msg or "3200" in msg:
time.sleep(1.5 * (attempt + 1))
continue
print(f" ERR: {msg:.160}")
return None
print(" ERR: retries exhausted (RU throttling)")
return None
def upsert_vertex(c, v: dict) -> bool:
"""Idempotent vertex upsert; 'pk' mirrors id (Cosmos partition key)."""
label = v["label"]
q = (
f"g.V().has('{label}','id',vid).fold()"
f".coalesce(unfold(), addV('{label}')"
f".property('id',vid).property('pk',vid))"
)
bindings = {"vid": v["id"]}
for i, (k, val) in enumerate(v["props"].items()):
q += f".property('{k}',p{i})"
bindings[f"p{i}"] = val
return submit(c, q, bindings) is not None
def upsert_edge(c, src: str, dst: str, label: str, props: dict) -> bool:
q = (
"g.V().has('id',src).as('s')"
".V().has('id',dst).as('d')"
".coalesce("
" select('s').outE(lbl).where(inV().as('d')),"
" addE(lbl).from('s').to('d')"
")"
)
bindings = {"src": src, "dst": dst, "lbl": label}
for i, (k, val) in enumerate(props.items()):
q += f".property('{k}',ep{i})"
bindings[f"ep{i}"] = val
return submit(c, q, bindings) is not None
# ── CLI ────────────────────────────────────────────────────────────────────
def main() -> int:
ap = argparse.ArgumentParser(
description="Load spectral codebook into Cosmos Gremlin (dry-run default)")
ap.add_argument("--codebook", type=Path, default=CODEBOOK_JSON)
ap.add_argument("--apply", action="store_true",
help="actually upsert (default: dry run, writes nothing)")
args = ap.parse_args()
doc = json.loads(args.codebook.read_text())
if "phi_corkscrew" not in doc:
sys.exit("codebook JSON lacks phi_corkscrew block — regenerate "
"(see module docstring)")
if "cartan_floor" not in doc:
print("NOTE: no cartan_floor block — within_cartan_floor edges skipped; "
"regenerate with --cartan-floor to include them")
plan = build_plan(doc)
by_label: dict[str, int] = {}
for v in plan["vertices"]:
by_label[v["label"]] = by_label.get(v["label"], 0) + 1
by_edge: dict[str, int] = {}
for _, _, lbl, _ in plan["edges"]:
by_edge[lbl] = by_edge.get(lbl, 0) + 1
print(f"Codebook {doc['generated']} → gremlin graph plan:")
print(f" vertices: {sum(by_label.values())} {by_label}")
print(f" edges: {sum(by_edge.values())} {by_edge}")
if not args.apply:
print("\nDRY RUN — nothing written. Sample vertex/edge:")
eq = next(v for v in plan["vertices"] if v["label"] == "equation")
print(f" {eq}")
print(f" {plan['edges'][0]}")
print("Use --apply to upsert into Cosmos (mathblob/concepts).")
return 0
c = make_client()
try:
before = submit(c, "g.V().count()")
print(f"\nConnected. Vertices before: {before}")
ok = fail = 0
for i, v in enumerate(plan["vertices"]):
if upsert_vertex(c, v):
ok += 1
else:
fail += 1
print(f" vertex failed: {v['id']}")
if (i + 1) % BATCH_PAUSE_EVERY == 0:
time.sleep(BATCH_PAUSE_SECS)
print(f" vertices upserted: {ok} ok, {fail} failed")
ok = fail = 0
for i, (src, dst, lbl, props) in enumerate(plan["edges"]):
if upsert_edge(c, src, dst, lbl, props):
ok += 1
else:
fail += 1
print(f" edge failed: {src} -{lbl}-> {dst}")
if (i + 1) % BATCH_PAUSE_EVERY == 0:
time.sleep(BATCH_PAUSE_SECS)
print(f" edges upserted: {ok} ok, {fail} failed")
after = submit(c, "g.V().count()")
eq_count = submit(c, "g.V().hasLabel('equation').count()")
print(f" vertices after: {after}; equation vertices: {eq_count}")
return 0
finally:
c.close()
if __name__ == "__main__":
sys.exit(main())