From 23a5fc77ffa7046679e573c2b57f3b22ccce1de1 Mon Sep 17 00:00:00 2001 From: Dimitri Pisarev Date: Fri, 20 Feb 2026 00:01:43 +0100 Subject: [PATCH] feat: Add event deduplication table and update daily_usage schema - Introduced `event_dedup` table for strict idempotency in event processing. - Modified `daily_usage` table schema to remove unnecessary `id` column and enforce non-negative `qr_count`. - Updated tests to ensure proper handling of deduplication logic and verify SQL execution. - Adjusted mock connection behavior in tests to simulate database interactions accurately. --- alembic/versions/001_initial_schema.py | 658 ++++---- .../repositories/event_store.py | 83 +- .../repositories/qr_history_repository.py | 10 +- db_schema.sql | 1381 ++++++++--------- tests/integration/conftest.py | 21 +- .../event_store/test_postgres_event_store.py | 40 +- 6 files changed, 991 insertions(+), 1202 deletions(-) diff --git a/alembic/versions/001_initial_schema.py b/alembic/versions/001_initial_schema.py index 3a979e6..5d6acfd 100644 --- a/alembic/versions/001_initial_schema.py +++ b/alembic/versions/001_initial_schema.py @@ -1,12 +1,12 @@ -"""Initial database schema with tariffs and extended users table +"""Initial database schema. Revision ID: 001_initial_schema Revises: Create Date: 2026-01-24 16:20:00.000000 - """ from collections.abc import Sequence +from datetime import date from alembic import op @@ -17,22 +17,78 @@ depends_on: str | Sequence[str] | None = None +_PARTITION_START = date(2026, 1, 1) +_PARTITION_MONTHS = 24 + + +def _next_month(dt: date) -> date: + if dt.month == 12: + return date(dt.year + 1, 1, 1) + return date(dt.year, dt.month + 1, 1) + + +def _iter_monthly_ranges(start: date, months: int) -> list[tuple[date, date]]: + ranges: list[tuple[date, date]] = [] + current = start + for _ in range(months): + end = _next_month(current) + ranges.append((current, end)) + current = end + return ranges + + +def _create_monthly_partitions( + *, + table_name: str, + start: date = _PARTITION_START, + months: int = _PARTITION_MONTHS, +) -> None: + for start_date, end_date in _iter_monthly_ranges(start, months): + partition_name = f"{table_name}_{start_date.year:04d}_{start_date.month:02d}" + op.execute( + f""" + CREATE TABLE IF NOT EXISTS {partition_name} + PARTITION OF {table_name} + FOR VALUES FROM ('{start_date.isoformat()}') TO ('{end_date.isoformat()}') + """ + ) + + +def _create_extension_if_possible(name: str) -> None: + op.execute( + f""" + DO $$ + BEGIN + EXECUTE 'CREATE EXTENSION IF NOT EXISTS {name}'; + EXCEPTION + WHEN insufficient_privilege THEN + RAISE NOTICE 'Skipping extension {name}: insufficient privileges'; + END + $$; + """ + ) + + def upgrade() -> None: - """Apply migration""" + """Apply migration.""" # ======================================== - # UTILITY FUNCTIONS (for indexes) + # UTILITY FUNCTIONS # ======================================== - # Create immutable date function for index expressions - op.execute(""" + op.execute( + """ CREATE OR REPLACE FUNCTION date_trunc_immutable(timestamp with time zone) RETURNS date AS $$ SELECT $1::date; $$ LANGUAGE SQL IMMUTABLE - """) + """ + ) - # Create tariffs table - op.execute(""" + # ======================================== + # CORE TABLES + # ======================================== + op.execute( + """ CREATE TABLE IF NOT EXISTS tariffs ( id SERIAL PRIMARY KEY, name VARCHAR(50) NOT NULL UNIQUE, @@ -55,15 +111,24 @@ def upgrade() -> None: created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP, updated_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP ) - """) - - # Create indexes + """ + ) op.execute("CREATE INDEX IF NOT EXISTS idx_tariffs_slug ON tariffs(slug)") op.execute("CREATE INDEX IF NOT EXISTS idx_tariffs_active ON tariffs(is_active)") - # Insert base tariffs (Trial first with id=1 for new users) - op.execute(""" - INSERT INTO tariffs (name, slug, daily_limit, monthly_limit, price_monthly, allowed_formats, history_retention_days, max_dynamic_qr, analytics_retention_days) + op.execute( + """ + INSERT INTO tariffs ( + name, + slug, + daily_limit, + monthly_limit, + price_monthly, + allowed_formats, + history_retention_days, + max_dynamic_qr, + analytics_retention_days + ) VALUES ('Trial', 'trial', 100, NULL, 0, ARRAY['png', 'svg'], 365, 10, 365), ('Freemium', 'freemium', 3, 90, 0, ARRAY['png'], 7, 0, 7), @@ -71,83 +136,56 @@ def upgrade() -> None: ('Premium', 'premium', 100, NULL, 999, ARRAY['png', 'svg'], 365, NULL, 365), ('Enterprise', 'enterprise', NULL, NULL, 2999, ARRAY['png', 'svg'], 365, NULL, 365) ON CONFLICT (slug) DO NOTHING - """) + """ + ) - # Create base users table - op.execute(""" + op.execute( + """ CREATE TABLE IF NOT EXISTS users ( telegram_id BIGINT PRIMARY KEY, - -- Profile language_code VARCHAR(10) NOT NULL DEFAULT 'en', username VARCHAR(255), first_name VARCHAR(255), last_name VARCHAR(255), - -- Metadata + tariff_id INTEGER NOT NULL DEFAULT 1 REFERENCES tariffs(id), + tariff_expires_at TIMESTAMPTZ, + + total_qr_generated INTEGER NOT NULL DEFAULT 0, + last_qr_generated_at TIMESTAMPTZ, + total_qr_deleted INTEGER NOT NULL DEFAULT 0, + last_qr_deleted_at TIMESTAMPTZ, + + is_active BOOLEAN NOT NULL DEFAULT TRUE, + is_banned BOOLEAN NOT NULL DEFAULT FALSE, + created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP, updated_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP ) - """) + """ + ) - # Extend users table (if columns don't exist) - op.execute(""" - DO $$ - BEGIN - IF NOT EXISTS (SELECT 1 FROM information_schema.columns WHERE table_name='users' AND column_name='tariff_id') THEN - ALTER TABLE users ADD COLUMN tariff_id INTEGER NOT NULL DEFAULT 1 REFERENCES tariffs(id); - END IF; - - IF NOT EXISTS (SELECT 1 FROM information_schema.columns WHERE table_name='users' AND column_name='tariff_expires_at') THEN - ALTER TABLE users ADD COLUMN tariff_expires_at TIMESTAMPTZ DEFAULT NULL; - END IF; - - IF NOT EXISTS (SELECT 1 FROM information_schema.columns WHERE table_name='users' AND column_name='total_qr_generated') THEN - ALTER TABLE users ADD COLUMN total_qr_generated INTEGER DEFAULT 0; - END IF; - - IF NOT EXISTS (SELECT 1 FROM information_schema.columns WHERE table_name='users' AND column_name='last_qr_generated_at') THEN - ALTER TABLE users ADD COLUMN last_qr_generated_at TIMESTAMPTZ; - END IF; - - IF NOT EXISTS (SELECT 1 FROM information_schema.columns WHERE table_name='users' AND column_name='total_qr_deleted') THEN - ALTER TABLE users ADD COLUMN total_qr_deleted INTEGER DEFAULT 0; - END IF; - - IF NOT EXISTS (SELECT 1 FROM information_schema.columns WHERE table_name='users' AND column_name='last_qr_deleted_at') THEN - ALTER TABLE users ADD COLUMN last_qr_deleted_at TIMESTAMPTZ; - END IF; - - IF NOT EXISTS (SELECT 1 FROM information_schema.columns WHERE table_name='users' AND column_name='is_active') THEN - ALTER TABLE users ADD COLUMN is_active BOOLEAN DEFAULT TRUE; - END IF; - - IF NOT EXISTS (SELECT 1 FROM information_schema.columns WHERE table_name='users' AND column_name='is_banned') THEN - ALTER TABLE users ADD COLUMN is_banned BOOLEAN DEFAULT FALSE; - END IF; - END $$ - """) - - # Create indexes for users - # O2-01: Optimized indexes for users - op.execute(""" + op.execute( + """ CREATE INDEX IF NOT EXISTS idx_users_active_tariff ON users (tariff_id, last_qr_generated_at DESC) WHERE is_active = TRUE AND is_banned = FALSE - """) - - op.execute(""" + """ + ) + op.execute( + """ CREATE INDEX IF NOT EXISTS idx_users_recently_active ON users (last_qr_generated_at DESC) INCLUDE (telegram_id, language_code, tariff_id) WHERE is_active = TRUE AND last_qr_generated_at IS NOT NULL - """) - + """ + ) op.execute("CREATE INDEX IF NOT EXISTS idx_users_tariff ON users(tariff_id)") op.execute("CREATE INDEX IF NOT EXISTS idx_users_active ON users(is_active)") - # Create qr_codes table - op.execute(""" + op.execute( + """ CREATE TABLE IF NOT EXISTS qr_codes ( id SERIAL PRIMARY KEY, user_id BIGINT NOT NULL REFERENCES users(telegram_id) ON DELETE CASCADE, @@ -165,27 +203,27 @@ def upgrade() -> None: short_code VARCHAR(20) UNIQUE, deleted_at TIMESTAMPTZ DEFAULT NULL, - created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP ) - """) + """ + ) - # Create indexes for qr_codes - # O2-01: Optimized indexes for dashboard and history queries - op.execute(""" + op.execute( + """ CREATE INDEX IF NOT EXISTS idx_qr_user_dashboard ON qr_codes (user_id, is_dynamic, deleted_at) INCLUDE (id, url, format, created_at, short_code) WHERE deleted_at IS NULL - """) - - op.execute(""" + """ + ) + op.execute( + """ CREATE INDEX IF NOT EXISTS idx_qr_user_history_active ON qr_codes (user_id, created_at DESC) INCLUDE (id, url, format, title) WHERE deleted_at IS NULL - """) - + """ + ) op.execute("CREATE INDEX IF NOT EXISTS idx_qr_created ON qr_codes(created_at DESC)") op.execute( "CREATE INDEX IF NOT EXISTS idx_qr_short_code ON qr_codes(short_code) WHERE short_code IS NOT NULL" @@ -194,8 +232,8 @@ def upgrade() -> None: "CREATE INDEX IF NOT EXISTS idx_qr_deleted ON qr_codes(deleted_at) WHERE deleted_at IS NOT NULL" ) - # Create qr_redirects table - op.execute(""" + op.execute( + """ CREATE TABLE IF NOT EXISTS qr_redirects ( id BIGSERIAL PRIMARY KEY, qr_code_id BIGINT NOT NULL UNIQUE REFERENCES qr_codes(id) ON DELETE CASCADE, @@ -212,20 +250,22 @@ def upgrade() -> None: created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP, updated_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP ) - """) + """ + ) - # O2-01: Optimized indexes for redirects - op.execute(""" + op.execute( + """ CREATE INDEX IF NOT EXISTS idx_redirects_lookup ON qr_redirects (short_code) INCLUDE (id, qr_code_id, target_url, is_active, clicks_count) - """) - - op.execute(""" + """ + ) + op.execute( + """ CREATE INDEX IF NOT EXISTS idx_redirects_user_dashboard ON qr_redirects (qr_code_id, is_active, created_at DESC) - """) - + """ + ) op.execute( "CREATE INDEX IF NOT EXISTS idx_redirects_qr_id ON qr_redirects(qr_code_id)" ) @@ -233,8 +273,8 @@ def upgrade() -> None: "CREATE INDEX IF NOT EXISTS idx_redirects_active ON qr_redirects(is_active)" ) - # Create qr_url_changes table - op.execute(""" + op.execute( + """ CREATE TABLE IF NOT EXISTS qr_url_changes ( id BIGSERIAL PRIMARY KEY, qr_redirect_id BIGINT NOT NULL REFERENCES qr_redirects(id) ON DELETE CASCADE, @@ -247,8 +287,8 @@ def upgrade() -> None: change_reason VARCHAR(255) ) - """) - + """ + ) op.execute( "CREATE INDEX IF NOT EXISTS idx_url_changes_redirect ON qr_url_changes(qr_redirect_id, changed_at DESC)" ) @@ -257,12 +297,10 @@ def upgrade() -> None: ) # ======================================== - # O2-02: PARTITIONED TABLES FOR SCALABILITY + # PARTITIONED HIGH-LOAD TABLES # ======================================== - - # Create qr_click_analytics table (PARTITIONED by clicked_at) - # Partition by month for optimal query performance at scale (10M+ rows) - op.execute(""" + op.execute( + """ CREATE TABLE IF NOT EXISTS qr_click_analytics ( id BIGSERIAL, qr_redirect_id BIGINT NOT NULL REFERENCES qr_redirects(id) ON DELETE CASCADE, @@ -271,204 +309,98 @@ def upgrade() -> None: user_agent TEXT, ip_address INET, - country_code VARCHAR(2), city VARCHAR(100), - referer TEXT, - device_type VARCHAR(20), os VARCHAR(50), browser VARCHAR(50), - -- PRIMARY KEY must include partition key for partitioned tables PRIMARY KEY (id, clicked_at) ) PARTITION BY RANGE (clicked_at) - """) - - # O2-02: Create partitions for qr_click_analytics (24 months: 2026-01 to 2028-01) - # Historical partitions (2026) - op.execute("""CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_01 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2026-01-01') TO ('2026-02-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_02 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2026-02-01') TO ('2026-03-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_03 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2026-03-01') TO ('2026-04-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_04 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2026-04-01') TO ('2026-05-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_05 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2026-05-01') TO ('2026-06-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_06 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2026-06-01') TO ('2026-07-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_07 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2026-07-01') TO ('2026-08-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_08 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2026-08-01') TO ('2026-09-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_09 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2026-09-01') TO ('2026-10-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_10 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2026-10-01') TO ('2026-11-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_11 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2026-11-01') TO ('2026-12-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_12 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2026-12-01') TO ('2027-01-01')""") - # Current and future partitions (2027) - op.execute("""CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_01 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2027-01-01') TO ('2027-02-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_02 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2027-02-01') TO ('2027-03-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_03 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2027-03-01') TO ('2027-04-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_04 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2027-04-01') TO ('2027-05-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_05 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2027-05-01') TO ('2027-06-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_06 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2027-06-01') TO ('2027-07-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_07 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2027-07-01') TO ('2027-08-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_08 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2027-08-01') TO ('2027-09-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_09 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2027-09-01') TO ('2027-10-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_10 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2027-10-01') TO ('2027-11-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_11 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2027-11-01') TO ('2027-12-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_12 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2027-12-01') TO ('2028-01-01')""") - - # O2-01: Optimized indexes for click analytics (automatically created on all partitions) + """ + ) + _create_monthly_partitions(table_name="qr_click_analytics") + op.execute( "CREATE INDEX IF NOT EXISTS idx_click_analytics_redirect_time ON qr_click_analytics(qr_redirect_id, clicked_at DESC)" ) - - op.execute(""" + op.execute( + """ CREATE INDEX IF NOT EXISTS idx_click_analytics_device_agg ON qr_click_analytics (qr_redirect_id, device_type) WHERE device_type IS NOT NULL - """) - - op.execute(""" + """ + ) + op.execute( + """ CREATE INDEX IF NOT EXISTS idx_click_analytics_country_agg ON qr_click_analytics (qr_redirect_id, country_code) WHERE country_code IS NOT NULL - """) - - op.execute(""" + """ + ) + op.execute( + """ CREATE INDEX IF NOT EXISTS idx_click_analytics_time_series ON qr_click_analytics (qr_redirect_id, date_trunc_immutable(clicked_at)) - """) - - op.execute(""" + """ + ) + op.execute( + """ CREATE INDEX IF NOT EXISTS idx_click_analytics_unique_check ON qr_click_analytics (qr_redirect_id, ip_address, date_trunc_immutable(clicked_at)) - """) - + """ + ) op.execute( "CREATE INDEX IF NOT EXISTS idx_click_analytics_ip ON qr_click_analytics(ip_address)" ) - - # GeoIP indexes for country analytics and backfill - op.execute(""" + op.execute( + """ CREATE INDEX IF NOT EXISTS idx_click_analytics_geo_backfill ON qr_click_analytics (ip_address, clicked_at) WHERE ip_address IS NOT NULL AND country_code IS NULL - """) - - op.execute(""" + """ + ) + op.execute( + """ CREATE INDEX IF NOT EXISTS idx_click_analytics_country_city ON qr_click_analytics (qr_redirect_id, country_code, city) WHERE country_code IS NOT NULL - """) + """ + ) - # Create daily_usage table (PARTITIONED by usage_date) - # O2-02: Partition by month for fast limit checks and efficient data management - op.execute(""" + op.execute( + """ CREATE TABLE IF NOT EXISTS daily_usage ( - id SERIAL, user_id BIGINT NOT NULL REFERENCES users(telegram_id) ON DELETE CASCADE, usage_date DATE NOT NULL, - qr_count INTEGER DEFAULT 0, + qr_count INTEGER NOT NULL DEFAULT 0 CHECK (qr_count >= 0), created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP, updated_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP, - -- PRIMARY KEY and UNIQUE must include partition key - PRIMARY KEY (id, usage_date), - UNIQUE(user_id, usage_date) + PRIMARY KEY (user_id, usage_date) ) PARTITION BY RANGE (usage_date) - """) - - # O2-02: Create partitions for daily_usage (24 months: 2026-01 to 2027-12) - # Current year partitions (2026) - op.execute("""CREATE TABLE IF NOT EXISTS daily_usage_2026_01 PARTITION OF daily_usage - FOR VALUES FROM ('2026-01-01') TO ('2026-02-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS daily_usage_2026_02 PARTITION OF daily_usage - FOR VALUES FROM ('2026-02-01') TO ('2026-03-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS daily_usage_2026_03 PARTITION OF daily_usage - FOR VALUES FROM ('2026-03-01') TO ('2026-04-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS daily_usage_2026_04 PARTITION OF daily_usage - FOR VALUES FROM ('2026-04-01') TO ('2026-05-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS daily_usage_2026_05 PARTITION OF daily_usage - FOR VALUES FROM ('2026-05-01') TO ('2026-06-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS daily_usage_2026_06 PARTITION OF daily_usage - FOR VALUES FROM ('2026-06-01') TO ('2026-07-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS daily_usage_2026_07 PARTITION OF daily_usage - FOR VALUES FROM ('2026-07-01') TO ('2026-08-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS daily_usage_2026_08 PARTITION OF daily_usage - FOR VALUES FROM ('2026-08-01') TO ('2026-09-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS daily_usage_2026_09 PARTITION OF daily_usage - FOR VALUES FROM ('2026-09-01') TO ('2026-10-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS daily_usage_2026_10 PARTITION OF daily_usage - FOR VALUES FROM ('2026-10-01') TO ('2026-11-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS daily_usage_2026_11 PARTITION OF daily_usage - FOR VALUES FROM ('2026-11-01') TO ('2026-12-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS daily_usage_2026_12 PARTITION OF daily_usage - FOR VALUES FROM ('2026-12-01') TO ('2027-01-01')""") - - # Future partitions (2027) - op.execute("""CREATE TABLE IF NOT EXISTS daily_usage_2027_01 PARTITION OF daily_usage - FOR VALUES FROM ('2027-01-01') TO ('2027-02-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS daily_usage_2027_02 PARTITION OF daily_usage - FOR VALUES FROM ('2027-02-01') TO ('2027-03-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS daily_usage_2027_03 PARTITION OF daily_usage - FOR VALUES FROM ('2027-03-01') TO ('2027-04-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS daily_usage_2027_04 PARTITION OF daily_usage - FOR VALUES FROM ('2027-04-01') TO ('2027-05-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS daily_usage_2027_05 PARTITION OF daily_usage - FOR VALUES FROM ('2027-05-01') TO ('2027-06-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS daily_usage_2027_06 PARTITION OF daily_usage - FOR VALUES FROM ('2027-06-01') TO ('2027-07-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS daily_usage_2027_07 PARTITION OF daily_usage - FOR VALUES FROM ('2027-07-01') TO ('2027-08-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS daily_usage_2027_08 PARTITION OF daily_usage - FOR VALUES FROM ('2027-08-01') TO ('2027-09-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS daily_usage_2027_09 PARTITION OF daily_usage - FOR VALUES FROM ('2027-09-01') TO ('2027-10-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS daily_usage_2027_10 PARTITION OF daily_usage - FOR VALUES FROM ('2027-10-01') TO ('2027-11-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS daily_usage_2027_11 PARTITION OF daily_usage - FOR VALUES FROM ('2027-11-01') TO ('2027-12-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS daily_usage_2027_12 PARTITION OF daily_usage - FOR VALUES FROM ('2027-12-01') TO ('2028-01-01')""") - - # O2-01: Optimized indexes for daily usage (critical for limit checks) - op.execute(""" + """ + ) + _create_monthly_partitions(table_name="daily_usage") + + op.execute( + """ CREATE INDEX IF NOT EXISTS idx_daily_usage_limit_check ON daily_usage (user_id, usage_date) INCLUDE (qr_count) - """) - - op.execute( - "CREATE INDEX IF NOT EXISTS idx_daily_usage_user_date ON daily_usage(user_id, usage_date DESC)" + """ ) op.execute( "CREATE INDEX IF NOT EXISTS idx_daily_usage_date ON daily_usage(usage_date DESC)" ) - # Create subscription_history table - op.execute(""" + # ======================================== + # SUBSCRIPTIONS + # ======================================== + op.execute( + """ CREATE TABLE IF NOT EXISTS subscription_history ( id SERIAL PRIMARY KEY, user_id BIGINT NOT NULL REFERENCES users(telegram_id) ON DELETE CASCADE, @@ -482,8 +414,8 @@ def upgrade() -> None: created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP ) - """) - + """ + ) op.execute( "CREATE INDEX IF NOT EXISTS idx_subscription_user ON subscription_history(user_id, start_date DESC)" ) @@ -491,38 +423,53 @@ def upgrade() -> None: "CREATE INDEX IF NOT EXISTS idx_subscription_tariff ON subscription_history(tariff_id)" ) - # Create update_updated_at_column function - op.execute(""" + # ======================================== + # TRIGGERS/FUNCTIONS + # ======================================== + op.execute( + """ CREATE OR REPLACE FUNCTION update_updated_at_column() RETURNS TRIGGER AS $$ BEGIN NEW.updated_at = CURRENT_TIMESTAMP; RETURN NEW; END; - $$ language 'plpgsql' - """) + $$ LANGUAGE plpgsql + """ + ) - # Create triggers op.execute("DROP TRIGGER IF EXISTS update_users_updated_at ON users") - op.execute(""" + op.execute( + """ CREATE TRIGGER update_users_updated_at BEFORE UPDATE ON users FOR EACH ROW EXECUTE FUNCTION update_updated_at_column() - """) + """ + ) op.execute("DROP TRIGGER IF EXISTS update_tariffs_updated_at ON tariffs") - op.execute(""" + op.execute( + """ CREATE TRIGGER update_tariffs_updated_at BEFORE UPDATE ON tariffs FOR EACH ROW EXECUTE FUNCTION update_updated_at_column() - """) + """ + ) + + op.execute("DROP TRIGGER IF EXISTS update_redirects_updated_at ON qr_redirects") + op.execute( + """ + CREATE TRIGGER update_redirects_updated_at + BEFORE UPDATE ON qr_redirects + FOR EACH ROW EXECUTE FUNCTION update_updated_at_column() + """ + ) # ======================================== - # FEATURE FLAGS TABLES + # FEATURE FLAGS # ======================================== - - # Create feature_flags table - op.execute(""" + op.execute( + """ CREATE TABLE IF NOT EXISTS feature_flags ( id SERIAL PRIMARY KEY, @@ -547,8 +494,8 @@ def upgrade() -> None: created_by BIGINT REFERENCES users(telegram_id), updated_by BIGINT REFERENCES users(telegram_id) ) - """) - + """ + ) op.execute("CREATE INDEX IF NOT EXISTS idx_ff_name ON feature_flags(name)") op.execute("CREATE INDEX IF NOT EXISTS idx_ff_status ON feature_flags(status)") op.execute( @@ -558,14 +505,16 @@ def upgrade() -> None: op.execute( "DROP TRIGGER IF EXISTS update_feature_flags_updated_at ON feature_flags" ) - op.execute(""" + op.execute( + """ CREATE TRIGGER update_feature_flags_updated_at BEFORE UPDATE ON feature_flags FOR EACH ROW EXECUTE FUNCTION update_updated_at_column() - """) + """ + ) - # Create feature_flag_history table - op.execute(""" + op.execute( + """ CREATE TABLE IF NOT EXISTS feature_flag_history ( id BIGSERIAL PRIMARY KEY, @@ -583,8 +532,8 @@ def upgrade() -> None: change_reason TEXT ) - """) - + """ + ) op.execute( "CREATE INDEX IF NOT EXISTS idx_ffh_flag_time ON feature_flag_history(flag_id, changed_at DESC)" ) @@ -599,11 +548,22 @@ def upgrade() -> None: ) # ======================================== - # DOMAIN_EVENTS (Event Store) - PARTITIONED + # EVENT STORE + GLOBAL DEDUP # ======================================== - # O2-02: Partition by occurred_at for efficient event sourcing at scale + op.execute( + """ + CREATE TABLE IF NOT EXISTS event_dedup ( + event_id UUID PRIMARY KEY, + created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP + ) + """ + ) + op.execute( + "CREATE INDEX IF NOT EXISTS idx_event_dedup_created_at ON event_dedup(created_at DESC)" + ) - op.execute(""" + op.execute( + """ CREATE TABLE IF NOT EXISTS domain_events ( id BIGSERIAL, event_id UUID NOT NULL, @@ -622,84 +582,31 @@ def upgrade() -> None: processed_at TIMESTAMPTZ, processing_error TEXT, - -- PRIMARY KEY and UNIQUE must include partition key PRIMARY KEY (id, occurred_at), UNIQUE (event_id, occurred_at) ) PARTITION BY RANGE (occurred_at) - """) - - # O2-02: Create partitions for domain_events (24 months: 2026-01 to 2027-12) - # Current year partitions (2026) - op.execute("""CREATE TABLE IF NOT EXISTS domain_events_2026_01 PARTITION OF domain_events - FOR VALUES FROM ('2026-01-01') TO ('2026-02-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS domain_events_2026_02 PARTITION OF domain_events - FOR VALUES FROM ('2026-02-01') TO ('2026-03-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS domain_events_2026_03 PARTITION OF domain_events - FOR VALUES FROM ('2026-03-01') TO ('2026-04-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS domain_events_2026_04 PARTITION OF domain_events - FOR VALUES FROM ('2026-04-01') TO ('2026-05-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS domain_events_2026_05 PARTITION OF domain_events - FOR VALUES FROM ('2026-05-01') TO ('2026-06-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS domain_events_2026_06 PARTITION OF domain_events - FOR VALUES FROM ('2026-06-01') TO ('2026-07-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS domain_events_2026_07 PARTITION OF domain_events - FOR VALUES FROM ('2026-07-01') TO ('2026-08-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS domain_events_2026_08 PARTITION OF domain_events - FOR VALUES FROM ('2026-08-01') TO ('2026-09-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS domain_events_2026_09 PARTITION OF domain_events - FOR VALUES FROM ('2026-09-01') TO ('2026-10-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS domain_events_2026_10 PARTITION OF domain_events - FOR VALUES FROM ('2026-10-01') TO ('2026-11-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS domain_events_2026_11 PARTITION OF domain_events - FOR VALUES FROM ('2026-11-01') TO ('2026-12-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS domain_events_2026_12 PARTITION OF domain_events - FOR VALUES FROM ('2026-12-01') TO ('2027-01-01')""") - - # Future partitions (2027) - op.execute("""CREATE TABLE IF NOT EXISTS domain_events_2027_01 PARTITION OF domain_events - FOR VALUES FROM ('2027-01-01') TO ('2027-02-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS domain_events_2027_02 PARTITION OF domain_events - FOR VALUES FROM ('2027-02-01') TO ('2027-03-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS domain_events_2027_03 PARTITION OF domain_events - FOR VALUES FROM ('2027-03-01') TO ('2027-04-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS domain_events_2027_04 PARTITION OF domain_events - FOR VALUES FROM ('2027-04-01') TO ('2027-05-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS domain_events_2027_05 PARTITION OF domain_events - FOR VALUES FROM ('2027-05-01') TO ('2027-06-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS domain_events_2027_06 PARTITION OF domain_events - FOR VALUES FROM ('2027-06-01') TO ('2027-07-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS domain_events_2027_07 PARTITION OF domain_events - FOR VALUES FROM ('2027-07-01') TO ('2027-08-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS domain_events_2027_08 PARTITION OF domain_events - FOR VALUES FROM ('2027-08-01') TO ('2027-09-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS domain_events_2027_09 PARTITION OF domain_events - FOR VALUES FROM ('2027-09-01') TO ('2027-10-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS domain_events_2027_10 PARTITION OF domain_events - FOR VALUES FROM ('2027-10-01') TO ('2027-11-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS domain_events_2027_11 PARTITION OF domain_events - FOR VALUES FROM ('2027-11-01') TO ('2027-12-01')""") - op.execute("""CREATE TABLE IF NOT EXISTS domain_events_2027_12 PARTITION OF domain_events - FOR VALUES FROM ('2027-12-01') TO ('2028-01-01')""") - - # Event store indexes - # O2-01: Optimized indexes for domain events + """ + ) + _create_monthly_partitions(table_name="domain_events") + op.execute( "CREATE INDEX IF NOT EXISTS idx_de_aggregate ON domain_events(aggregate_type, aggregate_id, id)" ) - - op.execute(""" + op.execute( + """ CREATE INDEX IF NOT EXISTS idx_de_aggregate_fast_replay ON domain_events (aggregate_type, aggregate_id, occurred_at) INCLUDE (event_name, event_version) WHERE processed = TRUE - """) - - op.execute(""" + """ + ) + op.execute( + """ CREATE INDEX IF NOT EXISTS idx_de_analytics ON domain_events (event_name, occurred_at) INCLUDE (aggregate_id, payload) - """) - + """ + ) op.execute( "CREATE INDEX IF NOT EXISTS idx_de_event_name ON domain_events(event_name, occurred_at DESC)" ) @@ -712,31 +619,22 @@ def upgrade() -> None: op.execute("CREATE INDEX IF NOT EXISTS idx_de_event_id ON domain_events(event_id)") # ======================================== - # O2-01: AUTOVACUUM CONFIGURATION + # VACUUM + EXTENSIONS + PARTITION CRON # ======================================== - # Configure autovacuum for high-traffic tables to prevent index bloat - # Note: Partitioned tables (qr_click_analytics, domain_events) don't support - # storage parameters - they must be set on individual partitions - - op.execute(""" + op.execute( + """ ALTER TABLE qr_codes SET ( autovacuum_vacuum_scale_factor = 0.1, autovacuum_analyze_scale_factor = 0.05 ) - """) - - # Enable pg_stat_statements for query performance monitoring - op.execute("CREATE EXTENSION IF NOT EXISTS pg_stat_statements") - - # Enable pgstattuple for index bloat monitoring - op.execute("CREATE EXTENSION IF NOT EXISTS pgstattuple") + """ + ) - # ======================================== - # O2-02: AUTOMATIC PARTITION MANAGEMENT - # ======================================== - # Function to automatically create next month's partitions for all partitioned tables + _create_extension_if_possible("pg_stat_statements") + _create_extension_if_possible("pgstattuple") - op.execute(""" + op.execute( + """ CREATE OR REPLACE FUNCTION create_monthly_partitions() RETURNS void AS $$ DECLARE @@ -745,11 +643,9 @@ def upgrade() -> None: start_date DATE; end_date DATE; BEGIN - -- Calculate next month start_date := date_trunc('month', CURRENT_DATE + INTERVAL '1 month'); end_date := start_date + INTERVAL '1 month'; - -- Create partitions for each partitioned table FOREACH table_name IN ARRAY ARRAY[ 'qr_click_analytics', 'daily_usage', @@ -757,9 +653,9 @@ def upgrade() -> None: ] LOOP partition_name := table_name || '_' || to_char(start_date, 'YYYY_MM'); - -- Check if partition exists IF NOT EXISTS ( - SELECT 1 FROM pg_class + SELECT 1 + FROM pg_class WHERE relname = partition_name ) THEN EXECUTE format( @@ -773,37 +669,38 @@ def upgrade() -> None: END IF; END LOOP; END; - $$ LANGUAGE plpgsql; - """) + $$ LANGUAGE plpgsql + """ + ) - # Create comment explaining usage - op.execute(""" + op.execute( + """ COMMENT ON FUNCTION create_monthly_partitions() IS 'Automatically creates next month partition for qr_click_analytics, daily_usage, and domain_events tables. Schedule this function to run monthly using pg_cron or external scheduler: Example: SELECT cron.schedule(''create-partitions'', ''0 0 25 * *'', ''SELECT create_monthly_partitions()'');' - """) + """ + ) def downgrade() -> None: - """Revert migration""" + """Revert migration.""" + + op.execute("DROP FUNCTION IF EXISTS create_monthly_partitions()") - # Drop domain_events table op.execute("DROP TABLE IF EXISTS domain_events CASCADE") + op.execute("DROP TABLE IF EXISTS event_dedup CASCADE") - # Drop feature flags tables op.execute( "DROP TRIGGER IF EXISTS update_feature_flags_updated_at ON feature_flags" ) op.execute("DROP TABLE IF EXISTS feature_flag_history CASCADE") op.execute("DROP TABLE IF EXISTS feature_flags CASCADE") - # Drop triggers + op.execute("DROP TRIGGER IF EXISTS update_redirects_updated_at ON qr_redirects") op.execute("DROP TRIGGER IF EXISTS update_tariffs_updated_at ON tariffs") op.execute("DROP TRIGGER IF EXISTS update_users_updated_at ON users") - op.execute("DROP FUNCTION IF EXISTS update_updated_at_column()") - # Drop tables op.execute("DROP TABLE IF EXISTS subscription_history CASCADE") op.execute("DROP TABLE IF EXISTS daily_usage CASCADE") op.execute("DROP TABLE IF EXISTS qr_click_analytics CASCADE") @@ -811,17 +708,8 @@ def downgrade() -> None: op.execute("DROP TABLE IF EXISTS qr_redirects CASCADE") op.execute("DROP TABLE IF EXISTS qr_codes CASCADE") - # Drop added columns from users - op.execute(""" - ALTER TABLE users - DROP COLUMN IF EXISTS is_banned, - DROP COLUMN IF EXISTS is_active, - DROP COLUMN IF EXISTS last_qr_deleted_at, - DROP COLUMN IF EXISTS total_qr_deleted, - DROP COLUMN IF EXISTS last_qr_generated_at, - DROP COLUMN IF EXISTS total_qr_generated, - DROP COLUMN IF EXISTS tariff_expires_at, - DROP COLUMN IF EXISTS tariff_id - """) - + op.execute("DROP TABLE IF EXISTS users CASCADE") op.execute("DROP TABLE IF EXISTS tariffs CASCADE") + + op.execute("DROP FUNCTION IF EXISTS update_updated_at_column()") + op.execute("DROP FUNCTION IF EXISTS date_trunc_immutable(timestamp with time zone)") diff --git a/bot/infrastructure/repositories/event_store.py b/bot/infrastructure/repositories/event_store.py index 6cba9ef..1eba894 100644 --- a/bot/infrastructure/repositories/event_store.py +++ b/bot/infrastructure/repositories/event_store.py @@ -83,6 +83,42 @@ def __init__(self, pool: asyncpg.Pool): """ self._pool = pool + async def _claim_event_id(self, conn: asyncpg.Connection, event_id: UUID) -> bool: + """Reserve event_id in global dedup table. + + Returns: + True if reservation created, False if event_id already exists. + """ + result = await conn.execute( + """ + INSERT INTO event_dedup (event_id) + VALUES ($1) + ON CONFLICT DO NOTHING + """, + event_id, + ) + return result == "INSERT 0 1" + + async def _insert_event(self, conn: asyncpg.Connection, event: BaseEvent) -> int: + """Insert event payload into partitioned event log.""" + row_id = await conn.fetchval( + """ + INSERT INTO domain_events + (event_id, event_name, event_version, + aggregate_type, aggregate_id, payload, occurred_at) + VALUES ($1, $2, $3, $4, $5, $6, $7) + RETURNING id + """, + event.event_id, + event.get_event_name(), + event.get_event_version(), + event.get_aggregate_type(), + event.aggregate_id, + json.dumps(event._get_payload()), + event.occurred_at, + ) + return int(row_id) + async def append(self, event: BaseEvent) -> int: """Append single event to the store. @@ -95,25 +131,13 @@ async def append(self, event: BaseEvent) -> int: Raises: DuplicateEventError: If event_id already exists """ - async with self._pool.acquire() as conn: + async with self._pool.acquire() as conn, conn.transaction(): try: - row_id = await conn.fetchval( - """ - INSERT INTO domain_events - (event_id, event_name, event_version, - aggregate_type, aggregate_id, payload, occurred_at) - VALUES ($1, $2, $3, $4, $5, $6, $7) - RETURNING id - """, - event.event_id, - event.get_event_name(), - event.get_event_version(), - event.get_aggregate_type(), - event.aggregate_id, - json.dumps(event._get_payload()), - event.occurred_at, - ) + claimed = await self._claim_event_id(conn, event.event_id) + if not claimed: + raise DuplicateEventError(event.event_id) + row_id = await self._insert_event(conn, event) logger.debug( "Stored event %s (id=%d, event_id=%s)", event.get_event_name(), @@ -121,8 +145,8 @@ async def append(self, event: BaseEvent) -> int: event.event_id, ) return row_id - except asyncpg.UniqueViolationError as e: + # Defensive mapping for legacy unique constraints. raise DuplicateEventError(event.event_id) from e async def append_batch(self, events: Sequence[BaseEvent]) -> list[int]: @@ -144,25 +168,14 @@ async def append_batch(self, events: Sequence[BaseEvent]) -> list[int]: return [] async with self._pool.acquire() as conn, conn.transaction(): - ids = [] + ids: list[int] = [] for event in events: try: - row_id = await conn.fetchval( - """ - INSERT INTO domain_events - (event_id, event_name, event_version, - aggregate_type, aggregate_id, payload, occurred_at) - VALUES ($1, $2, $3, $4, $5, $6, $7) - RETURNING id - """, - event.event_id, - event.get_event_name(), - event.get_event_version(), - event.get_aggregate_type(), - event.aggregate_id, - json.dumps(event._get_payload()), - event.occurred_at, - ) + claimed = await self._claim_event_id(conn, event.event_id) + if not claimed: + raise DuplicateEventError(event.event_id) + + row_id = await self._insert_event(conn, event) ids.append(row_id) except asyncpg.UniqueViolationError as e: raise DuplicateEventError(event.event_id) from e diff --git a/bot/infrastructure/repositories/qr_history_repository.py b/bot/infrastructure/repositories/qr_history_repository.py index 1733214..0575fb3 100644 --- a/bot/infrastructure/repositories/qr_history_repository.py +++ b/bot/infrastructure/repositories/qr_history_repository.py @@ -130,14 +130,8 @@ async def increment_daily_usage( raise ValueError("increment must be positive") query = """ - WITH seq AS ( - SELECT COALESCE( - pg_get_serial_sequence('daily_usage', 'id'), - 'daily_usage_id_seq' - ) AS seq_name - ) - INSERT INTO daily_usage (id, user_id, usage_date, qr_count) - VALUES (nextval((SELECT seq_name FROM seq)::regclass), $1, $2, $3) + INSERT INTO daily_usage (user_id, usage_date, qr_count) + VALUES ($1, $2, $3) ON CONFLICT (user_id, usage_date) DO UPDATE SET qr_count = daily_usage.qr_count + EXCLUDED.qr_count, diff --git a/db_schema.sql b/db_schema.sql index 2bf26a6..c236888 100644 --- a/db_schema.sql +++ b/db_schema.sql @@ -1,843 +1,738 @@ --- QR Code Bot Database Schema --- Идеальная архитектура с тарифами и лимитами +-- Generated from alembic/versions/001_initial_schema.py +-- Source of truth: Alembic revision 001_initial_schema --- ============================================ --- UTILITY FUNCTIONS (for indexes) --- ============================================ --- Create immutable date function for index expressions CREATE OR REPLACE FUNCTION date_trunc_immutable(timestamp with time zone) -RETURNS date AS $$ - SELECT $1::date; -$$ LANGUAGE SQL IMMUTABLE; + RETURNS date AS $$ + SELECT $1::date; + $$ LANGUAGE SQL IMMUTABLE; --- ============================================ --- 1. TARIFFS (Тарифные планы) --- ============================================ CREATE TABLE IF NOT EXISTS tariffs ( - id SERIAL PRIMARY KEY, - name VARCHAR(50) NOT NULL UNIQUE, - slug VARCHAR(50) NOT NULL UNIQUE, - - -- Лимиты (NULL = unlimited) - daily_limit INTEGER DEFAULT NULL, - monthly_limit INTEGER DEFAULT NULL, - max_url_length INTEGER NOT NULL DEFAULT 2048, - max_dynamic_qr INTEGER DEFAULT NULL, -- NULL = unlimited, >0 = limited, 0 = disabled - - -- Возможности - allowed_formats TEXT[] DEFAULT ARRAY['png', 'svg'], - history_retention_days INTEGER DEFAULT 30, -- сколько дней хранить историю - analytics_retention_days INTEGER DEFAULT 30, -- сколько дней хранить аналитику - priority_support BOOLEAN DEFAULT FALSE, - custom_branding BOOLEAN DEFAULT FALSE, - - -- Цена (в центах для точности) - price_monthly INTEGER DEFAULT 0, - - -- Метаданные - is_active BOOLEAN DEFAULT TRUE, - created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP, - updated_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP -); - --- Индексы + id SERIAL PRIMARY KEY, + name VARCHAR(50) NOT NULL UNIQUE, + slug VARCHAR(50) NOT NULL UNIQUE, + + daily_limit INTEGER DEFAULT NULL, + monthly_limit INTEGER DEFAULT NULL, + max_url_length INTEGER NOT NULL DEFAULT 2048, + max_dynamic_qr INTEGER DEFAULT NULL, + + allowed_formats TEXT[] DEFAULT ARRAY['png', 'svg'], + history_retention_days INTEGER DEFAULT 30, + analytics_retention_days INTEGER DEFAULT 30, + priority_support BOOLEAN DEFAULT FALSE, + custom_branding BOOLEAN DEFAULT FALSE, + + price_monthly INTEGER DEFAULT 0, + + is_active BOOLEAN DEFAULT TRUE, + created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP, + updated_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP + ); + CREATE INDEX IF NOT EXISTS idx_tariffs_slug ON tariffs(slug); -CREATE INDEX IF NOT EXISTS idx_tariffs_active ON tariffs(is_active); --- Базовые тарифы (цены в центах) --- Trial: 30-day full access for new users, auto-downgrades to Freemium -INSERT INTO tariffs (name, slug, daily_limit, monthly_limit, price_monthly, allowed_formats, history_retention_days, max_dynamic_qr, analytics_retention_days) VALUES - ('Trial', 'trial', 100, NULL, 0, ARRAY['png', 'svg'], 365, 10, 365), - ('Freemium', 'freemium', 3, 90, 0, ARRAY['png'], 7, 0, 7), - ('Basic', 'basic', 20, 600, 499, ARRAY['png', 'svg'], 30, 5, 30), - ('Premium', 'premium', 100, NULL, 999, ARRAY['png', 'svg'], 365, NULL, 365), - ('Enterprise', 'enterprise', NULL, NULL, 2999, ARRAY['png', 'svg'], 365, NULL, 365) -ON CONFLICT (slug) DO NOTHING; +CREATE INDEX IF NOT EXISTS idx_tariffs_active ON tariffs(is_active); +INSERT INTO tariffs ( + name, + slug, + daily_limit, + monthly_limit, + price_monthly, + allowed_formats, + history_retention_days, + max_dynamic_qr, + analytics_retention_days + ) + VALUES + ('Trial', 'trial', 100, NULL, 0, ARRAY['png', 'svg'], 365, 10, 365), + ('Freemium', 'freemium', 3, 90, 0, ARRAY['png'], 7, 0, 7), + ('Basic', 'basic', 20, 600, 499, ARRAY['png', 'svg'], 30, 5, 30), + ('Premium', 'premium', 100, NULL, 999, ARRAY['png', 'svg'], 365, NULL, 365), + ('Enterprise', 'enterprise', NULL, NULL, 2999, ARRAY['png', 'svg'], 365, NULL, 365) + ON CONFLICT (slug) DO NOTHING; --- ============================================ --- 2. USERS (обновленная таблица) --- ============================================ CREATE TABLE IF NOT EXISTS users ( - telegram_id BIGINT PRIMARY KEY, - - -- Профиль - language_code VARCHAR(10) NOT NULL DEFAULT 'en', - username VARCHAR(255), - first_name VARCHAR(255), - last_name VARCHAR(255), - - -- Метаданные - created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP, - updated_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP -); - --- Добавить новые колонки если их нет (миграция существующей таблицы) -DO $$ -BEGIN - -- Тариф - IF NOT EXISTS (SELECT 1 FROM information_schema.columns WHERE table_name='users' AND column_name='tariff_id') THEN - ALTER TABLE users ADD COLUMN tariff_id INTEGER NOT NULL DEFAULT 1 REFERENCES tariffs(id); - END IF; - - IF NOT EXISTS (SELECT 1 FROM information_schema.columns WHERE table_name='users' AND column_name='tariff_expires_at') THEN - ALTER TABLE users ADD COLUMN tariff_expires_at TIMESTAMPTZ DEFAULT NULL; - END IF; - - -- Статистика - IF NOT EXISTS (SELECT 1 FROM information_schema.columns WHERE table_name='users' AND column_name='total_qr_generated') THEN - ALTER TABLE users ADD COLUMN total_qr_generated INTEGER DEFAULT 0; - END IF; - - IF NOT EXISTS (SELECT 1 FROM information_schema.columns WHERE table_name='users' AND column_name='last_qr_generated_at') THEN - ALTER TABLE users ADD COLUMN last_qr_generated_at TIMESTAMPTZ; - END IF; - - IF NOT EXISTS (SELECT 1 FROM information_schema.columns WHERE table_name='users' AND column_name='total_qr_deleted') THEN - ALTER TABLE users ADD COLUMN total_qr_deleted INTEGER DEFAULT 0; - END IF; - - IF NOT EXISTS (SELECT 1 FROM information_schema.columns WHERE table_name='users' AND column_name='last_qr_deleted_at') THEN - ALTER TABLE users ADD COLUMN last_qr_deleted_at TIMESTAMPTZ; - END IF; - - -- Флаги - IF NOT EXISTS (SELECT 1 FROM information_schema.columns WHERE table_name='users' AND column_name='is_active') THEN - ALTER TABLE users ADD COLUMN is_active BOOLEAN DEFAULT TRUE; - END IF; - - IF NOT EXISTS (SELECT 1 FROM information_schema.columns WHERE table_name='users' AND column_name='is_banned') THEN - ALTER TABLE users ADD COLUMN is_banned BOOLEAN DEFAULT FALSE; - END IF; -END $$; - --- Индексы --- O2-01: Optimized indexes for users + telegram_id BIGINT PRIMARY KEY, + + language_code VARCHAR(10) NOT NULL DEFAULT 'en', + username VARCHAR(255), + first_name VARCHAR(255), + last_name VARCHAR(255), + + tariff_id INTEGER NOT NULL DEFAULT 1 REFERENCES tariffs(id), + tariff_expires_at TIMESTAMPTZ, + + total_qr_generated INTEGER NOT NULL DEFAULT 0, + last_qr_generated_at TIMESTAMPTZ, + total_qr_deleted INTEGER NOT NULL DEFAULT 0, + last_qr_deleted_at TIMESTAMPTZ, + + is_active BOOLEAN NOT NULL DEFAULT TRUE, + is_banned BOOLEAN NOT NULL DEFAULT FALSE, + + created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP, + updated_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP + ); + CREATE INDEX IF NOT EXISTS idx_users_active_tariff - ON users (tariff_id, last_qr_generated_at DESC) - WHERE is_active = TRUE AND is_banned = FALSE; + ON users (tariff_id, last_qr_generated_at DESC) + WHERE is_active = TRUE AND is_banned = FALSE; CREATE INDEX IF NOT EXISTS idx_users_recently_active - ON users (last_qr_generated_at DESC) - INCLUDE (telegram_id, language_code, tariff_id) - WHERE is_active = TRUE AND last_qr_generated_at IS NOT NULL; + ON users (last_qr_generated_at DESC) + INCLUDE (telegram_id, language_code, tariff_id) + WHERE is_active = TRUE AND last_qr_generated_at IS NOT NULL; -CREATE INDEX IF NOT EXISTS idx_users_language ON users(language_code); CREATE INDEX IF NOT EXISTS idx_users_tariff ON users(tariff_id); -CREATE INDEX IF NOT EXISTS idx_users_active ON users(is_active); +CREATE INDEX IF NOT EXISTS idx_users_active ON users(is_active); --- ============================================ --- 3. QR_CODES (История генераций) --- ============================================ CREATE TABLE IF NOT EXISTS qr_codes ( - id SERIAL PRIMARY KEY, - user_id BIGINT NOT NULL REFERENCES users(telegram_id) ON DELETE CASCADE, - - -- Данные QR кода - url TEXT NOT NULL, - url_hash VARCHAR(64) NOT NULL, -- SHA256 для дедупликации - format VARCHAR(10) NOT NULL DEFAULT 'png', - title VARCHAR(255) DEFAULT NULL, -- User-friendly title for the QR code - - -- Метаданные - file_size_bytes INTEGER, - generation_time_ms INTEGER, -- время генерации в миллисекундах - telegram_message_id BIGINT, -- ID сообщения в Telegram для возможности удаления - - -- Динамические QR коды (для будущего функционала) - is_dynamic BOOLEAN DEFAULT FALSE, - short_code VARCHAR(20) UNIQUE, - - -- Soft delete - deleted_at TIMESTAMPTZ DEFAULT NULL, - - created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP -); - --- Индексы для быстрого поиска --- O2-01: Optimized indexes for dashboard and history queries + id SERIAL PRIMARY KEY, + user_id BIGINT NOT NULL REFERENCES users(telegram_id) ON DELETE CASCADE, + + url TEXT NOT NULL, + url_hash VARCHAR(64) NOT NULL, + format VARCHAR(10) NOT NULL DEFAULT 'png', + title VARCHAR(255) DEFAULT NULL, + + file_size_bytes INTEGER, + generation_time_ms INTEGER, + telegram_message_id BIGINT, + + is_dynamic BOOLEAN DEFAULT FALSE, + short_code VARCHAR(20) UNIQUE, + + deleted_at TIMESTAMPTZ DEFAULT NULL, + created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP + ); + CREATE INDEX IF NOT EXISTS idx_qr_user_dashboard - ON qr_codes (user_id, is_dynamic, deleted_at) - INCLUDE (id, url, format, created_at, short_code) - WHERE deleted_at IS NULL; + ON qr_codes (user_id, is_dynamic, deleted_at) + INCLUDE (id, url, format, created_at, short_code) + WHERE deleted_at IS NULL; CREATE INDEX IF NOT EXISTS idx_qr_user_history_active - ON qr_codes (user_id, created_at DESC) - INCLUDE (id, url, format, title) - WHERE deleted_at IS NULL; + ON qr_codes (user_id, created_at DESC) + INCLUDE (id, url, format, title) + WHERE deleted_at IS NULL; CREATE INDEX IF NOT EXISTS idx_qr_created ON qr_codes(created_at DESC); + CREATE INDEX IF NOT EXISTS idx_qr_short_code ON qr_codes(short_code) WHERE short_code IS NOT NULL; + CREATE INDEX IF NOT EXISTS idx_qr_deleted ON qr_codes(deleted_at) WHERE deleted_at IS NOT NULL; --- Note: Removed idx_qr_user_date_only as it caused IMMUTABLE function issues --- We'll use daily_usage table for fast limit checks instead +CREATE TABLE IF NOT EXISTS qr_redirects ( + id BIGSERIAL PRIMARY KEY, + qr_code_id BIGINT NOT NULL UNIQUE REFERENCES qr_codes(id) ON DELETE CASCADE, + short_code VARCHAR(12) NOT NULL UNIQUE, + original_url TEXT NOT NULL, + target_url TEXT NOT NULL, + + clicks_count INTEGER DEFAULT 0, + unique_clicks_count INTEGER DEFAULT 0, + last_clicked_at TIMESTAMPTZ, + + is_active BOOLEAN DEFAULT TRUE, + created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP, + updated_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP + ); --- ============================================ --- 4. QR_REDIRECTS (Dynamic QR redirects) --- ============================================ -CREATE TABLE IF NOT EXISTS qr_redirects ( - id BIGSERIAL PRIMARY KEY, - qr_code_id BIGINT NOT NULL UNIQUE REFERENCES qr_codes(id) ON DELETE CASCADE, - short_code VARCHAR(12) NOT NULL UNIQUE, - - -- URLs - original_url TEXT NOT NULL, - target_url TEXT NOT NULL, - - -- Analytics - clicks_count INTEGER DEFAULT 0, - unique_clicks_count INTEGER DEFAULT 0, - last_clicked_at TIMESTAMPTZ, - - -- Status - is_active BOOLEAN DEFAULT TRUE, - created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP, - updated_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP -); - --- Indexes --- O2-01: Optimized indexes for redirects CREATE INDEX IF NOT EXISTS idx_redirects_lookup - ON qr_redirects (short_code) - INCLUDE (id, qr_code_id, target_url, is_active, clicks_count); + ON qr_redirects (short_code) + INCLUDE (id, qr_code_id, target_url, is_active, clicks_count); CREATE INDEX IF NOT EXISTS idx_redirects_user_dashboard - ON qr_redirects (qr_code_id, is_active, created_at DESC); + ON qr_redirects (qr_code_id, is_active, created_at DESC); CREATE INDEX IF NOT EXISTS idx_redirects_qr_id ON qr_redirects(qr_code_id); -CREATE INDEX IF NOT EXISTS idx_redirects_active ON qr_redirects(is_active); +CREATE INDEX IF NOT EXISTS idx_redirects_active ON qr_redirects(is_active); --- ============================================ --- 5. QR_URL_CHANGES (URL change history) --- ============================================ CREATE TABLE IF NOT EXISTS qr_url_changes ( - id BIGSERIAL PRIMARY KEY, - qr_redirect_id BIGINT NOT NULL REFERENCES qr_redirects(id) ON DELETE CASCADE, - - old_url TEXT NOT NULL, - new_url TEXT NOT NULL, - - changed_by BIGINT NOT NULL REFERENCES users(telegram_id), - changed_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP, - - change_reason VARCHAR(255) -); - --- Indexes + id BIGSERIAL PRIMARY KEY, + qr_redirect_id BIGINT NOT NULL REFERENCES qr_redirects(id) ON DELETE CASCADE, + + old_url TEXT NOT NULL, + new_url TEXT NOT NULL, + + changed_by BIGINT NOT NULL REFERENCES users(telegram_id), + changed_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP, + + change_reason VARCHAR(255) + ); + CREATE INDEX IF NOT EXISTS idx_url_changes_redirect ON qr_url_changes(qr_redirect_id, changed_at DESC); -CREATE INDEX IF NOT EXISTS idx_url_changes_user ON qr_url_changes(changed_by); +CREATE INDEX IF NOT EXISTS idx_url_changes_user ON qr_url_changes(changed_by); --- ============================================ --- 6. QR_CLICK_ANALYTICS (Click analytics) - PARTITIONED --- ============================================ --- O2-02: Partitioned by clicked_at (monthly) for optimal performance at scale (10M+ rows) CREATE TABLE IF NOT EXISTS qr_click_analytics ( - id BIGSERIAL, - qr_redirect_id BIGINT NOT NULL REFERENCES qr_redirects(id) ON DELETE CASCADE, - - clicked_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, - - user_agent TEXT, - ip_address INET, - - country_code VARCHAR(2), - city VARCHAR(100), - - referer TEXT, - - device_type VARCHAR(20), - os VARCHAR(50), - browser VARCHAR(50), - - -- PRIMARY KEY must include partition key - PRIMARY KEY (id, clicked_at) -) PARTITION BY RANGE (clicked_at); - --- O2-02: Create partitions (24 months: 2026-01 to 2027-12) --- Current year partitions (2026) -CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_01 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2026-01-01') TO ('2026-02-01'); -CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_02 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2026-02-01') TO ('2026-03-01'); -CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_03 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2026-03-01') TO ('2026-04-01'); -CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_04 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2026-04-01') TO ('2026-05-01'); -CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_05 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2026-05-01') TO ('2026-06-01'); -CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_06 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2026-06-01') TO ('2026-07-01'); -CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_07 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2026-07-01') TO ('2026-08-01'); -CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_08 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2026-08-01') TO ('2026-09-01'); -CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_09 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2026-09-01') TO ('2026-10-01'); -CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_10 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2026-10-01') TO ('2026-11-01'); -CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_11 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2026-11-01') TO ('2026-12-01'); -CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_12 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2026-12-01') TO ('2027-01-01'); - --- Future partitions (2027) -CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_01 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2027-01-01') TO ('2027-02-01'); -CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_02 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2027-02-01') TO ('2027-03-01'); -CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_03 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2027-03-01') TO ('2027-04-01'); -CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_04 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2027-04-01') TO ('2027-05-01'); -CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_05 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2027-05-01') TO ('2027-06-01'); -CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_06 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2027-06-01') TO ('2027-07-01'); -CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_07 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2027-07-01') TO ('2027-08-01'); -CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_08 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2027-08-01') TO ('2027-09-01'); -CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_09 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2027-09-01') TO ('2027-10-01'); -CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_10 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2027-10-01') TO ('2027-11-01'); -CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_11 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2027-11-01') TO ('2027-12-01'); -CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_12 PARTITION OF qr_click_analytics - FOR VALUES FROM ('2027-12-01') TO ('2028-01-01'); - --- Indexes --- O2-01: Optimized indexes for click analytics + id BIGSERIAL, + qr_redirect_id BIGINT NOT NULL REFERENCES qr_redirects(id) ON DELETE CASCADE, + + clicked_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, + + user_agent TEXT, + ip_address INET, + country_code VARCHAR(2), + city VARCHAR(100), + referer TEXT, + device_type VARCHAR(20), + os VARCHAR(50), + browser VARCHAR(50), + + PRIMARY KEY (id, clicked_at) + ) PARTITION BY RANGE (clicked_at); + +CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_01 + PARTITION OF qr_click_analytics + FOR VALUES FROM ('2026-01-01') TO ('2026-02-01'); + +CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_02 + PARTITION OF qr_click_analytics + FOR VALUES FROM ('2026-02-01') TO ('2026-03-01'); + +CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_03 + PARTITION OF qr_click_analytics + FOR VALUES FROM ('2026-03-01') TO ('2026-04-01'); + +CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_04 + PARTITION OF qr_click_analytics + FOR VALUES FROM ('2026-04-01') TO ('2026-05-01'); + +CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_05 + PARTITION OF qr_click_analytics + FOR VALUES FROM ('2026-05-01') TO ('2026-06-01'); + +CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_06 + PARTITION OF qr_click_analytics + FOR VALUES FROM ('2026-06-01') TO ('2026-07-01'); + +CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_07 + PARTITION OF qr_click_analytics + FOR VALUES FROM ('2026-07-01') TO ('2026-08-01'); + +CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_08 + PARTITION OF qr_click_analytics + FOR VALUES FROM ('2026-08-01') TO ('2026-09-01'); + +CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_09 + PARTITION OF qr_click_analytics + FOR VALUES FROM ('2026-09-01') TO ('2026-10-01'); + +CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_10 + PARTITION OF qr_click_analytics + FOR VALUES FROM ('2026-10-01') TO ('2026-11-01'); + +CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_11 + PARTITION OF qr_click_analytics + FOR VALUES FROM ('2026-11-01') TO ('2026-12-01'); + +CREATE TABLE IF NOT EXISTS qr_click_analytics_2026_12 + PARTITION OF qr_click_analytics + FOR VALUES FROM ('2026-12-01') TO ('2027-01-01'); + +CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_01 + PARTITION OF qr_click_analytics + FOR VALUES FROM ('2027-01-01') TO ('2027-02-01'); + +CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_02 + PARTITION OF qr_click_analytics + FOR VALUES FROM ('2027-02-01') TO ('2027-03-01'); + +CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_03 + PARTITION OF qr_click_analytics + FOR VALUES FROM ('2027-03-01') TO ('2027-04-01'); + +CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_04 + PARTITION OF qr_click_analytics + FOR VALUES FROM ('2027-04-01') TO ('2027-05-01'); + +CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_05 + PARTITION OF qr_click_analytics + FOR VALUES FROM ('2027-05-01') TO ('2027-06-01'); + +CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_06 + PARTITION OF qr_click_analytics + FOR VALUES FROM ('2027-06-01') TO ('2027-07-01'); + +CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_07 + PARTITION OF qr_click_analytics + FOR VALUES FROM ('2027-07-01') TO ('2027-08-01'); + +CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_08 + PARTITION OF qr_click_analytics + FOR VALUES FROM ('2027-08-01') TO ('2027-09-01'); + +CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_09 + PARTITION OF qr_click_analytics + FOR VALUES FROM ('2027-09-01') TO ('2027-10-01'); + +CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_10 + PARTITION OF qr_click_analytics + FOR VALUES FROM ('2027-10-01') TO ('2027-11-01'); + +CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_11 + PARTITION OF qr_click_analytics + FOR VALUES FROM ('2027-11-01') TO ('2027-12-01'); + +CREATE TABLE IF NOT EXISTS qr_click_analytics_2027_12 + PARTITION OF qr_click_analytics + FOR VALUES FROM ('2027-12-01') TO ('2028-01-01'); + CREATE INDEX IF NOT EXISTS idx_click_analytics_redirect_time ON qr_click_analytics(qr_redirect_id, clicked_at DESC); CREATE INDEX IF NOT EXISTS idx_click_analytics_device_agg - ON qr_click_analytics (qr_redirect_id, device_type) - WHERE device_type IS NOT NULL; + ON qr_click_analytics (qr_redirect_id, device_type) + WHERE device_type IS NOT NULL; CREATE INDEX IF NOT EXISTS idx_click_analytics_country_agg - ON qr_click_analytics (qr_redirect_id, country_code) - WHERE country_code IS NOT NULL; + ON qr_click_analytics (qr_redirect_id, country_code) + WHERE country_code IS NOT NULL; CREATE INDEX IF NOT EXISTS idx_click_analytics_time_series - ON qr_click_analytics (qr_redirect_id, date_trunc_immutable(clicked_at)); + ON qr_click_analytics (qr_redirect_id, date_trunc_immutable(clicked_at)); CREATE INDEX IF NOT EXISTS idx_click_analytics_unique_check - ON qr_click_analytics (qr_redirect_id, ip_address, date_trunc_immutable(clicked_at)); + ON qr_click_analytics (qr_redirect_id, ip_address, date_trunc_immutable(clicked_at)); CREATE INDEX IF NOT EXISTS idx_click_analytics_ip ON qr_click_analytics(ip_address); --- O2-03: GeoIP backfill index — find records with IP but no geo data CREATE INDEX IF NOT EXISTS idx_click_analytics_geo_backfill - ON qr_click_analytics (ip_address, clicked_at) - WHERE ip_address IS NOT NULL AND country_code IS NULL; + ON qr_click_analytics (ip_address, clicked_at) + WHERE ip_address IS NOT NULL AND country_code IS NULL; --- O2-03: Composite index for country + city breakdown queries CREATE INDEX IF NOT EXISTS idx_click_analytics_country_city - ON qr_click_analytics (qr_redirect_id, country_code, city) - WHERE country_code IS NOT NULL; - + ON qr_click_analytics (qr_redirect_id, country_code, city) + WHERE country_code IS NOT NULL; --- ============================================ --- 7. DAILY_USAGE (Дневная статистика) - PARTITIONED --- ============================================ --- O2-02: Partitioned by usage_date (monthly) for fast limit checks --- Агрегированная таблица для быстрой проверки лимитов CREATE TABLE IF NOT EXISTS daily_usage ( - id SERIAL, - user_id BIGINT NOT NULL REFERENCES users(telegram_id) ON DELETE CASCADE, - usage_date DATE NOT NULL, - qr_count INTEGER DEFAULT 0, - - created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP, - updated_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP, - - -- PRIMARY KEY and UNIQUE must include partition key - PRIMARY KEY (id, usage_date), - UNIQUE(user_id, usage_date) -) PARTITION BY RANGE (usage_date); - --- O2-02: Create partitions (24 months: 2026-01 to 2027-12) --- Current year partitions (2026) -CREATE TABLE IF NOT EXISTS daily_usage_2026_01 PARTITION OF daily_usage - FOR VALUES FROM ('2026-01-01') TO ('2026-02-01'); -CREATE TABLE IF NOT EXISTS daily_usage_2026_02 PARTITION OF daily_usage - FOR VALUES FROM ('2026-02-01') TO ('2026-03-01'); -CREATE TABLE IF NOT EXISTS daily_usage_2026_03 PARTITION OF daily_usage - FOR VALUES FROM ('2026-03-01') TO ('2026-04-01'); -CREATE TABLE IF NOT EXISTS daily_usage_2026_04 PARTITION OF daily_usage - FOR VALUES FROM ('2026-04-01') TO ('2026-05-01'); -CREATE TABLE IF NOT EXISTS daily_usage_2026_05 PARTITION OF daily_usage - FOR VALUES FROM ('2026-05-01') TO ('2026-06-01'); -CREATE TABLE IF NOT EXISTS daily_usage_2026_06 PARTITION OF daily_usage - FOR VALUES FROM ('2026-06-01') TO ('2026-07-01'); -CREATE TABLE IF NOT EXISTS daily_usage_2026_07 PARTITION OF daily_usage - FOR VALUES FROM ('2026-07-01') TO ('2026-08-01'); -CREATE TABLE IF NOT EXISTS daily_usage_2026_08 PARTITION OF daily_usage - FOR VALUES FROM ('2026-08-01') TO ('2026-09-01'); -CREATE TABLE IF NOT EXISTS daily_usage_2026_09 PARTITION OF daily_usage - FOR VALUES FROM ('2026-09-01') TO ('2026-10-01'); -CREATE TABLE IF NOT EXISTS daily_usage_2026_10 PARTITION OF daily_usage - FOR VALUES FROM ('2026-10-01') TO ('2026-11-01'); -CREATE TABLE IF NOT EXISTS daily_usage_2026_11 PARTITION OF daily_usage - FOR VALUES FROM ('2026-11-01') TO ('2026-12-01'); -CREATE TABLE IF NOT EXISTS daily_usage_2026_12 PARTITION OF daily_usage - FOR VALUES FROM ('2026-12-01') TO ('2027-01-01'); - --- Future partitions (2027) -CREATE TABLE IF NOT EXISTS daily_usage_2027_01 PARTITION OF daily_usage - FOR VALUES FROM ('2027-01-01') TO ('2027-02-01'); -CREATE TABLE IF NOT EXISTS daily_usage_2027_02 PARTITION OF daily_usage - FOR VALUES FROM ('2027-02-01') TO ('2027-03-01'); -CREATE TABLE IF NOT EXISTS daily_usage_2027_03 PARTITION OF daily_usage - FOR VALUES FROM ('2027-03-01') TO ('2027-04-01'); -CREATE TABLE IF NOT EXISTS daily_usage_2027_04 PARTITION OF daily_usage - FOR VALUES FROM ('2027-04-01') TO ('2027-05-01'); -CREATE TABLE IF NOT EXISTS daily_usage_2027_05 PARTITION OF daily_usage - FOR VALUES FROM ('2027-05-01') TO ('2027-06-01'); -CREATE TABLE IF NOT EXISTS daily_usage_2027_06 PARTITION OF daily_usage - FOR VALUES FROM ('2027-06-01') TO ('2027-07-01'); -CREATE TABLE IF NOT EXISTS daily_usage_2027_07 PARTITION OF daily_usage - FOR VALUES FROM ('2027-07-01') TO ('2027-08-01'); -CREATE TABLE IF NOT EXISTS daily_usage_2027_08 PARTITION OF daily_usage - FOR VALUES FROM ('2027-08-01') TO ('2027-09-01'); -CREATE TABLE IF NOT EXISTS daily_usage_2027_09 PARTITION OF daily_usage - FOR VALUES FROM ('2027-09-01') TO ('2027-10-01'); -CREATE TABLE IF NOT EXISTS daily_usage_2027_10 PARTITION OF daily_usage - FOR VALUES FROM ('2027-10-01') TO ('2027-11-01'); -CREATE TABLE IF NOT EXISTS daily_usage_2027_11 PARTITION OF daily_usage - FOR VALUES FROM ('2027-11-01') TO ('2027-12-01'); -CREATE TABLE IF NOT EXISTS daily_usage_2027_12 PARTITION OF daily_usage - FOR VALUES FROM ('2027-12-01') TO ('2028-01-01'); - --- Индексы --- O2-01: Optimized indexes for daily usage (critical for limit checks) + user_id BIGINT NOT NULL REFERENCES users(telegram_id) ON DELETE CASCADE, + usage_date DATE NOT NULL, + qr_count INTEGER NOT NULL DEFAULT 0 CHECK (qr_count >= 0), + + created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP, + updated_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP, + + PRIMARY KEY (user_id, usage_date) + ) PARTITION BY RANGE (usage_date); + +CREATE TABLE IF NOT EXISTS daily_usage_2026_01 + PARTITION OF daily_usage + FOR VALUES FROM ('2026-01-01') TO ('2026-02-01'); + +CREATE TABLE IF NOT EXISTS daily_usage_2026_02 + PARTITION OF daily_usage + FOR VALUES FROM ('2026-02-01') TO ('2026-03-01'); + +CREATE TABLE IF NOT EXISTS daily_usage_2026_03 + PARTITION OF daily_usage + FOR VALUES FROM ('2026-03-01') TO ('2026-04-01'); + +CREATE TABLE IF NOT EXISTS daily_usage_2026_04 + PARTITION OF daily_usage + FOR VALUES FROM ('2026-04-01') TO ('2026-05-01'); + +CREATE TABLE IF NOT EXISTS daily_usage_2026_05 + PARTITION OF daily_usage + FOR VALUES FROM ('2026-05-01') TO ('2026-06-01'); + +CREATE TABLE IF NOT EXISTS daily_usage_2026_06 + PARTITION OF daily_usage + FOR VALUES FROM ('2026-06-01') TO ('2026-07-01'); + +CREATE TABLE IF NOT EXISTS daily_usage_2026_07 + PARTITION OF daily_usage + FOR VALUES FROM ('2026-07-01') TO ('2026-08-01'); + +CREATE TABLE IF NOT EXISTS daily_usage_2026_08 + PARTITION OF daily_usage + FOR VALUES FROM ('2026-08-01') TO ('2026-09-01'); + +CREATE TABLE IF NOT EXISTS daily_usage_2026_09 + PARTITION OF daily_usage + FOR VALUES FROM ('2026-09-01') TO ('2026-10-01'); + +CREATE TABLE IF NOT EXISTS daily_usage_2026_10 + PARTITION OF daily_usage + FOR VALUES FROM ('2026-10-01') TO ('2026-11-01'); + +CREATE TABLE IF NOT EXISTS daily_usage_2026_11 + PARTITION OF daily_usage + FOR VALUES FROM ('2026-11-01') TO ('2026-12-01'); + +CREATE TABLE IF NOT EXISTS daily_usage_2026_12 + PARTITION OF daily_usage + FOR VALUES FROM ('2026-12-01') TO ('2027-01-01'); + +CREATE TABLE IF NOT EXISTS daily_usage_2027_01 + PARTITION OF daily_usage + FOR VALUES FROM ('2027-01-01') TO ('2027-02-01'); + +CREATE TABLE IF NOT EXISTS daily_usage_2027_02 + PARTITION OF daily_usage + FOR VALUES FROM ('2027-02-01') TO ('2027-03-01'); + +CREATE TABLE IF NOT EXISTS daily_usage_2027_03 + PARTITION OF daily_usage + FOR VALUES FROM ('2027-03-01') TO ('2027-04-01'); + +CREATE TABLE IF NOT EXISTS daily_usage_2027_04 + PARTITION OF daily_usage + FOR VALUES FROM ('2027-04-01') TO ('2027-05-01'); + +CREATE TABLE IF NOT EXISTS daily_usage_2027_05 + PARTITION OF daily_usage + FOR VALUES FROM ('2027-05-01') TO ('2027-06-01'); + +CREATE TABLE IF NOT EXISTS daily_usage_2027_06 + PARTITION OF daily_usage + FOR VALUES FROM ('2027-06-01') TO ('2027-07-01'); + +CREATE TABLE IF NOT EXISTS daily_usage_2027_07 + PARTITION OF daily_usage + FOR VALUES FROM ('2027-07-01') TO ('2027-08-01'); + +CREATE TABLE IF NOT EXISTS daily_usage_2027_08 + PARTITION OF daily_usage + FOR VALUES FROM ('2027-08-01') TO ('2027-09-01'); + +CREATE TABLE IF NOT EXISTS daily_usage_2027_09 + PARTITION OF daily_usage + FOR VALUES FROM ('2027-09-01') TO ('2027-10-01'); + +CREATE TABLE IF NOT EXISTS daily_usage_2027_10 + PARTITION OF daily_usage + FOR VALUES FROM ('2027-10-01') TO ('2027-11-01'); + +CREATE TABLE IF NOT EXISTS daily_usage_2027_11 + PARTITION OF daily_usage + FOR VALUES FROM ('2027-11-01') TO ('2027-12-01'); + +CREATE TABLE IF NOT EXISTS daily_usage_2027_12 + PARTITION OF daily_usage + FOR VALUES FROM ('2027-12-01') TO ('2028-01-01'); + CREATE INDEX IF NOT EXISTS idx_daily_usage_limit_check - ON daily_usage (user_id, usage_date) - INCLUDE (qr_count); + ON daily_usage (user_id, usage_date) + INCLUDE (qr_count); -CREATE INDEX IF NOT EXISTS idx_daily_usage_user_date ON daily_usage(user_id, usage_date DESC); CREATE INDEX IF NOT EXISTS idx_daily_usage_date ON daily_usage(usage_date DESC); - --- ============================================ --- 8. SUBSCRIPTION_HISTORY (История подписок) --- ============================================ CREATE TABLE IF NOT EXISTS subscription_history ( - id SERIAL PRIMARY KEY, - user_id BIGINT NOT NULL REFERENCES users(telegram_id) ON DELETE CASCADE, - tariff_id INTEGER NOT NULL REFERENCES tariffs(id), - - start_date TIMESTAMPTZ NOT NULL, - end_date TIMESTAMPTZ, - - price_paid INTEGER, - payment_method VARCHAR(50), - - created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP -); - --- Индексы -CREATE INDEX IF NOT EXISTS idx_subscription_user ON subscription_history(user_id, start_date DESC); -CREATE INDEX IF NOT EXISTS idx_subscription_tariff ON subscription_history(tariff_id); + id SERIAL PRIMARY KEY, + user_id BIGINT NOT NULL REFERENCES users(telegram_id) ON DELETE CASCADE, + tariff_id INTEGER NOT NULL REFERENCES tariffs(id), + start_date TIMESTAMPTZ NOT NULL, + end_date TIMESTAMPTZ, --- ============================================ --- TRIGGERS для автоматических обновлений --- ============================================ + price_paid INTEGER, + payment_method VARCHAR(50), --- Триггер для обновления updated_at -CREATE OR REPLACE FUNCTION update_updated_at_column() -RETURNS TRIGGER AS $$ -BEGIN - NEW.updated_at = CURRENT_TIMESTAMP; - RETURN NEW; -END; -$$ LANGUAGE plpgsql; + created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP + ); -CREATE TRIGGER update_users_updated_at BEFORE UPDATE ON users - FOR EACH ROW EXECUTE FUNCTION update_updated_at_column(); +CREATE INDEX IF NOT EXISTS idx_subscription_user ON subscription_history(user_id, start_date DESC); -CREATE TRIGGER update_tariffs_updated_at BEFORE UPDATE ON tariffs - FOR EACH ROW EXECUTE FUNCTION update_updated_at_column(); +CREATE INDEX IF NOT EXISTS idx_subscription_tariff ON subscription_history(tariff_id); + +CREATE OR REPLACE FUNCTION update_updated_at_column() + RETURNS TRIGGER AS $$ + BEGIN + NEW.updated_at = CURRENT_TIMESTAMP; + RETURN NEW; + END; + $$ LANGUAGE plpgsql; DROP TRIGGER IF EXISTS update_users_updated_at ON users; -CREATE TRIGGER update_users_updated_at -BEFORE UPDATE ON users -FOR EACH ROW EXECUTE FUNCTION update_updated_at_column(); + +CREATE TRIGGER update_users_updated_at + BEFORE UPDATE ON users + FOR EACH ROW EXECUTE FUNCTION update_updated_at_column(); DROP TRIGGER IF EXISTS update_tariffs_updated_at ON tariffs; -CREATE TRIGGER update_tariffs_updated_at -BEFORE UPDATE ON tariffs -FOR EACH ROW EXECUTE FUNCTION update_updated_at_column(); - - --- Триггер для автоматического обновления daily_usage (УДАЛЕН - используется событийная архитектура) --- Обновление статистики теперь происходит через event handlers - - --- ============================================ --- VIEWS для аналитики --- ============================================ - --- Текущая статистика пользователей -CREATE OR REPLACE VIEW v_user_stats AS -SELECT - u.telegram_id, - u.username, - u.first_name, - t.name as tariff_name, - t.daily_limit, - t.monthly_limit, - u.total_qr_generated, - u.total_qr_deleted, - u.last_qr_generated_at, - u.created_at as user_since -FROM users u -LEFT JOIN tariffs t ON u.tariff_id = t.id -WHERE u.is_active = TRUE; - - --- Топ пользователей по генерациям -CREATE OR REPLACE VIEW v_top_users AS -SELECT - u.telegram_id, - u.username, - u.total_qr_generated, - t.name as tariff_name, - u.created_at -FROM users u -LEFT JOIN tariffs t ON u.tariff_id = t.id -ORDER BY u.total_qr_generated DESC -LIMIT 100; - - --- ============================================ --- КОММЕНТАРИИ к таблицам --- ============================================ -COMMENT ON TABLE tariffs IS 'Тарифные планы с лимитами и возможностями'; -COMMENT ON TABLE users IS 'Пользователи Telegram с привязкой к тарифам'; -COMMENT ON TABLE qr_codes IS 'История всех сгенерированных QR кодов (с soft delete)'; -COMMENT ON TABLE daily_usage IS 'Агрегированная статистика использования по дням'; -COMMENT ON TABLE subscription_history IS 'История изменений подписок'; - -COMMENT ON COLUMN tariffs.history_retention_days IS 'Сколько дней хранить QR коды'; -COMMENT ON COLUMN users.tariff_expires_at IS 'NULL для freemium (не истекает)'; -COMMENT ON COLUMN qr_codes.deleted_at IS 'Soft delete - дата удаления QR кода'; -COMMENT ON COLUMN qr_codes.is_dynamic IS 'Динамический QR код с возможностью смены URL'; -COMMENT ON COLUMN qr_codes.short_code IS 'Короткий код для редиректа динамических QR'; - - --- ============================================ --- 9. FEATURE_FLAGS (Feature Flags Configuration) --- ============================================ + +CREATE TRIGGER update_tariffs_updated_at + BEFORE UPDATE ON tariffs + FOR EACH ROW EXECUTE FUNCTION update_updated_at_column(); + +DROP TRIGGER IF EXISTS update_redirects_updated_at ON qr_redirects; + +CREATE TRIGGER update_redirects_updated_at + BEFORE UPDATE ON qr_redirects + FOR EACH ROW EXECUTE FUNCTION update_updated_at_column(); + CREATE TABLE IF NOT EXISTS feature_flags ( - id SERIAL PRIMARY KEY, - - -- Идентификация - name VARCHAR(100) NOT NULL UNIQUE, - description TEXT NOT NULL DEFAULT '', - - -- Статус флага - status VARCHAR(20) NOT NULL DEFAULT 'disabled' - CHECK (status IN ('disabled', 'enabled', 'percentage', 'users', 'segments')), - - -- Rollout настройки - rollout_percentage INTEGER DEFAULT 0 - CHECK (rollout_percentage >= 0 AND rollout_percentage <= 100), - rollout_salt VARCHAR(50), - - -- Targeting (массивы для гибкости) - enabled_users BIGINT[] DEFAULT ARRAY[]::BIGINT[], - disabled_users BIGINT[] DEFAULT ARRAY[]::BIGINT[], - enabled_segments TEXT[] DEFAULT ARRAY[]::TEXT[], - - -- Metadata (произвольные данные) - metadata JSONB DEFAULT '{}', - - -- Audit - created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP, - updated_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP, - created_by BIGINT REFERENCES users(telegram_id), - updated_by BIGINT REFERENCES users(telegram_id) -); - --- Индексы для feature_flags + id SERIAL PRIMARY KEY, + + name VARCHAR(100) NOT NULL UNIQUE, + description TEXT NOT NULL DEFAULT '', + + status VARCHAR(20) NOT NULL DEFAULT 'disabled' + CHECK (status IN ('disabled', 'enabled', 'percentage', 'users', 'segments')), + + rollout_percentage INTEGER DEFAULT 0 + CHECK (rollout_percentage >= 0 AND rollout_percentage <= 100), + rollout_salt VARCHAR(50), + + enabled_users BIGINT[] DEFAULT ARRAY[]::BIGINT[], + disabled_users BIGINT[] DEFAULT ARRAY[]::BIGINT[], + enabled_segments TEXT[] DEFAULT ARRAY[]::TEXT[], + + metadata JSONB DEFAULT '{}', + + created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP, + updated_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP, + created_by BIGINT REFERENCES users(telegram_id), + updated_by BIGINT REFERENCES users(telegram_id) + ); + CREATE INDEX IF NOT EXISTS idx_ff_name ON feature_flags(name); + CREATE INDEX IF NOT EXISTS idx_ff_status ON feature_flags(status); + CREATE INDEX IF NOT EXISTS idx_ff_updated ON feature_flags(updated_at DESC); --- Trigger для updated_at DROP TRIGGER IF EXISTS update_feature_flags_updated_at ON feature_flags; + CREATE TRIGGER update_feature_flags_updated_at -BEFORE UPDATE ON feature_flags -FOR EACH ROW EXECUTE FUNCTION update_updated_at_column(); + BEFORE UPDATE ON feature_flags + FOR EACH ROW EXECUTE FUNCTION update_updated_at_column(); -COMMENT ON TABLE feature_flags IS 'Feature flags for gradual rollout and A/B testing'; -COMMENT ON COLUMN feature_flags.status IS 'disabled, enabled, percentage, users, segments'; -COMMENT ON COLUMN feature_flags.rollout_salt IS 'Salt for deterministic percentage rollout'; -COMMENT ON COLUMN feature_flags.enabled_segments IS 'User segments: beta_testers, premium_users, admin_users, etc.'; +CREATE TABLE IF NOT EXISTS feature_flag_history ( + id BIGSERIAL PRIMARY KEY, + flag_id INTEGER NOT NULL, + flag_name VARCHAR(100) NOT NULL, + + action VARCHAR(20) NOT NULL + CHECK (action IN ('created', 'updated', 'deleted')), + + old_value JSONB, + new_value JSONB, + + changed_by BIGINT NOT NULL REFERENCES users(telegram_id), + changed_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP, + + change_reason TEXT + ); --- ============================================ --- 10. FEATURE_FLAG_HISTORY (Audit Trail) --- ============================================ -CREATE TABLE IF NOT EXISTS feature_flag_history ( - id BIGSERIAL PRIMARY KEY, - - -- Ссылка на флаг (без CASCADE - храним историю удаленных) - flag_id INTEGER NOT NULL, - flag_name VARCHAR(100) NOT NULL, - - -- Действие - action VARCHAR(20) NOT NULL - CHECK (action IN ('created', 'updated', 'deleted')), - - -- Снимок до/после - old_value JSONB, - new_value JSONB, - - -- Кто и когда - changed_by BIGINT NOT NULL REFERENCES users(telegram_id), - changed_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP, - - -- Причина изменения - change_reason TEXT -); - --- Индексы для feature_flag_history CREATE INDEX IF NOT EXISTS idx_ffh_flag_time ON feature_flag_history(flag_id, changed_at DESC); + CREATE INDEX IF NOT EXISTS idx_ffh_name_time ON feature_flag_history(flag_name, changed_at DESC); + CREATE INDEX IF NOT EXISTS idx_ffh_user ON feature_flag_history(changed_by); -CREATE INDEX IF NOT EXISTS idx_ffh_action ON feature_flag_history(action); -COMMENT ON TABLE feature_flag_history IS 'Audit trail for all feature flag changes'; -COMMENT ON COLUMN feature_flag_history.action IS 'created, updated, deleted'; +CREATE INDEX IF NOT EXISTS idx_ffh_action ON feature_flag_history(action); +CREATE TABLE IF NOT EXISTS event_dedup ( + event_id UUID PRIMARY KEY, + created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP + ); --- ============================================ --- 11. DOMAIN_EVENTS (Event Store) - PARTITIONED --- ============================================ --- O2-02: Partitioned by occurred_at (monthly) for efficient event sourcing at scale --- Event Sourcing storage for all domain events. --- Provides full audit trail, replay capability, and debugging. --- --- Design decisions: --- - Append-only (events are immutable) --- - JSONB payload for flexible schema evolution --- - Separate indexes for aggregate queries and time-based queries --- - No foreign keys to users - events may outlive user records --- - Partitioned by month for better query performance and data management +CREATE INDEX IF NOT EXISTS idx_event_dedup_created_at ON event_dedup(created_at DESC); CREATE TABLE IF NOT EXISTS domain_events ( - -- Identity - id BIGSERIAL, - event_id UUID NOT NULL, - - -- Event type info (for deserialization) - event_name VARCHAR(100) NOT NULL, - event_version VARCHAR(20) NOT NULL DEFAULT '1.0.0', - aggregate_type VARCHAR(50) NOT NULL, - aggregate_id BIGINT NOT NULL, - - -- Event payload (complete event data as JSON) - payload JSONB NOT NULL, - - -- Timing - occurred_at TIMESTAMPTZ NOT NULL, - stored_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP, - - -- Processing status (for replay/reprocessing) - processed BOOLEAN DEFAULT TRUE, - processed_at TIMESTAMPTZ, - processing_error TEXT, - - -- PRIMARY KEY and UNIQUE must include partition key - PRIMARY KEY (id, occurred_at), - UNIQUE (event_id, occurred_at) -) PARTITION BY RANGE (occurred_at); - --- O2-02: Create partitions (24 months: 2026-01 to 2027-12) --- Current year partitions (2026) -CREATE TABLE IF NOT EXISTS domain_events_2026_01 PARTITION OF domain_events - FOR VALUES FROM ('2026-01-01') TO ('2026-02-01'); -CREATE TABLE IF NOT EXISTS domain_events_2026_02 PARTITION OF domain_events - FOR VALUES FROM ('2026-02-01') TO ('2026-03-01'); -CREATE TABLE IF NOT EXISTS domain_events_2026_03 PARTITION OF domain_events - FOR VALUES FROM ('2026-03-01') TO ('2026-04-01'); -CREATE TABLE IF NOT EXISTS domain_events_2026_04 PARTITION OF domain_events - FOR VALUES FROM ('2026-04-01') TO ('2026-05-01'); -CREATE TABLE IF NOT EXISTS domain_events_2026_05 PARTITION OF domain_events - FOR VALUES FROM ('2026-05-01') TO ('2026-06-01'); -CREATE TABLE IF NOT EXISTS domain_events_2026_06 PARTITION OF domain_events - FOR VALUES FROM ('2026-06-01') TO ('2026-07-01'); -CREATE TABLE IF NOT EXISTS domain_events_2026_07 PARTITION OF domain_events - FOR VALUES FROM ('2026-07-01') TO ('2026-08-01'); -CREATE TABLE IF NOT EXISTS domain_events_2026_08 PARTITION OF domain_events - FOR VALUES FROM ('2026-08-01') TO ('2026-09-01'); -CREATE TABLE IF NOT EXISTS domain_events_2026_09 PARTITION OF domain_events - FOR VALUES FROM ('2026-09-01') TO ('2026-10-01'); -CREATE TABLE IF NOT EXISTS domain_events_2026_10 PARTITION OF domain_events - FOR VALUES FROM ('2026-10-01') TO ('2026-11-01'); -CREATE TABLE IF NOT EXISTS domain_events_2026_11 PARTITION OF domain_events - FOR VALUES FROM ('2026-11-01') TO ('2026-12-01'); -CREATE TABLE IF NOT EXISTS domain_events_2026_12 PARTITION OF domain_events - FOR VALUES FROM ('2026-12-01') TO ('2027-01-01'); - --- Future partitions (2027) -CREATE TABLE IF NOT EXISTS domain_events_2027_01 PARTITION OF domain_events - FOR VALUES FROM ('2027-01-01') TO ('2027-02-01'); -CREATE TABLE IF NOT EXISTS domain_events_2027_02 PARTITION OF domain_events - FOR VALUES FROM ('2027-02-01') TO ('2027-03-01'); -CREATE TABLE IF NOT EXISTS domain_events_2027_03 PARTITION OF domain_events - FOR VALUES FROM ('2027-03-01') TO ('2027-04-01'); -CREATE TABLE IF NOT EXISTS domain_events_2027_04 PARTITION OF domain_events - FOR VALUES FROM ('2027-04-01') TO ('2027-05-01'); -CREATE TABLE IF NOT EXISTS domain_events_2027_05 PARTITION OF domain_events - FOR VALUES FROM ('2027-05-01') TO ('2027-06-01'); -CREATE TABLE IF NOT EXISTS domain_events_2027_06 PARTITION OF domain_events - FOR VALUES FROM ('2027-06-01') TO ('2027-07-01'); -CREATE TABLE IF NOT EXISTS domain_events_2027_07 PARTITION OF domain_events - FOR VALUES FROM ('2027-07-01') TO ('2027-08-01'); -CREATE TABLE IF NOT EXISTS domain_events_2027_08 PARTITION OF domain_events - FOR VALUES FROM ('2027-08-01') TO ('2027-09-01'); -CREATE TABLE IF NOT EXISTS domain_events_2027_09 PARTITION OF domain_events - FOR VALUES FROM ('2027-09-01') TO ('2027-10-01'); -CREATE TABLE IF NOT EXISTS domain_events_2027_10 PARTITION OF domain_events - FOR VALUES FROM ('2027-10-01') TO ('2027-11-01'); -CREATE TABLE IF NOT EXISTS domain_events_2027_11 PARTITION OF domain_events - FOR VALUES FROM ('2027-11-01') TO ('2027-12-01'); -CREATE TABLE IF NOT EXISTS domain_events_2027_12 PARTITION OF domain_events - FOR VALUES FROM ('2027-12-01') TO ('2028-01-01'); - --- Primary query patterns: --- 1. Get events by aggregate (event sourcing) -CREATE INDEX IF NOT EXISTS idx_de_aggregate - ON domain_events(aggregate_type, aggregate_id, id); - --- O2-01: Fast replay with covering index + id BIGSERIAL, + event_id UUID NOT NULL, + + event_name VARCHAR(100) NOT NULL, + event_version VARCHAR(20) NOT NULL DEFAULT '1.0.0', + aggregate_type VARCHAR(50) NOT NULL, + aggregate_id BIGINT NOT NULL, + + payload JSONB NOT NULL, + + occurred_at TIMESTAMPTZ NOT NULL, + stored_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP, + + processed BOOLEAN DEFAULT TRUE, + processed_at TIMESTAMPTZ, + processing_error TEXT, + + PRIMARY KEY (id, occurred_at), + UNIQUE (event_id, occurred_at) + ) PARTITION BY RANGE (occurred_at); + +CREATE TABLE IF NOT EXISTS domain_events_2026_01 + PARTITION OF domain_events + FOR VALUES FROM ('2026-01-01') TO ('2026-02-01'); + +CREATE TABLE IF NOT EXISTS domain_events_2026_02 + PARTITION OF domain_events + FOR VALUES FROM ('2026-02-01') TO ('2026-03-01'); + +CREATE TABLE IF NOT EXISTS domain_events_2026_03 + PARTITION OF domain_events + FOR VALUES FROM ('2026-03-01') TO ('2026-04-01'); + +CREATE TABLE IF NOT EXISTS domain_events_2026_04 + PARTITION OF domain_events + FOR VALUES FROM ('2026-04-01') TO ('2026-05-01'); + +CREATE TABLE IF NOT EXISTS domain_events_2026_05 + PARTITION OF domain_events + FOR VALUES FROM ('2026-05-01') TO ('2026-06-01'); + +CREATE TABLE IF NOT EXISTS domain_events_2026_06 + PARTITION OF domain_events + FOR VALUES FROM ('2026-06-01') TO ('2026-07-01'); + +CREATE TABLE IF NOT EXISTS domain_events_2026_07 + PARTITION OF domain_events + FOR VALUES FROM ('2026-07-01') TO ('2026-08-01'); + +CREATE TABLE IF NOT EXISTS domain_events_2026_08 + PARTITION OF domain_events + FOR VALUES FROM ('2026-08-01') TO ('2026-09-01'); + +CREATE TABLE IF NOT EXISTS domain_events_2026_09 + PARTITION OF domain_events + FOR VALUES FROM ('2026-09-01') TO ('2026-10-01'); + +CREATE TABLE IF NOT EXISTS domain_events_2026_10 + PARTITION OF domain_events + FOR VALUES FROM ('2026-10-01') TO ('2026-11-01'); + +CREATE TABLE IF NOT EXISTS domain_events_2026_11 + PARTITION OF domain_events + FOR VALUES FROM ('2026-11-01') TO ('2026-12-01'); + +CREATE TABLE IF NOT EXISTS domain_events_2026_12 + PARTITION OF domain_events + FOR VALUES FROM ('2026-12-01') TO ('2027-01-01'); + +CREATE TABLE IF NOT EXISTS domain_events_2027_01 + PARTITION OF domain_events + FOR VALUES FROM ('2027-01-01') TO ('2027-02-01'); + +CREATE TABLE IF NOT EXISTS domain_events_2027_02 + PARTITION OF domain_events + FOR VALUES FROM ('2027-02-01') TO ('2027-03-01'); + +CREATE TABLE IF NOT EXISTS domain_events_2027_03 + PARTITION OF domain_events + FOR VALUES FROM ('2027-03-01') TO ('2027-04-01'); + +CREATE TABLE IF NOT EXISTS domain_events_2027_04 + PARTITION OF domain_events + FOR VALUES FROM ('2027-04-01') TO ('2027-05-01'); + +CREATE TABLE IF NOT EXISTS domain_events_2027_05 + PARTITION OF domain_events + FOR VALUES FROM ('2027-05-01') TO ('2027-06-01'); + +CREATE TABLE IF NOT EXISTS domain_events_2027_06 + PARTITION OF domain_events + FOR VALUES FROM ('2027-06-01') TO ('2027-07-01'); + +CREATE TABLE IF NOT EXISTS domain_events_2027_07 + PARTITION OF domain_events + FOR VALUES FROM ('2027-07-01') TO ('2027-08-01'); + +CREATE TABLE IF NOT EXISTS domain_events_2027_08 + PARTITION OF domain_events + FOR VALUES FROM ('2027-08-01') TO ('2027-09-01'); + +CREATE TABLE IF NOT EXISTS domain_events_2027_09 + PARTITION OF domain_events + FOR VALUES FROM ('2027-09-01') TO ('2027-10-01'); + +CREATE TABLE IF NOT EXISTS domain_events_2027_10 + PARTITION OF domain_events + FOR VALUES FROM ('2027-10-01') TO ('2027-11-01'); + +CREATE TABLE IF NOT EXISTS domain_events_2027_11 + PARTITION OF domain_events + FOR VALUES FROM ('2027-11-01') TO ('2027-12-01'); + +CREATE TABLE IF NOT EXISTS domain_events_2027_12 + PARTITION OF domain_events + FOR VALUES FROM ('2027-12-01') TO ('2028-01-01'); + +CREATE INDEX IF NOT EXISTS idx_de_aggregate ON domain_events(aggregate_type, aggregate_id, id); + CREATE INDEX IF NOT EXISTS idx_de_aggregate_fast_replay - ON domain_events (aggregate_type, aggregate_id, occurred_at) - INCLUDE (event_name, event_version) - WHERE processed = TRUE; + ON domain_events (aggregate_type, aggregate_id, occurred_at) + INCLUDE (event_name, event_version) + WHERE processed = TRUE; --- O2-01: Analytics queries with time window CREATE INDEX IF NOT EXISTS idx_de_analytics - ON domain_events (event_name, occurred_at) - INCLUDE (aggregate_id, payload) - WHERE occurred_at > CURRENT_TIMESTAMP - INTERVAL '7 days'; - --- 2. Get events by type (analytics, handlers) -CREATE INDEX IF NOT EXISTS idx_de_event_name - ON domain_events(event_name, occurred_at DESC); - --- 3. Get events by time (debugging, audit) -CREATE INDEX IF NOT EXISTS idx_de_occurred_at - ON domain_events(occurred_at DESC); - --- 4. Find unprocessed events (replay) -CREATE INDEX IF NOT EXISTS idx_de_unprocessed - ON domain_events(processed) - WHERE processed = FALSE; - --- 5. Lookup by event_id -CREATE INDEX IF NOT EXISTS idx_de_event_id - ON domain_events(event_id); - --- Partitioning hint: for high-volume, partition by occurred_at month --- CREATE TABLE domain_events_2026_01 PARTITION OF domain_events --- FOR VALUES FROM ('2026-01-01') TO ('2026-02-01'); - -COMMENT ON TABLE domain_events IS 'Event store for all domain events (append-only)'; -COMMENT ON COLUMN domain_events.event_id IS 'UUID of the event instance'; -COMMENT ON COLUMN domain_events.event_name IS 'Event type name (e.g., qr_code.generated)'; -COMMENT ON COLUMN domain_events.event_version IS 'Event schema version for evolution'; -COMMENT ON COLUMN domain_events.aggregate_type IS 'Type of aggregate (User, QRCode, etc.)'; -COMMENT ON COLUMN domain_events.aggregate_id IS 'ID of the aggregate that emitted event'; -COMMENT ON COLUMN domain_events.payload IS 'Complete event data as JSON'; -COMMENT ON COLUMN domain_events.occurred_at IS 'When the event happened (from event)'; -COMMENT ON COLUMN domain_events.stored_at IS 'When the event was stored (server time)'; -COMMENT ON COLUMN domain_events.processed IS 'FALSE if event needs reprocessing'; -COMMENT ON COLUMN domain_events.processing_error IS 'Error message if processing failed'; - - --- ============================================ --- O2-01: AUTOVACUUM CONFIGURATION --- ============================================ --- Configure autovacuum for high-traffic tables to prevent index bloat - -ALTER TABLE qr_click_analytics SET ( - autovacuum_vacuum_scale_factor = 0.05, - autovacuum_analyze_scale_factor = 0.02, - autovacuum_vacuum_cost_delay = 10 -); + ON domain_events (event_name, occurred_at) + INCLUDE (aggregate_id, payload); -ALTER TABLE qr_codes SET ( - autovacuum_vacuum_scale_factor = 0.1, - autovacuum_analyze_scale_factor = 0.05 -); +CREATE INDEX IF NOT EXISTS idx_de_event_name ON domain_events(event_name, occurred_at DESC); -ALTER TABLE domain_events SET ( - autovacuum_vacuum_scale_factor = 0.05, - autovacuum_analyze_scale_factor = 0.02 -); +CREATE INDEX IF NOT EXISTS idx_de_occurred_at ON domain_events(occurred_at DESC); --- Enable extensions for monitoring -CREATE EXTENSION IF NOT EXISTS pg_stat_statements; -CREATE EXTENSION IF NOT EXISTS pgstattuple; +CREATE INDEX IF NOT EXISTS idx_de_unprocessed ON domain_events(processed) WHERE processed = FALSE; +CREATE INDEX IF NOT EXISTS idx_de_event_id ON domain_events(event_id); --- ============================================ --- O2-02: AUTOMATIC PARTITION MANAGEMENT --- ============================================ --- Function to automatically create next month's partitions for all partitioned tables +ALTER TABLE qr_codes SET ( + autovacuum_vacuum_scale_factor = 0.1, + autovacuum_analyze_scale_factor = 0.05 + ); + +DO $$ + BEGIN + EXECUTE 'CREATE EXTENSION IF NOT EXISTS pg_stat_statements'; + EXCEPTION + WHEN insufficient_privilege THEN + RAISE NOTICE 'Skipping extension pg_stat_statements: insufficient privileges'; + END + $$; + +DO $$ + BEGIN + EXECUTE 'CREATE EXTENSION IF NOT EXISTS pgstattuple'; + EXCEPTION + WHEN insufficient_privilege THEN + RAISE NOTICE 'Skipping extension pgstattuple: insufficient privileges'; + END + $$; CREATE OR REPLACE FUNCTION create_monthly_partitions() -RETURNS void AS $$ -DECLARE - table_name TEXT; - partition_name TEXT; - start_date DATE; - end_date DATE; -BEGIN - -- Calculate next month - start_date := date_trunc('month', CURRENT_DATE + INTERVAL '1 month'); - end_date := start_date + INTERVAL '1 month'; - - -- Create partitions for each partitioned table - FOREACH table_name IN ARRAY ARRAY[ - 'qr_click_analytics', - 'daily_usage', - 'domain_events' - ] LOOP - partition_name := table_name || '_' || to_char(start_date, 'YYYY_MM'); - - -- Check if partition exists - IF NOT EXISTS ( - SELECT 1 FROM pg_class - WHERE relname = partition_name - ) THEN - EXECUTE format( - 'CREATE TABLE %I PARTITION OF %I FOR VALUES FROM (%L) TO (%L)', - partition_name, - table_name, - start_date, - end_date - ); - RAISE NOTICE 'Created partition: %', partition_name; - END IF; - END LOOP; -END; -$$ LANGUAGE plpgsql; + RETURNS void AS $$ + DECLARE + table_name TEXT; + partition_name TEXT; + start_date DATE; + end_date DATE; + BEGIN + start_date := date_trunc('month', CURRENT_DATE + INTERVAL '1 month'); + end_date := start_date + INTERVAL '1 month'; + + FOREACH table_name IN ARRAY ARRAY[ + 'qr_click_analytics', + 'daily_usage', + 'domain_events' + ] LOOP + partition_name := table_name || '_' || to_char(start_date, 'YYYY_MM'); + + IF NOT EXISTS ( + SELECT 1 + FROM pg_class + WHERE relname = partition_name + ) THEN + EXECUTE format( + 'CREATE TABLE %I PARTITION OF %I FOR VALUES FROM (%L) TO (%L)', + partition_name, + table_name, + start_date, + end_date + ); + RAISE NOTICE 'Created partition: %', partition_name; + END IF; + END LOOP; + END; + $$ LANGUAGE plpgsql; COMMENT ON FUNCTION create_monthly_partitions() IS -'Automatically creates next month partition for qr_click_analytics, daily_usage, and domain_events tables. -Schedule this function to run monthly using pg_cron or external scheduler: -Example: SELECT cron.schedule(''create-partitions'', ''0 0 25 * *'', ''SELECT create_monthly_partitions()''); - -Or run manually before month end: -SELECT create_monthly_partitions();'; + 'Automatically creates next month partition for qr_click_analytics, daily_usage, and domain_events tables. + Schedule this function to run monthly using pg_cron or external scheduler: + Example: SELECT cron.schedule(''create-partitions'', ''0 0 25 * *'', ''SELECT create_monthly_partitions()'');'; diff --git a/tests/integration/conftest.py b/tests/integration/conftest.py index ca85243..c728127 100644 --- a/tests/integration/conftest.py +++ b/tests/integration/conftest.py @@ -100,6 +100,7 @@ async def _apply_schema(pool: asyncpg.Pool) -> None: async with pool.acquire() as conn: # Drop existing tables to start fresh await conn.execute(""" + DROP TABLE IF EXISTS event_dedup CASCADE; DROP TABLE IF EXISTS domain_events CASCADE; DROP TABLE IF EXISTS feature_flags CASCADE; DROP TABLE IF EXISTS click_analytics CASCADE; @@ -234,22 +235,15 @@ async def _apply_schema(pool: asyncpg.Pool) -> None: # Create daily_usage table - updated schema to match production await conn.execute(""" CREATE TABLE daily_usage ( - id INTEGER NOT NULL, user_id BIGINT NOT NULL REFERENCES users(telegram_id) ON DELETE CASCADE, usage_date DATE NOT NULL, - qr_count INTEGER DEFAULT 0, + qr_count INTEGER NOT NULL DEFAULT 0 CHECK (qr_count >= 0), created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP, updated_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP, - PRIMARY KEY (id, usage_date), - UNIQUE(user_id, usage_date) + PRIMARY KEY (user_id, usage_date) ) PARTITION BY RANGE (usage_date) """) - # Create daily_usage partitions for testing - await conn.execute(""" - CREATE SEQUENCE IF NOT EXISTS daily_usage_id_seq - """) - # Create partition for current month (2026-01) await conn.execute(""" CREATE TABLE IF NOT EXISTS daily_usage_2026_01 PARTITION OF daily_usage @@ -282,6 +276,14 @@ async def _apply_schema(pool: asyncpg.Pool) -> None: ) """) + # Global event dedup table for strict idempotency + await conn.execute(""" + CREATE TABLE event_dedup ( + event_id UUID PRIMARY KEY, + created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP + ) + """) + # Create domain_events table (event store) await conn.execute(""" CREATE TABLE domain_events ( @@ -319,6 +321,7 @@ async def _clean_tables(pool: asyncpg.Pool) -> None: # Order matters: respect foreign keys await conn.execute(""" TRUNCATE TABLE + event_dedup, domain_events, feature_flags, click_analytics, diff --git a/tests/unit/infrastructure/event_store/test_postgres_event_store.py b/tests/unit/infrastructure/event_store/test_postgres_event_store.py index 4956045..840d036 100644 --- a/tests/unit/infrastructure/event_store/test_postgres_event_store.py +++ b/tests/unit/infrastructure/event_store/test_postgres_event_store.py @@ -10,7 +10,6 @@ from unittest.mock import AsyncMock, MagicMock from uuid import uuid4 -import asyncpg import pytest from bot.domain.events.base import BaseEvent, EventMeta @@ -52,7 +51,15 @@ def mock_pool(): @pytest.fixture def mock_connection(): """Create a mocked asyncpg connection.""" - return AsyncMock() + conn = AsyncMock() + conn.execute.return_value = "INSERT 0 1" + + tx = AsyncMock() + tx.__aenter__ = AsyncMock(return_value=None) + tx.__aexit__ = AsyncMock(return_value=None) + conn.transaction = MagicMock(return_value=tx) + + return conn @pytest.fixture @@ -142,6 +149,9 @@ async def test_append_stores_event( result = await event_store.append(sample_event) assert result == 42 + mock_connection.execute.assert_called_once() + claim_sql = mock_connection.execute.call_args[0][0] + assert "INSERT INTO event_dedup" in claim_sql mock_connection.fetchval.assert_called_once() # Verify SQL contains correct columns @@ -172,15 +182,14 @@ async def test_append_sends_correct_values( async def test_append_raises_duplicate_error( self, event_store, mock_connection, sample_event ): - """Should raise DuplicateEventError on unique violation.""" - # Create UniqueViolationError without kwargs - error = asyncpg.UniqueViolationError() - mock_connection.fetchval.side_effect = error + """Should raise DuplicateEventError when dedup claim fails.""" + mock_connection.execute.return_value = "INSERT 0 0" with pytest.raises(DuplicateEventError) as exc_info: await event_store.append(sample_event) assert exc_info.value.event_id == sample_event.event_id + mock_connection.fetchval.assert_not_called() class TestPostgresEventStoreAppendBatch: @@ -201,15 +210,10 @@ async def test_append_batch_stores_all_events(self, event_store, mock_connection ] mock_connection.fetchval.side_effect = [10, 20, 30] - # Mock transaction context manager properly - mock_transaction = AsyncMock() - mock_transaction.__aenter__ = AsyncMock() - mock_transaction.__aexit__ = AsyncMock(return_value=None) - mock_connection.transaction = MagicMock(return_value=mock_transaction) - result = await event_store.append_batch(events) assert result == [10, 20, 30] + assert mock_connection.execute.call_count == 3 assert mock_connection.fetchval.call_count == 3 async def test_append_batch_raises_on_duplicate(self, event_store, mock_connection): @@ -218,16 +222,8 @@ async def test_append_batch_raises_on_duplicate(self, event_store, mock_connecti SampleEvent(user_id=1, data="first"), SampleEvent(user_id=2, data="second"), ] - mock_connection.fetchval.side_effect = [ - 10, - asyncpg.UniqueViolationError(), - ] - - # Mock transaction context manager properly - mock_transaction = AsyncMock() - mock_transaction.__aenter__ = AsyncMock() - mock_transaction.__aexit__ = AsyncMock(return_value=None) - mock_connection.transaction = MagicMock(return_value=mock_transaction) + mock_connection.execute.side_effect = ["INSERT 0 1", "INSERT 0 0"] + mock_connection.fetchval.return_value = 10 with pytest.raises(DuplicateEventError): await event_store.append_batch(events)