AutonomousAgent / bootstrap_db.py
dehua123's picture
Initial deploy: autonomous ingestion agent (plan→discover→dedupe→download→ingest→report)
11dde75 verified
Raw History Blame Contribute Delete
5.85 kB
r"""
bootstrap_db.py — one-shot fresh-database bootstrap, run at container start.
Purpose: let the agent Space target an EMPTY database that has the exact same
schema as the team's live database, without any manual DBA work. The sandbox
this Space was deployed from cannot reach Postgres, but the Space itself can —
so the clone happens here, on first boot.
Behaviour (idempotent, safe to run every boot):
1. No-op unless BOOTSTRAP_CLONE_FROM is set (the SOURCE database name,
e.g. "AIMDatabase"). DB_NAME is the TARGET (e.g. "aim_agent").
2. If the target database does not exist: CREATE DATABASE.
3. If the target already has the material tables: exit quietly (bootstrapped).
4. Otherwise clone schema only (pg_dump --schema-only --no-owner
--no-privileges | psql): every table, column, index and sequence —
zero rows.
5. Verify: material tables + sources present, all empty, and
pg_mirror.check_schema() passes against the target.
Never fatal: any failure prints a loud warning and the app still starts —
the orchestrator independently refuses to write to an unmigrated schema.
"""
from __future__ import annotations
import os
import subprocess
import sys
import psycopg
from psycopg import sql
MATERIAL_TABLES = ("Polymers", "Fibers", "Composites_materials")
def env(key: str, default: str | None = None) -> str | None:
v = os.environ.get(key)
return v if v not in (None, "") else default
def connect(dbname: str) -> psycopg.Connection:
return psycopg.connect(
host=env("DB_HOST"),
port=int(env("DB_PORT", "5432")),
dbname=dbname,
user=env("DB_USER"),
password=env("DB_PASSWORD"),
sslmode=env("DB_SSLMODE", "require"),
connect_timeout=20,
)
def has_material_tables(conn: psycopg.Connection) -> bool:
with conn.cursor() as cur:
cur.execute(
"SELECT count(*) FROM information_schema.tables "
"WHERE table_schema = 'public' AND table_name = ANY(%s)",
(list(MATERIAL_TABLES),),
)
return cur.fetchone()[0] == len(MATERIAL_TABLES)
def pg_tool_env() -> dict[str, str]:
e = dict(os.environ)
e["PGPASSWORD"] = env("DB_PASSWORD") or ""
e["PGSSLMODE"] = env("DB_SSLMODE", "require")
return e
def main() -> int:
source = env("BOOTSTRAP_CLONE_FROM")
target = env("DB_NAME")
if not source:
print("[bootstrap] BOOTSTRAP_CLONE_FROM not set — skipping.")
return 0
if not target or target == source:
print(f"[bootstrap] refusing: DB_NAME ({target!r}) must differ from "
f"BOOTSTRAP_CLONE_FROM ({source!r}).")
return 0
host = env("DB_HOST")
print(f"[bootstrap] target={target} source={source} host={host}")
# 1. Does the target database exist? (ask via the source DB)
# autocommit from the start: CREATE DATABASE cannot run inside a
# transaction, and psycopg3 otherwise opens one implicitly on first query.
with connect(source) as src:
src.autocommit = True
with src.cursor() as cur:
cur.execute("SELECT 1 FROM pg_database WHERE datname = %s", (target,))
target_exists = cur.fetchone() is not None
if not target_exists:
with src.cursor() as cur:
cur.execute(
sql.SQL("CREATE DATABASE {}").format(sql.Identifier(target)))
print(f"[bootstrap] created database {target}")
# 2. Already bootstrapped?
with connect(target) as tgt:
if has_material_tables(tgt):
print(f"[bootstrap] {target} already has material tables — done.")
return 0
# 3. Schema-only clone.
print(f"[bootstrap] cloning schema {source} -> {target} ...")
port = env("DB_PORT", "5432")
user = env("DB_USER")
dump = subprocess.run(
["pg_dump", "--schema-only", "--no-owner", "--no-privileges",
"-h", host, "-p", port, "-U", user, "-d", source],
env=pg_tool_env(), capture_output=True, text=True, timeout=300,
)
if dump.returncode != 0:
print(f"[bootstrap] WARNING: pg_dump failed:\n{dump.stderr[-2000:]}")
return 0
restore = subprocess.run(
["psql", "-h", host, "-p", port, "-U", user, "-d", target,
"-v", "ON_ERROR_STOP=0", "-q"],
env=pg_tool_env(), input=dump.stdout,
capture_output=True, text=True, timeout=300,
)
errors = [l for l in restore.stderr.splitlines() if "ERROR" in l]
if errors:
print(f"[bootstrap] psql reported {len(errors)} error(s), first few:")
for line in errors[:5]:
print(f" {line}")
# 4. Verify.
with connect(target) as tgt:
if not has_material_tables(tgt):
print("[bootstrap] WARNING: clone finished but material tables "
"missing — agent will fail safe (no writes).")
return 0
with tgt.cursor() as cur:
for t in MATERIAL_TABLES + ("sources",):
try:
cur.execute(
sql.SQL("SELECT count(*) FROM {}").format(sql.Identifier(t)))
print(f"[bootstrap] {t}: {cur.fetchone()[0]} rows")
except psycopg.Error:
tgt.rollback()
print(f"[bootstrap] {t}: MISSING")
try:
import pg_mirror
pg_mirror.check_schema(tgt)
print("[bootstrap] schema check PASSED — fresh empty DB ready.")
except Exception as exc: # noqa: BLE001
print(f"[bootstrap] WARNING: schema check failed: {exc}")
return 0
if __name__ == "__main__":
try:
sys.exit(main())
except Exception as exc: # noqa: BLE001 — never block app start
print(f"[bootstrap] WARNING: unexpected failure: {exc}")
sys.exit(0)