File size: 6,153 Bytes
6993919 d1b47a8 6993919 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 | """Catalog database schema — initial migration.
Creates the tables referenced by the v4.0 catalog models.
Idempotent: safe to run multiple times.
"""
from __future__ import annotations
import asyncio
import os
import asyncpg
from app.core.logging import get_logger
logger = get_logger(__name__)
SCHEMA = """
-- Tokens (T27)
CREATE TABLE IF NOT EXISTS tokens (
token_id TEXT PRIMARY KEY,
chain TEXT NOT NULL,
address TEXT NOT NULL,
symbol TEXT,
name TEXT,
decimals INT,
deployer_wallet_id TEXT,
deployed_at TIMESTAMPTZ NOT NULL,
initial_supply BIGINT,
current_supply BIGINT,
is_honeypot BOOLEAN,
is_mintable BOOLEAN,
is_proxy BOOLEAN,
tax_buy_bps INT,
tax_sell_bps INT,
risk_tier TEXT,
risk_score INT,
risk_factors TEXT[],
rag_embedding_id TEXT,
updated_at TIMESTAMPTZ DEFAULT NOW()
);
CREATE INDEX IF NOT EXISTS idx_tokens_chain ON tokens(chain);
CREATE INDEX IF NOT EXISTS idx_tokens_deployer ON tokens(deployer_wallet_id);
CREATE INDEX IF NOT EXISTS idx_tokens_risk ON tokens(risk_score DESC) WHERE risk_score IS NOT NULL;
-- Alerts
CREATE TABLE IF NOT EXISTS alerts (
alert_id TEXT PRIMARY KEY,
token_id TEXT,
wallet_id TEXT,
chain TEXT,
alert_type TEXT NOT NULL,
severity TEXT NOT NULL,
title TEXT NOT NULL,
description TEXT,
evidence JSONB DEFAULT '{}'::jsonb,
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
resolved_at TIMESTAMPTZ
);
CREATE INDEX IF NOT EXISTS idx_alerts_token ON alerts(token_id);
CREATE INDEX IF NOT EXISTS idx_alerts_wallet ON alerts(wallet_id);
CREATE INDEX IF NOT EXISTS idx_alerts_created ON alerts(created_at DESC);
CREATE INDEX IF NOT EXISTS idx_alerts_severity ON alerts(severity);
-- News items (T28)
CREATE TABLE IF NOT EXISTS news_items (
news_id TEXT PRIMARY KEY,
url TEXT NOT NULL,
title TEXT NOT NULL,
summary TEXT,
body_markdown TEXT,
source TEXT NOT NULL,
published_at TIMESTAMPTZ NOT NULL,
ingested_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
chains_mentioned TEXT[] DEFAULT ARRAY[]::TEXT[],
tokens_mentioned TEXT[] DEFAULT ARRAY[]::TEXT[],
wallets_mentioned TEXT[] DEFAULT ARRAY[]::TEXT[],
sentiment_score REAL,
ai_analysis TEXT,
rag_embedding_id TEXT,
body_tsv tsvector
);
CREATE INDEX IF NOT EXISTS idx_news_published ON news_items(published_at DESC);
CREATE INDEX IF NOT EXISTS idx_news_source ON news_items(source);
CREATE INDEX IF NOT EXISTS idx_news_sentiment ON news_items(sentiment_score);
CREATE INDEX IF NOT EXISTS idx_news_body_tsv ON news_items USING gin(body_tsv);
-- Auto-update the body_tsv column on insert/update
CREATE OR REPLACE FUNCTION news_items_tsv_update() RETURNS trigger AS $$
BEGIN
NEW.body_tsv := to_tsvector('english', COALESCE(NEW.title, '') || ' ' ||
COALESCE(NEW.summary, '') || ' ' ||
COALESCE(NEW.body_markdown, ''));
RETURN NEW;
END
$$ LANGUAGE plpgsql;
DROP TRIGGER IF EXISTS news_items_tsv_trigger ON news_items;
CREATE TRIGGER news_items_tsv_trigger
BEFORE INSERT OR UPDATE ON news_items
FOR EACH ROW EXECUTE FUNCTION news_items_tsv_update();
-- Scan reports (T29)
CREATE TABLE IF NOT EXISTS scan_reports (
report_id TEXT PRIMARY KEY,
subject_type TEXT NOT NULL,
subject_id TEXT NOT NULL,
generated_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
generated_by_model TEXT NOT NULL,
risk_score INT NOT NULL,
risk_tier TEXT NOT NULL,
sections JSONB DEFAULT '{}'::jsonb,
markdown_url TEXT,
paid_via_x402 TEXT
);
CREATE INDEX IF NOT EXISTS idx_reports_subject ON scan_reports(subject_type, subject_id);
CREATE INDEX IF NOT EXISTS idx_reports_generated ON scan_reports(generated_at DESC);
-- RAG findings metadata (Qdrant has the vectors; this is metadata for query)
CREATE TABLE IF NOT EXISTS rag_findings (
finding_id TEXT PRIMARY KEY,
source_type TEXT NOT NULL,
source_url TEXT,
source_token_id TEXT,
source_wallet_id TEXT,
claim TEXT NOT NULL,
confidence REAL NOT NULL,
extracted_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
qdrant_point_id TEXT NOT NULL
);
CREATE INDEX IF NOT EXISTS idx_findings_token ON rag_findings(source_token_id);
CREATE INDEX IF NOT EXISTS idx_findings_wallet ON rag_findings(source_wallet_id);
-- x402 receipts (T34)
CREATE TABLE IF NOT EXISTS x402_receipts (
tx_hash TEXT PRIMARY KEY,
agent_id TEXT,
tool TEXT NOT NULL,
amount_usd REAL NOT NULL,
chain TEXT,
paid_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
tier TEXT
);
CREATE INDEX IF NOT EXISTS idx_x402_paid_at ON x402_receipts(paid_at DESC);
CREATE INDEX IF NOT EXISTS idx_x402_tool ON x402_receipts(tool);
CREATE INDEX IF NOT EXISTS idx_x402_agent ON x402_receipts(agent_id);
"""
async def main():
"""Run the schema migration."""
cfg = {
"host": os.getenv("POSTGRES_HOST", "rmi-postgres"),
"port": int(os.getenv("POSTGRES_PORT", "5432")),
"user": os.getenv("POSTGRES_USER", "rmi"),
"password": os.getenv("POSTGRES_PASSWORD", "RMI_PROD_POSTGRES_2026"),
"database": os.getenv("POSTGRES_DB", "rmi"),
}
logger.info(f"Connecting to postgres at {cfg['host']}:{cfg['port']} db={cfg['database']}")
conn = await asyncpg.connect(**cfg)
try:
await conn.execute(SCHEMA)
logger.info("✓ Schema applied")
# Verify
tables = await conn.fetch(
"SELECT tablename FROM pg_tables WHERE schemaname='public' ORDER BY tablename"
)
logger.info(f"Tables ({len(tables)}):")
for t in tables:
logger.info(f" {t['tablename']}")
finally:
await conn.close()
if __name__ == "__main__":
asyncio.run(main())
|