358 lines
16 KiB
Python
358 lines
16 KiB
Python
"""Async PostgreSQL connection management using asyncpg."""
|
|
|
|
import asyncpg
|
|
from typing import Optional
|
|
|
|
from config import settings
|
|
|
|
_pool: Optional[asyncpg.Pool] = None
|
|
|
|
|
|
async def get_pool() -> asyncpg.Pool:
|
|
"""Get or create the asyncpg connection pool."""
|
|
global _pool
|
|
if _pool is None:
|
|
dsn = settings.DATABASE_URL
|
|
# Convert postgresql+asyncpg:// to asyncpg-compatible DSN
|
|
dsn = dsn.replace("postgresql+asyncpg://", "postgresql://")
|
|
_pool = await asyncpg.create_pool(
|
|
dsn=dsn,
|
|
min_size=2,
|
|
max_size=20,
|
|
command_timeout=60,
|
|
max_inactive_connection_lifetime=300,
|
|
)
|
|
return _pool
|
|
|
|
|
|
async def close_pool() -> None:
|
|
"""Close the connection pool."""
|
|
global _pool
|
|
if _pool is not None:
|
|
await _pool.close()
|
|
_pool = None
|
|
|
|
|
|
# Alias for compatibility
|
|
close_db = close_pool
|
|
|
|
|
|
async def get_connection():
|
|
"""Get a connection from the pool."""
|
|
pool = await get_pool()
|
|
return await pool.acquire()
|
|
|
|
|
|
async def release_connection(conn):
|
|
"""Return a connection to the pool."""
|
|
pool = await get_pool()
|
|
await pool.release(conn)
|
|
|
|
|
|
async def execute_query(sql: str, params: tuple = ()) -> list[dict]:
|
|
"""Execute a query and return results as list of dicts."""
|
|
conn = await get_connection()
|
|
try:
|
|
rows = await conn.fetch(sql, *params)
|
|
return [dict(row) for row in rows]
|
|
finally:
|
|
await release_connection(conn)
|
|
|
|
|
|
async def execute_command(sql: str, params: tuple = ()) -> None:
|
|
"""Execute a command (INSERT/UPDATE/DELETE) without returning rows."""
|
|
conn = await get_connection()
|
|
try:
|
|
await conn.execute(sql, *params)
|
|
finally:
|
|
await release_connection(conn)
|
|
|
|
|
|
async def execute_one(sql: str, params: tuple = ()) -> dict:
|
|
"""Execute a query and return a single row as dict."""
|
|
conn = await get_connection()
|
|
try:
|
|
row = await conn.fetchrow(sql, *params)
|
|
return dict(row) if row else {}
|
|
finally:
|
|
await release_connection(conn)
|
|
|
|
|
|
async def init_db() -> None:
|
|
"""Initialize database tables if they don't exist.
|
|
|
|
This is a safety net — in production, migrations should be managed
|
|
separately. We just create missing tables from the project schema.
|
|
"""
|
|
conn = await get_connection()
|
|
try:
|
|
await conn.execute("""
|
|
DO $$
|
|
BEGIN
|
|
-- stock_profiles
|
|
IF NOT EXISTS (SELECT FROM pg_tables WHERE schemaname='public' AND tablename='stock_profiles') THEN
|
|
CREATE TABLE stock_profiles (
|
|
ticker VARCHAR(20) PRIMARY KEY,
|
|
name VARCHAR(500),
|
|
exchange VARCHAR(20),
|
|
sector VARCHAR(100),
|
|
industry VARCHAR(200),
|
|
market_cap BIGINT,
|
|
description TEXT,
|
|
website TEXT,
|
|
ceo VARCHAR(255),
|
|
employees INTEGER,
|
|
pe_ratio DECIMAL(10,2),
|
|
eps DECIMAL(10,4),
|
|
dividend_yield DECIMAL(8,4),
|
|
beta DECIMAL(6,4),
|
|
last_updated TIMESTAMPTZ DEFAULT NOW()
|
|
);
|
|
END IF;
|
|
|
|
-- prices
|
|
IF NOT EXISTS (SELECT FROM pg_tables WHERE schemaname='public' AND tablename='prices') THEN
|
|
CREATE TABLE prices (
|
|
ticker VARCHAR(20) NOT NULL,
|
|
date TIMESTAMPTZ NOT NULL,
|
|
open DECIMAL(15,4),
|
|
high DECIMAL(15,4),
|
|
low DECIMAL(15,4),
|
|
close DECIMAL(15,4),
|
|
volume BIGINT,
|
|
adjusted_close DECIMAL(15,4),
|
|
PRIMARY KEY (ticker, date)
|
|
);
|
|
END IF;
|
|
|
|
-- watchlists
|
|
IF NOT EXISTS (SELECT FROM pg_tables WHERE schemaname='public' AND tablename='watchlists') THEN
|
|
CREATE TABLE watchlists (
|
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
|
user_id UUID NOT NULL,
|
|
name VARCHAR(255) NOT NULL,
|
|
description TEXT,
|
|
is_default BOOLEAN DEFAULT FALSE,
|
|
created_at TIMESTAMPTZ DEFAULT NOW(),
|
|
updated_at TIMESTAMPTZ DEFAULT NOW()
|
|
);
|
|
CREATE INDEX idx_watchlists_user ON watchlists(user_id);
|
|
END IF;
|
|
|
|
-- watchlist_items
|
|
IF NOT EXISTS (SELECT FROM pg_tables WHERE schemaname='public' AND tablename='watchlist_items') THEN
|
|
CREATE TABLE watchlist_items (
|
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
|
watchlist_id UUID NOT NULL REFERENCES watchlists(id) ON DELETE CASCADE,
|
|
ticker VARCHAR(20) NOT NULL,
|
|
type VARCHAR(20) NOT NULL CHECK (type IN ('stock', 'etf', 'index')),
|
|
custom_notes TEXT,
|
|
added_at TIMESTAMPTZ DEFAULT NOW()
|
|
);
|
|
CREATE UNIQUE INDEX idx_watchlist_items_unique ON watchlist_items(watchlist_id, ticker);
|
|
END IF;
|
|
|
|
-- sec_filings
|
|
IF NOT EXISTS (SELECT FROM pg_tables WHERE schemaname='public' AND tablename='sec_filings') THEN
|
|
CREATE TABLE sec_filings (
|
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
|
ticker VARCHAR(20) NOT NULL,
|
|
cik VARCHAR(20),
|
|
form_type VARCHAR(10) NOT NULL,
|
|
filing_date DATE NOT NULL,
|
|
report_date DATE,
|
|
accession_number VARCHAR(50),
|
|
url TEXT,
|
|
content_summary TEXT,
|
|
key_metrics JSONB,
|
|
sentiment_score DECIMAL(5,4),
|
|
tags TEXT[],
|
|
created_at TIMESTAMPTZ DEFAULT NOW()
|
|
);
|
|
CREATE INDEX idx_sec_filings_ticker ON sec_filings(ticker);
|
|
CREATE INDEX idx_sec_filings_form ON sec_filings(form_type);
|
|
CREATE INDEX idx_sec_filings_date ON sec_filings(filing_date);
|
|
END IF;
|
|
|
|
-- insider_trades
|
|
IF NOT EXISTS (SELECT FROM pg_tables WHERE schemaname='public' AND tablename='insider_trades') THEN
|
|
CREATE TABLE insider_trades (
|
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
|
ticker VARCHAR(20) NOT NULL,
|
|
insider_name VARCHAR(500),
|
|
insider_title VARCHAR(500),
|
|
transaction_date DATE NOT NULL,
|
|
transaction_type VARCHAR(10),
|
|
shares INTEGER,
|
|
price_per_share DECIMAL(10,4),
|
|
total_value DECIMAL(15,4),
|
|
shares_owned_after INTEGER,
|
|
filing_date DATE,
|
|
source_url TEXT,
|
|
created_at TIMESTAMPTZ DEFAULT NOW()
|
|
);
|
|
CREATE INDEX idx_insider_trades_ticker ON insider_trades(ticker);
|
|
CREATE INDEX idx_insider_trades_type ON insider_trades(transaction_type);
|
|
END IF;
|
|
|
|
-- strategies
|
|
IF NOT EXISTS (SELECT FROM pg_tables WHERE schemaname='public' AND tablename='strategies') THEN
|
|
CREATE TABLE strategies (
|
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
|
user_id UUID NOT NULL,
|
|
name VARCHAR(255) NOT NULL,
|
|
description TEXT,
|
|
type VARCHAR(20) NOT NULL CHECK (type IN ('technical', 'fundamental', 'hybrid')),
|
|
conditions JSONB NOT NULL,
|
|
backtest_results JSONB,
|
|
created_at TIMESTAMPTZ DEFAULT NOW(),
|
|
updated_at TIMESTAMPTZ DEFAULT NOW()
|
|
);
|
|
END IF;
|
|
|
|
-- alerts
|
|
IF NOT EXISTS (SELECT FROM pg_tables WHERE schemaname='public' AND tablename='alerts') THEN
|
|
CREATE TABLE alerts (
|
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
|
watchlist_id UUID NOT NULL REFERENCES watchlists(id) ON DELETE CASCADE,
|
|
strategy_id UUID REFERENCES strategies(id),
|
|
ticker VARCHAR(20),
|
|
type VARCHAR(30) NOT NULL,
|
|
trigger_type VARCHAR(50),
|
|
message TEXT NOT NULL,
|
|
severity VARCHAR(10) DEFAULT 'info' CHECK (severity IN ('info', 'warning', 'critical')),
|
|
status VARCHAR(20) DEFAULT 'active' CHECK (status IN ('active', 'resolved', 'dismissed')),
|
|
triggered_at TIMESTAMPTZ DEFAULT NOW(),
|
|
resolved_at TIMESTAMPTZ,
|
|
metadata JSONB DEFAULT '{}'
|
|
);
|
|
CREATE INDEX idx_alerts_watchlist ON alerts(watchlist_id);
|
|
CREATE INDEX idx_alerts_status ON alerts(status);
|
|
CREATE INDEX idx_alerts_triggered ON alerts(triggered_at);
|
|
END IF;
|
|
|
|
-- sector_rotations
|
|
IF NOT EXISTS (SELECT FROM pg_tables WHERE schemaname='public' AND tablename='sector_rotations') THEN
|
|
CREATE TABLE sector_rotations (
|
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
|
detection_date DATE NOT NULL,
|
|
sector_ticker VARCHAR(20) NOT NULL,
|
|
sector_name VARCHAR(100),
|
|
rank_now INTEGER,
|
|
rank_previous INTEGER,
|
|
rank_change INTEGER,
|
|
momentum_20d DECIMAL(8,4),
|
|
momentum_50d DECIMAL(8,4),
|
|
momentum_200d DECIMAL(8,4),
|
|
relative_strength DECIMAL(8,4),
|
|
rotation_signal VARCHAR(20),
|
|
macro_context JSONB,
|
|
analysis_summary TEXT,
|
|
created_at TIMESTAMPTZ DEFAULT NOW()
|
|
);
|
|
CREATE INDEX idx_sector_rotations_date ON sector_rotations(detection_date);
|
|
CREATE INDEX idx_sector_rotations_signal ON sector_rotations(rotation_signal);
|
|
END IF;
|
|
|
|
-- screeners
|
|
IF NOT EXISTS (SELECT FROM pg_tables WHERE schemaname='public' AND tablename='screeners') THEN
|
|
CREATE TABLE screeners (
|
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
|
user_id UUID NOT NULL,
|
|
name VARCHAR(255) NOT NULL,
|
|
description TEXT,
|
|
conditions JSONB NOT NULL,
|
|
results_count INTEGER DEFAULT 0,
|
|
last_run_at TIMESTAMPTZ,
|
|
created_at TIMESTAMPTZ DEFAULT NOW(),
|
|
updated_at TIMESTAMPTZ DEFAULT NOW()
|
|
);
|
|
END IF;
|
|
|
|
-- screener_results
|
|
IF NOT EXISTS (SELECT FROM pg_tables WHERE schemaname='public' AND tablename='screener_results') THEN
|
|
CREATE TABLE screener_results (
|
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
|
screener_id UUID NOT NULL REFERENCES screeners(id) ON DELETE CASCADE,
|
|
ticker VARCHAR(20) NOT NULL,
|
|
match_score DECIMAL(5,4),
|
|
ranked_position INTEGER,
|
|
result_data JSONB,
|
|
generated_at TIMESTAMPTZ DEFAULT NOW()
|
|
);
|
|
CREATE INDEX idx_screener_results_screener ON screener_results(screener_id);
|
|
END IF;
|
|
|
|
-- peer_groups
|
|
IF NOT EXISTS (SELECT FROM pg_tables WHERE schemaname='public' AND tablename='peer_groups') THEN
|
|
CREATE TABLE peer_groups (
|
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
|
ticker VARCHAR(20) NOT NULL,
|
|
peer_ticker VARCHAR(20) NOT NULL,
|
|
similarity_score DECIMAL(5,4),
|
|
created_at TIMESTAMPTZ DEFAULT NOW(),
|
|
UNIQUE (ticker, peer_ticker)
|
|
);
|
|
CREATE INDEX idx_peer_groups_ticker ON peer_groups(ticker);
|
|
END IF;
|
|
|
|
-- users
|
|
IF NOT EXISTS (SELECT FROM pg_tables WHERE schemaname='public' AND tablename='users') THEN
|
|
CREATE TABLE users (
|
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
|
email VARCHAR(255) UNIQUE NOT NULL,
|
|
password_hash VARCHAR(255),
|
|
name VARCHAR(255),
|
|
avatar_url TEXT,
|
|
timezone VARCHAR(50) DEFAULT 'UTC',
|
|
settings JSONB DEFAULT '{}',
|
|
created_at TIMESTAMPTZ DEFAULT NOW(),
|
|
updated_at TIMESTAMPTZ DEFAULT NOW()
|
|
);
|
|
END IF;
|
|
|
|
-- watchlist_strategies junction table
|
|
IF NOT EXISTS (SELECT FROM pg_tables WHERE schemaname='public' AND tablename='watchlist_strategies') THEN
|
|
CREATE TABLE watchlist_strategies (
|
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
|
watchlist_id UUID NOT NULL REFERENCES watchlists(id) ON DELETE CASCADE,
|
|
strategy_id UUID NOT NULL REFERENCES strategies(id) ON DELETE CASCADE,
|
|
is_active BOOLEAN DEFAULT TRUE,
|
|
created_at TIMESTAMPTZ DEFAULT NOW()
|
|
);
|
|
CREATE UNIQUE INDEX idx_watchlist_strategies_unique ON watchlist_strategies(watchlist_id, strategy_id);
|
|
END IF;
|
|
|
|
-- rotation_alert_preferences
|
|
IF NOT EXISTS (SELECT FROM pg_tables WHERE schemaname='public' AND tablename='rotation_alert_preferences') THEN
|
|
CREATE TABLE rotation_alert_preferences (
|
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
|
user_id UUID NOT NULL REFERENCES users(id) ON DELETE CASCADE,
|
|
sectors TEXT[] NOT NULL DEFAULT '{}',
|
|
min_rank_change INTEGER DEFAULT 2,
|
|
email_enabled BOOLEAN DEFAULT TRUE,
|
|
push_enabled BOOLEAN DEFAULT TRUE,
|
|
created_at TIMESTAMPTZ DEFAULT NOW(),
|
|
updated_at TIMESTAMPTZ DEFAULT NOW()
|
|
);
|
|
END IF;
|
|
|
|
-- password_reset_tokens
|
|
IF NOT EXISTS (SELECT FROM pg_tables WHERE schemaname='public' AND tablename='password_reset_tokens') THEN
|
|
CREATE TABLE password_reset_tokens (
|
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
|
user_id UUID NOT NULL REFERENCES users(id) ON DELETE CASCADE,
|
|
token_hash VARCHAR(64) UNIQUE NOT NULL,
|
|
expires_at TIMESTAMPTZ NOT NULL,
|
|
used BOOLEAN DEFAULT FALSE,
|
|
created_at TIMESTAMPTZ DEFAULT NOW()
|
|
);
|
|
CREATE INDEX idx_password_reset_tokens_user ON password_reset_tokens(user_id);
|
|
CREATE INDEX idx_password_reset_tokens_hash ON password_reset_tokens(token_hash);
|
|
CREATE INDEX idx_password_reset_tokens_expires ON password_reset_tokens(expires_at);
|
|
END IF;
|
|
END $$;
|
|
""")
|
|
finally:
|
|
await release_connection(conn)
|