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())