bankml / sAGI /connectors.py
Gregory-L's picture
bankML: the whole source (github.com/cryptoAGI/bankml @ 12ae409) and its page, with the bankML persona; the live engine (Dockerfile, hf/start.sh) ready for Docker hardware
28c70af verified
Raw History Blame Contribute Delete
13.2 kB
#!/usr/bin/env python3
# SPDX-License-Identifier: MIT OR Apache-2.0
"""Connectors for bankml agents: PostgreSQL (pgvector / pgvectorscale) — publish an agent, load it back, verified.
Through the `psql` client and the standard library (no driver). Values never enter SQL text: they travel as COPY
data (text format, escaped) into a temporary table, and the statements read them from there.
Publish writes the agent's public parts — persona, prompt, card, ledger, THOT manifest — and, only when asked
(`include_private=True`), its .history and .memory lines. Without that, the database holds the commitments (the THOT
manifest's digests of history and memory) and not the content: proof of data without the data.
Load rebuilds the agent from the row and refuses it unless every public file re-hashes to the digest the stored
THOT manifest commits to (and the manifest's own identity recomputes).
BANKML_PG_DSN libpq connection string (default "dbname=bankml": the local socket, the operator's role)
Schema: bankml_agents, bankml_exchanges (with embedding vector(1024) — bge-m3's width, as mindX uses — when the
vector extension is present; DiskANN index when pgvectorscale is, HNSW otherwise), bankml_memory.
"""
from __future__ import annotations
import hashlib
import json
import os
import subprocess
import agents
import thot
DSN = os.environ.get("BANKML_PG_DSN", "dbname=bankml")
SETUP = ('sudo -u postgres psql -c "CREATE ROLE $USER LOGIN CREATEDB" -c "CREATE DATABASE bankml OWNER $USER" '
'&& sudo -u postgres psql -d bankml -c "CREATE EXTENSION IF NOT EXISTS vector"')
EMBED_DIM = 1024
def _esc(v) -> str:
"""COPY text format: backslash, tab, newline and carriage return escaped; None is \\N."""
if v is None:
return r"\N"
return str(v).replace("\\", "\\\\").replace("\t", "\\t").replace("\n", "\\n").replace("\r", "\\r")
def psql(script: str, rows: list | None = None, dsn: str | None = None, atomic: bool = False) -> str:
"""Run a script; `rows` ([(k, v), …]) are COPY-ed into temp table _in(k text, v text) first. Returns stdout.
atomic: the whole script, COPY included, is one transaction (psql --single-transaction): all of it or none."""
head = ""
if rows is not None:
head = "CREATE TEMP TABLE _in(k text, v text);\nCOPY _in FROM STDIN;\n" + "".join(f"{_esc(k)}\t{_esc(v)}\n" for k, v in rows) + "\\.\n"
p = subprocess.run(["psql", "-X", "-q", "-A", "-t", "-v", "ON_ERROR_STOP=1"] + (["--single-transaction"] if atomic else []) + ["-d", dsn or DSN, "-f", "-"],
input=head + script, capture_output=True, text=True, timeout=120)
if p.returncode:
raise RuntimeError(p.stderr.strip() or f"psql exit {p.returncode}")
return p.stdout
def status(dsn: str | None = None) -> dict:
try:
out = psql("SELECT json_build_object('user', current_user, 'db', current_database(), 'server', current_setting('server_version'),"
" 'vector', (SELECT extversion FROM pg_extension WHERE extname='vector'),"
" 'vectorscale', (SELECT extversion FROM pg_extension WHERE extname='vectorscale'),"
" 'vector_available', EXISTS (SELECT 1 FROM pg_available_extensions WHERE name='vector'),"
" 'vectorscale_available', EXISTS (SELECT 1 FROM pg_available_extensions WHERE name='vectorscale'));", dsn=dsn)
return {"ok": True, **json.loads(out.strip())}
except (RuntimeError, OSError, subprocess.TimeoutExpired, ValueError) as e:
return {"ok": False, "error": str(e).splitlines()[-1] if str(e) else repr(e), "setup": SETUP}
def ensure_schema(dsn: str | None = None) -> dict:
st = status(dsn)
if not st["ok"]:
return st
psql("""
CREATE TABLE IF NOT EXISTS bankml_agents (slug text PRIMARY KEY, name text NOT NULL, persona jsonb NOT NULL, persona_raw text NOT NULL, prompt text NOT NULL,
card jsonb NOT NULL, commitments jsonb NOT NULL, manifest jsonb NOT NULL, thot_cid text NOT NULL, content_root text NOT NULL,
generation int NOT NULL, private_included boolean NOT NULL, published_at timestamptz NOT NULL DEFAULT now());
CREATE TABLE IF NOT EXISTS bankml_exchanges (agent text NOT NULL REFERENCES bankml_agents(slug) ON DELETE CASCADE, seq int NOT NULL,
line_sha256 text NOT NULL, line text NOT NULL, rec jsonb NOT NULL, PRIMARY KEY (agent, seq));
CREATE TABLE IF NOT EXISTS bankml_memory (agent text NOT NULL REFERENCES bankml_agents(slug) ON DELETE CASCADE, seq int NOT NULL,
line_sha256 text NOT NULL, line text NOT NULL, rec jsonb NOT NULL, PRIMARY KEY (agent, seq));
""", dsn=dsn)
if st.get("vector"):
psql(f"ALTER TABLE bankml_exchanges ADD COLUMN IF NOT EXISTS embedding vector({EMBED_DIM});", dsn=dsn)
if st.get("vectorscale"):
psql("CREATE INDEX IF NOT EXISTS bankml_exchanges_emb ON bankml_exchanges USING diskann (embedding vector_cosine_ops);", dsn=dsn)
else:
psql("CREATE INDEX IF NOT EXISTS bankml_exchanges_emb ON bankml_exchanges USING hnsw (embedding vector_cosine_ops);", dsn=dsn)
return {**status(dsn), "schema": "ready", "index": "diskann (pgvectorscale)" if st.get("vectorscale") else ("hnsw (pgvector)" if st.get("vector") else "none")}
def _lines(p):
try:
return [l for l in p.read_bytes().split(b"\n") if l.strip()]
except OSError:
return []
def publish(slug: str, include_private: bool = False, dsn: str | None = None) -> dict:
"""Upsert the agent (and, if asked, its history and memory lines). The THOT manifest is rebuilt first."""
sch = ensure_schema(dsn)
f = agents.files(slug)
m = thot.build(slug)
bad = thot.verify(slug)
if bad:
raise ValueError(f"the agent does not verify, not published: {bad}")
led = json.loads(f["commitments.json"].read_text(encoding="utf-8"))
rows = [("slug", slug), ("name", m["bundle"]["officer"]), ("persona", f["persona"].read_text(encoding="utf-8")),
("prompt", f["prompt"].read_text(encoding="utf-8")), ("card", f["agentcard.json"].read_text(encoding="utf-8")),
("commitments", json.dumps(led, ensure_ascii=False)), ("manifest", json.dumps(m, ensure_ascii=False)),
("thot_cid", m["identity"]["cid"]), ("content_root", m["identity"]["contentRoot"]), ("generation", str(m["bundle"]["generation"])),
("private", "true" if include_private else "false")]
get = lambda k: f"(SELECT v FROM _in WHERE k='{k}')" # noqa: E731 — fixed column names only; values stay in _in
hist = _lines(f["history"]) if include_private else []
mem = _lines(f["memory"]) if include_private else []
# k = "h:<position>" / "m:<position>"; v = the line's exact text (the bytes its sha256 and the Merkle proofs are over)
rows += [(f"h:{i}", l.decode("utf-8")) for i, l in enumerate(hist)] + [(f"m:{i}", l.decode("utf-8")) for i, l in enumerate(mem)]
emb_note, n_emb = "no private lines published", 0
if hist and sch.get("vector"): # meaning-search in the database, as mindX's memories: bge-m3, 1024-d
import embed
try:
recs = [json.loads(l) for l in hist]
embed.index(f["history"], recs)
have = embed.cached(f["history"])
for i, r in enumerate(recs):
v = have.get(hashlib.sha256(embed.exchange_text(r).encode()).hexdigest())
if v is not None:
rows.append((f"e:{i}", "[" + ",".join(f"{x:.7g}" for x in v) + "]"))
n_emb += 1
emb_note = f"{n_emb} of {len(hist)} exchanges embedded with {embed.MODEL} ({embed.provenance().get('digest', '')[:12]}…)"
except (embed.Unavailable, ValueError) as e:
emb_note = f"not embedded: {e}"
elif hist:
emb_note = "not embedded: the database has no pgvector"
lines_sql = "".join(f"""
DELETE FROM {table} WHERE agent = {get('slug')};
INSERT INTO {table} (agent, seq, line_sha256, line, rec)
SELECT {get('slug')}, substr(k, 3)::int, encode(sha256(convert_to(v, 'UTF8')), 'hex'), v, v::jsonb FROM _in WHERE k LIKE '{p}:%';
""" for table, p in (("bankml_exchanges", "h"), ("bankml_memory", "m"))) + ("""
UPDATE bankml_exchanges x SET embedding = i.v::vector FROM _in i
WHERE i.k LIKE 'e:%' AND x.agent = (SELECT v FROM _in WHERE k='slug') AND x.seq = substr(i.k, 3)::int;
""" if n_emb else "")
# one transaction: the agent row and its lines change together or not at all
psql(f"""
INSERT INTO bankml_agents (slug, name, persona, persona_raw, prompt, card, commitments, manifest, thot_cid, content_root, generation, private_included)
VALUES ({get('slug')}, {get('name')}, {get('persona')}::jsonb, {get('persona')}, {get('prompt')}, {get('card')}::jsonb, {get('commitments')}::jsonb,
{get('manifest')}::jsonb, {get('thot_cid')}, {get('content_root')}, {get('generation')}::int, {get('private')}::boolean)
ON CONFLICT (slug) DO UPDATE SET name=EXCLUDED.name, persona=EXCLUDED.persona, persona_raw=EXCLUDED.persona_raw, prompt=EXCLUDED.prompt, card=EXCLUDED.card,
commitments=EXCLUDED.commitments, manifest=EXCLUDED.manifest, thot_cid=EXCLUDED.thot_cid, content_root=EXCLUDED.content_root,
generation=EXCLUDED.generation, private_included=EXCLUDED.private_included, published_at=now();
{lines_sql}""", rows, dsn, atomic=True)
n_ex, n_mem = len(hist), len(mem)
return {"slug": slug, "thot_cid": m["identity"]["cid"], "content_root": m["identity"]["contentRoot"], "generation": m["bundle"]["generation"],
"private_included": include_private, "exchanges": n_ex, "memory": n_mem, "embedded": n_emb, "embedding": emb_note}
def published(dsn: str | None = None) -> list:
out = psql("SELECT coalesce(json_agg(json_build_object('slug', slug, 'name', name, 'generation', generation, 'thot_cid', thot_cid,"
" 'private_included', private_included, 'published_at', published_at) ORDER BY published_at DESC), '[]') FROM bankml_agents;", dsn=dsn)
return json.loads(out.strip() or "[]")
def load(slug: str, as_slug: str | None = None, dsn: str | None = None, expect_content_root: str | None = None) -> dict:
"""Rebuild a published agent locally — refused unless every public file matches the stored THOT manifest.
The manifest comes from the same row, so on its own that is self-consistency: whoever can write the database can
write a consistent row. Pass `expect_content_root` (an iNFT's contentRoot, or one published elsewhere) to anchor
the manifest outside the database; the result says which kind of check was made."""
out = psql("SELECT row_to_json(a) FROM bankml_agents a WHERE slug = (SELECT v FROM _in WHERE k='slug');", [("slug", slug)], dsn).strip()
if not out:
raise KeyError(f"no published agent {slug!r}")
row = json.loads(out)
m = row["manifest"]
findings = thot.check_structure(m)
if expect_content_root and (m.get("identity") or {}).get("contentRoot", "").lower() != expect_content_root.lower():
raise ValueError(f"refused: the stored manifest's contentRoot is not the expected {expect_content_root}")
persona_b = row["persona_raw"].encode() # the exact bytes; the jsonb copy is for queries
want = {x["facet"]: x["sha256"] for x in m["facets"]}
got = {"persona": hashlib.sha256(persona_b).hexdigest(), "prompt": hashlib.sha256(row["prompt"].encode()).hexdigest()}
findings += [f"{k}: stored bytes do not match the manifest" for k in got if got[k] != want.get(k)]
if findings:
raise ValueError(f"refused: {findings}")
# every private blob fetched and checked BEFORE anything is written: a refusal leaves the target untouched
private = {}
if row.get("private_included"):
for table, key, facet in (("bankml_exchanges", "history", "x-bankml.history"), ("bankml_memory", "memory", "x-bankml.memory")):
lines = json.loads(psql(f"SELECT coalesce(json_agg(line ORDER BY seq), '[]') FROM {table} WHERE agent = (SELECT v FROM _in WHERE k='slug');",
[("slug", slug)], dsn).strip())
b = "".join(l + "\n" for l in lines).encode()
if hashlib.sha256(b).hexdigest() != want.get(facet):
raise ValueError(f"refused: the stored {key} does not match the manifest's {facet} digest")
private[key] = (b, len(lines))
target = as_slug or slug
d = agents.agent_dir(target)
d.mkdir(parents=True, exist_ok=True)
f = agents.files(target)
f["persona"].write_bytes(persona_b)
f["prompt"].write_text(row["prompt"], encoding="utf-8")
(d / f"{target}.thot.json").write_text(json.dumps(m, indent=1, ensure_ascii=False) + "\n", encoding="utf-8")
agents.rebind(target, derived_from=row["commitments"].get("derived_from"))
n = 0
for key, (b, k) in private.items():
f[key].write_bytes(b)
n += k
return {"slug": target, "from": slug, "thot_cid": row["thot_cid"], "generation": row["generation"], "verified": True, "private_lines": n,
"anchored": bool(expect_content_root), "check": "against the expected contentRoot" if expect_content_root else
"self-consistency (files against the manifest stored beside them)"}