3230 lines
118 KiB
Python
3230 lines
118 KiB
Python
"""PostgreSQL БД бота: пользователи, баланс, покупки, инвойсы, сессии."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import json
|
|
import re
|
|
import secrets
|
|
from dataclasses import dataclass
|
|
from datetime import date, datetime, timedelta, timezone
|
|
from typing import Any
|
|
|
|
from psycopg import errors as pg_errors
|
|
from psycopg.rows import dict_row
|
|
from psycopg_pool import ConnectionPool
|
|
|
|
from services.passwords import hash_password, new_session_token, verify_password
|
|
|
|
LOGIN_RE = re.compile(r"^[a-zA-Z0-9_]{3,32}$")
|
|
EMAIL_RE = re.compile(r"^[a-zA-Z0-9._%+\-]+@[a-zA-Z0-9.\-]+\.[a-zA-Z]{2,}$")
|
|
QR_LOGIN_TTL = timedelta(minutes=5)
|
|
WEBAUTHN_CHALLENGE_TTL = timedelta(minutes=5)
|
|
SESSION_DAYS = 30
|
|
# Блокировка миграций схемы (bot + webapp стартуют вместе)
|
|
_SCHEMA_ADVISORY_LOCK = 872_014_401
|
|
|
|
|
|
def _utc_now() -> datetime:
|
|
return datetime.now(timezone.utc)
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class BotUser:
|
|
telegram_id: int
|
|
username: str | None
|
|
full_name: str | None
|
|
balance: float
|
|
created_at: datetime | str
|
|
login: str | None = None
|
|
has_password: bool = False
|
|
is_telegram: bool = True
|
|
extra_devices: int = 0
|
|
device_blocked: bool = False
|
|
extra_traffic_packs: int = 0
|
|
extras_next_billing_at: datetime | None = None
|
|
extras_active: bool = True
|
|
auto_renew: bool = True
|
|
email: str | None = None
|
|
banned: bool = False
|
|
referral_code: str | None = None
|
|
referred_by: int | None = None
|
|
trial_used_at: datetime | None = None
|
|
notify_telegram: bool = True
|
|
notify_email: bool = True
|
|
|
|
@property
|
|
def display_name(self) -> str:
|
|
if self.full_name:
|
|
return self.full_name
|
|
if self.login:
|
|
return self.login
|
|
if self.username:
|
|
return f"@{self.username}"
|
|
return f"ID {self.telegram_id}"
|
|
|
|
|
|
class AuthError(ValueError):
|
|
pass
|
|
|
|
|
|
class Database:
|
|
def __init__(self, database_url: str) -> None:
|
|
if not database_url:
|
|
raise ValueError("DATABASE_URL пуст")
|
|
self.database_url = database_url
|
|
self.pool = ConnectionPool(
|
|
conninfo=database_url,
|
|
min_size=1,
|
|
max_size=10,
|
|
kwargs={"row_factory": dict_row},
|
|
open=True,
|
|
)
|
|
self.init_schema()
|
|
|
|
def close(self) -> None:
|
|
self.pool.close()
|
|
|
|
def init_schema(self) -> None:
|
|
statements = [
|
|
"""
|
|
CREATE TABLE IF NOT EXISTS users (
|
|
telegram_id BIGINT PRIMARY KEY,
|
|
username TEXT,
|
|
full_name TEXT,
|
|
balance NUMERIC(12, 2) NOT NULL DEFAULT 0,
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
|
|
)
|
|
""",
|
|
"""
|
|
ALTER TABLE users ADD COLUMN IF NOT EXISTS login TEXT
|
|
""",
|
|
"""
|
|
ALTER TABLE users ADD COLUMN IF NOT EXISTS password_hash TEXT
|
|
""",
|
|
"""
|
|
ALTER TABLE users ADD COLUMN IF NOT EXISTS session_token TEXT
|
|
""",
|
|
"""
|
|
ALTER TABLE users ADD COLUMN IF NOT EXISTS session_expires TIMESTAMPTZ
|
|
""",
|
|
"""
|
|
ALTER TABLE users ADD COLUMN IF NOT EXISTS extra_devices INTEGER NOT NULL DEFAULT 0
|
|
""",
|
|
"""
|
|
ALTER TABLE users ADD COLUMN IF NOT EXISTS device_blocked BOOLEAN NOT NULL DEFAULT FALSE
|
|
""",
|
|
"""
|
|
ALTER TABLE users ADD COLUMN IF NOT EXISTS extra_traffic_packs INTEGER NOT NULL DEFAULT 0
|
|
""",
|
|
"""
|
|
ALTER TABLE users ADD COLUMN IF NOT EXISTS extras_next_billing_at TIMESTAMPTZ
|
|
""",
|
|
"""
|
|
ALTER TABLE users ADD COLUMN IF NOT EXISTS extras_active BOOLEAN NOT NULL DEFAULT TRUE
|
|
""",
|
|
"""
|
|
ALTER TABLE users ADD COLUMN IF NOT EXISTS expiry_reminder_sent_on DATE
|
|
""",
|
|
"""
|
|
ALTER TABLE users ADD COLUMN IF NOT EXISTS auto_renew BOOLEAN NOT NULL DEFAULT TRUE
|
|
""",
|
|
"""
|
|
ALTER TABLE users ADD COLUMN IF NOT EXISTS email TEXT
|
|
""",
|
|
"""
|
|
CREATE UNIQUE INDEX IF NOT EXISTS idx_users_login_unique
|
|
ON users (login) WHERE login IS NOT NULL
|
|
""",
|
|
"""
|
|
CREATE UNIQUE INDEX IF NOT EXISTS idx_users_email_unique
|
|
ON users (lower(email)) WHERE email IS NOT NULL AND email <> ''
|
|
""",
|
|
"""
|
|
CREATE UNIQUE INDEX IF NOT EXISTS idx_users_session_token
|
|
ON users (session_token) WHERE session_token IS NOT NULL
|
|
""",
|
|
"""
|
|
CREATE TABLE IF NOT EXISTS login_qr (
|
|
token TEXT PRIMARY KEY,
|
|
status TEXT NOT NULL DEFAULT 'pending',
|
|
telegram_id BIGINT,
|
|
session_token TEXT,
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
|
expires_at TIMESTAMPTZ NOT NULL,
|
|
approved_at TIMESTAMPTZ
|
|
)
|
|
""",
|
|
"""
|
|
CREATE INDEX IF NOT EXISTS idx_login_qr_expires
|
|
ON login_qr (expires_at)
|
|
""",
|
|
"""
|
|
CREATE TABLE IF NOT EXISTS webauthn_credentials (
|
|
credential_id TEXT PRIMARY KEY,
|
|
telegram_id BIGINT NOT NULL REFERENCES users(telegram_id) ON DELETE CASCADE,
|
|
public_key BYTEA NOT NULL,
|
|
sign_count BIGINT NOT NULL DEFAULT 0,
|
|
transports TEXT NOT NULL DEFAULT '[]',
|
|
device_name TEXT,
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
|
last_used_at TIMESTAMPTZ
|
|
)
|
|
""",
|
|
"""
|
|
CREATE INDEX IF NOT EXISTS idx_webauthn_credentials_user
|
|
ON webauthn_credentials (telegram_id)
|
|
""",
|
|
"""
|
|
CREATE TABLE IF NOT EXISTS webauthn_challenges (
|
|
challenge TEXT PRIMARY KEY,
|
|
telegram_id BIGINT,
|
|
purpose TEXT NOT NULL,
|
|
expires_at TIMESTAMPTZ NOT NULL,
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
|
|
)
|
|
""",
|
|
"""
|
|
CREATE INDEX IF NOT EXISTS idx_webauthn_challenges_expires
|
|
ON webauthn_challenges (expires_at)
|
|
""",
|
|
"""
|
|
ALTER TABLE users ADD COLUMN IF NOT EXISTS banned BOOLEAN NOT NULL DEFAULT FALSE
|
|
""",
|
|
"""
|
|
ALTER TABLE users ADD COLUMN IF NOT EXISTS referral_code TEXT
|
|
""",
|
|
"""
|
|
ALTER TABLE users ADD COLUMN IF NOT EXISTS referred_by BIGINT
|
|
""",
|
|
"""
|
|
ALTER TABLE users ADD COLUMN IF NOT EXISTS trial_used_at TIMESTAMPTZ
|
|
""",
|
|
"""
|
|
ALTER TABLE users ADD COLUMN IF NOT EXISTS notify_telegram BOOLEAN NOT NULL DEFAULT TRUE
|
|
""",
|
|
"""
|
|
ALTER TABLE users ADD COLUMN IF NOT EXISTS notify_email BOOLEAN NOT NULL DEFAULT TRUE
|
|
""",
|
|
"""
|
|
CREATE UNIQUE INDEX IF NOT EXISTS idx_users_referral_code
|
|
ON users (referral_code) WHERE referral_code IS NOT NULL
|
|
""",
|
|
"""
|
|
CREATE TABLE IF NOT EXISTS user_sessions (
|
|
token TEXT PRIMARY KEY,
|
|
telegram_id BIGINT NOT NULL REFERENCES users(telegram_id) ON DELETE CASCADE,
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
|
expires_at TIMESTAMPTZ NOT NULL,
|
|
user_agent TEXT,
|
|
ip TEXT
|
|
)
|
|
""",
|
|
"""
|
|
CREATE INDEX IF NOT EXISTS idx_user_sessions_user
|
|
ON user_sessions (telegram_id, expires_at DESC)
|
|
""",
|
|
"""
|
|
CREATE TABLE IF NOT EXISTS promo_codes (
|
|
id BIGSERIAL PRIMARY KEY,
|
|
code TEXT NOT NULL,
|
|
kind TEXT NOT NULL,
|
|
value NUMERIC(12, 2) NOT NULL,
|
|
max_uses INTEGER,
|
|
used_count INTEGER NOT NULL DEFAULT 0,
|
|
tariff_id TEXT,
|
|
active BOOLEAN NOT NULL DEFAULT TRUE,
|
|
expires_at TIMESTAMPTZ,
|
|
note TEXT,
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
|
|
)
|
|
""",
|
|
"""
|
|
CREATE UNIQUE INDEX IF NOT EXISTS idx_promo_codes_code
|
|
ON promo_codes (lower(code))
|
|
""",
|
|
"""
|
|
CREATE TABLE IF NOT EXISTS promo_redemptions (
|
|
id BIGSERIAL PRIMARY KEY,
|
|
promo_id BIGINT NOT NULL REFERENCES promo_codes(id) ON DELETE CASCADE,
|
|
telegram_id BIGINT NOT NULL REFERENCES users(telegram_id) ON DELETE CASCADE,
|
|
amount_saved NUMERIC(12, 2) NOT NULL DEFAULT 0,
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
|
UNIQUE (promo_id, telegram_id)
|
|
)
|
|
""",
|
|
"""
|
|
CREATE TABLE IF NOT EXISTS password_reset_tokens (
|
|
token TEXT PRIMARY KEY,
|
|
telegram_id BIGINT NOT NULL REFERENCES users(telegram_id) ON DELETE CASCADE,
|
|
expires_at TIMESTAMPTZ NOT NULL,
|
|
used_at TIMESTAMPTZ,
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
|
|
)
|
|
""",
|
|
"""
|
|
CREATE TABLE IF NOT EXISTS payment_events (
|
|
id BIGSERIAL PRIMARY KEY,
|
|
provider TEXT NOT NULL,
|
|
event_type TEXT NOT NULL,
|
|
telegram_id BIGINT,
|
|
amount NUMERIC(12, 2),
|
|
invoice_ref TEXT,
|
|
status TEXT NOT NULL DEFAULT 'ok',
|
|
meta TEXT,
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
|
|
)
|
|
""",
|
|
"""
|
|
CREATE INDEX IF NOT EXISTS idx_payment_events_created
|
|
ON payment_events (created_at DESC)
|
|
""",
|
|
"""
|
|
CREATE TABLE IF NOT EXISTS referral_rewards (
|
|
id BIGSERIAL PRIMARY KEY,
|
|
referrer_id BIGINT NOT NULL REFERENCES users(telegram_id),
|
|
referee_id BIGINT NOT NULL REFERENCES users(telegram_id),
|
|
amount NUMERIC(12, 2) NOT NULL,
|
|
kind TEXT NOT NULL,
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
|
UNIQUE (referee_id, kind)
|
|
)
|
|
""",
|
|
"""
|
|
CREATE SEQUENCE IF NOT EXISTS web_user_id_seq START WITH 1 INCREMENT BY 1
|
|
""",
|
|
"""
|
|
CREATE TABLE IF NOT EXISTS transactions (
|
|
id BIGSERIAL PRIMARY KEY,
|
|
telegram_id BIGINT NOT NULL REFERENCES users(telegram_id),
|
|
amount NUMERIC(12, 2) NOT NULL,
|
|
kind TEXT NOT NULL,
|
|
meta TEXT,
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
|
|
)
|
|
""",
|
|
"""
|
|
CREATE INDEX IF NOT EXISTS idx_transactions_user_id
|
|
ON transactions (telegram_id, id DESC)
|
|
""",
|
|
"""
|
|
CREATE TABLE IF NOT EXISTS purchases (
|
|
id BIGSERIAL PRIMARY KEY,
|
|
telegram_id BIGINT NOT NULL REFERENCES users(telegram_id),
|
|
tariff_id TEXT NOT NULL,
|
|
price NUMERIC(12, 2) NOT NULL,
|
|
remnawave_uuid TEXT,
|
|
subscription_url TEXT,
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
|
|
)
|
|
""",
|
|
"""
|
|
CREATE TABLE IF NOT EXISTS invoices (
|
|
id BIGSERIAL PRIMARY KEY,
|
|
telegram_id BIGINT NOT NULL REFERENCES users(telegram_id),
|
|
invoice_id BIGINT NOT NULL,
|
|
kind TEXT NOT NULL,
|
|
amount_rub NUMERIC(12, 2) NOT NULL,
|
|
tariff_id TEXT,
|
|
status TEXT NOT NULL DEFAULT 'active',
|
|
pay_url TEXT,
|
|
payload TEXT,
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
|
paid_at TIMESTAMPTZ
|
|
)
|
|
""",
|
|
"""
|
|
CREATE SEQUENCE IF NOT EXISTS local_invoice_id_seq START WITH 900000000001
|
|
""",
|
|
"""
|
|
ALTER TABLE invoices ADD COLUMN IF NOT EXISTS provider TEXT NOT NULL DEFAULT 'cryptobot'
|
|
""",
|
|
"""
|
|
ALTER TABLE invoices ADD COLUMN IF NOT EXISTS order_id TEXT
|
|
""",
|
|
"""
|
|
ALTER TABLE invoices ADD COLUMN IF NOT EXISTS external_uuid TEXT
|
|
""",
|
|
"""
|
|
CREATE UNIQUE INDEX IF NOT EXISTS idx_invoices_provider_invoice
|
|
ON invoices (provider, invoice_id)
|
|
""",
|
|
"""
|
|
CREATE UNIQUE INDEX IF NOT EXISTS idx_invoices_order_id
|
|
ON invoices (order_id) WHERE order_id IS NOT NULL
|
|
""",
|
|
"""
|
|
CREATE TABLE IF NOT EXISTS digiseller_payments (
|
|
inv BIGINT PRIMARY KEY,
|
|
unique_code TEXT NOT NULL,
|
|
telegram_id BIGINT NOT NULL REFERENCES users(telegram_id),
|
|
amount_rub NUMERIC(12, 2) NOT NULL,
|
|
product_id BIGINT,
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
|
|
)
|
|
""",
|
|
"""
|
|
CREATE UNIQUE INDEX IF NOT EXISTS idx_digiseller_unique_code
|
|
ON digiseller_payments (unique_code)
|
|
""",
|
|
"""
|
|
CREATE TABLE IF NOT EXISTS topup_codes (
|
|
code TEXT PRIMARY KEY,
|
|
amount_rub NUMERIC(12, 2) NOT NULL,
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
|
created_by TEXT,
|
|
used_at TIMESTAMPTZ,
|
|
used_by BIGINT REFERENCES users(telegram_id)
|
|
)
|
|
""",
|
|
"""
|
|
CREATE INDEX IF NOT EXISTS idx_topup_codes_used
|
|
ON topup_codes (used_at NULLS FIRST, created_at DESC)
|
|
""",
|
|
"""
|
|
CREATE TABLE IF NOT EXISTS tickets (
|
|
id BIGSERIAL PRIMARY KEY,
|
|
telegram_id BIGINT NOT NULL REFERENCES users(telegram_id),
|
|
subject TEXT NOT NULL,
|
|
status TEXT NOT NULL DEFAULT 'open',
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
|
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
|
|
)
|
|
""",
|
|
"""
|
|
CREATE INDEX IF NOT EXISTS idx_tickets_user
|
|
ON tickets (telegram_id, updated_at DESC)
|
|
""",
|
|
"""
|
|
CREATE INDEX IF NOT EXISTS idx_tickets_status
|
|
ON tickets (status, updated_at DESC)
|
|
""",
|
|
"""
|
|
CREATE TABLE IF NOT EXISTS ticket_messages (
|
|
id BIGSERIAL PRIMARY KEY,
|
|
ticket_id BIGINT NOT NULL REFERENCES tickets(id) ON DELETE CASCADE,
|
|
author_role TEXT NOT NULL,
|
|
author_id BIGINT,
|
|
body TEXT NOT NULL,
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
|
|
)
|
|
""",
|
|
"""
|
|
CREATE INDEX IF NOT EXISTS idx_ticket_messages_ticket
|
|
ON ticket_messages (ticket_id, id ASC)
|
|
""",
|
|
"""
|
|
CREATE TABLE IF NOT EXISTS app_settings (
|
|
key TEXT PRIMARY KEY,
|
|
value TEXT NOT NULL DEFAULT '',
|
|
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
|
|
)
|
|
""",
|
|
"""
|
|
CREATE TABLE IF NOT EXISTS xui_servers (
|
|
id BIGSERIAL PRIMARY KEY,
|
|
name TEXT NOT NULL,
|
|
api_url TEXT NOT NULL,
|
|
username TEXT NOT NULL DEFAULT '',
|
|
password TEXT NOT NULL DEFAULT '',
|
|
api_key TEXT NOT NULL DEFAULT '',
|
|
inbound_id INTEGER,
|
|
inbound_remark TEXT,
|
|
subscription_url TEXT NOT NULL DEFAULT '',
|
|
enabled BOOLEAN NOT NULL DEFAULT TRUE,
|
|
verify_ssl BOOLEAN NOT NULL DEFAULT TRUE,
|
|
last_check_at TIMESTAMPTZ,
|
|
last_check_ok BOOLEAN,
|
|
last_check_msg TEXT,
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
|
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
|
|
)
|
|
""",
|
|
"""
|
|
ALTER TABLE xui_servers ADD COLUMN IF NOT EXISTS api_key TEXT NOT NULL DEFAULT ''
|
|
""",
|
|
"""
|
|
CREATE TABLE IF NOT EXISTS xui_wg_configs (
|
|
id BIGSERIAL PRIMARY KEY,
|
|
telegram_id BIGINT NOT NULL REFERENCES users(telegram_id),
|
|
server_id BIGINT NOT NULL REFERENCES xui_servers(id) ON DELETE CASCADE,
|
|
email TEXT NOT NULL,
|
|
sub_id TEXT NOT NULL,
|
|
subscription_url TEXT NOT NULL,
|
|
remark TEXT,
|
|
enabled BOOLEAN NOT NULL DEFAULT TRUE,
|
|
expire_at TIMESTAMPTZ,
|
|
traffic_bytes BIGINT NOT NULL DEFAULT 0,
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
|
UNIQUE (server_id, email)
|
|
)
|
|
""",
|
|
"""
|
|
ALTER TABLE xui_wg_configs ADD COLUMN IF NOT EXISTS expire_at TIMESTAMPTZ
|
|
""",
|
|
"""
|
|
ALTER TABLE xui_wg_configs ADD COLUMN IF NOT EXISTS traffic_bytes BIGINT NOT NULL DEFAULT 0
|
|
""",
|
|
"""
|
|
CREATE INDEX IF NOT EXISTS idx_xui_wg_configs_user
|
|
ON xui_wg_configs (telegram_id, id DESC)
|
|
""",
|
|
"""
|
|
CREATE INDEX IF NOT EXISTS idx_xui_servers_enabled
|
|
ON xui_servers (enabled, id)
|
|
""",
|
|
"""
|
|
CREATE TABLE IF NOT EXISTS awg_servers (
|
|
id BIGSERIAL PRIMARY KEY,
|
|
name TEXT NOT NULL,
|
|
api_url TEXT NOT NULL,
|
|
api_token TEXT NOT NULL DEFAULT '',
|
|
panel_server_id INTEGER NOT NULL DEFAULT 0,
|
|
protocol TEXT NOT NULL DEFAULT 'awg2',
|
|
enabled BOOLEAN NOT NULL DEFAULT TRUE,
|
|
verify_ssl BOOLEAN NOT NULL DEFAULT TRUE,
|
|
last_check_at TIMESTAMPTZ,
|
|
last_check_ok BOOLEAN,
|
|
last_check_msg TEXT,
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
|
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
|
|
)
|
|
""",
|
|
"""
|
|
CREATE TABLE IF NOT EXISTS awg_configs (
|
|
id BIGSERIAL PRIMARY KEY,
|
|
telegram_id BIGINT NOT NULL REFERENCES users(telegram_id),
|
|
server_id BIGINT NOT NULL REFERENCES awg_servers(id) ON DELETE CASCADE,
|
|
client_id TEXT NOT NULL,
|
|
client_name TEXT NOT NULL DEFAULT '',
|
|
remark TEXT,
|
|
config_text TEXT NOT NULL DEFAULT '',
|
|
vpn_link TEXT NOT NULL DEFAULT '',
|
|
protocol TEXT NOT NULL DEFAULT 'awg2',
|
|
enabled BOOLEAN NOT NULL DEFAULT TRUE,
|
|
expire_at TIMESTAMPTZ,
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
|
UNIQUE (server_id, client_id)
|
|
)
|
|
""",
|
|
"""
|
|
CREATE INDEX IF NOT EXISTS idx_awg_configs_user
|
|
ON awg_configs (telegram_id, id DESC)
|
|
""",
|
|
"""
|
|
CREATE INDEX IF NOT EXISTS idx_awg_servers_enabled
|
|
ON awg_servers (enabled, id)
|
|
""",
|
|
]
|
|
with self.pool.connection() as conn:
|
|
# Один процесс за раз — иначе гонка CREATE IF NOT EXISTS → 42P07
|
|
conn.execute("SELECT pg_advisory_lock(%s)", (_SCHEMA_ADVISORY_LOCK,))
|
|
try:
|
|
for stmt in statements:
|
|
try:
|
|
with conn.transaction():
|
|
conn.execute(stmt)
|
|
except (pg_errors.DuplicateTable, pg_errors.DuplicateObject):
|
|
# Параллельный старт / повторный apply — таблица уже есть
|
|
continue
|
|
# Старый UNIQUE(invoice_id) мешает разным провайдерам — убираем, если есть
|
|
with conn.transaction():
|
|
conn.execute(
|
|
"""
|
|
DO $$
|
|
BEGIN
|
|
IF EXISTS (
|
|
SELECT 1 FROM pg_constraint
|
|
WHERE conname = 'invoices_invoice_id_key'
|
|
) THEN
|
|
ALTER TABLE invoices DROP CONSTRAINT invoices_invoice_id_key;
|
|
END IF;
|
|
END $$;
|
|
"""
|
|
)
|
|
with conn.transaction():
|
|
conn.execute(
|
|
"""
|
|
UPDATE users
|
|
SET extras_next_billing_at = NOW() + INTERVAL '30 days'
|
|
WHERE extras_next_billing_at IS NULL
|
|
AND extra_devices > 0
|
|
"""
|
|
)
|
|
conn.commit()
|
|
finally:
|
|
conn.execute("SELECT pg_advisory_unlock(%s)", (_SCHEMA_ADVISORY_LOCK,))
|
|
conn.commit()
|
|
|
|
def ensure_user(
|
|
self,
|
|
telegram_id: int,
|
|
*,
|
|
username: str | None = None,
|
|
full_name: str | None = None,
|
|
) -> BotUser:
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"SELECT * FROM users WHERE telegram_id = %s",
|
|
(telegram_id,),
|
|
).fetchone()
|
|
if row:
|
|
conn.execute(
|
|
"""
|
|
UPDATE users
|
|
SET username = COALESCE(%s, username),
|
|
full_name = COALESCE(%s, full_name)
|
|
WHERE telegram_id = %s
|
|
""",
|
|
(username, full_name, telegram_id),
|
|
)
|
|
else:
|
|
conn.execute(
|
|
"""
|
|
INSERT INTO users (telegram_id, username, full_name, balance, created_at)
|
|
VALUES (%s, %s, %s, 0, %s)
|
|
""",
|
|
(telegram_id, username, full_name, _utc_now()),
|
|
)
|
|
conn.commit()
|
|
row = conn.execute(
|
|
"SELECT * FROM users WHERE telegram_id = %s",
|
|
(telegram_id,),
|
|
).fetchone()
|
|
assert row is not None
|
|
return self._to_user(row)
|
|
|
|
def get_user(self, telegram_id: int) -> BotUser | None:
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"SELECT * FROM users WHERE telegram_id = %s",
|
|
(telegram_id,),
|
|
).fetchone()
|
|
return self._to_user(row) if row else None
|
|
|
|
def get_user_by_login(self, login: str) -> BotUser | None:
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"SELECT * FROM users WHERE lower(login) = lower(%s)",
|
|
(login.strip(),),
|
|
).fetchone()
|
|
return self._to_user(row) if row else None
|
|
|
|
def get_user_by_session(self, token: str | None) -> BotUser | None:
|
|
if not token:
|
|
return None
|
|
now = _utc_now()
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"""
|
|
SELECT u.*
|
|
FROM user_sessions s
|
|
JOIN users u ON u.telegram_id = s.telegram_id
|
|
WHERE s.token = %s AND s.expires_at > %s
|
|
""",
|
|
(token, now),
|
|
).fetchone()
|
|
if row:
|
|
user = self._to_user(row)
|
|
if user.banned:
|
|
return None
|
|
return user
|
|
row = conn.execute(
|
|
"""
|
|
SELECT * FROM users
|
|
WHERE session_token = %s
|
|
AND session_expires IS NOT NULL
|
|
AND session_expires > %s
|
|
""",
|
|
(token, now),
|
|
).fetchone()
|
|
if not row:
|
|
return None
|
|
user = self._to_user(row)
|
|
if user.banned:
|
|
return None
|
|
return user
|
|
|
|
def create_session(
|
|
self,
|
|
telegram_id: int,
|
|
*,
|
|
user_agent: str | None = None,
|
|
ip: str | None = None,
|
|
) -> str:
|
|
token = new_session_token()
|
|
expires = _utc_now() + timedelta(days=SESSION_DAYS)
|
|
with self.pool.connection() as conn:
|
|
conn.execute(
|
|
"""
|
|
INSERT INTO user_sessions (token, telegram_id, created_at, expires_at, user_agent, ip)
|
|
VALUES (%s, %s, %s, %s, %s, %s)
|
|
""",
|
|
(
|
|
token,
|
|
telegram_id,
|
|
_utc_now(),
|
|
expires,
|
|
(user_agent or "")[:300] or None,
|
|
(ip or "")[:80] or None,
|
|
),
|
|
)
|
|
# legacy single-token column (совместимость)
|
|
conn.execute(
|
|
"""
|
|
UPDATE users
|
|
SET session_token = %s, session_expires = %s
|
|
WHERE telegram_id = %s
|
|
""",
|
|
(token, expires, telegram_id),
|
|
)
|
|
conn.commit()
|
|
return token
|
|
|
|
def clear_session(self, token: str | None) -> None:
|
|
if not token:
|
|
return
|
|
with self.pool.connection() as conn:
|
|
conn.execute("DELETE FROM user_sessions WHERE token = %s", (token,))
|
|
conn.execute(
|
|
"""
|
|
UPDATE users
|
|
SET session_token = NULL, session_expires = NULL
|
|
WHERE session_token = %s
|
|
""",
|
|
(token,),
|
|
)
|
|
conn.commit()
|
|
|
|
def clear_all_sessions(self, telegram_id: int, *, keep_token: str | None = None) -> int:
|
|
with self.pool.connection() as conn:
|
|
if keep_token:
|
|
row = conn.execute(
|
|
"""
|
|
DELETE FROM user_sessions
|
|
WHERE telegram_id = %s AND token <> %s
|
|
""",
|
|
(telegram_id, keep_token),
|
|
)
|
|
else:
|
|
row = conn.execute(
|
|
"DELETE FROM user_sessions WHERE telegram_id = %s",
|
|
(telegram_id,),
|
|
)
|
|
deleted = row.rowcount if row is not None else 0
|
|
if not keep_token:
|
|
conn.execute(
|
|
"""
|
|
UPDATE users
|
|
SET session_token = NULL, session_expires = NULL
|
|
WHERE telegram_id = %s
|
|
""",
|
|
(telegram_id,),
|
|
)
|
|
conn.commit()
|
|
return int(deleted or 0)
|
|
|
|
def list_sessions(self, telegram_id: int) -> list[dict[str, Any]]:
|
|
with self.pool.connection() as conn:
|
|
rows = conn.execute(
|
|
"""
|
|
SELECT token, created_at, expires_at, user_agent, ip
|
|
FROM user_sessions
|
|
WHERE telegram_id = %s AND expires_at > %s
|
|
ORDER BY created_at DESC
|
|
LIMIT 20
|
|
""",
|
|
(telegram_id, _utc_now()),
|
|
).fetchall()
|
|
return [dict(r) for r in rows]
|
|
|
|
def set_user_email(self, telegram_id: int, email: str | None) -> BotUser:
|
|
raw = (email or "").strip().lower()
|
|
if raw and not EMAIL_RE.fullmatch(raw):
|
|
raise AuthError("Некорректный email")
|
|
value = raw or None
|
|
with self.pool.connection() as conn:
|
|
if value:
|
|
clash = conn.execute(
|
|
"""
|
|
SELECT telegram_id FROM users
|
|
WHERE lower(email) = %s AND telegram_id <> %s
|
|
LIMIT 1
|
|
""",
|
|
(value, telegram_id),
|
|
).fetchone()
|
|
if clash:
|
|
raise AuthError("Этот email уже привязан к другому аккаунту")
|
|
conn.execute(
|
|
"UPDATE users SET email = %s WHERE telegram_id = %s",
|
|
(value, telegram_id),
|
|
)
|
|
conn.commit()
|
|
user = self.get_user(telegram_id)
|
|
if not user:
|
|
raise AuthError("Пользователь не найден")
|
|
return user
|
|
|
|
def create_login_qr(self) -> dict[str, Any]:
|
|
token = secrets.token_urlsafe(24)
|
|
now = _utc_now()
|
|
expires = now + QR_LOGIN_TTL
|
|
with self.pool.connection() as conn:
|
|
conn.execute(
|
|
"""
|
|
INSERT INTO login_qr (token, status, created_at, expires_at)
|
|
VALUES (%s, 'pending', %s, %s)
|
|
""",
|
|
(token, now, expires),
|
|
)
|
|
conn.commit()
|
|
return {"token": token, "expires_at": expires, "expires_in": int(QR_LOGIN_TTL.total_seconds())}
|
|
|
|
def approve_login_qr(self, token: str, telegram_id: int) -> None:
|
|
token = (token or "").strip()
|
|
if not token:
|
|
raise AuthError("Некорректный QR-токен")
|
|
now = _utc_now()
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"SELECT * FROM login_qr WHERE token = %s",
|
|
(token,),
|
|
).fetchone()
|
|
if not row:
|
|
raise AuthError("QR-сессия не найдена")
|
|
if row["status"] != "pending":
|
|
raise AuthError("QR уже использован или подтверждён")
|
|
if row["expires_at"] <= now:
|
|
conn.execute(
|
|
"UPDATE login_qr SET status = 'expired' WHERE token = %s",
|
|
(token,),
|
|
)
|
|
conn.commit()
|
|
raise AuthError("QR истёк — обновите код на сайте")
|
|
updated = conn.execute(
|
|
"""
|
|
UPDATE login_qr
|
|
SET status = 'approved',
|
|
telegram_id = %s,
|
|
approved_at = %s
|
|
WHERE token = %s AND status = 'pending' AND expires_at > %s
|
|
RETURNING token
|
|
""",
|
|
(telegram_id, now, token, now),
|
|
).fetchone()
|
|
conn.commit()
|
|
if not updated:
|
|
raise AuthError("Не удалось подтвердить QR")
|
|
|
|
def poll_login_qr(self, token: str) -> dict[str, Any]:
|
|
token = (token or "").strip()
|
|
if not token:
|
|
raise AuthError("Некорректный QR-токен")
|
|
now = _utc_now()
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"SELECT * FROM login_qr WHERE token = %s",
|
|
(token,),
|
|
).fetchone()
|
|
if not row:
|
|
raise AuthError("QR-сессия не найдена")
|
|
if row["expires_at"] <= now and row["status"] == "pending":
|
|
conn.execute(
|
|
"UPDATE login_qr SET status = 'expired' WHERE token = %s",
|
|
(token,),
|
|
)
|
|
conn.commit()
|
|
return {"status": "expired"}
|
|
if row["status"] == "pending":
|
|
return {"status": "pending"}
|
|
if row["status"] in {"expired", "consumed"}:
|
|
return {"status": row["status"]}
|
|
if row["status"] != "approved":
|
|
return {"status": str(row["status"])}
|
|
|
|
telegram_id = int(row["telegram_id"])
|
|
claimed = conn.execute(
|
|
"""
|
|
UPDATE login_qr
|
|
SET status = 'consumed'
|
|
WHERE token = %s AND status = 'approved'
|
|
RETURNING telegram_id
|
|
""",
|
|
(token,),
|
|
).fetchone()
|
|
conn.commit()
|
|
if not claimed:
|
|
return {"status": "consumed"}
|
|
|
|
session = self.create_session(telegram_id)
|
|
with self.pool.connection() as conn:
|
|
conn.execute(
|
|
"UPDATE login_qr SET session_token = %s WHERE token = %s",
|
|
(session, token),
|
|
)
|
|
user_row = conn.execute(
|
|
"SELECT * FROM users WHERE telegram_id = %s",
|
|
(telegram_id,),
|
|
).fetchone()
|
|
conn.commit()
|
|
if not user_row:
|
|
raise AuthError("Пользователь не найден")
|
|
return {
|
|
"status": "approved",
|
|
"token": session,
|
|
"user": self._to_user(user_row),
|
|
}
|
|
|
|
def save_webauthn_challenge(
|
|
self,
|
|
*,
|
|
challenge_b64: str,
|
|
purpose: str,
|
|
telegram_id: int | None = None,
|
|
) -> None:
|
|
expires = _utc_now() + WEBAUTHN_CHALLENGE_TTL
|
|
with self.pool.connection() as conn:
|
|
conn.execute(
|
|
"DELETE FROM webauthn_challenges WHERE expires_at < %s",
|
|
(_utc_now(),),
|
|
)
|
|
conn.execute(
|
|
"""
|
|
INSERT INTO webauthn_challenges (challenge, telegram_id, purpose, expires_at)
|
|
VALUES (%s, %s, %s, %s)
|
|
ON CONFLICT (challenge) DO UPDATE
|
|
SET telegram_id = EXCLUDED.telegram_id,
|
|
purpose = EXCLUDED.purpose,
|
|
expires_at = EXCLUDED.expires_at
|
|
""",
|
|
(challenge_b64, telegram_id, purpose, expires),
|
|
)
|
|
conn.commit()
|
|
|
|
def take_webauthn_challenge(
|
|
self, challenge_b64: str, *, purpose: str
|
|
) -> dict[str, Any] | None:
|
|
now = _utc_now()
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"""
|
|
DELETE FROM webauthn_challenges
|
|
WHERE challenge = %s AND purpose = %s AND expires_at > %s
|
|
RETURNING challenge, telegram_id, purpose
|
|
""",
|
|
(challenge_b64, purpose, now),
|
|
).fetchone()
|
|
conn.commit()
|
|
return dict(row) if row else None
|
|
|
|
def list_webauthn_credentials(self, telegram_id: int) -> list[dict[str, Any]]:
|
|
with self.pool.connection() as conn:
|
|
rows = conn.execute(
|
|
"""
|
|
SELECT credential_id, telegram_id, public_key, sign_count,
|
|
transports, device_name, created_at, last_used_at
|
|
FROM webauthn_credentials
|
|
WHERE telegram_id = %s
|
|
ORDER BY created_at DESC
|
|
""",
|
|
(telegram_id,),
|
|
).fetchall()
|
|
return [dict(r) for r in rows]
|
|
|
|
def get_webauthn_credential(self, credential_id: str) -> dict[str, Any] | None:
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"""
|
|
SELECT credential_id, telegram_id, public_key, sign_count,
|
|
transports, device_name, created_at, last_used_at
|
|
FROM webauthn_credentials
|
|
WHERE credential_id = %s
|
|
""",
|
|
(credential_id,),
|
|
).fetchone()
|
|
return dict(row) if row else None
|
|
|
|
def add_webauthn_credential(
|
|
self,
|
|
*,
|
|
telegram_id: int,
|
|
credential_id: str,
|
|
public_key: bytes,
|
|
sign_count: int,
|
|
transports: list[str] | None = None,
|
|
device_name: str | None = None,
|
|
) -> None:
|
|
with self.pool.connection() as conn:
|
|
conn.execute(
|
|
"""
|
|
INSERT INTO webauthn_credentials (
|
|
credential_id, telegram_id, public_key, sign_count,
|
|
transports, device_name, created_at
|
|
) VALUES (%s, %s, %s, %s, %s, %s, %s)
|
|
""",
|
|
(
|
|
credential_id,
|
|
telegram_id,
|
|
public_key,
|
|
int(sign_count),
|
|
json.dumps(transports or []),
|
|
(device_name or "").strip() or None,
|
|
_utc_now(),
|
|
),
|
|
)
|
|
conn.commit()
|
|
|
|
def update_webauthn_sign_count(self, credential_id: str, sign_count: int) -> None:
|
|
with self.pool.connection() as conn:
|
|
conn.execute(
|
|
"""
|
|
UPDATE webauthn_credentials
|
|
SET sign_count = %s, last_used_at = %s
|
|
WHERE credential_id = %s
|
|
""",
|
|
(int(sign_count), _utc_now(), credential_id),
|
|
)
|
|
conn.commit()
|
|
|
|
def delete_webauthn_credential(self, telegram_id: int, credential_id: str) -> bool:
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"""
|
|
DELETE FROM webauthn_credentials
|
|
WHERE telegram_id = %s AND credential_id = %s
|
|
RETURNING credential_id
|
|
""",
|
|
(telegram_id, credential_id),
|
|
).fetchone()
|
|
conn.commit()
|
|
return row is not None
|
|
|
|
def register_with_password(
|
|
self,
|
|
*,
|
|
login: str,
|
|
password: str,
|
|
full_name: str | None = None,
|
|
auto_renew: bool = True,
|
|
) -> tuple[BotUser, str]:
|
|
login_norm = login.strip()
|
|
if not LOGIN_RE.fullmatch(login_norm):
|
|
raise AuthError("Логин: 3–32 символа (латиница, цифры, _)")
|
|
if len(password) < 6:
|
|
raise AuthError("Пароль минимум 6 символов")
|
|
if self.get_user_by_login(login_norm):
|
|
raise AuthError("Такой логин уже занят")
|
|
|
|
pwd_hash = hash_password(password)
|
|
with self.pool.connection() as conn:
|
|
web_id = conn.execute("SELECT nextval('web_user_id_seq') AS id").fetchone()
|
|
assert web_id is not None
|
|
telegram_id = -int(web_id["id"])
|
|
try:
|
|
conn.execute(
|
|
"""
|
|
INSERT INTO users (
|
|
telegram_id, username, full_name, balance, created_at,
|
|
login, password_hash, auto_renew
|
|
) VALUES (%s, %s, %s, 0, %s, %s, %s, %s)
|
|
""",
|
|
(
|
|
telegram_id,
|
|
login_norm,
|
|
full_name or login_norm,
|
|
_utc_now(),
|
|
login_norm,
|
|
pwd_hash,
|
|
bool(auto_renew),
|
|
),
|
|
)
|
|
except Exception as exc: # noqa: BLE001
|
|
conn.rollback()
|
|
raise AuthError("Не удалось зарегистрировать") from exc
|
|
conn.commit()
|
|
|
|
token = self.create_session(telegram_id)
|
|
user = self.get_user(telegram_id)
|
|
assert user is not None
|
|
return user, token
|
|
|
|
def login_with_password(self, login: str, password: str) -> tuple[BotUser, str]:
|
|
login_norm = login.strip()
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"SELECT * FROM users WHERE lower(login) = lower(%s)",
|
|
(login_norm,),
|
|
).fetchone()
|
|
if not row or not row.get("password_hash"):
|
|
raise AuthError("Неверный логин или пароль")
|
|
if not verify_password(password, row["password_hash"]):
|
|
raise AuthError("Неверный логин или пароль")
|
|
|
|
user = self._to_user(row)
|
|
token = self.create_session(user.telegram_id)
|
|
return user, token
|
|
|
|
def login_or_register_telegram(
|
|
self,
|
|
telegram_id: int,
|
|
*,
|
|
username: str | None = None,
|
|
full_name: str | None = None,
|
|
link_from: int | None = None,
|
|
) -> tuple[BotUser, str]:
|
|
if telegram_id <= 0:
|
|
raise AuthError("Некорректный Telegram ID")
|
|
|
|
if link_from is not None and link_from < 0 and link_from != telegram_id:
|
|
user = self._merge_web_into_telegram(
|
|
web_id=link_from,
|
|
telegram_id=telegram_id,
|
|
username=username,
|
|
full_name=full_name,
|
|
)
|
|
else:
|
|
user = self.ensure_user(
|
|
telegram_id,
|
|
username=username,
|
|
full_name=full_name,
|
|
)
|
|
token = self.create_session(user.telegram_id)
|
|
return user, token
|
|
|
|
def _merge_web_into_telegram(
|
|
self,
|
|
*,
|
|
web_id: int,
|
|
telegram_id: int,
|
|
username: str | None,
|
|
full_name: str | None,
|
|
) -> BotUser:
|
|
with self.pool.connection() as conn:
|
|
web = conn.execute(
|
|
"SELECT * FROM users WHERE telegram_id = %s FOR UPDATE",
|
|
(web_id,),
|
|
).fetchone()
|
|
if not web:
|
|
raise AuthError("Веб-аккаунт не найден")
|
|
|
|
tg = conn.execute(
|
|
"SELECT * FROM users WHERE telegram_id = %s FOR UPDATE",
|
|
(telegram_id,),
|
|
).fetchone()
|
|
|
|
if tg is None:
|
|
# Переносим PK: создаём tg-строку и перекидываем FK
|
|
conn.execute(
|
|
"""
|
|
INSERT INTO users (
|
|
telegram_id, username, full_name, balance, created_at,
|
|
login, password_hash
|
|
) VALUES (%s, %s, %s, %s, %s, %s, %s)
|
|
""",
|
|
(
|
|
telegram_id,
|
|
username or web["username"],
|
|
full_name or web["full_name"],
|
|
web["balance"],
|
|
web["created_at"],
|
|
web["login"],
|
|
web["password_hash"],
|
|
),
|
|
)
|
|
for table in ("transactions", "purchases", "invoices"):
|
|
conn.execute(
|
|
f"UPDATE {table} SET telegram_id = %s WHERE telegram_id = %s",
|
|
(telegram_id, web_id),
|
|
)
|
|
conn.execute("DELETE FROM users WHERE telegram_id = %s", (web_id,))
|
|
else:
|
|
# Telegram уже есть — сливаем баланс и логин в него
|
|
new_balance = float(tg["balance"]) + float(web["balance"])
|
|
login = tg["login"] or web["login"]
|
|
pwd = tg["password_hash"] or web["password_hash"]
|
|
conn.execute(
|
|
"""
|
|
UPDATE users
|
|
SET balance = %s,
|
|
login = %s,
|
|
password_hash = %s,
|
|
username = COALESCE(%s, username),
|
|
full_name = COALESCE(%s, full_name)
|
|
WHERE telegram_id = %s
|
|
""",
|
|
(
|
|
new_balance,
|
|
login,
|
|
pwd,
|
|
username,
|
|
full_name,
|
|
telegram_id,
|
|
),
|
|
)
|
|
for table in ("transactions", "purchases", "invoices"):
|
|
conn.execute(
|
|
f"UPDATE {table} SET telegram_id = %s WHERE telegram_id = %s",
|
|
(telegram_id, web_id),
|
|
)
|
|
conn.execute("DELETE FROM users WHERE telegram_id = %s", (web_id,))
|
|
|
|
conn.commit()
|
|
row = conn.execute(
|
|
"SELECT * FROM users WHERE telegram_id = %s",
|
|
(telegram_id,),
|
|
).fetchone()
|
|
assert row is not None
|
|
return self._to_user(row)
|
|
|
|
def get_balance(self, telegram_id: int) -> float:
|
|
return self.ensure_user(telegram_id).balance
|
|
|
|
def add_balance(
|
|
self,
|
|
telegram_id: int,
|
|
amount: float,
|
|
*,
|
|
kind: str = "topup",
|
|
meta: str | None = None,
|
|
) -> float:
|
|
if amount == 0:
|
|
return self.get_balance(telegram_id)
|
|
with self.pool.connection() as conn:
|
|
self._ensure_user_conn(conn, telegram_id)
|
|
conn.execute(
|
|
"UPDATE users SET balance = balance + %s WHERE telegram_id = %s",
|
|
(amount, telegram_id),
|
|
)
|
|
conn.execute(
|
|
"""
|
|
INSERT INTO transactions (telegram_id, amount, kind, meta, created_at)
|
|
VALUES (%s, %s, %s, %s, %s)
|
|
""",
|
|
(telegram_id, amount, kind, meta, _utc_now()),
|
|
)
|
|
row = conn.execute(
|
|
"SELECT balance FROM users WHERE telegram_id = %s",
|
|
(telegram_id,),
|
|
).fetchone()
|
|
conn.commit()
|
|
assert row is not None
|
|
return float(row["balance"])
|
|
|
|
def try_charge(
|
|
self,
|
|
telegram_id: int,
|
|
amount: float,
|
|
*,
|
|
kind: str = "purchase",
|
|
meta: str | None = None,
|
|
) -> tuple[bool, float]:
|
|
if amount < 0:
|
|
raise ValueError("amount must be >= 0")
|
|
with self.pool.connection() as conn:
|
|
self._ensure_user_conn(conn, telegram_id)
|
|
row = conn.execute(
|
|
"SELECT balance FROM users WHERE telegram_id = %s FOR UPDATE",
|
|
(telegram_id,),
|
|
).fetchone()
|
|
assert row is not None
|
|
balance = float(row["balance"])
|
|
if balance < amount:
|
|
conn.commit()
|
|
return False, balance
|
|
conn.execute(
|
|
"UPDATE users SET balance = balance - %s WHERE telegram_id = %s",
|
|
(amount, telegram_id),
|
|
)
|
|
conn.execute(
|
|
"""
|
|
INSERT INTO transactions (telegram_id, amount, kind, meta, created_at)
|
|
VALUES (%s, %s, %s, %s, %s)
|
|
""",
|
|
(telegram_id, -amount, kind, meta, _utc_now()),
|
|
)
|
|
conn.commit()
|
|
return True, balance - amount
|
|
|
|
def record_purchase(
|
|
self,
|
|
telegram_id: int,
|
|
*,
|
|
tariff_id: str,
|
|
price: float,
|
|
remnawave_uuid: str | None,
|
|
subscription_url: str | None,
|
|
) -> None:
|
|
with self.pool.connection() as conn:
|
|
conn.execute(
|
|
"""
|
|
INSERT INTO purchases (
|
|
telegram_id, tariff_id, price, remnawave_uuid, subscription_url, created_at
|
|
) VALUES (%s, %s, %s, %s, %s, %s)
|
|
""",
|
|
(
|
|
telegram_id,
|
|
tariff_id,
|
|
price,
|
|
remnawave_uuid,
|
|
subscription_url,
|
|
_utc_now(),
|
|
),
|
|
)
|
|
conn.commit()
|
|
|
|
def recent_transactions(self, telegram_id: int, limit: int = 5) -> list[dict[str, Any]]:
|
|
items, _total = self.list_transactions(telegram_id, limit=limit, offset=0)
|
|
return items
|
|
|
|
def list_transactions(
|
|
self,
|
|
telegram_id: int,
|
|
*,
|
|
limit: int = 10,
|
|
offset: int = 0,
|
|
) -> tuple[list[dict[str, Any]], int]:
|
|
limit = max(1, min(int(limit), 50))
|
|
offset = max(0, int(offset))
|
|
with self.pool.connection() as conn:
|
|
total_row = conn.execute(
|
|
"SELECT COUNT(*) AS c FROM transactions WHERE telegram_id = %s",
|
|
(telegram_id,),
|
|
).fetchone()
|
|
total = int(total_row["c"] if total_row else 0)
|
|
rows = conn.execute(
|
|
"""
|
|
SELECT amount, kind, meta, created_at
|
|
FROM transactions
|
|
WHERE telegram_id = %s
|
|
ORDER BY id DESC
|
|
LIMIT %s OFFSET %s
|
|
""",
|
|
(telegram_id, limit, offset),
|
|
).fetchall()
|
|
return [dict(r) for r in rows], total
|
|
|
|
def list_users_admin(
|
|
self,
|
|
*,
|
|
query: str = "",
|
|
limit: int = 50,
|
|
offset: int = 0,
|
|
) -> tuple[list[dict[str, Any]], int]:
|
|
q = (query or "").strip()
|
|
limit = max(1, min(int(limit), 200))
|
|
offset = max(0, int(offset))
|
|
with self.pool.connection() as conn:
|
|
if q:
|
|
pattern = f"%{q}%"
|
|
tid: int | None = None
|
|
if q.lstrip("-").isdigit():
|
|
try:
|
|
tid = int(q)
|
|
except ValueError:
|
|
tid = None
|
|
if tid is not None:
|
|
rows = conn.execute(
|
|
"""
|
|
SELECT telegram_id, username, full_name, login, balance, created_at, banned
|
|
FROM users
|
|
WHERE telegram_id = %s
|
|
OR login ILIKE %s
|
|
OR username ILIKE %s
|
|
OR full_name ILIKE %s
|
|
ORDER BY created_at DESC
|
|
LIMIT %s OFFSET %s
|
|
""",
|
|
(tid, pattern, pattern, pattern, limit, offset),
|
|
).fetchall()
|
|
count_row = conn.execute(
|
|
"""
|
|
SELECT COUNT(*) AS c FROM users
|
|
WHERE telegram_id = %s
|
|
OR login ILIKE %s
|
|
OR username ILIKE %s
|
|
OR full_name ILIKE %s
|
|
""",
|
|
(tid, pattern, pattern, pattern),
|
|
).fetchone()
|
|
else:
|
|
rows = conn.execute(
|
|
"""
|
|
SELECT telegram_id, username, full_name, login, balance, created_at, banned
|
|
FROM users
|
|
WHERE login ILIKE %s
|
|
OR username ILIKE %s
|
|
OR full_name ILIKE %s
|
|
OR CAST(telegram_id AS TEXT) LIKE %s
|
|
ORDER BY created_at DESC
|
|
LIMIT %s OFFSET %s
|
|
""",
|
|
(pattern, pattern, pattern, pattern, limit, offset),
|
|
).fetchall()
|
|
count_row = conn.execute(
|
|
"""
|
|
SELECT COUNT(*) AS c FROM users
|
|
WHERE login ILIKE %s
|
|
OR username ILIKE %s
|
|
OR full_name ILIKE %s
|
|
OR CAST(telegram_id AS TEXT) LIKE %s
|
|
""",
|
|
(pattern, pattern, pattern, pattern),
|
|
).fetchone()
|
|
else:
|
|
rows = conn.execute(
|
|
"""
|
|
SELECT telegram_id, username, full_name, login, balance, created_at, banned
|
|
FROM users
|
|
ORDER BY created_at DESC
|
|
LIMIT %s OFFSET %s
|
|
""",
|
|
(limit, offset),
|
|
).fetchall()
|
|
count_row = conn.execute("SELECT COUNT(*) AS c FROM users").fetchone()
|
|
total = int(count_row["c"] or 0) if count_row else 0
|
|
return [dict(r) for r in rows], total
|
|
|
|
def admin_credit_balance(
|
|
self,
|
|
telegram_id: int,
|
|
amount: float,
|
|
*,
|
|
comment: str,
|
|
admin_note: str | None = None,
|
|
) -> float:
|
|
"""Пополнение с баланса админом; comment виден пользователю в истории."""
|
|
if amount <= 0:
|
|
raise ValueError("amount must be > 0")
|
|
text = (comment or "").strip()[:200] or "Зачисление администратором"
|
|
meta = text
|
|
if admin_note:
|
|
meta = f"{text} | admin:{admin_note[:80]}"
|
|
return self.add_balance(
|
|
telegram_id,
|
|
amount,
|
|
kind="admin_credit",
|
|
meta=meta,
|
|
)
|
|
|
|
def next_local_invoice_id(self) -> int:
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute("SELECT nextval('local_invoice_id_seq') AS id").fetchone()
|
|
conn.commit()
|
|
assert row is not None
|
|
return int(row["id"])
|
|
|
|
def create_invoice_record(
|
|
self,
|
|
*,
|
|
telegram_id: int,
|
|
invoice_id: int,
|
|
kind: str,
|
|
amount_rub: float,
|
|
pay_url: str | None,
|
|
payload: str | None,
|
|
tariff_id: str | None = None,
|
|
provider: str = "cryptobot",
|
|
order_id: str | None = None,
|
|
external_uuid: str | None = None,
|
|
) -> None:
|
|
self.ensure_user(telegram_id)
|
|
with self.pool.connection() as conn:
|
|
conn.execute(
|
|
"""
|
|
INSERT INTO invoices (
|
|
telegram_id, invoice_id, kind, amount_rub, tariff_id,
|
|
status, pay_url, payload, created_at,
|
|
provider, order_id, external_uuid
|
|
) VALUES (%s, %s, %s, %s, %s, 'active', %s, %s, %s, %s, %s, %s)
|
|
""",
|
|
(
|
|
telegram_id,
|
|
invoice_id,
|
|
kind,
|
|
amount_rub,
|
|
tariff_id,
|
|
pay_url,
|
|
payload,
|
|
_utc_now(),
|
|
provider,
|
|
order_id,
|
|
external_uuid,
|
|
),
|
|
)
|
|
conn.commit()
|
|
|
|
def get_invoice(
|
|
self,
|
|
invoice_id: int,
|
|
*,
|
|
provider: str | None = None,
|
|
) -> dict[str, Any] | None:
|
|
with self.pool.connection() as conn:
|
|
if provider:
|
|
row = conn.execute(
|
|
"""
|
|
SELECT * FROM invoices
|
|
WHERE invoice_id = %s AND provider = %s
|
|
""",
|
|
(invoice_id, provider),
|
|
).fetchone()
|
|
else:
|
|
row = conn.execute(
|
|
"SELECT * FROM invoices WHERE invoice_id = %s",
|
|
(invoice_id,),
|
|
).fetchone()
|
|
return dict(row) if row else None
|
|
|
|
def get_invoice_by_order_id(self, order_id: str) -> dict[str, Any] | None:
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"SELECT * FROM invoices WHERE order_id = %s",
|
|
(order_id,),
|
|
).fetchone()
|
|
return dict(row) if row else None
|
|
|
|
def digiseller_inv_used(self, inv: int) -> bool:
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"SELECT 1 FROM digiseller_payments WHERE inv = %s",
|
|
(inv,),
|
|
).fetchone()
|
|
return row is not None
|
|
|
|
def digiseller_code_used(self, unique_code: str) -> bool:
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"SELECT 1 FROM digiseller_payments WHERE unique_code = %s",
|
|
(unique_code,),
|
|
).fetchone()
|
|
return row is not None
|
|
|
|
def record_digiseller_payment(
|
|
self,
|
|
*,
|
|
inv: int,
|
|
unique_code: str,
|
|
telegram_id: int,
|
|
amount_rub: float,
|
|
product_id: int | None,
|
|
) -> bool:
|
|
"""True если записали впервые, False если уже был такой inv/code."""
|
|
self.ensure_user(telegram_id)
|
|
with self.pool.connection() as conn:
|
|
try:
|
|
conn.execute(
|
|
"""
|
|
INSERT INTO digiseller_payments (
|
|
inv, unique_code, telegram_id, amount_rub, product_id, created_at
|
|
) VALUES (%s, %s, %s, %s, %s, %s)
|
|
""",
|
|
(
|
|
inv,
|
|
unique_code,
|
|
telegram_id,
|
|
amount_rub,
|
|
product_id,
|
|
_utc_now(),
|
|
),
|
|
)
|
|
conn.commit()
|
|
return True
|
|
except Exception: # noqa: BLE001
|
|
conn.rollback()
|
|
return False
|
|
|
|
def create_topup_codes(
|
|
self,
|
|
*,
|
|
amount_rub: float,
|
|
count: int,
|
|
created_by: str = "",
|
|
codes: list[str],
|
|
) -> int:
|
|
if amount_rub <= 0 or not codes:
|
|
return 0
|
|
inserted = 0
|
|
with self.pool.connection() as conn:
|
|
for code in codes:
|
|
try:
|
|
conn.execute(
|
|
"""
|
|
INSERT INTO topup_codes (code, amount_rub, created_by, created_at)
|
|
VALUES (%s, %s, %s, %s)
|
|
""",
|
|
(code, amount_rub, created_by or None, _utc_now()),
|
|
)
|
|
inserted += 1
|
|
except Exception: # noqa: BLE001
|
|
continue
|
|
conn.commit()
|
|
return inserted
|
|
|
|
def redeem_topup_code(
|
|
self,
|
|
*,
|
|
code: str,
|
|
telegram_id: int,
|
|
) -> tuple[float, bool]:
|
|
"""(сумма, уже_использован). NotFound -> ValueError."""
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"""
|
|
SELECT code, amount_rub, used_at
|
|
FROM topup_codes
|
|
WHERE code = %s
|
|
FOR UPDATE
|
|
""",
|
|
(code,),
|
|
).fetchone()
|
|
if not row:
|
|
raise ValueError("Код не найден")
|
|
amount = float(row["amount_rub"])
|
|
if row["used_at"] is not None:
|
|
conn.commit()
|
|
return amount, True
|
|
conn.execute(
|
|
"""
|
|
UPDATE topup_codes
|
|
SET used_at = %s, used_by = %s
|
|
WHERE code = %s AND used_at IS NULL
|
|
""",
|
|
(_utc_now(), telegram_id, code),
|
|
)
|
|
conn.commit()
|
|
return amount, False
|
|
|
|
def list_topup_codes(
|
|
self,
|
|
*,
|
|
status: str = "all",
|
|
limit: int = 50,
|
|
offset: int = 0,
|
|
) -> list[dict[str, Any]]:
|
|
limit = max(1, min(int(limit), 200))
|
|
offset = max(0, int(offset))
|
|
where = ""
|
|
if status == "unused":
|
|
where = "WHERE used_at IS NULL"
|
|
elif status == "used":
|
|
where = "WHERE used_at IS NOT NULL"
|
|
with self.pool.connection() as conn:
|
|
rows = conn.execute(
|
|
f"""
|
|
SELECT code, amount_rub, created_at, created_by, used_at, used_by
|
|
FROM topup_codes
|
|
{where}
|
|
ORDER BY created_at DESC
|
|
LIMIT %s OFFSET %s
|
|
""",
|
|
(limit, offset),
|
|
).fetchall()
|
|
return [dict(r) for r in rows]
|
|
|
|
def topup_code_stats(self) -> dict[str, int]:
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"""
|
|
SELECT
|
|
COUNT(*) FILTER (WHERE used_at IS NULL) AS unused,
|
|
COUNT(*) FILTER (WHERE used_at IS NOT NULL) AS used,
|
|
COUNT(*) AS total
|
|
FROM topup_codes
|
|
"""
|
|
).fetchone()
|
|
return {
|
|
"unused": int(row["unused"] or 0),
|
|
"used": int(row["used"] or 0),
|
|
"total": int(row["total"] or 0),
|
|
}
|
|
|
|
def mark_invoice_paid(self, invoice_id: int, *, provider: str | None = None) -> bool:
|
|
with self.pool.connection() as conn:
|
|
if provider:
|
|
row = conn.execute(
|
|
"""
|
|
UPDATE invoices
|
|
SET status = 'paid', paid_at = %s
|
|
WHERE invoice_id = %s AND provider = %s
|
|
AND status IS DISTINCT FROM 'paid'
|
|
RETURNING invoice_id
|
|
""",
|
|
(_utc_now(), invoice_id, provider),
|
|
).fetchone()
|
|
else:
|
|
row = conn.execute(
|
|
"""
|
|
UPDATE invoices
|
|
SET status = 'paid', paid_at = %s
|
|
WHERE invoice_id = %s AND status IS DISTINCT FROM 'paid'
|
|
RETURNING invoice_id
|
|
""",
|
|
(_utc_now(), invoice_id),
|
|
).fetchone()
|
|
conn.commit()
|
|
return row is not None
|
|
|
|
@staticmethod
|
|
def _to_user(row: dict[str, Any]) -> BotUser:
|
|
tid = int(row["telegram_id"])
|
|
return BotUser(
|
|
telegram_id=tid,
|
|
username=row["username"],
|
|
full_name=row["full_name"],
|
|
balance=float(row["balance"]),
|
|
created_at=row["created_at"],
|
|
login=row.get("login"),
|
|
has_password=bool(row.get("password_hash")),
|
|
is_telegram=tid > 0,
|
|
extra_devices=int(row.get("extra_devices") or 0),
|
|
device_blocked=bool(row.get("device_blocked")),
|
|
extra_traffic_packs=int(row.get("extra_traffic_packs") or 0),
|
|
extras_next_billing_at=row.get("extras_next_billing_at"),
|
|
extras_active=bool(row.get("extras_active", True)),
|
|
auto_renew=bool(row["auto_renew"]) if row.get("auto_renew") is not None else True,
|
|
email=(str(row["email"]).strip() if row.get("email") else None) or None,
|
|
banned=bool(row.get("banned")),
|
|
referral_code=(str(row["referral_code"]).strip() if row.get("referral_code") else None)
|
|
or None,
|
|
referred_by=int(row["referred_by"]) if row.get("referred_by") is not None else None,
|
|
trial_used_at=row.get("trial_used_at"),
|
|
notify_telegram=bool(row["notify_telegram"])
|
|
if row.get("notify_telegram") is not None
|
|
else True,
|
|
notify_email=bool(row["notify_email"])
|
|
if row.get("notify_email") is not None
|
|
else True,
|
|
)
|
|
|
|
def set_device_blocked(self, telegram_id: int, blocked: bool) -> None:
|
|
with self.pool.connection() as conn:
|
|
self._ensure_user_conn(conn, telegram_id)
|
|
conn.execute(
|
|
"UPDATE users SET device_blocked = %s WHERE telegram_id = %s",
|
|
(bool(blocked), telegram_id),
|
|
)
|
|
conn.commit()
|
|
|
|
def add_extra_devices(self, telegram_id: int, slots: int) -> int:
|
|
if slots <= 0:
|
|
return self.ensure_user(telegram_id).extra_devices
|
|
with self.pool.connection() as conn:
|
|
self._ensure_user_conn(conn, telegram_id)
|
|
row = conn.execute(
|
|
"""
|
|
UPDATE users
|
|
SET extra_devices = extra_devices + %s
|
|
WHERE telegram_id = %s
|
|
RETURNING extra_devices
|
|
""",
|
|
(int(slots), telegram_id),
|
|
).fetchone()
|
|
conn.commit()
|
|
assert row is not None
|
|
return int(row["extra_devices"])
|
|
|
|
def add_extra_traffic_packs(self, telegram_id: int, packs: int) -> int:
|
|
if packs <= 0:
|
|
return self.ensure_user(telegram_id).extra_traffic_packs
|
|
with self.pool.connection() as conn:
|
|
self._ensure_user_conn(conn, telegram_id)
|
|
row = conn.execute(
|
|
"""
|
|
UPDATE users
|
|
SET extra_traffic_packs = extra_traffic_packs + %s
|
|
WHERE telegram_id = %s
|
|
RETURNING extra_traffic_packs
|
|
""",
|
|
(int(packs), telegram_id),
|
|
).fetchone()
|
|
conn.commit()
|
|
assert row is not None
|
|
return int(row["extra_traffic_packs"])
|
|
|
|
def set_extras_active(self, telegram_id: int, active: bool) -> None:
|
|
with self.pool.connection() as conn:
|
|
self._ensure_user_conn(conn, telegram_id)
|
|
conn.execute(
|
|
"UPDATE users SET extras_active = %s WHERE telegram_id = %s",
|
|
(bool(active), telegram_id),
|
|
)
|
|
conn.commit()
|
|
|
|
def set_extras_billing(
|
|
self,
|
|
telegram_id: int,
|
|
*,
|
|
next_billing_at: datetime | None,
|
|
active: bool,
|
|
) -> None:
|
|
with self.pool.connection() as conn:
|
|
self._ensure_user_conn(conn, telegram_id)
|
|
conn.execute(
|
|
"""
|
|
UPDATE users
|
|
SET extras_next_billing_at = %s,
|
|
extras_active = %s
|
|
WHERE telegram_id = %s
|
|
""",
|
|
(next_billing_at, bool(active), telegram_id),
|
|
)
|
|
conn.commit()
|
|
|
|
def list_users_for_extras_billing(self, *, limit: int = 200) -> list[int]:
|
|
with self.pool.connection() as conn:
|
|
rows = conn.execute(
|
|
"""
|
|
SELECT telegram_id FROM users
|
|
WHERE (extra_devices > 0 OR extra_traffic_packs > 0)
|
|
AND (
|
|
extras_active = FALSE
|
|
OR extras_next_billing_at IS NULL
|
|
OR extras_next_billing_at <= NOW()
|
|
)
|
|
ORDER BY extras_next_billing_at NULLS FIRST
|
|
LIMIT %s
|
|
""",
|
|
(int(limit),),
|
|
).fetchall()
|
|
return [int(r["telegram_id"]) for r in rows]
|
|
|
|
def list_telegram_ids_with_purchases(self, *, limit: int = 500) -> list[int]:
|
|
with self.pool.connection() as conn:
|
|
rows = conn.execute(
|
|
"""
|
|
SELECT DISTINCT p.telegram_id
|
|
FROM purchases p
|
|
WHERE p.telegram_id > 0
|
|
ORDER BY p.telegram_id
|
|
LIMIT %s
|
|
""",
|
|
(int(limit),),
|
|
).fetchall()
|
|
return [int(r["telegram_id"]) for r in rows]
|
|
|
|
def set_auto_renew(self, telegram_id: int, enabled: bool) -> BotUser:
|
|
with self.pool.connection() as conn:
|
|
self._ensure_user_conn(conn, telegram_id)
|
|
row = conn.execute(
|
|
"""
|
|
UPDATE users
|
|
SET auto_renew = %s
|
|
WHERE telegram_id = %s
|
|
RETURNING *
|
|
""",
|
|
(bool(enabled), int(telegram_id)),
|
|
).fetchone()
|
|
conn.commit()
|
|
assert row is not None
|
|
return self._to_user(row)
|
|
|
|
def get_last_tariff_id(self, telegram_id: int) -> str | None:
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"""
|
|
SELECT tariff_id FROM purchases
|
|
WHERE telegram_id = %s
|
|
ORDER BY id DESC
|
|
LIMIT 1
|
|
""",
|
|
(telegram_id,),
|
|
).fetchone()
|
|
return str(row["tariff_id"]) if row else None
|
|
|
|
def expiry_reminder_sent_on(self, telegram_id: int) -> date | None:
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"SELECT expiry_reminder_sent_on FROM users WHERE telegram_id = %s",
|
|
(telegram_id,),
|
|
).fetchone()
|
|
if not row:
|
|
return None
|
|
val = row.get("expiry_reminder_sent_on")
|
|
if val is None:
|
|
return None
|
|
if isinstance(val, date):
|
|
return val
|
|
return val.date() if hasattr(val, "date") else None
|
|
|
|
def mark_expiry_reminder_sent(self, telegram_id: int, on: date) -> None:
|
|
with self.pool.connection() as conn:
|
|
self._ensure_user_conn(conn, telegram_id)
|
|
conn.execute(
|
|
"""
|
|
UPDATE users SET expiry_reminder_sent_on = %s
|
|
WHERE telegram_id = %s
|
|
""",
|
|
(on, telegram_id),
|
|
)
|
|
conn.commit()
|
|
|
|
def count_open_tickets(self, telegram_id: int) -> int:
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"""
|
|
SELECT COUNT(*) AS c FROM tickets
|
|
WHERE telegram_id = %s AND status <> 'closed'
|
|
""",
|
|
(telegram_id,),
|
|
).fetchone()
|
|
return int(row["c"]) if row else 0
|
|
|
|
def create_ticket(
|
|
self,
|
|
telegram_id: int,
|
|
*,
|
|
subject: str,
|
|
body: str,
|
|
) -> dict[str, Any]:
|
|
subject = (subject or "").strip()
|
|
body = (body or "").strip()
|
|
if not subject or not body:
|
|
raise ValueError("Тема и сообщение обязательны")
|
|
with self.pool.connection() as conn:
|
|
self._ensure_user_conn(conn, telegram_id)
|
|
row = conn.execute(
|
|
"""
|
|
INSERT INTO tickets (telegram_id, subject, status, created_at, updated_at)
|
|
VALUES (%s, %s, 'open', NOW(), NOW())
|
|
RETURNING *
|
|
""",
|
|
(telegram_id, subject),
|
|
).fetchone()
|
|
assert row is not None
|
|
ticket_id = int(row["id"])
|
|
conn.execute(
|
|
"""
|
|
INSERT INTO ticket_messages (ticket_id, author_role, author_id, body)
|
|
VALUES (%s, 'user', %s, %s)
|
|
""",
|
|
(ticket_id, telegram_id, body),
|
|
)
|
|
conn.commit()
|
|
return dict(row)
|
|
|
|
def list_tickets_for_user(
|
|
self,
|
|
telegram_id: int,
|
|
*,
|
|
limit: int = 30,
|
|
offset: int = 0,
|
|
include_closed: bool = True,
|
|
) -> list[dict[str, Any]]:
|
|
limit = max(1, min(int(limit), 100))
|
|
offset = max(0, int(offset))
|
|
with self.pool.connection() as conn:
|
|
if include_closed:
|
|
rows = conn.execute(
|
|
"""
|
|
SELECT t.*,
|
|
(
|
|
SELECT COUNT(*) FROM ticket_messages m
|
|
WHERE m.ticket_id = t.id
|
|
) AS message_count
|
|
FROM tickets t
|
|
WHERE t.telegram_id = %s
|
|
ORDER BY t.updated_at DESC
|
|
LIMIT %s OFFSET %s
|
|
""",
|
|
(telegram_id, limit, offset),
|
|
).fetchall()
|
|
else:
|
|
rows = conn.execute(
|
|
"""
|
|
SELECT t.*,
|
|
(
|
|
SELECT COUNT(*) FROM ticket_messages m
|
|
WHERE m.ticket_id = t.id
|
|
) AS message_count
|
|
FROM tickets t
|
|
WHERE t.telegram_id = %s AND t.status <> 'closed'
|
|
ORDER BY t.updated_at DESC
|
|
LIMIT %s OFFSET %s
|
|
""",
|
|
(telegram_id, limit, offset),
|
|
).fetchall()
|
|
return [dict(r) for r in rows]
|
|
|
|
def list_tickets_admin(
|
|
self,
|
|
*,
|
|
status: str = "open",
|
|
limit: int = 50,
|
|
offset: int = 0,
|
|
) -> tuple[list[dict[str, Any]], int]:
|
|
limit = max(1, min(int(limit), 100))
|
|
offset = max(0, int(offset))
|
|
status = (status or "open").strip().lower()
|
|
with self.pool.connection() as conn:
|
|
if status == "all":
|
|
total_row = conn.execute("SELECT COUNT(*) AS c FROM tickets").fetchone()
|
|
rows = conn.execute(
|
|
"""
|
|
SELECT t.*,
|
|
u.username, u.full_name, u.login,
|
|
(
|
|
SELECT COUNT(*) FROM ticket_messages m
|
|
WHERE m.ticket_id = t.id
|
|
) AS message_count
|
|
FROM tickets t
|
|
LEFT JOIN users u ON u.telegram_id = t.telegram_id
|
|
ORDER BY
|
|
CASE t.status
|
|
WHEN 'open' THEN 0
|
|
WHEN 'answered' THEN 1
|
|
ELSE 2
|
|
END,
|
|
t.updated_at DESC
|
|
LIMIT %s OFFSET %s
|
|
""",
|
|
(limit, offset),
|
|
).fetchall()
|
|
else:
|
|
total_row = conn.execute(
|
|
"SELECT COUNT(*) AS c FROM tickets WHERE status = %s",
|
|
(status,),
|
|
).fetchone()
|
|
rows = conn.execute(
|
|
"""
|
|
SELECT t.*,
|
|
u.username, u.full_name, u.login,
|
|
(
|
|
SELECT COUNT(*) FROM ticket_messages m
|
|
WHERE m.ticket_id = t.id
|
|
) AS message_count
|
|
FROM tickets t
|
|
LEFT JOIN users u ON u.telegram_id = t.telegram_id
|
|
WHERE t.status = %s
|
|
ORDER BY t.updated_at DESC
|
|
LIMIT %s OFFSET %s
|
|
""",
|
|
(status, limit, offset),
|
|
).fetchall()
|
|
total = int(total_row["c"]) if total_row else 0
|
|
return [dict(r) for r in rows], total
|
|
|
|
def ticket_stats(self) -> dict[str, int]:
|
|
with self.pool.connection() as conn:
|
|
rows = conn.execute(
|
|
"""
|
|
SELECT status, COUNT(*) AS c
|
|
FROM tickets
|
|
GROUP BY status
|
|
"""
|
|
).fetchall()
|
|
stats = {"open": 0, "answered": 0, "closed": 0, "total": 0}
|
|
for r in rows:
|
|
st = str(r["status"])
|
|
c = int(r["c"])
|
|
if st in stats:
|
|
stats[st] = c
|
|
stats["total"] += c
|
|
return stats
|
|
|
|
def get_ticket(self, ticket_id: int) -> dict[str, Any] | None:
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"""
|
|
SELECT t.*, u.username, u.full_name, u.login
|
|
FROM tickets t
|
|
LEFT JOIN users u ON u.telegram_id = t.telegram_id
|
|
WHERE t.id = %s
|
|
""",
|
|
(int(ticket_id),),
|
|
).fetchone()
|
|
return dict(row) if row else None
|
|
|
|
def list_ticket_messages(
|
|
self,
|
|
ticket_id: int,
|
|
*,
|
|
limit: int = 200,
|
|
) -> list[dict[str, Any]]:
|
|
limit = max(1, min(int(limit), 500))
|
|
with self.pool.connection() as conn:
|
|
rows = conn.execute(
|
|
"""
|
|
SELECT * FROM ticket_messages
|
|
WHERE ticket_id = %s
|
|
ORDER BY id ASC
|
|
LIMIT %s
|
|
""",
|
|
(int(ticket_id), limit),
|
|
).fetchall()
|
|
return [dict(r) for r in rows]
|
|
|
|
def add_ticket_message(
|
|
self,
|
|
ticket_id: int,
|
|
*,
|
|
author_role: str,
|
|
author_id: int | None,
|
|
body: str,
|
|
new_status: str | None = None,
|
|
) -> dict[str, Any]:
|
|
body = (body or "").strip()
|
|
if not body:
|
|
raise ValueError("Пустое сообщение")
|
|
role = (author_role or "").strip().lower()
|
|
if role not in ("user", "admin", "ai"):
|
|
raise ValueError("Некорректная роль")
|
|
with self.pool.connection() as conn:
|
|
ticket = conn.execute(
|
|
"SELECT * FROM tickets WHERE id = %s FOR UPDATE",
|
|
(int(ticket_id),),
|
|
).fetchone()
|
|
if not ticket:
|
|
raise ValueError("Тикет не найден")
|
|
if str(ticket["status"]) == "closed":
|
|
raise ValueError("Тикет закрыт")
|
|
status = new_status or str(ticket["status"])
|
|
if new_status is None:
|
|
status = "answered" if role in ("admin", "ai") else "open"
|
|
conn.execute(
|
|
"""
|
|
INSERT INTO ticket_messages (ticket_id, author_role, author_id, body)
|
|
VALUES (%s, %s, %s, %s)
|
|
""",
|
|
(int(ticket_id), role, author_id, body),
|
|
)
|
|
conn.execute(
|
|
"""
|
|
UPDATE tickets
|
|
SET status = %s, updated_at = NOW()
|
|
WHERE id = %s
|
|
""",
|
|
(status, int(ticket_id)),
|
|
)
|
|
conn.commit()
|
|
row = conn.execute(
|
|
"SELECT * FROM tickets WHERE id = %s",
|
|
(int(ticket_id),),
|
|
).fetchone()
|
|
assert row is not None
|
|
return dict(row)
|
|
|
|
def set_ticket_status(self, ticket_id: int, status: str) -> dict[str, Any] | None:
|
|
status = (status or "").strip().lower()
|
|
if status not in ("open", "answered", "closed"):
|
|
raise ValueError("Некорректный статус")
|
|
with self.pool.connection() as conn:
|
|
conn.execute(
|
|
"""
|
|
UPDATE tickets
|
|
SET status = %s, updated_at = NOW()
|
|
WHERE id = %s
|
|
""",
|
|
(status, int(ticket_id)),
|
|
)
|
|
conn.commit()
|
|
row = conn.execute(
|
|
"SELECT * FROM tickets WHERE id = %s",
|
|
(int(ticket_id),),
|
|
).fetchone()
|
|
return dict(row) if row else None
|
|
|
|
def get_setting(self, key: str, default: str = "") -> str:
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"SELECT value FROM app_settings WHERE key = %s",
|
|
(key,),
|
|
).fetchone()
|
|
if not row:
|
|
return default
|
|
return str(row["value"] if row["value"] is not None else default)
|
|
|
|
def set_setting(self, key: str, value: str) -> None:
|
|
with self.pool.connection() as conn:
|
|
conn.execute(
|
|
"""
|
|
INSERT INTO app_settings (key, value, updated_at)
|
|
VALUES (%s, %s, NOW())
|
|
ON CONFLICT (key) DO UPDATE
|
|
SET value = EXCLUDED.value, updated_at = NOW()
|
|
""",
|
|
(key, value if value is not None else ""),
|
|
)
|
|
conn.commit()
|
|
|
|
def get_settings_map(self, keys: list[str]) -> dict[str, str]:
|
|
if not keys:
|
|
return {}
|
|
with self.pool.connection() as conn:
|
|
rows = conn.execute(
|
|
"SELECT key, value FROM app_settings WHERE key = ANY(%s)",
|
|
(list(keys),),
|
|
).fetchall()
|
|
out = {k: "" for k in keys}
|
|
for r in rows:
|
|
out[str(r["key"])] = str(r["value"] if r["value"] is not None else "")
|
|
return out
|
|
|
|
def set_settings_map(self, data: dict[str, str]) -> None:
|
|
if not data:
|
|
return
|
|
with self.pool.connection() as conn:
|
|
for key, value in data.items():
|
|
conn.execute(
|
|
"""
|
|
INSERT INTO app_settings (key, value, updated_at)
|
|
VALUES (%s, %s, NOW())
|
|
ON CONFLICT (key) DO UPDATE
|
|
SET value = EXCLUDED.value, updated_at = NOW()
|
|
""",
|
|
(str(key), "" if value is None else str(value)),
|
|
)
|
|
conn.commit()
|
|
|
|
# ─── 3x-ui servers ─────────────────────────────────────────────────────
|
|
|
|
def list_xui_servers(self, *, enabled_only: bool = False) -> list[dict[str, Any]]:
|
|
where = "WHERE enabled = TRUE" if enabled_only else ""
|
|
with self.pool.connection() as conn:
|
|
rows = conn.execute(
|
|
f"""
|
|
SELECT id, name, api_url, username, password, api_key, inbound_id, inbound_remark,
|
|
subscription_url, enabled, verify_ssl,
|
|
last_check_at, last_check_ok, last_check_msg,
|
|
created_at, updated_at
|
|
FROM xui_servers
|
|
{where}
|
|
ORDER BY id ASC
|
|
"""
|
|
).fetchall()
|
|
return [dict(r) for r in rows]
|
|
|
|
def get_xui_server(self, server_id: int) -> dict[str, Any] | None:
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"""
|
|
SELECT id, name, api_url, username, password, api_key, inbound_id, inbound_remark,
|
|
subscription_url, enabled, verify_ssl,
|
|
last_check_at, last_check_ok, last_check_msg,
|
|
created_at, updated_at
|
|
FROM xui_servers
|
|
WHERE id = %s
|
|
""",
|
|
(int(server_id),),
|
|
).fetchone()
|
|
return dict(row) if row else None
|
|
|
|
def create_xui_server(
|
|
self,
|
|
*,
|
|
name: str,
|
|
api_url: str,
|
|
username: str = "",
|
|
password: str = "",
|
|
api_key: str = "",
|
|
subscription_url: str = "",
|
|
inbound_id: int | None = None,
|
|
inbound_remark: str | None = None,
|
|
enabled: bool = True,
|
|
verify_ssl: bool = True,
|
|
) -> dict[str, Any]:
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"""
|
|
INSERT INTO xui_servers (
|
|
name, api_url, username, password, api_key, subscription_url,
|
|
inbound_id, inbound_remark, enabled, verify_ssl
|
|
)
|
|
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
|
|
RETURNING id, name, api_url, username, password, api_key, inbound_id, inbound_remark,
|
|
subscription_url, enabled, verify_ssl,
|
|
last_check_at, last_check_ok, last_check_msg,
|
|
created_at, updated_at
|
|
""",
|
|
(
|
|
name.strip(),
|
|
api_url.strip(),
|
|
(username or "").strip(),
|
|
password or "",
|
|
(api_key or "").strip(),
|
|
(subscription_url or "").strip(),
|
|
inbound_id,
|
|
(inbound_remark or "").strip() or None,
|
|
bool(enabled),
|
|
bool(verify_ssl),
|
|
),
|
|
).fetchone()
|
|
conn.commit()
|
|
return dict(row)
|
|
|
|
def update_xui_server(self, server_id: int, **fields: Any) -> dict[str, Any] | None:
|
|
allowed = {
|
|
"name",
|
|
"api_url",
|
|
"username",
|
|
"password",
|
|
"api_key",
|
|
"subscription_url",
|
|
"inbound_id",
|
|
"inbound_remark",
|
|
"enabled",
|
|
"verify_ssl",
|
|
"last_check_at",
|
|
"last_check_ok",
|
|
"last_check_msg",
|
|
}
|
|
updates = {k: v for k, v in fields.items() if k in allowed}
|
|
if not updates:
|
|
return self.get_xui_server(server_id)
|
|
|
|
cols = []
|
|
vals: list[Any] = []
|
|
for key, value in updates.items():
|
|
cols.append(f"{key} = %s")
|
|
if key in {"name", "api_url", "username", "subscription_url", "api_key"} and isinstance(
|
|
value, str
|
|
):
|
|
vals.append(value.strip())
|
|
elif key == "inbound_remark":
|
|
vals.append((value or "").strip() or None if isinstance(value, str) else value)
|
|
else:
|
|
vals.append(value)
|
|
cols.append("updated_at = NOW()")
|
|
vals.append(int(server_id))
|
|
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
f"""
|
|
UPDATE xui_servers
|
|
SET {", ".join(cols)}
|
|
WHERE id = %s
|
|
RETURNING id, name, api_url, username, password, api_key, inbound_id, inbound_remark,
|
|
subscription_url, enabled, verify_ssl,
|
|
last_check_at, last_check_ok, last_check_msg,
|
|
created_at, updated_at
|
|
""",
|
|
tuple(vals),
|
|
).fetchone()
|
|
conn.commit()
|
|
return dict(row) if row else None
|
|
|
|
def delete_xui_server(self, server_id: int) -> bool:
|
|
with self.pool.connection() as conn:
|
|
cur = conn.execute(
|
|
"DELETE FROM xui_servers WHERE id = %s",
|
|
(int(server_id),),
|
|
)
|
|
conn.commit()
|
|
return cur.rowcount > 0
|
|
|
|
def set_xui_server_check(
|
|
self,
|
|
server_id: int,
|
|
*,
|
|
ok: bool,
|
|
message: str,
|
|
) -> None:
|
|
with self.pool.connection() as conn:
|
|
conn.execute(
|
|
"""
|
|
UPDATE xui_servers
|
|
SET last_check_at = NOW(),
|
|
last_check_ok = %s,
|
|
last_check_msg = %s,
|
|
updated_at = NOW()
|
|
WHERE id = %s
|
|
""",
|
|
(bool(ok), (message or "")[:500], int(server_id)),
|
|
)
|
|
conn.commit()
|
|
|
|
# ─── 3x-ui WireGuard configs ───────────────────────────────────────────
|
|
|
|
def count_xui_wg_configs(self, telegram_id: int) -> int:
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"SELECT COUNT(*) AS c FROM xui_wg_configs WHERE telegram_id = %s",
|
|
(int(telegram_id),),
|
|
).fetchone()
|
|
return int(row["c"] if row else 0)
|
|
|
|
def list_xui_wg_configs(self, telegram_id: int) -> list[dict[str, Any]]:
|
|
with self.pool.connection() as conn:
|
|
rows = conn.execute(
|
|
"""
|
|
SELECT c.id, c.telegram_id, c.server_id, c.email, c.sub_id,
|
|
c.subscription_url, c.remark, c.enabled, c.expire_at,
|
|
c.traffic_bytes, c.created_at,
|
|
s.name AS server_name
|
|
FROM xui_wg_configs c
|
|
LEFT JOIN xui_servers s ON s.id = c.server_id
|
|
WHERE c.telegram_id = %s
|
|
ORDER BY c.id DESC
|
|
""",
|
|
(int(telegram_id),),
|
|
).fetchall()
|
|
return [dict(r) for r in rows]
|
|
|
|
def get_xui_wg_config(
|
|
self, config_id: int, *, telegram_id: int | None = None
|
|
) -> dict[str, Any] | None:
|
|
with self.pool.connection() as conn:
|
|
if telegram_id is None:
|
|
row = conn.execute(
|
|
"""
|
|
SELECT id, telegram_id, server_id, email, sub_id,
|
|
subscription_url, remark, enabled, expire_at,
|
|
traffic_bytes, created_at
|
|
FROM xui_wg_configs
|
|
WHERE id = %s
|
|
""",
|
|
(int(config_id),),
|
|
).fetchone()
|
|
else:
|
|
row = conn.execute(
|
|
"""
|
|
SELECT id, telegram_id, server_id, email, sub_id,
|
|
subscription_url, remark, enabled, expire_at,
|
|
traffic_bytes, created_at
|
|
FROM xui_wg_configs
|
|
WHERE id = %s AND telegram_id = %s
|
|
""",
|
|
(int(config_id), int(telegram_id)),
|
|
).fetchone()
|
|
return dict(row) if row else None
|
|
|
|
def create_xui_wg_config(
|
|
self,
|
|
*,
|
|
telegram_id: int,
|
|
server_id: int,
|
|
email: str,
|
|
sub_id: str,
|
|
subscription_url: str,
|
|
remark: str = "",
|
|
enabled: bool = True,
|
|
expire_at: datetime | None = None,
|
|
traffic_bytes: int = 0,
|
|
) -> dict[str, Any]:
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"""
|
|
INSERT INTO xui_wg_configs (
|
|
telegram_id, server_id, email, sub_id, subscription_url, remark,
|
|
enabled, expire_at, traffic_bytes
|
|
)
|
|
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s)
|
|
RETURNING id, telegram_id, server_id, email, sub_id,
|
|
subscription_url, remark, enabled, expire_at,
|
|
traffic_bytes, created_at
|
|
""",
|
|
(
|
|
int(telegram_id),
|
|
int(server_id),
|
|
email,
|
|
sub_id,
|
|
subscription_url,
|
|
(remark or "").strip() or None,
|
|
bool(enabled),
|
|
expire_at,
|
|
int(traffic_bytes or 0),
|
|
),
|
|
).fetchone()
|
|
conn.commit()
|
|
return dict(row)
|
|
|
|
def set_xui_wg_config_enabled(
|
|
self, config_id: int, telegram_id: int, enabled: bool
|
|
) -> dict[str, Any] | None:
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"""
|
|
UPDATE xui_wg_configs
|
|
SET enabled = %s
|
|
WHERE id = %s AND telegram_id = %s
|
|
RETURNING id, telegram_id, server_id, email, sub_id,
|
|
subscription_url, remark, enabled, expire_at,
|
|
traffic_bytes, created_at
|
|
""",
|
|
(bool(enabled), int(config_id), int(telegram_id)),
|
|
).fetchone()
|
|
conn.commit()
|
|
return dict(row) if row else None
|
|
|
|
def update_xui_wg_config_limits(
|
|
self,
|
|
config_id: int,
|
|
*,
|
|
expire_at: datetime | None = None,
|
|
traffic_bytes: int | None = None,
|
|
) -> None:
|
|
sets: list[str] = []
|
|
vals: list[Any] = []
|
|
if expire_at is not None:
|
|
sets.append("expire_at = %s")
|
|
vals.append(expire_at)
|
|
if traffic_bytes is not None:
|
|
sets.append("traffic_bytes = %s")
|
|
vals.append(int(traffic_bytes))
|
|
if not sets:
|
|
return
|
|
vals.append(int(config_id))
|
|
with self.pool.connection() as conn:
|
|
conn.execute(
|
|
f"UPDATE xui_wg_configs SET {', '.join(sets)} WHERE id = %s",
|
|
tuple(vals),
|
|
)
|
|
conn.commit()
|
|
|
|
def delete_xui_wg_config(self, config_id: int, telegram_id: int) -> bool:
|
|
with self.pool.connection() as conn:
|
|
cur = conn.execute(
|
|
"DELETE FROM xui_wg_configs WHERE id = %s AND telegram_id = %s",
|
|
(int(config_id), int(telegram_id)),
|
|
)
|
|
conn.commit()
|
|
return cur.rowcount > 0
|
|
|
|
# ─── AmneziaWG (AWG 2.0) via Amnezia Web Panel ─────────────────────────
|
|
|
|
def list_awg_servers(self, *, enabled_only: bool = False) -> list[dict[str, Any]]:
|
|
where = "WHERE enabled = TRUE" if enabled_only else ""
|
|
with self.pool.connection() as conn:
|
|
rows = conn.execute(
|
|
f"""
|
|
SELECT id, name, api_url, api_token, panel_server_id, protocol,
|
|
enabled, verify_ssl,
|
|
last_check_at, last_check_ok, last_check_msg,
|
|
created_at, updated_at
|
|
FROM awg_servers
|
|
{where}
|
|
ORDER BY id ASC
|
|
"""
|
|
).fetchall()
|
|
return [dict(r) for r in rows]
|
|
|
|
def get_awg_server(self, server_id: int) -> dict[str, Any] | None:
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"""
|
|
SELECT id, name, api_url, api_token, panel_server_id, protocol,
|
|
enabled, verify_ssl,
|
|
last_check_at, last_check_ok, last_check_msg,
|
|
created_at, updated_at
|
|
FROM awg_servers
|
|
WHERE id = %s
|
|
""",
|
|
(int(server_id),),
|
|
).fetchone()
|
|
return dict(row) if row else None
|
|
|
|
def create_awg_server(
|
|
self,
|
|
*,
|
|
name: str,
|
|
api_url: str,
|
|
api_token: str = "",
|
|
panel_server_id: int = 0,
|
|
protocol: str = "awg2",
|
|
enabled: bool = True,
|
|
verify_ssl: bool = True,
|
|
) -> dict[str, Any]:
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"""
|
|
INSERT INTO awg_servers (
|
|
name, api_url, api_token, panel_server_id, protocol,
|
|
enabled, verify_ssl
|
|
)
|
|
VALUES (%s, %s, %s, %s, %s, %s, %s)
|
|
RETURNING id, name, api_url, api_token, panel_server_id, protocol,
|
|
enabled, verify_ssl,
|
|
last_check_at, last_check_ok, last_check_msg,
|
|
created_at, updated_at
|
|
""",
|
|
(
|
|
name.strip(),
|
|
api_url.strip(),
|
|
(api_token or "").strip(),
|
|
int(panel_server_id),
|
|
(protocol or "awg2").strip() or "awg2",
|
|
bool(enabled),
|
|
bool(verify_ssl),
|
|
),
|
|
).fetchone()
|
|
conn.commit()
|
|
return dict(row)
|
|
|
|
def update_awg_server(self, server_id: int, **fields: Any) -> dict[str, Any] | None:
|
|
allowed = {
|
|
"name",
|
|
"api_url",
|
|
"api_token",
|
|
"panel_server_id",
|
|
"protocol",
|
|
"enabled",
|
|
"verify_ssl",
|
|
"last_check_at",
|
|
"last_check_ok",
|
|
"last_check_msg",
|
|
}
|
|
updates = {k: v for k, v in fields.items() if k in allowed}
|
|
if not updates:
|
|
return self.get_awg_server(server_id)
|
|
cols = []
|
|
vals: list[Any] = []
|
|
for k, v in updates.items():
|
|
cols.append(f"{k} = %s")
|
|
vals.append(v)
|
|
cols.append("updated_at = NOW()")
|
|
vals.append(int(server_id))
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
f"""
|
|
UPDATE awg_servers
|
|
SET {', '.join(cols)}
|
|
WHERE id = %s
|
|
RETURNING id, name, api_url, api_token, panel_server_id, protocol,
|
|
enabled, verify_ssl,
|
|
last_check_at, last_check_ok, last_check_msg,
|
|
created_at, updated_at
|
|
""",
|
|
tuple(vals),
|
|
).fetchone()
|
|
conn.commit()
|
|
return dict(row) if row else None
|
|
|
|
def delete_awg_server(self, server_id: int) -> bool:
|
|
with self.pool.connection() as conn:
|
|
cur = conn.execute(
|
|
"DELETE FROM awg_servers WHERE id = %s",
|
|
(int(server_id),),
|
|
)
|
|
conn.commit()
|
|
return cur.rowcount > 0
|
|
|
|
def set_awg_server_check(
|
|
self,
|
|
server_id: int,
|
|
*,
|
|
ok: bool,
|
|
message: str,
|
|
) -> None:
|
|
with self.pool.connection() as conn:
|
|
conn.execute(
|
|
"""
|
|
UPDATE awg_servers
|
|
SET last_check_at = NOW(),
|
|
last_check_ok = %s,
|
|
last_check_msg = %s,
|
|
updated_at = NOW()
|
|
WHERE id = %s
|
|
""",
|
|
(bool(ok), (message or "")[:500], int(server_id)),
|
|
)
|
|
conn.commit()
|
|
|
|
def count_awg_configs(self, telegram_id: int) -> int:
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"SELECT COUNT(*) AS c FROM awg_configs WHERE telegram_id = %s",
|
|
(int(telegram_id),),
|
|
).fetchone()
|
|
return int(row["c"] if row else 0)
|
|
|
|
def list_awg_configs(self, telegram_id: int) -> list[dict[str, Any]]:
|
|
with self.pool.connection() as conn:
|
|
rows = conn.execute(
|
|
"""
|
|
SELECT c.id, c.telegram_id, c.server_id, c.client_id, c.client_name,
|
|
c.remark, c.config_text, c.vpn_link, c.protocol, c.enabled,
|
|
c.expire_at, c.created_at,
|
|
s.name AS server_name
|
|
FROM awg_configs c
|
|
LEFT JOIN awg_servers s ON s.id = c.server_id
|
|
WHERE c.telegram_id = %s
|
|
ORDER BY c.id DESC
|
|
""",
|
|
(int(telegram_id),),
|
|
).fetchall()
|
|
return [dict(r) for r in rows]
|
|
|
|
def get_awg_config(
|
|
self, config_id: int, *, telegram_id: int | None = None
|
|
) -> dict[str, Any] | None:
|
|
with self.pool.connection() as conn:
|
|
if telegram_id is None:
|
|
row = conn.execute(
|
|
"""
|
|
SELECT id, telegram_id, server_id, client_id, client_name,
|
|
remark, config_text, vpn_link, protocol, enabled,
|
|
expire_at, created_at
|
|
FROM awg_configs
|
|
WHERE id = %s
|
|
""",
|
|
(int(config_id),),
|
|
).fetchone()
|
|
else:
|
|
row = conn.execute(
|
|
"""
|
|
SELECT id, telegram_id, server_id, client_id, client_name,
|
|
remark, config_text, vpn_link, protocol, enabled,
|
|
expire_at, created_at
|
|
FROM awg_configs
|
|
WHERE id = %s AND telegram_id = %s
|
|
""",
|
|
(int(config_id), int(telegram_id)),
|
|
).fetchone()
|
|
return dict(row) if row else None
|
|
|
|
def create_awg_config(
|
|
self,
|
|
*,
|
|
telegram_id: int,
|
|
server_id: int,
|
|
client_id: str,
|
|
client_name: str = "",
|
|
remark: str = "",
|
|
config_text: str = "",
|
|
vpn_link: str = "",
|
|
protocol: str = "awg2",
|
|
enabled: bool = True,
|
|
expire_at: datetime | None = None,
|
|
) -> dict[str, Any]:
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"""
|
|
INSERT INTO awg_configs (
|
|
telegram_id, server_id, client_id, client_name, remark,
|
|
config_text, vpn_link, protocol, enabled, expire_at
|
|
)
|
|
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
|
|
RETURNING id, telegram_id, server_id, client_id, client_name,
|
|
remark, config_text, vpn_link, protocol, enabled,
|
|
expire_at, created_at
|
|
""",
|
|
(
|
|
int(telegram_id),
|
|
int(server_id),
|
|
client_id,
|
|
(client_name or "").strip(),
|
|
(remark or "").strip() or None,
|
|
config_text or "",
|
|
vpn_link or "",
|
|
(protocol or "awg2").strip() or "awg2",
|
|
bool(enabled),
|
|
expire_at,
|
|
),
|
|
).fetchone()
|
|
conn.commit()
|
|
return dict(row)
|
|
|
|
def update_awg_config(self, config_id: int, **fields: Any) -> dict[str, Any] | None:
|
|
allowed = {
|
|
"client_name",
|
|
"remark",
|
|
"config_text",
|
|
"vpn_link",
|
|
"protocol",
|
|
"enabled",
|
|
"expire_at",
|
|
}
|
|
updates = {k: v for k, v in fields.items() if k in allowed}
|
|
if not updates:
|
|
return self.get_awg_config(config_id)
|
|
cols = []
|
|
vals: list[Any] = []
|
|
for k, v in updates.items():
|
|
cols.append(f"{k} = %s")
|
|
vals.append(v)
|
|
vals.append(int(config_id))
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
f"""
|
|
UPDATE awg_configs
|
|
SET {', '.join(cols)}
|
|
WHERE id = %s
|
|
RETURNING id, telegram_id, server_id, client_id, client_name,
|
|
remark, config_text, vpn_link, protocol, enabled,
|
|
expire_at, created_at
|
|
""",
|
|
tuple(vals),
|
|
).fetchone()
|
|
conn.commit()
|
|
return dict(row) if row else None
|
|
|
|
def delete_awg_config(self, config_id: int, telegram_id: int) -> bool:
|
|
with self.pool.connection() as conn:
|
|
cur = conn.execute(
|
|
"DELETE FROM awg_configs WHERE id = %s AND telegram_id = %s",
|
|
(int(config_id), int(telegram_id)),
|
|
)
|
|
conn.commit()
|
|
return cur.rowcount > 0
|
|
|
|
def list_telegram_ids_with_enabled_vpn_configs(self, *, limit: int = 500) -> list[int]:
|
|
"""Пользователи с включёнными WG или AWG конфигами (для фонового enforce)."""
|
|
with self.pool.connection() as conn:
|
|
rows = conn.execute(
|
|
"""
|
|
SELECT DISTINCT telegram_id FROM (
|
|
SELECT telegram_id FROM xui_wg_configs WHERE enabled = TRUE
|
|
UNION
|
|
SELECT telegram_id FROM awg_configs WHERE enabled = TRUE
|
|
) t
|
|
ORDER BY telegram_id
|
|
LIMIT %s
|
|
""",
|
|
(int(limit),),
|
|
).fetchall()
|
|
return [int(r["telegram_id"]) for r in rows]
|
|
|
|
# ─── Growth: ban / referral / promo / trial / password reset / audit ───
|
|
|
|
def set_banned(self, telegram_id: int, banned: bool) -> BotUser:
|
|
with self.pool.connection() as conn:
|
|
self._ensure_user_conn(conn, telegram_id)
|
|
conn.execute(
|
|
"UPDATE users SET banned = %s WHERE telegram_id = %s",
|
|
(bool(banned), telegram_id),
|
|
)
|
|
if banned:
|
|
conn.execute(
|
|
"DELETE FROM user_sessions WHERE telegram_id = %s",
|
|
(telegram_id,),
|
|
)
|
|
conn.execute(
|
|
"""
|
|
UPDATE users SET session_token = NULL, session_expires = NULL
|
|
WHERE telegram_id = %s
|
|
""",
|
|
(telegram_id,),
|
|
)
|
|
conn.commit()
|
|
user = self.get_user(telegram_id)
|
|
assert user is not None
|
|
return user
|
|
|
|
def delete_user(self, telegram_id: int) -> bool:
|
|
"""Полностью удалить пользователя и связанные записи. False если не найден."""
|
|
tid = int(telegram_id)
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"SELECT telegram_id FROM users WHERE telegram_id = %s",
|
|
(tid,),
|
|
).fetchone()
|
|
if not row:
|
|
return False
|
|
|
|
# Снять ссылки с других пользователей
|
|
conn.execute(
|
|
"UPDATE users SET referred_by = NULL WHERE referred_by = %s",
|
|
(tid,),
|
|
)
|
|
conn.execute(
|
|
"UPDATE topup_codes SET used_by = NULL WHERE used_by = %s",
|
|
(tid,),
|
|
)
|
|
conn.execute(
|
|
"UPDATE login_qr SET telegram_id = NULL WHERE telegram_id = %s",
|
|
(tid,),
|
|
)
|
|
|
|
# Зависимые таблицы без ON DELETE CASCADE
|
|
conn.execute(
|
|
"DELETE FROM referral_rewards WHERE referrer_id = %s OR referee_id = %s",
|
|
(tid, tid),
|
|
)
|
|
conn.execute("DELETE FROM tickets WHERE telegram_id = %s", (tid,))
|
|
conn.execute("DELETE FROM xui_wg_configs WHERE telegram_id = %s", (tid,))
|
|
conn.execute("DELETE FROM awg_configs WHERE telegram_id = %s", (tid,))
|
|
conn.execute("DELETE FROM transactions WHERE telegram_id = %s", (tid,))
|
|
conn.execute("DELETE FROM purchases WHERE telegram_id = %s", (tid,))
|
|
conn.execute("DELETE FROM invoices WHERE telegram_id = %s", (tid,))
|
|
conn.execute("DELETE FROM digiseller_payments WHERE telegram_id = %s", (tid,))
|
|
conn.execute("DELETE FROM payment_events WHERE telegram_id = %s", (tid,))
|
|
conn.execute("DELETE FROM webauthn_challenges WHERE telegram_id = %s", (tid,))
|
|
conn.execute("DELETE FROM user_sessions WHERE telegram_id = %s", (tid,))
|
|
conn.execute("DELETE FROM webauthn_credentials WHERE telegram_id = %s", (tid,))
|
|
conn.execute("DELETE FROM promo_redemptions WHERE telegram_id = %s", (tid,))
|
|
conn.execute("DELETE FROM password_reset_tokens WHERE telegram_id = %s", (tid,))
|
|
conn.execute("DELETE FROM users WHERE telegram_id = %s", (tid,))
|
|
conn.commit()
|
|
return True
|
|
|
|
def set_notify_prefs(
|
|
self,
|
|
telegram_id: int,
|
|
*,
|
|
notify_telegram: bool | None = None,
|
|
notify_email: bool | None = None,
|
|
) -> BotUser:
|
|
user = self.ensure_user(telegram_id)
|
|
tg = user.notify_telegram if notify_telegram is None else bool(notify_telegram)
|
|
em = user.notify_email if notify_email is None else bool(notify_email)
|
|
with self.pool.connection() as conn:
|
|
conn.execute(
|
|
"""
|
|
UPDATE users
|
|
SET notify_telegram = %s, notify_email = %s
|
|
WHERE telegram_id = %s
|
|
""",
|
|
(tg, em, telegram_id),
|
|
)
|
|
conn.commit()
|
|
updated = self.get_user(telegram_id)
|
|
assert updated is not None
|
|
return updated
|
|
|
|
def ensure_referral_code(self, telegram_id: int) -> str:
|
|
user = self.ensure_user(telegram_id)
|
|
if user.referral_code:
|
|
return user.referral_code
|
|
for _ in range(8):
|
|
code = secrets.token_urlsafe(6).replace("-", "").replace("_", "")[:10].lower()
|
|
with self.pool.connection() as conn:
|
|
try:
|
|
conn.execute(
|
|
"""
|
|
UPDATE users SET referral_code = %s
|
|
WHERE telegram_id = %s AND referral_code IS NULL
|
|
""",
|
|
(code, telegram_id),
|
|
)
|
|
conn.commit()
|
|
except Exception: # noqa: BLE001
|
|
conn.rollback()
|
|
continue
|
|
fresh = self.get_user(telegram_id)
|
|
if fresh and fresh.referral_code:
|
|
return fresh.referral_code
|
|
raise AuthError("Не удалось создать реферальный код")
|
|
|
|
def get_user_by_referral_code(self, code: str) -> BotUser | None:
|
|
raw = (code or "").strip().lower()
|
|
if not raw:
|
|
return None
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"SELECT * FROM users WHERE lower(referral_code) = %s",
|
|
(raw,),
|
|
).fetchone()
|
|
return self._to_user(row) if row else None
|
|
|
|
def attach_referrer(self, telegram_id: int, referrer_code: str) -> bool:
|
|
"""Привязать реферера только новым пользователям без referred_by."""
|
|
referrer = self.get_user_by_referral_code(referrer_code)
|
|
if not referrer or referrer.telegram_id == telegram_id:
|
|
return False
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"""
|
|
UPDATE users
|
|
SET referred_by = %s
|
|
WHERE telegram_id = %s
|
|
AND referred_by IS NULL
|
|
AND created_at > NOW() - INTERVAL '7 days'
|
|
RETURNING telegram_id
|
|
""",
|
|
(referrer.telegram_id, telegram_id),
|
|
).fetchone()
|
|
conn.commit()
|
|
return row is not None
|
|
|
|
def grant_referral_rewards_on_first_purchase(
|
|
self,
|
|
referee_id: int,
|
|
*,
|
|
referrer_bonus: float,
|
|
referee_bonus: float,
|
|
) -> None:
|
|
user = self.get_user(referee_id)
|
|
if not user or not user.referred_by:
|
|
return
|
|
referrer_id = int(user.referred_by)
|
|
with self.pool.connection() as conn:
|
|
exists = conn.execute(
|
|
"""
|
|
SELECT 1 FROM referral_rewards
|
|
WHERE referee_id = %s AND kind = 'first_purchase'
|
|
""",
|
|
(referee_id,),
|
|
).fetchone()
|
|
if exists:
|
|
return
|
|
try:
|
|
conn.execute(
|
|
"""
|
|
INSERT INTO referral_rewards (referrer_id, referee_id, amount, kind)
|
|
VALUES (%s, %s, %s, 'first_purchase')
|
|
""",
|
|
(referrer_id, referee_id, float(referrer_bonus)),
|
|
)
|
|
except Exception: # noqa: BLE001
|
|
conn.rollback()
|
|
return
|
|
conn.commit()
|
|
if referrer_bonus > 0:
|
|
self.add_balance(
|
|
referrer_id,
|
|
referrer_bonus,
|
|
kind="referral_bonus",
|
|
meta=f"ref:{referee_id}",
|
|
)
|
|
if referee_bonus > 0:
|
|
self.add_balance(
|
|
referee_id,
|
|
referee_bonus,
|
|
kind="referral_welcome",
|
|
meta=f"from:{referrer_id}",
|
|
)
|
|
|
|
def mark_trial_used(self, telegram_id: int) -> None:
|
|
with self.pool.connection() as conn:
|
|
conn.execute(
|
|
"""
|
|
UPDATE users SET trial_used_at = %s
|
|
WHERE telegram_id = %s AND trial_used_at IS NULL
|
|
""",
|
|
(_utc_now(), telegram_id),
|
|
)
|
|
conn.commit()
|
|
|
|
def create_promo_code(
|
|
self,
|
|
*,
|
|
code: str,
|
|
kind: str,
|
|
value: float,
|
|
max_uses: int | None = None,
|
|
tariff_id: str | None = None,
|
|
expires_at: datetime | None = None,
|
|
note: str | None = None,
|
|
) -> dict[str, Any]:
|
|
code_norm = (code or "").strip().upper()
|
|
if len(code_norm) < 3:
|
|
raise AuthError("Код слишком короткий")
|
|
kind_n = (kind or "").strip().lower()
|
|
if kind_n not in {"percent", "fixed", "days"}:
|
|
raise AuthError("kind: percent | fixed | days")
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"""
|
|
INSERT INTO promo_codes (
|
|
code, kind, value, max_uses, tariff_id, expires_at, note
|
|
) VALUES (%s, %s, %s, %s, %s, %s, %s)
|
|
RETURNING *
|
|
""",
|
|
(
|
|
code_norm,
|
|
kind_n,
|
|
float(value),
|
|
max_uses,
|
|
tariff_id or None,
|
|
expires_at,
|
|
(note or "")[:200] or None,
|
|
),
|
|
).fetchone()
|
|
conn.commit()
|
|
assert row is not None
|
|
return dict(row)
|
|
|
|
def list_promo_codes(self, *, limit: int = 100) -> list[dict[str, Any]]:
|
|
with self.pool.connection() as conn:
|
|
rows = conn.execute(
|
|
"""
|
|
SELECT * FROM promo_codes
|
|
ORDER BY id DESC
|
|
LIMIT %s
|
|
""",
|
|
(max(1, min(limit, 200)),),
|
|
).fetchall()
|
|
return [dict(r) for r in rows]
|
|
|
|
def set_promo_active(self, promo_id: int, active: bool) -> None:
|
|
with self.pool.connection() as conn:
|
|
conn.execute(
|
|
"UPDATE promo_codes SET active = %s WHERE id = %s",
|
|
(bool(active), int(promo_id)),
|
|
)
|
|
conn.commit()
|
|
|
|
def apply_promo_to_price(
|
|
self,
|
|
*,
|
|
code: str,
|
|
telegram_id: int,
|
|
tariff_id: str,
|
|
price: float,
|
|
days: int,
|
|
) -> tuple[float, int, dict[str, Any]]:
|
|
"""Вернуть (new_price, bonus_days, promo_row)."""
|
|
code_norm = (code or "").strip().upper()
|
|
if not code_norm:
|
|
raise AuthError("Укажите промокод")
|
|
now = _utc_now()
|
|
with self.pool.connection() as conn:
|
|
promo = conn.execute(
|
|
"SELECT * FROM promo_codes WHERE lower(code) = lower(%s)",
|
|
(code_norm,),
|
|
).fetchone()
|
|
if not promo or not promo.get("active"):
|
|
raise AuthError("Промокод не найден или выключен")
|
|
if promo.get("expires_at") and promo["expires_at"] <= now:
|
|
raise AuthError("Промокод истёк")
|
|
if promo.get("tariff_id") and str(promo["tariff_id"]) != str(tariff_id):
|
|
raise AuthError("Промокод не для этого тарифа")
|
|
max_uses = promo.get("max_uses")
|
|
if max_uses is not None and int(promo.get("used_count") or 0) >= int(max_uses):
|
|
raise AuthError("Лимит использований промокода исчерпан")
|
|
used = conn.execute(
|
|
"""
|
|
SELECT 1 FROM promo_redemptions
|
|
WHERE promo_id = %s AND telegram_id = %s
|
|
""",
|
|
(promo["id"], telegram_id),
|
|
).fetchone()
|
|
if used:
|
|
raise AuthError("Вы уже использовали этот промокод")
|
|
|
|
kind = str(promo["kind"])
|
|
value = float(promo["value"])
|
|
new_price = float(price)
|
|
bonus_days = 0
|
|
if kind == "percent":
|
|
new_price = max(0.0, price * (1.0 - min(100.0, value) / 100.0))
|
|
elif kind == "fixed":
|
|
new_price = max(0.0, price - value)
|
|
elif kind == "days":
|
|
bonus_days = max(0, int(value))
|
|
return round(new_price, 2), bonus_days, dict(promo)
|
|
|
|
def redeem_promo(self, *, promo_id: int, telegram_id: int, amount_saved: float) -> None:
|
|
with self.pool.connection() as conn:
|
|
conn.execute(
|
|
"""
|
|
INSERT INTO promo_redemptions (promo_id, telegram_id, amount_saved)
|
|
VALUES (%s, %s, %s)
|
|
ON CONFLICT (promo_id, telegram_id) DO NOTHING
|
|
""",
|
|
(int(promo_id), telegram_id, float(amount_saved)),
|
|
)
|
|
conn.execute(
|
|
"""
|
|
UPDATE promo_codes SET used_count = used_count + 1
|
|
WHERE id = %s
|
|
""",
|
|
(int(promo_id),),
|
|
)
|
|
conn.commit()
|
|
|
|
def create_password_reset_token(self, telegram_id: int) -> str:
|
|
token = secrets.token_urlsafe(32)
|
|
expires = _utc_now() + timedelta(hours=1)
|
|
with self.pool.connection() as conn:
|
|
conn.execute(
|
|
"DELETE FROM password_reset_tokens WHERE telegram_id = %s OR expires_at < %s",
|
|
(telegram_id, _utc_now()),
|
|
)
|
|
conn.execute(
|
|
"""
|
|
INSERT INTO password_reset_tokens (token, telegram_id, expires_at)
|
|
VALUES (%s, %s, %s)
|
|
""",
|
|
(token, telegram_id, expires),
|
|
)
|
|
conn.commit()
|
|
return token
|
|
|
|
def consume_password_reset_token(self, token: str) -> int:
|
|
token = (token or "").strip()
|
|
if not token:
|
|
raise AuthError("Некорректный токен")
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"""
|
|
UPDATE password_reset_tokens
|
|
SET used_at = %s
|
|
WHERE token = %s AND used_at IS NULL AND expires_at > %s
|
|
RETURNING telegram_id
|
|
""",
|
|
(_utc_now(), token, _utc_now()),
|
|
).fetchone()
|
|
conn.commit()
|
|
if not row:
|
|
raise AuthError("Ссылка недействительна или истекла")
|
|
return int(row["telegram_id"])
|
|
|
|
def set_password_hash(self, telegram_id: int, password: str) -> None:
|
|
if len(password) < 6:
|
|
raise AuthError("Пароль минимум 6 символов")
|
|
pwd_hash = hash_password(password)
|
|
with self.pool.connection() as conn:
|
|
conn.execute(
|
|
"UPDATE users SET password_hash = %s WHERE telegram_id = %s",
|
|
(pwd_hash, telegram_id),
|
|
)
|
|
conn.commit()
|
|
|
|
def get_user_by_email(self, email: str) -> BotUser | None:
|
|
raw = (email or "").strip().lower()
|
|
if not raw:
|
|
return None
|
|
with self.pool.connection() as conn:
|
|
row = conn.execute(
|
|
"SELECT * FROM users WHERE lower(email) = %s",
|
|
(raw,),
|
|
).fetchone()
|
|
return self._to_user(row) if row else None
|
|
|
|
def log_payment_event(
|
|
self,
|
|
*,
|
|
provider: str,
|
|
event_type: str,
|
|
telegram_id: int | None = None,
|
|
amount: float | None = None,
|
|
invoice_ref: str | None = None,
|
|
status: str = "ok",
|
|
meta: str | None = None,
|
|
) -> None:
|
|
with self.pool.connection() as conn:
|
|
conn.execute(
|
|
"""
|
|
INSERT INTO payment_events (
|
|
provider, event_type, telegram_id, amount, invoice_ref, status, meta
|
|
) VALUES (%s, %s, %s, %s, %s, %s, %s)
|
|
""",
|
|
(
|
|
(provider or "")[:40],
|
|
(event_type or "")[:60],
|
|
telegram_id,
|
|
amount,
|
|
(invoice_ref or "")[:120] or None,
|
|
(status or "ok")[:40],
|
|
(meta or "")[:500] or None,
|
|
),
|
|
)
|
|
conn.commit()
|
|
|
|
def list_payment_events(self, *, limit: int = 50) -> list[dict[str, Any]]:
|
|
with self.pool.connection() as conn:
|
|
rows = conn.execute(
|
|
"""
|
|
SELECT * FROM payment_events
|
|
ORDER BY id DESC
|
|
LIMIT %s
|
|
""",
|
|
(max(1, min(limit, 200)),),
|
|
).fetchall()
|
|
return [dict(r) for r in rows]
|
|
|
|
@staticmethod
|
|
def _ensure_user_conn(conn, telegram_id: int) -> None:
|
|
row = conn.execute(
|
|
"SELECT 1 FROM users WHERE telegram_id = %s",
|
|
(telegram_id,),
|
|
).fetchone()
|
|
if not row:
|
|
conn.execute(
|
|
"""
|
|
INSERT INTO users (telegram_id, username, full_name, balance, created_at)
|
|
VALUES (%s, NULL, NULL, 0, %s)
|
|
ON CONFLICT (telegram_id) DO NOTHING
|
|
""",
|
|
(telegram_id, _utc_now()),
|
|
)
|