import asyncpg import os import json from typing import Optional DATABASE_URL = os.getenv( "DATABASE_URL", f"postgresql://{os.getenv('POSTGRES_USER', 'postgres')}:{os.getenv('POSTGRES_PASSWORD', 'postgres')}@{os.getenv('POSTGRES_HOST', 'localhost')}:{os.getenv('POSTGRES_PORT', '5432')}/{os.getenv('POSTGRES_DB', 'copykar')}" ) SCHEMA = """ CREATE TABLE IF NOT EXISTS sources ( id SERIAL PRIMARY KEY, channel_id BIGINT UNIQUE NOT NULL, username VARCHAR(255), title VARCHAR(255), is_active BOOLEAN DEFAULT TRUE, created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP ); CREATE TABLE IF NOT EXISTS targets ( id SERIAL PRIMARY KEY, channel_id BIGINT UNIQUE NOT NULL, title VARCHAR(255), username VARCHAR(255), post_interval_min INT DEFAULT 30, personality TEXT DEFAULT '', custom_footer TEXT DEFAULT '', custom_prompt TEXT DEFAULT '', sleep_start_hour INT DEFAULT 0, sleep_end_hour INT DEFAULT 0, is_sleep_enabled BOOLEAN DEFAULT FALSE, auto_source_ids BIGINT[] DEFAULT '{}', language VARCHAR(32) DEFAULT 'fa', last_post_time TIMESTAMPTZ, is_active BOOLEAN DEFAULT TRUE, created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP ); CREATE TABLE IF NOT EXISTS posts ( id BIGSERIAL PRIMARY KEY, source_channel_id BIGINT NOT NULL, source_message_id BIGINT NOT NULL, raw_text TEXT, media_path TEXT, media_type VARCHAR(64), content_hash VARCHAR(128), tags TEXT[] DEFAULT '{}', is_duplicate BOOLEAN DEFAULT FALSE, duplicate_of_id BIGINT REFERENCES posts(id) ON DELETE SET NULL, similarity_reason TEXT, subject VARCHAR(255), ai_text TEXT, suggested_target_id INT REFERENCES targets(id) ON DELETE SET NULL, target_channel_id INT REFERENCES targets(id) ON DELETE SET NULL, status VARCHAR(32) DEFAULT 'pending_review', rejection_reason TEXT DEFAULT '', published_to JSONB DEFAULT '[]'::jsonb, review_message_id BIGINT, scheduled_at TIMESTAMPTZ, published_at TIMESTAMPTZ, created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP, CONSTRAINT unique_source_message UNIQUE (source_channel_id, source_message_id) ); CREATE INDEX IF NOT EXISTS idx_posts_content_hash ON posts(content_hash); CREATE INDEX IF NOT EXISTS idx_posts_status ON posts(status); CREATE INDEX IF NOT EXISTS idx_posts_target_status ON posts(target_channel_id, status); CREATE INDEX IF NOT EXISTS idx_posts_tags ON posts USING GIN (tags); CREATE INDEX IF NOT EXISTS idx_posts_created_at ON posts(created_at DESC); CREATE TABLE IF NOT EXISTS settings ( key VARCHAR(128) PRIMARY KEY, value TEXT NOT NULL, description TEXT ); CREATE TABLE IF NOT EXISTS error_logs ( id BIGSERIAL PRIMARY KEY, service_name VARCHAR(64) NOT NULL, error_type VARCHAR(128) NOT NULL, error_message TEXT NOT NULL, traceback TEXT, context JSONB DEFAULT '{}'::jsonb, resolved BOOLEAN DEFAULT FALSE, resolved_at TIMESTAMPTZ, resolved_note TEXT, created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP ); CREATE INDEX IF NOT EXISTS idx_error_logs_created_at ON error_logs(created_at DESC); CREATE INDEX IF NOT EXISTS idx_error_logs_service ON error_logs(service_name); CREATE TABLE IF NOT EXISTS ai_logs ( id BIGSERIAL PRIMARY KEY, action_name VARCHAR(64), provider VARCHAR(32), model VARCHAR(128), prompt TEXT, system_prompt TEXT, response_text TEXT, duration_sec FLOAT DEFAULT 0.0, status VARCHAR(16) DEFAULT 'success', error_message TEXT, created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP ); CREATE INDEX IF NOT EXISTS idx_ai_logs_created_at ON ai_logs(created_at DESC); CREATE TABLE IF NOT EXISTS ai_providers ( id BIGSERIAL PRIMARY KEY, name VARCHAR(128) NOT NULL, provider_type VARCHAR(32) NOT NULL, base_url TEXT DEFAULT '', api_key TEXT DEFAULT '', model VARCHAR(128) NOT NULL, reasoning_effort VARCHAR(32) DEFAULT '', is_active BOOLEAN DEFAULT FALSE, fallback_provider_id BIGINT REFERENCES ai_providers(id) ON DELETE SET NULL, supports_vision BOOLEAN DEFAULT FALSE, created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP ); CREATE TABLE IF NOT EXISTS channel_categories ( id SERIAL PRIMARY KEY, name VARCHAR(128) NOT NULL, type VARCHAR(32) DEFAULT 'both', description TEXT DEFAULT '', created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP ); CREATE INDEX IF NOT EXISTS idx_ai_providers_active ON ai_providers(is_active); -- Migration safety for existing tables ALTER TABLE ai_providers ADD COLUMN IF NOT EXISTS fallback_provider_id BIGINT REFERENCES ai_providers(id) ON DELETE SET NULL; ALTER TABLE ai_providers ADD COLUMN IF NOT EXISTS supports_vision BOOLEAN DEFAULT FALSE; ALTER TABLE sources ADD COLUMN IF NOT EXISTS category_id INT REFERENCES channel_categories(id) ON DELETE SET NULL; ALTER TABLE targets ADD COLUMN IF NOT EXISTS category_id INT REFERENCES channel_categories(id) ON DELETE SET NULL; ALTER TABLE sources ADD COLUMN IF NOT EXISTS is_active BOOLEAN DEFAULT TRUE; ALTER TABLE targets ADD COLUMN IF NOT EXISTS is_active BOOLEAN DEFAULT TRUE; ALTER TABLE sources ADD COLUMN IF NOT EXISTS created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP; ALTER TABLE targets ADD COLUMN IF NOT EXISTS created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP; CREATE INDEX IF NOT EXISTS idx_sources_category ON sources(category_id); CREATE INDEX IF NOT EXISTS idx_targets_category ON targets(category_id); ALTER TABLE targets ADD COLUMN IF NOT EXISTS personality TEXT DEFAULT ''; ALTER TABLE targets ADD COLUMN IF NOT EXISTS custom_footer TEXT DEFAULT ''; ALTER TABLE targets ADD COLUMN IF NOT EXISTS sleep_start_hour INT DEFAULT 0; ALTER TABLE targets ADD COLUMN IF NOT EXISTS sleep_end_hour INT DEFAULT 0; ALTER TABLE targets ADD COLUMN IF NOT EXISTS is_sleep_enabled BOOLEAN DEFAULT FALSE; ALTER TABLE targets ADD COLUMN IF NOT EXISTS auto_source_ids BIGINT[] DEFAULT '{}'; ALTER TABLE targets ADD COLUMN IF NOT EXISTS language VARCHAR(32) DEFAULT 'fa'; ALTER TABLE targets ADD COLUMN IF NOT EXISTS custom_prompt TEXT DEFAULT ''; ALTER TABLE posts ADD COLUMN IF NOT EXISTS published_to JSONB DEFAULT '[]'::jsonb; ALTER TABLE posts ADD COLUMN IF NOT EXISTS is_deleted BOOLEAN DEFAULT FALSE; ALTER TABLE posts ADD COLUMN IF NOT EXISTS rejection_reason TEXT DEFAULT ''; ALTER TABLE error_logs ADD COLUMN IF NOT EXISTS resolved BOOLEAN DEFAULT FALSE; ALTER TABLE error_logs ADD COLUMN IF NOT EXISTS resolved_at TIMESTAMPTZ; ALTER TABLE error_logs ADD COLUMN IF NOT EXISTS resolved_note TEXT; CREATE INDEX IF NOT EXISTS idx_error_logs_open ON error_logs(resolved, created_at DESC); -- 'pending_ai' and 'approved' belong to an earlier workflow that no longer exists. -- Posts left in those states are invisible to every current code path, so return -- them to the review queue. UPDATE posts SET status = 'pending_review' WHERE status IN ('pending_ai', 'approved') AND is_deleted = FALSE; """ _pool: Optional[asyncpg.Pool] = None async def get_db_pool(dsn: str = DATABASE_URL) -> asyncpg.Pool: global _pool if _pool is None or _pool._closed: _pool = await asyncpg.create_pool(dsn=dsn, min_size=2, max_size=10) return _pool async def close_db_pool(): global _pool if _pool is not None and not _pool._closed: await _pool.close() _pool = None async def init_db(dsn: str = DATABASE_URL): pool = await get_db_pool(dsn) async with pool.acquire() as conn: await conn.execute(SCHEMA)