diff --git "a/app.py" "b/app.py" new file mode 100644--- /dev/null +++ "b/app.py" @@ -0,0 +1,2966 @@ +""" +SalGram — мессенджер. Этапы 1–4: аккаунты, профили, переписка, модерация. + +Бэкенд: Flask + flask-sock (WebSocket) + SQLite. Запуск: python app.py. + +Особые аккаунты: + @salgr — администратор: видит любые аватарки и контакты, удаляет аккаунты, + его нельзя заблокировать; рядом с ником — галочка. + @salgram — официальный аккаунт: от его имени приходят уведомления о банах + и рассылки. Писать в него может только @salgr — такое сообщение + рассылается всем пользователям от имени @salgram. + @reports — автоматизированный: принимает жалобы, пересылает переписку админу. + @notes — автоматизированный: личные заметки (каждый пишет туда сам себе). + @addgif — бот: принимает GIF в общую библиотеку (поиск по ключевым словам). + В системные аккаунты нельзя войти, их нельзя заблокировать или удалить. + +Забаненный аккаунт (banned=1) может войти и читать переписку, но всё в режиме +«только чтение»: писать, изменять и удалять что-либо нельзя. + +Шифрование: текст сообщений и файлы фото хранятся зашифрованными (Fernet, +AES-128-CBC + HMAC); ключ — в /data/secret.key. Аккаунты помечаются удалёнными +(deleted=1), а не стираются — ID не переиспользуются. + +HTTP API (все, кроме register/login, требуют cookie LoggedAccount): + POST /api/register | /api/login | /api/logout + GET /api/me POST /api/me/ + GET|POST /api/me/contacts DELETE /api/me/contacts/ + GET /api/avatar/ GET /api/users/ GET /api/peer/ + POST /api/users//block | /unblock | /report + POST /api/admin/users//delete (только админ) + POST /api/admin/reports// (только админ) + GET /api/search?q&type&mode GET /api/chats + POST /api/chats//delete (удаление переписки у обоих) + POST /api/messages/photo GET /api/photos/ + POST /api/messages/file GET /api/files/?download=0|1 + POST /api/messages/gif GET /api/gifs/ + POST /api/messages/audio GET /api/audio/ (голосовые/аудиофайлы) + GET /api/push/key POST /api/push/subscribe | /unsubscribe + GET|POST /api/gifs (поиск / добавление GIF в общую библиотеку @addgif) + POST /api/me/reactions (закреплённые реакции) + GET|POST|DELETE /api/me/wallpaper (обои чата: фото или видео) + POST /api/groups GET /api/groups/ GET /api/groups/hash/ + POST /api/groups//join | /leave | /delete | /avatar (аватарка — владелец) + POST /api/groups//add (добавить пользователя — участником) + POST /api/groups/join/ (вступление по ссылке-приглашению) + POST /api/channels//link | /unlink (привязка группы комментариев — владелец) + POST /api/channels//comments (вступление в группу комментариев) + POST /api/groups//moderation (настройки модерации — модератор) + POST /api/groups//members// (модерация) + GET /api/group_avatar/ + +Вложения дедуплицируются: имя файла — SHA-256 содержимого, повторная отправка +того же фото/файла ссылается на уже сохранённый файл; файл удаляется с диска, +только когда на него не остаётся ни одной ссылки из сообщений. + +Группы: публичные находятся поиском, в приватные можно попасть только по +пригласительной ссылке /#join=. Сообщения групп лежат в той же таблице +messages (group_id заполнен, recipient_id = 0). + +WebSocket /ws (JSON; «group» вместо «to»/«with» — для групповых чатов): + -> {type: "send", to|group: , text} <- {type: "message", message} + -> {type: "history", with|group, before?} <- {type: "history", with|group, messages, pinned} + -> {type: "edit", id, text} <- {type: "message_edited", message} + -> {type: "delete", id} <- {type: "message_deleted", id, from, to, group} + -> {type: "react", id, emoji} <- {type: "reaction", id, from, to, group, reactions} + -> {type: "pin", id, pinned} <- {type: "message_pinned", message} + <- {type: "chat_deleted", with} + <- {type: "group_removed", group} + <- {type: "group_added", group} + <- {type: "moderation", group, text} + <- {type: "error", error} +""" +import base64 +import hashlib +import io +import json +import os +import re +import secrets +import sqlite3 +import threading +import urllib.request +import zipfile +from collections import defaultdict +from functools import wraps + +from cryptography.fernet import Fernet, InvalidToken +from flask import Flask, g, jsonify, request, send_file +from flask_sock import Sock +from PIL import Image, ImageOps +from werkzeug.middleware.proxy_fix import ProxyFix +from werkzeug.security import check_password_hash, generate_password_hash + +import hf_backup # резервное копирование БД в приватный репозиторий HF (см. модуль) + +# Каталог данных: SALGRAM_DATA_DIR (напр. для Hugging Face Spaces), иначе +# постоянное хранилище /data (Amvera/HF persistent storage), иначе — рядом с кодом. +DATA_DIR = (os.environ.get("SALGRAM_DATA_DIR") + or ("/data" if os.path.isdir("/data") else os.path.dirname(__file__))) +os.makedirs(DATA_DIR, exist_ok=True) +# До загрузки ключа и init_db: если диск пуст (свежий контейнер) — восстановить +# БД, ключ и медиа из приватного бэкапа. No-op, если бэкап не настроен. +hf_backup.restore_on_boot(DATA_DIR) +DB_PATH = os.path.join(DATA_DIR, "salgram.db") +AVATAR_DIR = os.path.join(DATA_DIR, "avatars") +PHOTO_DIR = os.path.join(DATA_DIR, "photos") +FILE_DIR = os.path.join(DATA_DIR, "files") +WALLPAPER_DIR = os.path.join(DATA_DIR, "wallpapers") +GIF_DIR = os.path.join(DATA_DIR, "gifs") +AUDIO_DIR = os.path.join(DATA_DIR, "audio") +KEY_PATH = os.path.join(DATA_DIR, "secret.key") +VAPID_KEY_PATH = os.path.join(DATA_DIR, "vapid.json") # ключи Web Push (VAPID) + +ADMIN_USERNAME = "salgr" # администратор +SALGRAM_USERNAME = "salgram" # официальный аккаунт: рассылки и уведомления о банах +REPORTS_USERNAME = "reports" # автоматизированный приём жалоб +NOTES_USERNAME = "notes" # личные заметки +ADDGIF_USERNAME = "addgif" # бот: приём GIF в общую библиотеку +NEWS_USERNAME = "salgram_news" # официальный канал: автоподписка всех новых пользователей + +USERNAME_RE = re.compile(r"^[A-Za-z0-9_]{3,32}$") +PRIVACY_VALUES = {"everyone", "contacts", "nobody"} +SLOW_SCOPES = {"nobody", "noncontacts", "everyone"} +WHITELIST_MODES = {"any", "all"} # белый список: хотя бы одно слово / все слова +MOD_WORDS_MAX = 100 # максимум записей в списке слов модерации +MOD_SLOW_MAX = 86400 # медленный режим группы — до суток + +# Реакции: полный набор (порядок — как в выборе на фронтенде). +REACTIONS = [ + "😁", "🤣", "😍", "😘", "🥰", "🥲", "🤩", "🤔", "🫡", "🤨", "😐", "😑", "😶", + "🙄", "😯", "🥱", "😴", "🤤", "🤑", "🙁", "😭", "😨", "🤯", "😬", "😱", "🤬", + "🤮", "🥺", "🥹", "🤡", "🤓", "😈", "👿", "💀", "💩", "👀", "🤞", "🫸", "🤝", + "👌", "👍", "👎", "🤌", "✊", "👋", "👏", "🎉", "💋", "🏆", "🗿", "📝", "❤️", + "❌", "✅", +] +# Для проверки принимаем и реакции из старого набора (могли остаться в БД). +REACTION_OK = set(REACTIONS) | {"😂", "😮", "😢", "🔥"} +DEFAULT_PINNED = ["🤝", "👍", "👎", "✅", "❌"] # закреплённые по умолчанию + +AVATAR_SIDE = 1024 +PHOTO_SIDE = 1600 +MSG_MAX = 4096 +GROUP_TITLE_MAX = 64 +FILE_MAX = 25 * 1024 * 1024 # файлы — до 25 МБ +GIF_MAX = 15 * 1024 * 1024 # GIF-анимации — до 15 МБ +AUDIO_MAX = 25 * 1024 * 1024 # аудио/голосовые — до 25 МБ +# Допустимые форматы аудио: имя файла = SHA-256 содержимого + расширение из MIME. +AUDIO_EXT = {"audio/webm": "webm", "audio/ogg": "ogg", "audio/oga": "ogg", + "audio/mpeg": "mp3", "audio/mp3": "mp3", "audio/mp4": "m4a", + "audio/x-m4a": "m4a", "audio/aac": "aac", "audio/wav": "wav", + "audio/x-wav": "wav", "audio/wave": "wav", "audio/flac": "flac"} +AUDIO_NAME_RE = re.compile(r"[0-9a-f]{64}\.(webm|ogg|mp3|m4a|aac|wav|flac)") +AVATAR_VIDEO_MAX = 10 * 1024 * 1024 # видео-аватарка — до 10 МБ (длина ≤10 с — проверка в браузере) +GIF_DAILY_LIMIT = 5 # сколько GIF в сутки можно добавить в библиотеку +KEYWORDS_MAX = 200 # длина строки ключевых слов GIF +PREVIEW_CHARS = 400 # длина текстового предпросмотра документов +# Расширение файла аватарки зависит от её типа (картинка/гиф/видео). +AVATAR_EXT = {"image": "webp", "gif": "gif", "video": "mp4"} + +app = Flask(__name__, static_folder="static", static_url_path="") +app.config["MAX_CONTENT_LENGTH"] = FILE_MAX + 1024 * 1024 # + запас на multipart +# За реверс-прокси хостинга: реальный IP клиента и схема https из X-Forwarded-*, +# чтобы в логах были адреса клиентов, а cookie получала флаг Secure. +app.wsgi_app = ProxyFix(app.wsgi_app, x_for=1, x_proto=1) +sock = Sock(app) + +# Активные WebSocket-соединения: user_id -> множество соединений. +connections: dict[int, set] = defaultdict(set) + + +# ---------- Шифрование (Fernet: AES-128-CBC + HMAC) ---------- + +def _load_key() -> Fernet: + if not os.path.exists(KEY_PATH): + with open(KEY_PATH, "wb") as f: + f.write(Fernet.generate_key()) + with open(KEY_PATH, "rb") as f: + return Fernet(f.read()) + + +fernet = _load_key() + + +def enc(text: str) -> str: + return fernet.encrypt(text.encode()).decode() + + +def dec(token: str) -> str: + """Расшифровка; незашифрованные строки (старые БД) возвращаются как есть.""" + try: + return fernet.decrypt(token.encode()).decode() + except (InvalidToken, ValueError): + return token + + +# ---------- Web Push (VAPID) ---------- +# Уведомления в фоне для PWA/браузера/десктопа. Если библиотека недоступна — +# приложение работает как раньше, просто без фоновых push (остаются уведомления +# при открытом приложении по WebSocket). + +try: + from pywebpush import webpush, WebPushException + from cryptography.hazmat.primitives import serialization + from cryptography.hazmat.primitives.asymmetric import ec + _WEBPUSH_OK = True +except Exception: # pragma: no cover — библиотека не установлена + _WEBPUSH_OK = False + +VAPID_SUB = "mailto:admin@salgram.app" # контакт для push-сервисов (требование VAPID) +VAPID_PRIVATE_PEM = "" # приватный ключ (PEM) для подписи push +VAPID_PUBLIC_B64 = "" # публичный ключ (applicationServerKey для браузера) + + +def _load_vapid(): + """Загружает (или генерирует один раз) пару ключей VAPID для Web Push.""" + global VAPID_PRIVATE_PEM, VAPID_PUBLIC_B64 + if not _WEBPUSH_OK: + return + if os.path.exists(VAPID_KEY_PATH): + with open(VAPID_KEY_PATH) as f: + data = json.load(f) + VAPID_PRIVATE_PEM, VAPID_PUBLIC_B64 = data["private_pem"], data["public_b64"] + return + priv = ec.generate_private_key(ec.SECP256R1()) + VAPID_PRIVATE_PEM = priv.private_bytes( + serialization.Encoding.PEM, serialization.PrivateFormat.PKCS8, + serialization.NoEncryption()).decode() + raw = priv.public_key().public_bytes( + serialization.Encoding.X962, serialization.PublicFormat.UncompressedPoint) + VAPID_PUBLIC_B64 = base64.urlsafe_b64encode(raw).rstrip(b"=").decode() + with open(VAPID_KEY_PATH, "w") as f: + json.dump({"private_pem": VAPID_PRIVATE_PEM, "public_b64": VAPID_PUBLIC_B64}, f) + + +def _push_worker(subs: list[tuple[str, dict]], payload: str): + """Фоновая отправка push (сеть медленная — не блокируем обработчик сообщений). + Мёртвые подписки (404/410) удаляются через отдельное соединение с БД.""" + dead = [] + for endpoint, info in subs: + try: + webpush(subscription_info=info, data=payload, + vapid_private_key=VAPID_PRIVATE_PEM, + vapid_claims={"sub": VAPID_SUB}, timeout=10) + except WebPushException as e: + resp = getattr(e, "response", None) + if resp is not None and resp.status_code in (404, 410): + dead.append(endpoint) + except Exception: + pass + if dead: + try: + with sqlite3.connect(DB_PATH, timeout=15) as c: + c.executemany("DELETE FROM push_subscriptions WHERE endpoint = ?", + [(e,) for e in dead]) + except Exception: + pass + + +def _spawn_push(subs: list[tuple[str, dict]], title: str, body: str, data: dict): + if not (_WEBPUSH_OK and VAPID_PRIVATE_PEM and subs): + return + payload = json.dumps({"title": title, "body": body, **data}, ensure_ascii=False) + threading.Thread(target=_push_worker, args=(subs, payload), daemon=True).start() + + +def _notify_preview(kind: str, content: str, extra: dict | None) -> str: + """Короткий текст сообщения для уведомления (content для text — ещё не шифрован).""" + if kind == "text": + return content[:120] + ex = extra or {} + return {"photo": "📷 Фото", "gif": "GIF", "call": "📞 Звонок", + "file": "📎 " + (ex.get("name") or "Файл"), + "audio": "🎤 Голосовое сообщение" if ex.get("voice") else "🎵 Аудио", + }.get(kind, "Сообщение") + + +def push_notify_user(user_id: int, title: str, body: str, data: dict): + """Шлёт Web Push подпискам пользователя (в фоне). Если он онлайн (есть WS), + push не шлём — уведомит само открытое приложение, иначе будет дубль.""" + if connections.get(user_id): + return + rows = db().execute( + "SELECT endpoint, data FROM push_subscriptions WHERE user_id = ?", + (user_id,)).fetchall() + _spawn_push([(r["endpoint"], json.loads(r["data"])) for r in rows], title, body, data) + + +def push_notify_group(gid: int, sender_id: int, title: str, body: str, data: dict): + """Шлёт Web Push участникам группы, кроме автора и тех, кто сейчас онлайн.""" + rows = db().execute( + """SELECT ps.endpoint, ps.data, ps.user_id FROM push_subscriptions ps + JOIN group_members gm ON gm.user_id = ps.user_id + WHERE gm.group_id = ? AND ps.user_id != ?""", (gid, sender_id)).fetchall() + subs = [(r["endpoint"], json.loads(r["data"])) + for r in rows if not connections.get(r["user_id"])] + _spawn_push(subs, title, body, data) + + +# ---------- База данных ---------- + +def db() -> sqlite3.Connection: + """Соединение с БД, одно на запрос/WS-сессию (хранится в flask.g).""" + if "db" not in g: + g.db = sqlite3.connect(DB_PATH, timeout=15) + g.db.row_factory = sqlite3.Row + return g.db + + +@app.teardown_appcontext +def close_db(_exc): + if (conn := g.pop("db", None)) is not None: + conn.close() + + +def init_db(): + """Создаёт таблицы, применяет миграции, заводит системные аккаунты.""" + os.makedirs(AVATAR_DIR, exist_ok=True) + os.makedirs(PHOTO_DIR, exist_ok=True) + os.makedirs(FILE_DIR, exist_ok=True) + os.makedirs(WALLPAPER_DIR, exist_ok=True) + os.makedirs(GIF_DIR, exist_ok=True) + os.makedirs(AUDIO_DIR, exist_ok=True) + with sqlite3.connect(DB_PATH) as conn: + conn.execute("PRAGMA journal_mode=WAL") + conn.executescript(""" + CREATE TABLE IF NOT EXISTS users ( + id INTEGER PRIMARY KEY, + display_name TEXT NOT NULL, + username TEXT NOT NULL UNIQUE COLLATE NOCASE, + password_hash TEXT NOT NULL, + created_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP + ); + CREATE TABLE IF NOT EXISTS sessions ( + token TEXT PRIMARY KEY, + user_id INTEGER NOT NULL REFERENCES users(id), + created_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP + ); + CREATE TABLE IF NOT EXISTS contacts ( + user_id INTEGER NOT NULL REFERENCES users(id), + contact_id INTEGER NOT NULL REFERENCES users(id), + PRIMARY KEY (user_id, contact_id) + ); + CREATE TABLE IF NOT EXISTS blocks ( + blocker_id INTEGER NOT NULL REFERENCES users(id), + blocked_id INTEGER NOT NULL REFERENCES users(id), + PRIMARY KEY (blocker_id, blocked_id) + ); + CREATE TABLE IF NOT EXISTS messages ( + id INTEGER PRIMARY KEY, + sender_id INTEGER NOT NULL, + recipient_id INTEGER NOT NULL, + kind TEXT NOT NULL DEFAULT 'text', + content TEXT NOT NULL, -- текст (зашифрован) или файл фото + created_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP -- UTC (ЧП +0) + ); + CREATE TABLE IF NOT EXISTS reports ( + id INTEGER PRIMARY KEY, + reporter_id INTEGER NOT NULL REFERENCES users(id), + target_id INTEGER NOT NULL REFERENCES users(id), + status TEXT NOT NULL DEFAULT 'open', -- open | banned | dismissed + created_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP + ); + CREATE TABLE IF NOT EXISTS reactions ( + message_id INTEGER NOT NULL REFERENCES messages(id), + user_id INTEGER NOT NULL REFERENCES users(id), + emoji TEXT NOT NULL, + PRIMARY KEY (message_id, user_id) + ); + CREATE TABLE IF NOT EXISTS groups ( + id INTEGER PRIMARY KEY, + title TEXT NOT NULL, + username TEXT, -- юзернейм публичной группы + owner_id INTEGER NOT NULL REFERENCES users(id), + public INTEGER NOT NULL DEFAULT 1, -- 0 = вход только по ссылке + invite_hash TEXT NOT NULL UNIQUE, -- хэш пригласительной ссылки + created_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP + ); + CREATE TABLE IF NOT EXISTS group_members ( + group_id INTEGER NOT NULL REFERENCES groups(id), + user_id INTEGER NOT NULL REFERENCES users(id), + joined_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP, + PRIMARY KEY (group_id, user_id) + ); + CREATE TABLE IF NOT EXISTS group_bans ( + group_id INTEGER NOT NULL REFERENCES groups(id), + user_id INTEGER NOT NULL REFERENCES users(id), + PRIMARY KEY (group_id, user_id) + ); + CREATE TABLE IF NOT EXISTS gifs ( + id INTEGER PRIMARY KEY, + file TEXT NOT NULL, -- SHA-256 имя файла в GIF_DIR (.gif) + keywords TEXT NOT NULL, -- ключевые слова (нижний регистр) для поиска + uploader_id INTEGER NOT NULL REFERENCES users(id), + created_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP + ); + CREATE TABLE IF NOT EXISTS push_subscriptions ( + endpoint TEXT PRIMARY KEY, -- URL подписки браузера (уникален) + user_id INTEGER NOT NULL REFERENCES users(id), + data TEXT NOT NULL, -- JSON PushSubscription (keys p256dh/auth) + created_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP + ); + CREATE TABLE IF NOT EXISTS reads ( + reader_id INTEGER NOT NULL REFERENCES users(id), + peer_id INTEGER NOT NULL REFERENCES users(id), + last_read_id INTEGER NOT NULL DEFAULT 0, -- макс. id прочитанного сообщения собеседника + PRIMARY KEY (reader_id, peer_id) + ); + CREATE INDEX IF NOT EXISTS idx_push_user ON push_subscriptions (user_id); + CREATE INDEX IF NOT EXISTS idx_gifs_uploader ON gifs (uploader_id, created_at); + CREATE INDEX IF NOT EXISTS idx_msg_pair ON messages (sender_id, recipient_id, id); + CREATE INDEX IF NOT EXISTS idx_msg_recipient ON messages (recipient_id, id); + """) + # Миграции: недостающие колонки. + users_cols = {row[1] for row in conn.execute("PRAGMA table_info(users)")} + for col, ddl in { + "description": "TEXT NOT NULL DEFAULT ''", + "privacy_avatar": "TEXT NOT NULL DEFAULT 'everyone'", + "privacy_contacts": "TEXT NOT NULL DEFAULT 'everyone'", + "privacy_group_invite": "TEXT NOT NULL DEFAULT 'everyone'", # кто может добавлять в группы + "last_seen": "TEXT", + "tz_offset": "INTEGER NOT NULL DEFAULT 0", + "slow_seconds": "INTEGER NOT NULL DEFAULT 0", + "slow_scope": "TEXT NOT NULL DEFAULT 'nobody'", + "deleted": "INTEGER NOT NULL DEFAULT 0", # пометка вместо удаления + "banned": "INTEGER NOT NULL DEFAULT 0", + "system": "INTEGER NOT NULL DEFAULT 0", # @reports, @notes + "pinned_reactions": "TEXT NOT NULL DEFAULT ''", # JSON-список эмодзи + "wallpaper": "TEXT NOT NULL DEFAULT ''", # JSON {file, kind, mime} + "avatar_kind": "TEXT NOT NULL DEFAULT ''", # '' (legacy webp)|image|gif|video + }.items(): + if col not in users_cols: + conn.execute(f"ALTER TABLE users ADD COLUMN {col} {ddl}") + msg_cols = {row[1] for row in conn.execute("PRAGMA table_info(messages)")} + if "extra" not in msg_cols: + conn.execute("ALTER TABLE messages ADD COLUMN extra TEXT") # JSON-метаданные + if "group_id" not in msg_cols: + conn.execute("ALTER TABLE messages ADD COLUMN group_id INTEGER") # NULL = личное + conn.execute("CREATE INDEX IF NOT EXISTS idx_msg_group ON messages (group_id, id)") + if "reply_to" not in msg_cols: # id сообщения-родителя (ответ/комментарий) + conn.execute("ALTER TABLE messages ADD COLUMN reply_to INTEGER") + if "thread_root" not in msg_cols: # id поста-зеркала, под которым висит комментарий + conn.execute("ALTER TABLE messages ADD COLUMN thread_root INTEGER") + conn.execute("CREATE INDEX IF NOT EXISTS idx_msg_thread ON messages (thread_root)") + grp_cols = {row[1] for row in conn.execute("PRAGMA table_info(groups)")} + if "username" not in grp_cols: + conn.execute("ALTER TABLE groups ADD COLUMN username TEXT") + if "is_channel" not in grp_cols: + conn.execute("ALTER TABLE groups ADD COLUMN is_channel INTEGER NOT NULL DEFAULT 0") + if "linked_group_id" not in grp_cols: + # id привязанной для комментариев группы (только у каналов; NULL = нет) + conn.execute("ALTER TABLE groups ADD COLUMN linked_group_id INTEGER") + for col, ddl in { # настройки модерации группы/канала + "mod_enabled": "INTEGER NOT NULL DEFAULT 0", + "banned_words": "TEXT NOT NULL DEFAULT ''", # JSON-список слов/фраз + "whitelist_words": "TEXT NOT NULL DEFAULT ''", # JSON-список слов/фраз + "whitelist_mode": "TEXT NOT NULL DEFAULT 'any'", # any | all + "slow_seconds": "INTEGER NOT NULL DEFAULT 0", # медленный режим группы + }.items(): + if col not in grp_cols: + conn.execute(f"ALTER TABLE groups ADD COLUMN {col} {ddl}") + gm_cols = {row[1] for row in conn.execute("PRAGMA table_info(group_members)")} + if "is_admin" not in gm_cols: # назначенный админ группы/канала + conn.execute("ALTER TABLE group_members ADD COLUMN is_admin INTEGER NOT NULL DEFAULT 0") + if "muted" not in gm_cols: # мут: участник не может писать + conn.execute("ALTER TABLE group_members ADD COLUMN muted INTEGER NOT NULL DEFAULT 0") + # Юзернеймы групп уникальны без учёта регистра (у приватных их нет). + conn.execute("""CREATE UNIQUE INDEX IF NOT EXISTS idx_groups_username + ON groups (username COLLATE NOCASE) WHERE username IS NOT NULL""") + + # Миграция: шифруем тексты, сохранённые до появления шифрования. + for mid, content in conn.execute( + "SELECT id, content FROM messages WHERE kind = 'text'").fetchall(): + try: + fernet.decrypt(content.encode()) + except (InvalidToken, ValueError): + conn.execute("UPDATE messages SET content = ? WHERE id = ?", + (enc(content), mid)) + + # Системные аккаунты (войти в них нельзя — случайный пароль никому не известен). + for username, name in ((REPORTS_USERNAME, "Жалобы"), (NOTES_USERNAME, "Заметки"), + (SALGRAM_USERNAME, "SalGram"), (ADDGIF_USERNAME, "GIF")): + if not conn.execute("SELECT 1 FROM users WHERE username = ?", (username,)).fetchone(): + conn.execute( + "INSERT INTO users (display_name, username, password_hash, system) " + "VALUES (?, ?, ?, 1)", + (name, username, generate_password_hash(secrets.token_hex(32))), + ) + + # Официальный канал @salgram_news: создаётся один раз. Владелец — админ @salgr + # (если уже зарегистрирован), иначе официальный системный @salgram (админ всё + # равно может им управлять). Существующие реальные пользователи подписываются + # сразу, новые — при регистрации; отписаться можно как от обычного канала. + if not conn.execute("SELECT 1 FROM groups WHERE username = ? COLLATE NOCASE", + (NEWS_USERNAME,)).fetchone(): + owner = conn.execute( + "SELECT id FROM users WHERE username IN (?, ?) AND deleted = 0 " + "ORDER BY (username = ?) DESC LIMIT 1", + (ADMIN_USERNAME, SALGRAM_USERNAME, ADMIN_USERNAME)).fetchone() + if owner: + owner_id = owner[0] # соединение init_db без row_factory — строки-кортежи + cid = conn.execute( + "INSERT INTO groups (title, username, owner_id, public, invite_hash, is_channel) " + "VALUES (?, ?, ?, 1, ?, 1)", + ("Новости SalGram", NEWS_USERNAME, owner_id, + secrets.token_urlsafe(12))).lastrowid + conn.execute("INSERT OR IGNORE INTO group_members (group_id, user_id) " + "SELECT ?, id FROM users WHERE deleted = 0 AND system = 0", (cid,)) + conn.execute("INSERT OR IGNORE INTO group_members (group_id, user_id) " + "VALUES (?, ?)", (cid, owner_id)) + + +# ---------- Вспомогательные функции ---------- + +def error(message: str, status: int = 400): + return jsonify(error=message), status + + +def auth(fn): + """Декоратор: требует валидную сессию и передаёт user первым аргументом. + + Забаненный аккаунт работает в режиме «только чтение»: GET-запросы + разрешены (почитать переписку можно), любые изменения (POST/DELETE) — нет. + """ + @wraps(fn) + def wrapper(*args, **kwargs): + if not (user := current_user()): + return error("Не авторизован", 401) + if user["banned"] and request.method != "GET": + return error("Аккаунт заблокирован", 403) + return fn(user, *args, **kwargs) + return wrapper + + +def admin_only(fn): + """Декоратор поверх auth: только для администратора.""" + @wraps(fn) + @auth + def wrapper(user, *args, **kwargs): + if not user_is_admin(user): + return error("Нет прав", 403) + return fn(user, *args, **kwargs) + return wrapper + + +def current_user() -> sqlite3.Row | None: + if not (token := request.cookies.get("LoggedAccount")): + return None + # Забаненные не отсеиваются: им разрешён вход в режиме «только чтение». + return db().execute( + """SELECT u.* FROM sessions s JOIN users u ON u.id = s.user_id + WHERE s.token = ? AND u.deleted = 0""", + (token,), + ).fetchone() + + +def user_is_admin(u: sqlite3.Row | None) -> bool: + return bool(u) and not u["deleted"] and u["username"].lower() == ADMIN_USERNAME + + +def by_id(user_id: int) -> sqlite3.Row | None: + return db().execute("SELECT * FROM users WHERE id = ?", (user_id,)).fetchone() + + +def by_msg(message_id: int) -> sqlite3.Row | None: + return db().execute("SELECT * FROM messages WHERE id = ?", (message_id,)).fetchone() + + +def by_username(username: str) -> sqlite3.Row | None: + return db().execute( + "SELECT * FROM users WHERE username = ? AND deleted = 0", (username,) + ).fetchone() + + +def start_session(user_id: int, status: int = 200): + """Создаёт сессию и ставит cookie LoggedAccount с флагами безопасности.""" + token = secrets.token_hex(32) # 256 бит энтропии + conn = db() + conn.execute("INSERT INTO sessions (token, user_id) VALUES (?, ?)", (token, user_id)) + conn.execute("UPDATE users SET last_seen = datetime('now') WHERE id = ?", (user_id,)) + conn.commit() + + resp = jsonify(ok=True) + resp.set_cookie( + "LoggedAccount", token, + max_age=60 * 60 * 24 * 365, + httponly=True, # недоступна из JS — защита от XSS + samesite="Strict", # не отправляется с чужих сайтов — защита от CSRF + secure=request.is_secure, # на проде (HTTPS) — только по защищённому каналу + ) + return resp, status + + +def update_user(user_id: int, field: str, value): + conn = db() + conn.execute(f"UPDATE users SET {field} = ? WHERE id = ?", (value, user_id)) + conn.commit() + + +def touch_seen(user_id: int): + """Фиксирует время захода в сеть (хранится в UTC, ЧП +0).""" + conn = db() + conn.execute("UPDATE users SET last_seen = datetime('now') WHERE id = ?", (user_id,)) + conn.commit() + + +def soft_delete(user_id: int): + """Помечает аккаунт удалённым: строка остаётся, ID не переиспользуется.""" + conn = db() + conn.execute( + """UPDATE users SET deleted = 1, display_name = 'Удалённый аккаунт', + username = 'deleted_' || id || '_' || ?, password_hash = '!' WHERE id = ?""", + (secrets.token_hex(4), user_id), # юзернейм освобождается для регистрации + ) + conn.execute("DELETE FROM sessions WHERE user_id = ?", (user_id,)) + conn.execute("DELETE FROM contacts WHERE user_id = ? OR contact_id = ?", (user_id,) * 2) + conn.execute("DELETE FROM blocks WHERE blocker_id = ? OR blocked_id = ?", (user_id,) * 2) + conn.execute("DELETE FROM group_members WHERE user_id = ?", (user_id,)) + conn.execute("DELETE FROM reads WHERE reader_id = ? OR peer_id = ?", (user_id,) * 2) + conn.commit() + remove_avatar_files(user_id) + + +def is_contact(owner_id: int, other_id: int) -> bool: + return db().execute( + "SELECT 1 FROM contacts WHERE user_id = ? AND contact_id = ?", + (owner_id, other_id), + ).fetchone() is not None + + +def is_blocked(blocker_id: int, blocked_id: int) -> bool: + return db().execute( + "SELECT 1 FROM blocks WHERE blocker_id = ? AND blocked_id = ?", + (blocker_id, blocked_id), + ).fetchone() is not None + + +def avatar_file(user_id: int, kind: str = "image") -> str: + return os.path.join(AVATAR_DIR, f"{user_id}.{AVATAR_EXT.get(kind, 'webp')}") + + +def avatar_kind_of(u: sqlite3.Row) -> str: + """Тип аватарки: '' (legacy) трактуется как обычная картинка.""" + return u["avatar_kind"] or "image" + + +def user_avatar_path(u: sqlite3.Row) -> str: + return avatar_file(u["id"], avatar_kind_of(u)) + + +def has_avatar(u: sqlite3.Row) -> bool: + return os.path.exists(user_avatar_path(u)) + + +def remove_avatar_files(user_id: int): + """Удаляет аватарку любого типа (webp/gif/mp4) — при смене или удалении аккаунта.""" + for ext in set(AVATAR_EXT.values()): + path = os.path.join(AVATAR_DIR, f"{user_id}.{ext}") + if os.path.exists(path): + os.remove(path) + + +# ---------- Группы: вспомогательные функции ---------- + +def by_gid(gid: int) -> sqlite3.Row | None: + return db().execute("SELECT * FROM groups WHERE id = ?", (gid,)).fetchone() + + +def group_by_username(name: str) -> sqlite3.Row | None: + return db().execute( + "SELECT * FROM groups WHERE username = ? COLLATE NOCASE", (name,)).fetchone() + + +def group_avatar_file(gid: int) -> str: + return os.path.join(AVATAR_DIR, f"g{gid}.webp") + + +def group_member(gid: int, uid: int) -> bool: + return db().execute( + "SELECT 1 FROM group_members WHERE group_id = ? AND user_id = ?", (gid, uid) + ).fetchone() is not None + + +def group_member_count(gid: int) -> int: + return db().execute( + "SELECT COUNT(*) FROM group_members WHERE group_id = ?", (gid,)).fetchone()[0] + + +def push_group(gid: int, payload: dict): + """Рассылает событие всем участникам группы.""" + for r in db().execute("SELECT user_id FROM group_members WHERE group_id = ?", (gid,)): + push(r["user_id"], payload) + + +def can_add_to_group(adder: sqlite3.Row, target: sqlite3.Row) -> bool: + """Разрешает ли приватность target, чтобы adder добавил его в группу. + everyone — любой; contacts — только контакты target; nobody — никто. Админ — всегда.""" + if user_is_admin(adder): + return True + match target["privacy_group_invite"]: + case "everyone": + return True + case "contacts": + return is_contact(target["id"], adder["id"]) + case _: + return False + + +# ---------- Группы: модерация ---------- + +def group_role(gid: int, uid: int, grp: sqlite3.Row | None = None) -> str | None: + """Роль пользователя в группе: 'owner' | 'admin' | 'member' | None.""" + grp = grp or by_gid(gid) + if not grp: + return None + if grp["owner_id"] == uid: + return "owner" + row = db().execute("SELECT is_admin FROM group_members WHERE group_id=? AND user_id=?", + (gid, uid)).fetchone() + if not row: + return None + return "admin" if row["is_admin"] else "member" + + +def can_moderate(user: sqlite3.Row, grp: sqlite3.Row) -> bool: + """Может ли модерировать группу: владелец, назначенный админ или глоб. админ @salgr.""" + return user_is_admin(user) or group_role(grp["id"], user["id"], grp) in ("owner", "admin") + + +def group_is_banned(gid: int, uid: int) -> bool: + return db().execute("SELECT 1 FROM group_bans WHERE group_id=? AND user_id=?", + (gid, uid)).fetchone() is not None + + +def group_is_muted(gid: int, uid: int) -> bool: + row = db().execute("SELECT muted FROM group_members WHERE group_id=? AND user_id=?", + (gid, uid)).fetchone() + return bool(row and row["muted"]) + + +def group_slow_wait(grp: sqlite3.Row, uid: int) -> int: + """Сколько секу��д ждать по медленному режиму группы (0 — можно слать).""" + if grp["slow_seconds"] <= 0: + return 0 + elapsed = db().execute( + """SELECT CAST((julianday('now') - julianday(MAX(created_at))) * 86400 AS INT) + FROM messages WHERE group_id = ? AND sender_id = ?""", + (grp["id"], uid)).fetchone()[0] + return 0 if elapsed is None else max(0, grp["slow_seconds"] - elapsed) + + +def moderate_words(grp: sqlite3.Row, text: str) -> str | None: + """Проверка текста по бан-листу и белому списку (записи — слова или фразы, + сопоставление без учёта регистра). Возвращает причину нарушения или None.""" + if not grp["mod_enabled"]: + return None + low = text.lower() + for w in (json.loads(grp["banned_words"]) if grp["banned_words"] else []): + if w and w.lower() in low: + return f"запрещённое выражение «{w}»" + white = json.loads(grp["whitelist_words"]) if grp["whitelist_words"] else [] + if white: + if grp["whitelist_mode"] == "all": + for w in white: + if w and w.lower() not in low: + return f"в сообщении нет обязательного выражения «{w}»" + elif not any(w and w.lower() in low for w in white): + return "в сообщении нет ни одного разрешённого выражения" + return None + + +def moderation_guard(user: sqlite3.Row, gid: int, target_id: int, need_owner: bool = False): + """Общие проверки действий модерации над участником. Возвращает + (grp, target_member_row|None, error_response|None).""" + grp = by_gid(gid) + if not grp: + return None, None, error("Группа не найдена", 404) + if not can_moderate(user, grp): + return None, None, error("Нет прав модерации", 403) + is_owner = grp["owner_id"] == user["id"] or user_is_admin(user) + if need_owner and not is_owner: + return None, None, error("Это действие доступно только владельцу", 403) + if target_id == grp["owner_id"]: + return None, None, error("Нельзя применить к владельцу группы") + tm = db().execute("SELECT * FROM group_members WHERE group_id=? AND user_id=?", + (gid, target_id)).fetchone() + # Админа группы может модерировать только владелец/глобальный админ. + if tm and tm["is_admin"] and not is_owner: + return None, None, error("Модерировать админа может только владелец", 403) + return grp, tm, None + + +def group_payload(viewer: sqlite3.Row, grp: sqlite3.Row) -> dict: + member = group_member(grp["id"], viewer["id"]) + is_channel = bool(grp["is_channel"]) + p = { + "id": grp["id"], "title": grp["title"], "username": grp["username"], + "public": bool(grp["public"]), "is_channel": is_channel, + "is_member": member, "is_owner": grp["owner_id"] == viewer["id"], + "members": group_member_count(grp["id"]), + "has_avatar": os.path.exists(group_avatar_file(grp["id"])), + } + # У канала — привязанная для комментариев группа (если подключена). + if is_channel and grp["linked_group_id"] and (linked := by_gid(grp["linked_group_id"])): + p["linked_group_id"] = linked["id"] + p["linked_group_title"] = linked["title"] + if member: # пригласительная ссылка и список участников — только участникам + p["invite_hash"] = grp["invite_hash"] + p["my_role"] = group_role(grp["id"], viewer["id"], grp) + p["my_muted"] = group_is_muted(grp["id"], viewer["id"]) + # Список подписчиков канала виден только владельцу и админу (как в Telegram). + if not is_channel or grp["owner_id"] == viewer["id"] or user_is_admin(viewer): + p["member_list"] = [dict(r) for r in db().execute( + """SELECT u.id, u.username, u.display_name, gm.is_admin, gm.muted, + (u.id = :own) AS is_owner + FROM group_members gm JOIN users u ON u.id = gm.user_id + WHERE gm.group_id = :g + ORDER BY (u.id = :own) DESC, gm.is_admin DESC, + u.display_name COLLATE NOCASE""", + {"g": grp["id"], "own": grp["owner_id"]})] + # Настройки модерации — только модераторам (владелец, админы группы, @salgr). + moderator = member and can_moderate(viewer, grp) + p["can_moderate"] = bool(moderator) + if moderator: + p["moderation"] = { + "enabled": bool(grp["mod_enabled"]), + "banned_words": json.loads(grp["banned_words"]) if grp["banned_words"] else [], + "whitelist_words": json.loads(grp["whitelist_words"]) if grp["whitelist_words"] else [], + "whitelist_mode": grp["whitelist_mode"], + "slow_seconds": grp["slow_seconds"], + } + return p + + +def visible_to(viewer: sqlite3.Row | None, owner: sqlite3.Row, setting: str) -> bool: + """Приватность: everyone | contacts | nobody. Админ видит всё.""" + if viewer and (viewer["id"] == owner["id"] or user_is_admin(viewer)): + return True + match owner[setting]: + case "everyone": + return True + case "contacts": + return viewer is not None and is_contact(owner["id"], viewer["id"]) + case _: + return False + + +def profile_payload(viewer: sqlite3.Row, target: sqlite3.Row) -> dict: + """Чужой профиль глазами viewer — с учётом блокировок и приватности.""" + blocked_me = is_blocked(target["id"], viewer["id"]) + hidden = blocked_me and not user_is_admin(viewer) + payload = { + "id": target["id"], + "username": target["username"], + "display_name": target["display_name"], + "description": target["description"], + "has_avatar": has_avatar(target), + "avatar_kind": avatar_kind_of(target), + "is_admin": user_is_admin(target), + "is_system": bool(target["system"]), + "banned": bool(target["banned"]), + "is_contact": is_contact(viewer["id"], target["id"]), + "i_blocked": is_blocked(viewer["id"], target["id"]), + "blocked_me": blocked_me, + # Заблокированному не видно, в сети ли пользователь и когда заходил. + "online": not hidden and bool(connections.get(target["id"])), + "last_seen": None if hidden else target["last_seen"], + "contacts": None, + } + if not hidden and visible_to(viewer, target, "privacy_contacts"): + payload["contacts"] = [dict(r) for r in db().execute( + """SELECT u.username, u.display_name FROM contacts c + JOIN users u ON u.id = c.contact_id + WHERE c.user_id = ? ORDER BY u.display_name COLLATE NOCASE""", + (target["id"],), + )] + return payload + + +# ---------- Сообщения: общая логика ---------- + +# Форматирование текста: типы пометок и их допустимые значения. Пометка — это +# диапазон символов {s, e, t, v?}; рендер на клиенте строит из них стили. +FMT_TYPES = {"b", "i", "u", "s", "code", "quote", "size", "color"} +FMT_SIZES = {"huge", "big", "normal", "small", "tiny"} +FMT_COLOR_RE = re.compile(r"^#[0-9a-fA-F]{6}$") +FMT_MARKS_MAX = 200 # защита: не храним бесконечный список пометок + + +def clean_fmt(marks, text_len: int): + """Проверяет и нормализует список пометок форматирования от клиента. + Возвращает безопасный список (или None, если форматирования нет).""" + if not isinstance(marks, list): + return None + out = [] + for mk in marks[:FMT_MARKS_MAX]: + if not isinstance(mk, dict): + continue + t = mk.get("t") + if t not in FMT_TYPES: + continue + try: + s, e = int(mk.get("s")), int(mk.get("e")) + except (TypeError, ValueError): + continue + s, e = max(0, s), min(text_len, e) + if s >= e: + continue + clean = {"s": s, "e": e, "t": t} + if t == "size": + if mk.get("v") not in FMT_SIZES or mk.get("v") == "normal": + continue + clean["v"] = mk["v"] + elif t == "color": + if not (isinstance(mk.get("v"), str) and FMT_COLOR_RE.match(mk["v"])): + continue + clean["v"] = mk["v"] + out.append(clean) + return out or None + + +def reply_preview(reply_to: int | None) -> dict | None: + """Короткая карточка сообщения-родителя для отрисовки цитаты ответа.""" + if not reply_to: + return None + parent = db().execute("SELECT * FROM messages WHERE id = ?", (reply_to,)).fetchone() + if not parent: + return None + sender = by_id(parent["sender_id"]) + if parent["kind"] == "text": + text = dec(parent["content"])[:120] + else: + pextra = json.loads(parent["extra"]) if parent["extra"] else None + text = _notify_preview(parent["kind"], "", pextra) + return {"id": parent["id"], "from": parent["sender_id"], + "from_name": sender["display_name"] if sender else "?", "text": text} + + +def message_payload(m: sqlite3.Row) -> dict: + extra = json.loads(m["extra"]) if m["extra"] else None + p = { + "id": m["id"], "from": m["sender_id"], "to": m["recipient_id"], + "kind": m["kind"], + "content": dec(m["content"]) if m["kind"] == "text" else m["content"], + "created_at": m["created_at"], + "extra": extra, + "reactions": message_reactions(m["id"]), + "reply_to": m["reply_to"], + "reply_to_preview": reply_preview(m["reply_to"]), + "thread_root": m["thread_root"], + } + sender = by_id(m["sender_id"]) # имя автора — для подписи и уведомлений + p["from_name"] = sender["display_name"] if sender else "?" + if m["group_id"]: # в группе подписывается автор сообщения (+ аватарка слева) + p["group"] = m["group_id"] + if sender: + p["from_username"] = sender["username"] + p["from_avatar_kind"] = sender["avatar_kind"] + # Число комментариев: у поста канала (extra.disc → зеркало) или у самого + # зеркала поста в группе-обсуждении (extra.channel_post → его id и есть корень). + root = extra.get("disc") if extra and extra.get("disc") else ( + m["id"] if extra and extra.get("channel_post") else None) + if root: + p["comments"] = db().execute( + "SELECT COUNT(*) FROM messages WHERE thread_root = ?", (root,)).fetchone()[0] + return p + + +def message_reactions(message_id: int) -> list[dict]: + return [{"user_id": r["user_id"], "emoji": r["emoji"]} for r in db().execute( + "SELECT user_id, emoji FROM reactions WHERE message_id = ?", (message_id,))] + + +def remove_attachment_file(kind: str, fname: str): + """Удаляет файл вложения (фото, файл или GIF), если на него больше не ссылается + ни одно сообщение (вложение может разделяться сообщениями при рассылке). + GIF из общей библиотеки (@addgif) не удаляется — он принадлежит библиотеке.""" + if kind not in ("photo", "file", "gif", "audio"): + return + if db().execute("SELECT 1 FROM messages WHERE kind = ? AND content = ?", + (kind, fname)).fetchone(): + return + if kind == "gif": + if db().execute("SELECT 1 FROM gifs WHERE file = ?", (fname,)).fetchone(): + return + path = os.path.join(GIF_DIR, fname) + elif kind == "audio": + path = os.path.join(AUDIO_DIR, fname) + else: + path = os.path.join(PHOTO_DIR if kind == "photo" else FILE_DIR, fname) + if os.path.exists(path): + os.remove(path) + + +def _xml_text(xml: bytes) -> str: + """Грубое извлечение текста из XML: теги заменяются пробелами.""" + text = re.sub(rb"<[^>]*>", b" ", xml).decode("utf-8", "ignore") + return re.sub(r"\s+", " ", text).strip() + + +def file_preview(name: str, data: bytes) -> str | None: + """Текстовый предпросмотр: документы (docx/odt), презентации (pptx/odp), + таблицы (xlsx/ods/csv) и простой текст. Office-форматы — это zip с XML, + поэтому текст достаётся без сторонних библиотек.""" + ext = name.rsplit(".", 1)[-1].lower() if "." in name else "" + try: + if ext in ("txt", "csv", "md", "log", "json"): + return data[:8192].decode("utf-8", "ignore")[:PREVIEW_CHARS].strip() or None + if ext in ("docx", "pptx", "xlsx", "odt", "odp", "ods"): + with zipfile.ZipFile(io.BytesIO(data)) as z: + names = set(z.namelist()) + if ext == "docx": + parts = [z.read("word/document.xml")] + elif ext == "pptx": + slides = sorted(n for n in names + if re.fullmatch(r"ppt/slides/slide\d+\.xml", n)) + parts = [z.read(n) for n in slides[:3]] + elif ext == "xlsx": + parts = ([z.read("xl/sharedStrings.xml")] + if "xl/sharedStrings.xml" in names else []) + else: # OpenDocument: odt / odp / ods + parts = [z.read("content.xml")] if "content.xml" in names else [] + text = " ".join(filter(None, (_xml_text(p) for p in parts))) + return text[:PREVIEW_CHARS].strip() or None + except Exception: + pass # предпросмотр — не критичен: битый файл просто остаётся без него + return None + + +def push(user_id: int, payload: dict): + """Рассылает событие во все WebSocket-соединения пользователя.""" + data = json.dumps(payload, ensure_ascii=False) + for conn in list(connections.get(user_id, ())): + try: + conn.send(data) + except Exception: + pass + + +def store_message(sender_id: int, recipient_id: int, kind: str, content: str, + extra: dict | None = None, reply_to: int | None = None) -> dict: + """Шифрует, сохраняет и рассылает сообщение по WS (без проверок доступа).""" + conn = db() + cur = conn.execute( + "INSERT INTO messages (sender_id, recipient_id, kind, content, extra, reply_to) " + "VALUES (?, ?, ?, ?, ?, ?)", + (sender_id, recipient_id, kind, enc(content) if kind == "text" else content, + json.dumps(extra, ensure_ascii=False) if extra else None, reply_to), + ) + conn.commit() + row = conn.execute("SELECT * FROM messages WHERE id = ?", (cur.lastrowid,)).fetchone() + payload = {"type": "message", "message": message_payload(row)} + push(sender_id, payload) + push(recipient_id, payload) + if recipient_id != sender_id: # фоновое уведомление получателю + sender = by_id(sender_id) + push_notify_user(recipient_id, sender["display_name"] if sender else "SalGram", + _notify_preview(kind, content, extra), {"peer": sender_id}) + return payload + + +def store_group_message(sender_id: int, gid: int, kind: str, content: str, + extra: dict | None = None, reply_to: int | None = None, + thread_root: int | None = None) -> dict: + """Сохраняет групповое сообщение и рассылает его всем участникам. + Пост канала дополнительно зеркалится в привязанную группу-обсуждение.""" + conn = db() + cur = conn.execute( + "INSERT INTO messages (sender_id, recipient_id, group_id, kind, content, extra, " + "reply_to, thread_root) VALUES (?, 0, ?, ?, ?, ?, ?, ?)", + (sender_id, gid, kind, enc(content) if kind == "text" else content, + json.dumps(extra, ensure_ascii=False) if extra else None, reply_to, thread_root), + ) + conn.commit() + mid = cur.lastrowid + grp = by_gid(gid) + if grp and grp["is_channel"]: # пост канала → дублируем в обсуждение + mirror_post_to_discussion(grp, mid, kind) + row = conn.execute("SELECT * FROM messages WHERE id = ?", (mid,)).fetchone() + payload = {"type": "message", "message": message_payload(row)} + push_group(gid, payload) + sender = by_id(sender_id) # фоновое уведомление участникам, кроме автора + who = sender["display_name"] if sender else "?" + title = f"{who} • {grp['title']}" if grp else who + push_notify_group(gid, sender_id, title, _notify_preview(kind, content, extra), + {"group": gid}) + return payload + + +def mirror_post_to_discussion(grp: sqlite3.Row, post_id: int, kind: str): + """Дублирует пост канала в привязанную группу-обсуждение как корень треда + комментариев и записывает id зеркала в extra исходного поста (extra.disc).""" + if not grp["linked_group_id"] or kind not in ("text", "photo", "gif", "file"): + return + linked = by_gid(grp["linked_group_id"]) + if not linked: + return + conn = db() + post = conn.execute("SELECT * FROM messages WHERE id = ?", (post_id,)).fetchone() + # Зеркало хранит то же содержимое (content уже зашифрован/имя файла — берём как есть). + mextra = json.loads(post["extra"]) if post["extra"] else {} + mextra.pop("pinned", None) + mextra["channel_post"] = {"cid": grp["id"], "mid": post_id, "title": grp["title"]} + cur = conn.execute( + "INSERT INTO messages (sender_id, recipient_id, group_id, kind, content, extra) " + "VALUES (?, 0, ?, ?, ?, ?)", + (post["sender_id"], linked["id"], post["kind"], post["content"], + json.dumps(mextra, ensure_ascii=False)), + ) + mirror_id = cur.lastrowid + # Связываем пост с его зеркалом: по нему считаем комментарии и открываем тред. + post_extra = json.loads(post["extra"]) if post["extra"] else {} + post_extra["disc"] = mirror_id + post_extra["disc_group"] = linked["id"] + conn.execute("UPDATE messages SET extra = ? WHERE id = ?", + (json.dumps(post_extra, ensure_ascii=False), post_id)) + conn.commit() + row = conn.execute("SELECT * FROM messages WHERE id = ?", (mirror_id,)).fetchone() + push_group(linked["id"], {"type": "message", "message": message_payload(row)}) + + +def deliver_group_message(sender: sqlite3.Row, gid: int, kind: str, content: str, + extra: dict | None = None, reply_to: int | None = None, + thread_root: int | None = None): + """Проверка членства + сохранение + рассылка группового сообщения. + В канал писать может только владелец (и админ); остальные — подписчики.""" + grp = by_gid(gid) + if not grp or not group_member(gid, sender["id"]): + return None, "Вы не состоите в этой группе" + if grp["is_channel"] and grp["owner_id"] != sender["id"] and not user_is_admin(sender): + return None, "Публиковать в канал может только его владелец" + # Модерация (мут, медленный режим, фильтры слов) — кроме модераторов группы. + if not can_moderate(sender, grp): + if group_is_muted(gid, sender["id"]): + return None, "Вы не можете писать в этой группе (мут)" + if (wait := group_slow_wait(grp, sender["id"])) > 0: + return None, f"Медленный режим: подождите {wait} сек." + if kind == "text" and (reason := moderate_words(grp, content)): + return None, f"Сообщение удалено модерацией: {reason}" + return store_group_message(sender["id"], gid, kind, content, extra, + reply_to, thread_root), None + + +def message_visible_to(m: sqlite3.Row, uid: int) -> bool: + """Видит ли пользователь сообщение: участник группы либо одной из сторон.""" + if m["group_id"]: + return group_member(m["group_id"], uid) + return uid in (m["sender_id"], m["recipient_id"]) + + +def resolve_thread_root(parent: sqlite3.Row) -> int | None: + """Корень треда комментариев для ответа на parent: сам пост-зеркало (если parent — + зеркало поста канала) или уже известный thread_root родителя; иначе не комментарий.""" + if parent["thread_root"]: + return parent["thread_root"] + pextra = json.loads(parent["extra"]) if parent["extra"] else None + if pextra and pextra.get("channel_post"): + return parent["id"] + return None + + +def push_message_event(m: sqlite3.Row, payload: dict): + """Шлёт событие всем, кому видно сообщение: группе или паре собеседников.""" + if m["group_id"]: + push_group(m["group_id"], payload) + else: + push(m["sender_id"], payload) + push(m["recipient_id"], payload) + + +def slow_wait(sender: sqlite3.Row, recipient: sqlite3.Row) -> int: + """Сколько секунд ждать по медленному режиму получателя (0 — можно слать).""" + if recipient["slow_seconds"] <= 0 or recipient["slow_scope"] == "nobody": + return 0 + if recipient["slow_scope"] == "noncontacts" and is_contact(recipient["id"], sender["id"]): + return 0 + elapsed = db().execute( + """SELECT CAST((julianday('now') - julianday(MAX(created_at))) * 86400 AS INT) + FROM messages WHERE sender_id = ? AND recipient_id = ?""", + (sender["id"], recipient["id"]), + ).fetchone()[0] + return 0 if elapsed is None else max(0, recipient["slow_seconds"] - elapsed) + + +def forward_to_admin(sender: sqlite3.Row, recipient: sqlite3.Row, kind: str, text: str): + """Пока жалоба на пару открыта — копия каждого сообщения уходит админу от @reports.""" + rep = db().execute( + """SELECT id FROM reports WHERE status = 'open' + AND ((reporter_id = :a AND target_id = :b) OR (reporter_id = :b AND target_id = :a))""", + {"a": sender["id"], "b": recipient["id"]}, + ).fetchone() + if not rep: + return + admin, reports_acc = by_username(ADMIN_USERNAME), by_username(REPORTS_USERNAME) + if admin and reports_acc: + body = text if kind == "text" else \ + {"photo": "[Фото]", "gif": "[GIF]", "call": "[Звонок]", + "audio": "[Аудио]"}.get(kind, "[Файл]") + store_message(reports_acc["id"], admin["id"], + "text", f"Жалоба #{rep['id']} — {sender['display_name']} " + f"(@{sender['username']}): {body}") + + +def deliver_message(sender: sqlite3.Row, recipient_id: int, kind: str, content: str, + extra: dict | None = None, reply_to: int | None = None): + """Все проверки + сохранение + рассылка. Возвращает (payload, None) | (None, ошибка).""" + target = by_id(recipient_id) + if not target or target["deleted"]: + return None, "Пользователь не найден" + if target["banned"]: + return None, "Аккаунт получателя заблокирован" + if target["id"] == sender["id"]: + return None, "Нельзя писать самому себе (для заметок есть @notes)" + # В @salgram пишет только админ — его сообщение рассылается всем пользователям. + is_broadcast = (target["system"] and target["username"] == SALGRAM_USERNAME + and user_is_admin(sender)) + if target["system"] and target["username"] != NOTES_USERNAME and not is_broadcast: + return None, "Этому аккаунту нельзя писать" + if is_blocked(target["id"], sender["id"]): + return None, "Пользователь заблокировал вас" + if is_blocked(sender["id"], target["id"]): + return None, "Вы заблокировали пользователя" + if (wait := slow_wait(sender, target)) > 0: + return None, f"Медленный режим: подождите {wait} сек." + # Ответ возможен только на существующее сообщение этой же переписки. + if reply_to: + parent = by_msg(reply_to) + if not (parent and parent["group_id"] is None and not parent["thread_root"] + and {parent["sender_id"], parent["recipient_id"]} == {sender["id"], target["id"]}): + reply_to = None + + payload = store_message(sender["id"], recipient_id, kind, content, extra, reply_to) + if is_broadcast: + broadcast_from_salgram(target, sender, kind, content, extra) + forward_to_admin(sender, target, kind, content) + return payload, None + + +def broadcast_from_salgram(salgram: sqlite3.Row, admin: sqlite3.Row, kind: str, + content: str, extra: dict | None = None): + """Рассылка: сообщение админа в @salgram уходит всем пользователям от @salgram.""" + rows = db().execute( + "SELECT id FROM users WHERE deleted = 0 AND system = 0 AND banned = 0 AND id <> ?", + (admin["id"],), + ).fetchall() + for r in rows: + store_message(salgram["id"], r["id"], kind, content, extra) + + +# ---------- Аутентификация ---------- + +@app.get("/") +def index(): + return app.send_static_file("index.html") + + +@app.get("/.well-known/assetlinks.json") +def assetlinks(): + """Digital Asset Links для Android-приложения (TWA/APK): связывает домен с + подписью apk, чтобы приложение открывалось без адресной строки. Содержимое + кладётся в assetlinks.json (рядом с кодом или в DATA_DIR) — отпечаток ключа + выдаёт PWABuilder/Bubblewrap при сборке APK.""" + for path in (os.path.join(DATA_DIR, "assetlinks.json"), + os.path.join(os.path.dirname(os.path.abspath(__file__)), "assetlinks.json")): + if os.path.exists(path): + return send_file(path, mimetype="application/json") + return jsonify([]), 404 + + +@app.post("/api/register") +def register(): + d = request.get_json(silent=True) or {} + display_name = str(d.get("display_name", "")).strip() + username = str(d.get("username", "")).strip() + password = str(d.get("password", "")) + + if not 1 <= len(display_name) <= 50: + return error("Ник должен быть от 1 до 50 символов") + if not USERNAME_RE.fullmatch(username): + return error("Юзернейм: 3–32 символа, только латиница, цифры и _") + if len(password) < 8: + return error("Пароль должен быть не короче 8 символов") + if not d.get("tos"): + return error("Необходимо принять условия использования") + if group_by_username(username): # общее пространство имён с группами + return error("Этот юзернейм уже занят") + + try: + cur = db().execute( + "INSERT INTO users (display_name, username, password_hash) VALUES (?, ?, ?)", + (display_name, username, generate_password_hash(password)), + ) + except sqlite3.IntegrityError: + return error("Этот юзернейм уже занят") + # Автоподписка на официальный канал @salgram_news (отписаться можно потом). + if news := group_by_username(NEWS_USERNAME): + conn = db() + conn.execute("INSERT OR IGNORE INTO group_members (group_id, user_id) VALUES (?, ?)", + (news["id"], cur.lastrowid)) + conn.commit() + return start_session(cur.lastrowid, 201) + + +@app.post("/api/login") +def login(): + d = request.get_json(silent=True) or {} + user = db().execute( + "SELECT * FROM users WHERE username = ? AND deleted = 0", + (str(d.get("username", "")).strip(),), + ).fetchone() + + if not user or not check_password_hash(user["password_hash"], str(d.get("password", ""))): + return error("Неверный юзернейм или пароль", 401) + if user["system"]: + return error("В этот аккаунт нельзя войти", 403) + # Забаненный может войти — но только читать (режим «только чтение»). + return start_session(user["id"]) + + +@app.post("/api/logout") +def logout(): + if token := request.cookies.get("LoggedAccount"): + conn = db() + conn.execute("DELETE FROM sessions WHERE token = ?", (token,)) + conn.commit() + resp = jsonify(ok=True) + resp.delete_cookie("LoggedAccount") + return resp + + +# ---------- Собственный профиль и настройки ---------- + +@app.get("/api/me") +@auth +def me(user): + return jsonify( + id=user["id"], + display_name=user["display_name"], + username=user["username"], + description=user["description"], + has_avatar=has_avatar(user), + avatar_kind=avatar_kind_of(user), + is_admin=user_is_admin(user), + banned=bool(user["banned"]), + privacy={"avatar": user["privacy_avatar"], "contacts": user["privacy_contacts"], + "group_invite": user["privacy_group_invite"]}, + tz_offset=user["tz_offset"], + slow_seconds=user["slow_seconds"], + slow_scope=user["slow_scope"], + pinned_reactions=(json.loads(user["pinned_reactions"]) + if user["pinned_reactions"] else DEFAULT_PINNED), + wallpaper=({"kind": json.loads(user["wallpaper"])["kind"]} + if user["wallpaper"] else None), + ) + + +@app.post("/api/me/reactions") +@auth +def set_pinned_reactions(user): + """Закреплённые реакции (показываются первыми, максимум 5).""" + pinned = (request.get_json(silent=True) or {}).get("pinned", []) + if (not isinstance(pinned, list) or len(pinned) > 5 + or len(set(pinned)) != len(pinned) + or any(e not in REACTION_OK for e in pinned)): + return error("Недопустимый список реакций") + update_user(user["id"], "pinned_reactions", json.dumps(pinned, ensure_ascii=False)) + return jsonify(ok=True) + + +@app.post("/api/me/display_name") +@auth +def set_display_name(user): + name = str((request.get_json(silent=True) or {}).get("display_name", "")).strip() + if not 1 <= len(name) <= 50: + return error("Ник должен быть от 1 до 50 символов") + update_user(user["id"], "display_name", name) + return jsonify(ok=True) + + +@app.post("/api/me/username") +@auth +def set_username(user): + username = str((request.get_json(silent=True) or {}).get("username", "")).strip() + if not USERNAME_RE.fullmatch(username): + return error("Юзернейм: 3–32 символа, только латиница, цифры и _") + if group_by_username(username): # общее пространство имён с группами + return error("Этот юзернейм уже занят") + try: + update_user(user["id"], "username", username) + except sqlite3.IntegrityError: + return error("Этот юзернейм уже занят") + return jsonify(ok=True) + + +@app.post("/api/me/password") +@auth +def set_password(user): + d = request.get_json(silent=True) or {} + if not check_password_hash(user["password_hash"], str(d.get("old_password", ""))): + return error("Старый пароль неверен", 401) + new = str(d.get("new_password", "")) + if len(new) < 8: + return error("Новый пароль должен быть не короче 8 символов") + + conn = db() + conn.execute("UPDATE users SET password_hash = ? WHERE id = ?", + (generate_password_hash(new), user["id"])) + conn.execute("DELETE FROM sessions WHERE user_id = ? AND token <> ?", + (user["id"], request.cookies["LoggedAccount"])) + conn.commit() + return jsonify(ok=True) + + +@app.post("/api/me/description") +@auth +def set_description(user): + text = str((request.get_json(silent=True) or {}).get("description", "")).strip() + if len(text) > 512: + return error("Описание не длиннее 512 символов") + update_user(user["id"], "description", text) + return jsonify(ok=True) + + +@app.post("/api/me/privacy") +@auth +def set_privacy(user): + d = request.get_json(silent=True) or {} + for key, col in (("avatar", "privacy_avatar"), ("contacts", "privacy_contacts"), + ("group_invite", "privacy_group_invite")): + if key in d: + if d[key] not in PRIVACY_VALUES: + return error("Недопустимое значение приватности") + update_user(user["id"], col, d[key]) + return jsonify(ok=True) + + +@app.post("/api/me/settings") +@auth +def set_settings(user): + """Часовой пояс («Прочее») и медленный режим («Чаты»).""" + d = request.get_json(silent=True) or {} + if "tz_offset" in d: + tz = int(d["tz_offset"]) + if not -12 <= tz <= 14: + return error("Часовой пояс: от -12 до +14") + update_user(user["id"], "tz_offset", tz) + if "slow_seconds" in d: + sec = int(d["slow_seconds"]) + if not 0 <= sec <= 3600: + return error("Интервал: от 0 до 3600 секунд") + update_user(user["id"], "slow_seconds", sec) + if "slow_scope" in d: + if d["slow_scope"] not in SLOW_SCOPES: + return error("Недопустимая область медленного режима") + update_user(user["id"], "slow_scope", d["slow_scope"]) + return jsonify(ok=True) + + +@app.post("/api/me/delete") +@auth +def delete_account(user): + d = request.get_json(silent=True) or {} + if not check_password_hash(user["password_hash"], str(d.get("password", ""))): + return error("Неверный пароль", 401) + soft_delete(user["id"]) + resp = jsonify(ok=True) + resp.delete_cookie("LoggedAccount") + return resp + + +# ---------- Аватарка ---------- + +@app.post("/api/me/avatar") +@auth +def upload_avatar(user): + """Аватарка: обычная картинка (центр-кроп до 1024², WebP), анимированный GIF + (сохраняется как есть) или видео ≤10 с (как есть; длительность проверяет браузер).""" + if not (file := request.files.get("avatar")): + return error("Файл не выбран") + if (file.mimetype or "").startswith("video/"): + data = file.read() + if len(data) > AVATAR_VIDEO_MAX: + return error("Видео-аватарка больше 10 МБ") + kind = "video" + else: + raw = file.read() + try: + src = Image.open(io.BytesIO(raw)) # формат/анимацию проверяем до exif_transpose + if src.format == "GIF" and getattr(src, "is_animated", False): + if len(raw) > GIF_MAX: + return error("GIF больше 15 МБ") + data, kind = raw, "gif" # сохраняем анимацию как есть + else: + img = ImageOps.exif_transpose(src) + side = min(AVATAR_SIDE, *img.size) + img = ImageOps.fit(img.convert("RGBA"), (side, side)) + buf = io.BytesIO() + img.save(buf, "WEBP", quality=85) + data, kind = buf.getvalue(), "image" + except Exception: + return error("Не удалось обработать изображение") + remove_avatar_files(user["id"]) # убрать прежнюю аватарку любого типа + with open(avatar_file(user["id"], kind), "wb") as f: + f.write(data) + update_user(user["id"], "avatar_kind", kind) + return jsonify(ok=True, avatar_kind=kind) + + +AVATAR_MIME = {"image": "image/webp", "gif": "image/gif", "video": "video/mp4"} + + +@app.get("/api/avatar/") +def get_avatar(username): + target = by_username(username) + if not target or not has_avatar(target): + return error("Аватарка не найдена", 404) + if not visible_to(current_user(), target, "privacy_avatar"): + return error("Нет доступа", 403) + return send_file(user_avatar_path(target), + mimetype=AVATAR_MIME[avatar_kind_of(target)], max_age=0) + + +# ---------- Обои чата ---------- + +def remove_wallpaper_file(user: sqlite3.Row): + if user["wallpaper"]: + path = os.path.join(WALLPAPER_DIR, json.loads(user["wallpaper"])["file"]) + if os.path.exists(path): + os.remove(path) + + +@app.post("/api/me/wallpaper") +@auth +def upload_wallpaper(user): + """Обои переписки: фото (сжимается до 1600px, WebP) или видео до 25 МБ. + Хранятся зашифрованными, как и все вложения.""" + if not (file := request.files.get("wallpaper")): + return error("Файл не выбран") + mime = file.mimetype or "" + if mime.startswith("image/"): + raw = file.read() + try: + src = Image.open(io.BytesIO(raw)) # формат/анимацию проверяем до exif_transpose + if src.format == "GIF" and getattr(src, "is_animated", False): + if len(raw) > GIF_MAX: + return error("GIF больше 15 МБ") + data, kind, mime = raw, "photo", "image/gif" # анимация фоном (CSS) + else: + img = ImageOps.exif_transpose(src) + img.thumbnail((PHOTO_SIDE, PHOTO_SIDE)) + buf = io.BytesIO() + img.convert("RGB").save(buf, "WEBP", quality=82) + data, kind, mime = buf.getvalue(), "photo", "image/webp" + except Exception: + return error("Не удалось обработать изображение") + elif mime.startswith("video/"): + data = file.read() + if len(data) > FILE_MAX: + return error("Видео больше 25 МБ") + kind = "video" + else: + return error("Обои — это фото или видео") + + remove_wallpaper_file(user) + fname = secrets.token_hex(16) + ".bin" + with open(os.path.join(WALLPAPER_DIR, fname), "wb") as f: + f.write(fernet.encrypt(data)) # шифрование на диске + update_user(user["id"], "wallpaper", + json.dumps({"file": fname, "kind": kind, "mime": mime})) + return jsonify(wallpaper={"kind": kind}) + + +@app.delete("/api/me/wallpaper") +@auth +def delete_wallpaper(user): + remove_wallpaper_file(user) + update_user(user["id"], "wallpaper", "") + return jsonify(ok=True) + + +@app.get("/api/me/wallpaper") +@auth +def get_wallpaper(user): + """Обои видны только их владельцу; расшифровка на лету.""" + if not user["wallpaper"]: + return error("Обои не установлены", 404) + w = json.loads(user["wallpaper"]) + with open(os.path.join(WALLPAPER_DIR, w["file"]), "rb") as f: + raw = fernet.decrypt(f.read()) + return send_file(io.BytesIO(raw), mimetype=w["mime"], max_age=0) + + +# ---------- Контакты ---------- + +@app.get("/api/me/contacts") +@auth +def contacts_list(user): + rows = db().execute( + """SELECT u.username, u.display_name FROM contacts c + JOIN users u ON u.id = c.contact_id + WHERE c.user_id = ? ORDER BY u.display_name COLLATE NOCASE""", + (user["id"],), + ).fetchall() + return jsonify([dict(r) for r in rows]) + + +@app.post("/api/me/contacts") +@auth +def contacts_add(user): + username = str((request.get_json(silent=True) or {}).get("username", "")).strip() + target = by_username(username) + + if not target: + return error("Пользователь не найден", 404) + if target["id"] == user["id"]: + return error("Нельзя добавить самого себя") + try: + conn = db() + conn.execute("INSERT INTO contacts (user_id, contact_id) VALUES (?, ?)", + (user["id"], target["id"])) + conn.commit() + except sqlite3.IntegrityError: + return error("Уже в списке контактов") + return jsonify(username=target["username"], display_name=target["display_name"]) + + +@app.delete("/api/me/contacts/") +@auth +def contacts_remove(user, username): + conn = db() + conn.execute( + """DELETE FROM contacts WHERE user_id = ? + AND contact_id = (SELECT id FROM users WHERE username = ?)""", + (user["id"], username), + ) + conn.commit() + return jsonify(ok=True) + + +# ---------- Чужие профили, блокировки, жалобы ---------- + +@app.get("/api/users/") +@auth +def user_profile(viewer, username): + target = by_username(username) + if not target: + return error("Пользователь не найден", 404) + return jsonify(profile_payload(viewer, target)) + + +@app.get("/api/peer/") +@auth +def peer_profile(viewer, peer_id): + target = by_id(peer_id) + if not target or target["deleted"]: + return jsonify(id=peer_id, deleted=True) + return jsonify(profile_payload(viewer, target)) + + +@app.post("/api/users//block") +@auth +def block_user(user, target_id): + target = by_id(target_id) + if not target or target["deleted"]: + return error("Пользователь не найден", 404) + if target_id == user["id"]: + return error("Нельзя заблокировать себя") + if user_is_admin(target): + return error("Нельзя заблокировать администратора") + if target["system"]: + return error("Нельзя заблокировать системный аккаунт") + conn = db() + conn.execute("INSERT OR IGNORE INTO blocks (blocker_id, blocked_id) VALUES (?, ?)", + (user["id"], target_id)) + conn.commit() + return jsonify(ok=True) + + +@app.post("/api/users//unblock") +@auth +def unblock_user(user, target_id): + conn = db() + conn.execute("DELETE FROM blocks WHERE blocker_id = ? AND blocked_id = ?", + (user["id"], target_id)) + conn.commit() + return jsonify(ok=True) + + +@app.post("/api/users//report") +@auth +def report_user(user, target_id): + """Жалоба: переписка с пользователем пересылается админу от имени @reports.""" + target = by_id(target_id) + if not target or target["deleted"]: + return error("Пользователь не найден", 404) + if target["system"] or user_is_admin(target): + return error("На этот аккаунт нельзя пожаловаться") + admin, reports_acc = by_username(ADMIN_USERNAME), by_username(REPORTS_USERNAME) + if not admin or not reports_acc: + return error("Администратор недоступен") + + conn = db() + rid = conn.execute("INSERT INTO reports (reporter_id, target_id) VALUES (?, ?)", + (user["id"], target_id)).lastrowid + conn.commit() + + # Собираем переписку (последние 200 сообщений) в текстовые блоки. + msgs = db().execute( + """SELECT * FROM messages + WHERE (sender_id = :a AND recipient_id = :b) + OR (sender_id = :b AND recipient_id = :a) + ORDER BY id DESC LIMIT 200""", + {"a": user["id"], "b": target_id}, + ).fetchall() + names = {user["id"]: f"{user['display_name']} (@{user['username']})", + target_id: f"{target['display_name']} (@{target['username']})"} + lines = [f"[{m['created_at']}] {names[m['sender_id']]}: " + f"{dec(m['content']) if m['kind'] == 'text' else ('[Фото]' if m['kind'] == 'photo' else '[Файл]')}" + for m in reversed(msgs)] or ["(переписка пуста)"] + + chunk = "" + for line in lines: # режем на сообщения <= 3500 символов + if chunk and len(chunk) + len(line) > 3500: + store_message(reports_acc["id"], admin["id"], "text", chunk) + chunk = "" + chunk += ("\n" if chunk else "") + line + store_message(reports_acc["id"], admin["id"], "text", chunk) + + store_message( + reports_acc["id"], admin["id"], "text", + f"Жалоба #{rid}: @{user['username']} пожаловался(ась) на @{target['username']}. " + f"Переписка выше; новые сообщения пары будут пересылаться, пока жалоба открыта.", + extra={"report": rid}, # по этой метке фронтенд рисует кнопки Забанить/Отклонить + ) + return jsonify(ok=True) + + +# ---------- Администрирование ---------- + +@app.post("/api/admin/users//delete") +@admin_only +def admin_delete_user(_admin, target_id): + target = by_id(target_id) + if not target or target["deleted"]: + return error("Пользователь не найден", 404) + if target["system"] or user_is_admin(target): + return error("Этот аккаунт нельзя удалить") + soft_delete(target_id) + return jsonify(ok=True) + + +@app.post("/api/admin/reports//") +@admin_only +def admin_report_action(admin, rid, action): + """Кнопки в чате с @reports: ban — забанить нарушителя, dismiss — отклонить.""" + if action not in ("ban", "dismiss"): + return error("Не найдено", 404) + rep = db().execute("SELECT * FROM reports WHERE id = ?", (rid,)).fetchone() + if not rep: + return error("Жалоба не найдена", 404) + if rep["status"] != "open": + return error("Жалоба уже рассмотрена") + + conn = db() + if action == "ban": + target = by_id(rep["target_id"]) + if target and not target["deleted"]: + conn.execute("UPDATE users SET banned = 1 WHERE id = ?", (target["id"],)) + conn.execute("DELETE FROM sessions WHERE user_id = ?", (target["id"],)) + conn.commit() + # Уведомления о бане приходят от официального аккаунта @salgram. + notifier = by_username(SALGRAM_USERNAME) or admin + store_message(notifier["id"], rep["reporter_id"], "text", + f"Пользователь @{target['username']} был заблокирован " + f"по вашей жалобе.") + store_message(notifier["id"], target["id"], "text", + "Ваш аккаунт был заблокирован за нарушение правил.") + conn.execute("UPDATE reports SET status = ? WHERE id = ?", + ("banned" if action == "ban" else "dismissed", rid)) + # Помечаем сообщение с кнопками закрытым, чтобы они исчезли из чата. + conn.execute("UPDATE messages SET extra = ? WHERE extra = ?", + (json.dumps({"report": rid, "closed": True}), + json.dumps({"report": rid}))) + conn.commit() + return jsonify(ok=True) + + +# ---------- Группы ---------- + +@app.post("/api/groups") +@auth +def create_group(user): + """Создание группы или канала (is_channel). Публичным обязателен юзернейм + (ищется по нему и названию), приватным — без юзернейма, вход только по ссылке. + В канале пишет только владелец (и админ), остальные — подписчики.""" + d = request.get_json(silent=True) or {} + is_channel = bool(d.get("is_channel", False)) + kindword = "канала" if is_channel else "группы" + title = str(d.get("title", "")).strip() + if not 1 <= len(title) <= GROUP_TITLE_MAX: + return error(f"Название {kindword}: от 1 до {GROUP_TITLE_MAX} символов") + public = bool(d.get("public", True)) + username = None + if public: + username = str(d.get("username", "")).strip().lstrip("@") + if not USERNAME_RE.fullmatch(username): + return error(f"Юзернейм {kindword}: 3–32 символа, только латиница, цифры и _") + # Юзернеймы пользователей, групп и каналов — общее пространство имён. + if by_username(username) or group_by_username(username): + return error("Этот юзернейм уже занят") + conn = db() + try: + gid = conn.execute( + "INSERT INTO groups (title, username, owner_id, public, invite_hash, is_channel) " + "VALUES (?, ?, ?, ?, ?, ?)", + (title, username, user["id"], 1 if public else 0, + secrets.token_urlsafe(12), 1 if is_channel else 0), + ).lastrowid + except sqlite3.IntegrityError: + return error("Этот юзернейм уже занят") + conn.execute("INSERT INTO group_members (group_id, user_id) VALUES (?, ?)", + (gid, user["id"])) + conn.commit() + return jsonify(group_payload(user, by_gid(gid))), 201 + + +@app.get("/api/groups/") +@auth +def group_info(user, gid): + grp = by_gid(gid) + # Приватная группа не видна не-участникам: попасть в неё можно только по ссылке. + if not grp or (not grp["public"] and not group_member(gid, user["id"])): + return error("Группа не найдена", 404) + return jsonify(group_payload(user, grp)) + + +@app.get("/api/groups/hash/") +@auth +def group_by_hash(user, invite): + """Предпросмотр группы по пригласительной ссылке (знание хэша = приглашение).""" + grp = db().execute("SELECT * FROM groups WHERE invite_hash = ?", (invite,)).fetchone() + if not grp: + return error("Ссылка недействительна", 404) + return jsonify(group_payload(user, grp)) + + +def _group_join(user, grp): + if group_is_banned(grp["id"], user["id"]): + return error("Вы заблокированы в этой группе", 403) + conn = db() + conn.execute("INSERT OR IGNORE INTO group_members (group_id, user_id) VALUES (?, ?)", + (grp["id"], user["id"])) + conn.commit() + return jsonify(group_payload(user, grp)) + + +@app.post("/api/groups//join") +@auth +def group_join(user, gid): + grp = by_gid(gid) + if not grp or not grp["public"]: # в приватную группу — только по ссылке + return error("Группа не найдена", 404) + return _group_join(user, grp) + + +@app.post("/api/groups/join/") +@auth +def group_join_by_hash(user, invite): + grp = db().execute("SELECT * FROM groups WHERE invite_hash = ?", (invite,)).fetchone() + if not grp: + return error("Ссылка недействительна", 404) + return _group_join(user, grp) + + +@app.post("/api/groups//add") +@auth +def group_add_member(user, gid): + """Добавление пользователя в группу её участником (минуя ссылку-приглашение). + Учитывает приватность приглашений добавляемого (privacy_group_invite).""" + grp = by_gid(gid) + if not grp or not group_member(gid, user["id"]): + return error("Группа не найдена", 404) + username = str((request.get_json(silent=True) or {}).get("username", "")).strip().lstrip("@") + target = by_username(username) + if not target or target["deleted"]: + return error("Пользователь не найден", 404) + if target["system"]: + return error("Нельзя добавить системный аккаунт") + if target["banned"]: + return error("Нельзя добавить заблокированный аккаунт") + if group_member(gid, target["id"]): + return error("Пользователь уже в группе") + if group_is_banned(gid, target["id"]): + return error("Пользователь заблокирован в этой группе") + if not can_add_to_group(user, target): + return error("Пользователя нельзя добавить: так настроена его приватность") + conn = db() + conn.execute("INSERT OR IGNORE INTO group_members (group_id, user_id) VALUES (?, ?)", + (gid, target["id"])) + conn.commit() + push(target["id"], {"type": "group_added", "group": gid}) # обновить список чатов + return jsonify(group_payload(user, by_gid(gid))) + + +@app.post("/api/groups//avatar") +@auth +def upload_group_avatar(user, gid): + """Аватарка группы (только владелец): центр-кроп до 1024², WebP.""" + grp = by_gid(gid) + if not grp: + return error("Группа не найдена", 404) + if grp["owner_id"] != user["id"] and not user_is_admin(user): + return error("Менять аватарку может только владелец группы", 403) + if not (file := request.files.get("avatar")): + return error("Файл не выбран") + try: + img = ImageOps.exif_transpose(Image.open(file.stream)) + side = min(AVATAR_SIDE, *img.size) + img = ImageOps.fit(img.convert("RGBA"), (side, side)) + img.save(group_avatar_file(gid), "WEBP", quality=85) + except Exception: + return error("Не удалось обработать изображение") + return jsonify(ok=True) + + +@app.get("/api/group_avatar/") +@auth +def get_group_avatar(_user, gid): + if not os.path.exists(group_avatar_file(gid)): + return error("Аватарка не найдена", 404) + return send_file(group_avatar_file(gid), mimetype="image/webp", max_age=0) + + +@app.post("/api/groups//leave") +@auth +def group_leave(user, gid): + grp = by_gid(gid) + if not grp or not group_member(gid, user["id"]): + return error("Группа не найдена", 404) + if grp["owner_id"] == user["id"]: + return error("Владелец не может выйти из группы — удалите её") + conn = db() + conn.execute("DELETE FROM group_members WHERE group_id = ? AND user_id = ?", + (gid, user["id"])) + conn.commit() + push(user["id"], {"type": "group_removed", "group": gid}) # другие вкладки + return jsonify(ok=True) + + +@app.post("/api/groups//delete") +@auth +def group_delete(user, gid): + """Удаление группы (владелец или админ): сообщения, вложения, участники.""" + grp = by_gid(gid) + if not grp: + return error("Группа не найдена", 404) + if grp["owner_id"] != user["id"] and not user_is_admin(user): + return error("Удалить группу может только её владелец") + conn = db() + msgs = conn.execute("SELECT id, kind, content FROM messages WHERE group_id = ?", + (gid,)).fetchall() + member_ids = [r["user_id"] for r in conn.execute( + "SELECT user_id FROM group_members WHERE group_id = ?", (gid,))] + if msgs: + qmarks = ",".join("?" * len(msgs)) + ids = [m["id"] for m in msgs] + conn.execute(f"DELETE FROM reactions WHERE message_id IN ({qmarks})", ids) + conn.execute(f"DELETE FROM messages WHERE id IN ({qmarks})", ids) + conn.execute("DELETE FROM group_members WHERE group_id = ?", (gid,)) + conn.execute("DELETE FROM groups WHERE id = ?", (gid,)) + conn.commit() + for m in msgs: + remove_attachment_file(m["kind"], m["content"]) + if os.path.exists(group_avatar_file(gid)): + os.remove(group_avatar_file(gid)) + for uid in member_ids: + push(uid, {"type": "group_removed", "group": gid}) + return jsonify(ok=True) + + +# ---------- Каналы: привязка группы для комментариев ---------- + +@app.post("/api/channels//link") +@auth +def channel_link_group(user, cid): + """Привязка группы к каналу для комментариев (владелец канала и владелец группы). + Группу нельзя привязать, если она канал или уже привязана к другому каналу.""" + chan = by_gid(cid) + if not chan or not chan["is_channel"]: + return error("Канал не найден", 404) + if chan["owner_id"] != user["id"] and not user_is_admin(user): + return error("Привязывать группу может только владелец канала", 403) + group_id = int((request.get_json(silent=True) or {}).get("group_id", 0)) + grp = by_gid(group_id) + if not grp or grp["is_channel"]: + return error("Группа не найдена", 404) + if grp["owner_id"] != user["id"] and not user_is_admin(user): + return error("Привязать можно только свою группу", 403) + other = db().execute( + "SELECT id FROM groups WHERE linked_group_id = ? AND id <> ?", (group_id, cid)).fetchone() + if other: + return error("Эта группа уже привязана к другому каналу") + conn = db() + conn.execute("UPDATE groups SET linked_group_id = ? WHERE id = ?", (group_id, cid)) + conn.commit() + return jsonify(group_payload(user, by_gid(cid))) + + +@app.post("/api/channels//unlink") +@auth +def channel_unlink_group(user, cid): + chan = by_gid(cid) + if not chan or not chan["is_channel"]: + return error("Канал не найден", 404) + if chan["owner_id"] != user["id"] and not user_is_admin(user): + return error("Отвязать группу может только владелец канала", 403) + conn = db() + conn.execute("UPDATE groups SET linked_group_id = NULL WHERE id = ?", (cid,)) + conn.commit() + return jsonify(group_payload(user, by_gid(cid))) + + +@app.post("/api/channels//comments") +@auth +def channel_open_comments(user, cid): + """Открытие комментариев: подписчик канала автоматически вступает в привязанную + группу, чтобы читать и писать комментарии. Возвращает эту группу.""" + chan = by_gid(cid) + if not chan or not chan["is_channel"] or not chan["linked_group_id"]: + return error("У канала нет комментариев", 404) + if not group_member(cid, user["id"]): + return error("Сначала подпишитесь на канал", 403) + grp = by_gid(chan["linked_group_id"]) + if not grp: + return error("Группа комментариев недоступна", 404) + conn = db() + conn.execute("INSERT OR IGNORE INTO group_members (group_id, user_id) VALUES (?, ?)", + (grp["id"], user["id"])) + conn.commit() + return jsonify(group_payload(user, grp)) + + +@app.get("/api/me/groups") +@auth +def my_linkable_groups(user): + """Группы пользователя, которые можно привязать к каналу: его собственные, + не являющиеся каналом и ещё не привязанные ни к одному каналу.""" + rows = db().execute( + """SELECT id, title FROM groups + WHERE owner_id = ? AND is_channel = 0 + AND id NOT IN (SELECT linked_group_id FROM groups + WHERE linked_group_id IS NOT NULL) + ORDER BY title COLLATE NOCASE""", (user["id"],)).fetchall() + return jsonify([dict(r) for r in rows]) + + +# ---------- Группы: модерация (эндпоинты) ---------- + +@app.post("/api/groups//moderation") +@auth +def set_group_moderation(user, gid): + """Настройки модерации (владелец/админ группы): бан-лист, белый список и его + режим (any/all), медленный режим, общий выключатель.""" + grp = by_gid(gid) + if not grp: + return error("Группа не найдена", 404) + if not can_moderate(user, grp): + return error("Нет прав модерации", 403) + d = request.get_json(silent=True) or {} + conn = db() + if "enabled" in d: + conn.execute("UPDATE groups SET mod_enabled=? WHERE id=?", + (1 if d["enabled"] else 0, gid)) + for key in ("banned_words", "whitelist_words"): + if key in d: + words = d[key] + if not isinstance(words, list) or len(words) > MOD_WORDS_MAX: + return error(f"Список слов: не больше {MOD_WORDS_MAX} записей") + clean = [s for w in words if (s := str(w).strip())][:MOD_WORDS_MAX] + conn.execute(f"UPDATE groups SET {key}=? WHERE id=?", + (json.dumps(clean, ensure_ascii=False), gid)) + if "whitelist_mode" in d: + if d["whitelist_mode"] not in WHITELIST_MODES: + return error("Недопустимый режим белого списка") + conn.execute("UPDATE groups SET whitelist_mode=? WHERE id=?", (d["whitelist_mode"], gid)) + if "slow_seconds" in d: + sec = int(d["slow_seconds"]) + if not 0 <= sec <= MOD_SLOW_MAX: + return error(f"Медленный режим: от 0 до {MOD_SLOW_MAX} секунд") + conn.execute("UPDATE groups SET slow_seconds=? WHERE id=?", (sec, gid)) + conn.commit() + return jsonify(group_payload(user, by_gid(gid))) + + +@app.post("/api/groups//members//mute") +@auth +def group_mute(user, gid, uid): + grp, tm, err = moderation_guard(user, gid, uid) + if err: + return err + if not tm: + return error("Пользователь не в группе", 404) + muted = 1 if (request.get_json(silent=True) or {}).get("muted", True) else 0 + conn = db() + conn.execute("UPDATE group_members SET muted=? WHERE group_id=? AND user_id=?", + (muted, gid, uid)) + conn.commit() + notify_moderation(uid, grp, "замучены — вы не можете писать" if muted else "размучены") + return jsonify(group_payload(user, grp)) + + +@app.post("/api/groups//members//kick") +@auth +def group_kick(user, gid, uid): + grp, tm, err = moderation_guard(user, gid, uid) + if err: + return err + if not tm: + return error("Пользователь не в группе", 404) + conn = db() + conn.execute("DELETE FROM group_members WHERE group_id=? AND user_id=?", (gid, uid)) + conn.commit() + push(uid, {"type": "group_removed", "group": gid}) + notify_moderation(uid, grp, "исключены") + return jsonify(group_payload(user, grp)) + + +@app.post("/api/groups//members//ban") +@auth +def group_ban(user, gid, uid): + grp, _tm, err = moderation_guard(user, gid, uid) + if err: + return err + conn = db() + conn.execute("INSERT OR IGNORE INTO group_bans (group_id, user_id) VALUES (?, ?)", (gid, uid)) + conn.execute("DELETE FROM group_members WHERE group_id=? AND user_id=?", (gid, uid)) + conn.commit() + push(uid, {"type": "group_removed", "group": gid}) + notify_moderation(uid, grp, "заблокированы") + return jsonify(group_payload(user, grp)) + + +@app.post("/api/groups//members//unban") +@auth +def group_unban(user, gid, uid): + grp = by_gid(gid) + if not grp: + return error("Группа не найдена", 404) + if not can_moderate(user, grp): + return error("Нет прав модерации", 403) + conn = db() + conn.execute("DELETE FROM group_bans WHERE group_id=? AND user_id=?", (gid, uid)) + conn.commit() + return jsonify(ok=True) + + +@app.post("/api/groups//members//admin") +@auth +def group_set_admin(user, gid, uid): + """Назначение/снятие админа группы — только владелец (и глобальный @salgr).""" + grp, tm, err = moderation_guard(user, gid, uid, need_owner=True) + if err: + return err + if not tm: + return error("Пользователь не в группе", 404) + admin = 1 if (request.get_json(silent=True) or {}).get("admin", True) else 0 + conn = db() + conn.execute("UPDATE group_members SET is_admin=? WHERE group_id=? AND user_id=?", + (admin, gid, uid)) + conn.commit() + notify_moderation(uid, grp, "назначены администратором" if admin else "сняты с админа") + return jsonify(group_payload(user, grp)) + + +def notify_moderation(uid: int, grp: sqlite3.Row, what: str): + """WS-уведомление участнику о модерационном действии (тост в интерфейсе).""" + kind = "канале" if grp["is_channel"] else "группе" + push(uid, {"type": "moderation", "group": grp["id"], + "text": f"Вы {what} в {kind} «{grp['title']}»"}) + + +# ---------- Поиск ---------- + +@app.get("/api/search") +@auth +def search(user): + q = request.args.get("q", "").strip() + exact = request.args.get("mode") == "exact" + if not q: + return jsonify([]) + + if request.args.get("type") == "messages": + # Содержимое зашифровано — расшифровываем и фильтруем на стороне Python. + rows = db().execute( + """SELECT m.id, m.content, m.created_at, m.sender_id, + CASE WHEN m.sender_id = :me THEN m.recipient_id ELSE m.sender_id END AS peer_id, + u.display_name AS peer_name, u.username AS peer_username, + u.avatar_kind AS peer_avatar_kind, + COALESCE(u.deleted, 1) AS peer_deleted + FROM messages m + LEFT JOIN users u ON u.id = CASE WHEN m.sender_id = :me + THEN m.recipient_id ELSE m.sender_id END + WHERE (m.sender_id = :me OR m.recipient_id = :me) AND m.kind = 'text' + AND m.group_id IS NULL + ORDER BY m.id DESC""", {"me": user["id"]}).fetchall() + results = [] + for r in rows: + text = dec(r["content"]) + if (text == q) if exact else (q.lower() in text.lower()): + results.append({**dict(r), "content": text}) + if len(results) >= 50: + break + return jsonify(results) + + pat = "%" + q.replace("\\", "\\\\").replace("%", "\\%").replace("_", "\\_") + "%" + + if request.args.get("type") == "groups": + # Публичные группы: по названию и юзернейму (можно с ведущим @). + uq = q.lstrip("@") + upat = "%" + uq.replace("\\", "\\\\").replace("%", "\\%").replace("_", "\\_") + "%" + gcond = ("(title = :q OR username = :uq COLLATE NOCASE)" if exact + else "(title LIKE :pat ESCAPE '\\' OR username LIKE :upat ESCAPE '\\')") + grows = db().execute( + f"""SELECT id, title, username, is_channel, + (SELECT COUNT(*) FROM group_members gm + WHERE gm.group_id = groups.id) AS members + FROM groups WHERE public = 1 AND {gcond} + ORDER BY title COLLATE NOCASE LIMIT 20""", + {"q": q, "uq": uq, "pat": pat, "upat": upat}).fetchall() + return jsonify([{"group": True, **dict(r)} for r in grows]) + + # Пользователи: ник и юзернейм одним запросом — дублей не бывает + # (один пользователь = одна строка таблицы). Удалённые и @reports скрыты. + cond = ("(username = :q OR display_name = :q)" if exact + else "(username LIKE :pat ESCAPE '\\' OR display_name LIKE :pat ESCAPE '\\')") + rows = db().execute( + f"""SELECT id, username, display_name, avatar_kind FROM users + WHERE {cond} AND deleted = 0 AND username <> :reports + ORDER BY display_name COLLATE NOCASE LIMIT 20""", + {"q": q, "pat": pat, "reports": REPORTS_USERNAME}).fetchall() + return jsonify([dict(r) for r in rows]) + + +# ---------- Чаты и сообщения ---------- + +@app.get("/api/chats") +@auth +def chats(user): + """Список чатов: последнее сообщение каждой переписки и группы, новые сверху.""" + rows = db().execute(""" + SELECT m.id, m.kind, m.content, m.created_at, m.sender_id, m.extra, + CASE WHEN m.sender_id = :me THEN m.recipient_id ELSE m.sender_id END AS peer_id, + u.display_name AS peer_name, u.username AS peer_username, + u.avatar_kind AS peer_avatar_kind, + COALESCE(u.deleted, 1) AS peer_deleted + FROM messages m + JOIN (SELECT MAX(id) AS mid FROM messages + WHERE (sender_id = :me OR recipient_id = :me) AND group_id IS NULL + GROUP BY CASE WHEN sender_id = :me THEN recipient_id ELSE sender_id END) t + ON t.mid = m.id + LEFT JOIN users u ON u.id = CASE WHEN m.sender_id = :me THEN m.recipient_id ELSE m.sender_id END + ORDER BY m.id DESC""", {"me": user["id"]}).fetchall() + items = [{ + **dict(r), + "content": dec(r["content"]) if r["kind"] == "text" else r["content"], + "extra": json.loads(r["extra"]) if r["extra"] else None, + "from_me": r["sender_id"] == user["id"], + } for r in rows] + + # Группы пользователя: последнее сообщение либо пустая новая группа. + grows = db().execute(""" + SELECT g.id AS group_id, g.title AS group_title, g.created_at AS g_created, + g.is_channel AS is_channel, + m.id, m.kind, m.content, m.created_at, m.sender_id, m.extra, + su.display_name AS sender_name + FROM group_members gm + JOIN groups g ON g.id = gm.group_id + LEFT JOIN messages m ON m.id = (SELECT MAX(id) FROM messages WHERE group_id = g.id) + LEFT JOIN users su ON su.id = m.sender_id + WHERE gm.user_id = :me""", {"me": user["id"]}).fetchall() + for r in grows: + items.append({ + "group_id": r["group_id"], "group_title": r["group_title"], + "is_channel": r["is_channel"], + "kind": r["kind"], + "content": ((dec(r["content"]) if r["kind"] == "text" else r["content"]) + if r["id"] else ""), + "created_at": r["created_at"] or r["g_created"], + "extra": json.loads(r["extra"]) if r["extra"] else None, + "from_me": r["sender_id"] == user["id"], + "sender_name": r["sender_name"], + }) + # Личные и групповые чаты сортируются вместе по времени последнего события. + items.sort(key=lambda x: x["created_at"] or "", reverse=True) + return jsonify(items) + + +@app.post("/api/chats//delete") +@auth +def delete_chat(user, peer_id): + """Удаляет переписку целиком — у обоих участников, с реакциями и фото.""" + conn = db() + msgs = conn.execute( + """SELECT id, kind, content FROM messages + WHERE ((sender_id = :a AND recipient_id = :b) + OR (sender_id = :b AND recipient_id = :a)) + AND group_id IS NULL""", + {"a": user["id"], "b": peer_id}, + ).fetchall() + if not msgs: + return error("Переписка не найдена", 404) + qmarks = ",".join("?" * len(msgs)) + ids = [m["id"] for m in msgs] + conn.execute(f"DELETE FROM reactions WHERE message_id IN ({qmarks})", ids) + conn.execute(f"DELETE FROM messages WHERE id IN ({qmarks})", ids) + conn.commit() + for m in msgs: + remove_attachment_file(m["kind"], m["content"]) + push(user["id"], {"type": "chat_deleted", "with": peer_id}) + push(peer_id, {"type": "chat_deleted", "with": user["id"]}) + return jsonify(ok=True) + + +@app.post("/api/messages/photo") +@auth +def send_photo(user): + """Фото-сообщение: сжатие до 1600px, WebP, файл хранится зашифрованным.""" + if not (file := request.files.get("photo")): + return error("Файл не выбран") + try: + to = int(request.form.get("to", 0)) + gid = int(request.form.get("group", 0)) + except ValueError: + return error("Некорректный получатель") + + raw = file.read() + try: + src = Image.open(io.BytesIO(raw)) # формат/анимацию проверяем до exif_transpose + if src.format == "GIF" and getattr(src, "is_animated", False): + # Анимированный GIF отправляется как анимация (не сжимается в статику). + if len(raw) > GIF_MAX: + return error("GIF больше 15 МБ") + kind, data = "gif", raw + fname = hashlib.sha256(raw).hexdigest() + ".gif" + path = os.path.join(GIF_DIR, fname) + else: + kind = "photo" + img = ImageOps.exif_transpose(src) + img.thumbnail((PHOTO_SIDE, PHOTO_SIDE)) + buf = io.BytesIO() + img.convert("RGBA").save(buf, "WEBP", quality=82) + data = buf.getvalue() + fname = hashlib.sha256(data).hexdigest() + ".webp" + path = os.path.join(PHOTO_DIR, fname) + # Дедупликация: имя файла — хэш содержимого, повторная отправка ссылается + # на уже сохранённый файл. + if not os.path.exists(path): + with open(path, "wb") as f: + f.write(fernet.encrypt(data)) # шифрование на диске + except Exception: + return error("Не удалось обработать изображение") + + _, err = (deliver_group_message(user, gid, kind, fname) if gid + else deliver_message(user, to, kind, fname)) + if err: + # Удалится, только если файл не используется другими сообщениями. + remove_attachment_file(kind, fname) + return error(err) + touch_seen(user["id"]) + return jsonify(ok=True) + + +@app.get("/api/photos/") +@auth +def get_photo(user, fname): + """Фото из сообщения — только участникам переписки; расшифровка на лету.""" + # 32 символа — старые случайные имена, 64 — SHA-256 (дедупликация). + if not re.fullmatch(r"[0-9a-f]{32,64}\.webp", fname): + return error("Не найдено", 404) + allowed = db().execute( + """SELECT 1 FROM messages m WHERE m.kind = 'photo' AND m.content = :f + AND (m.sender_id = :u OR m.recipient_id = :u + OR (m.group_id IS NOT NULL AND EXISTS ( + SELECT 1 FROM group_members gm + WHERE gm.group_id = m.group_id AND gm.user_id = :u)))""", + {"f": fname, "u": user["id"]}, + ).fetchone() + if not allowed: + return error("Не найдено", 404) + with open(os.path.join(PHOTO_DIR, fname), "rb") as f: + raw = f.read() + try: + raw = fernet.decrypt(raw) + except InvalidToken: + pass # файл из БД до включения шифрования + return send_file(io.BytesIO(raw), mimetype="image/webp", max_age=31536000) + + +# ---------- Файлы ---------- + +@app.post("/api/messages/file") +@auth +def send_file_msg(user): + """Файл-сообщение (любой тип, ≤25 МБ): хранится зашифрованным; для документов, + презентаций и таблиц извлекается текстовый предпросмотр.""" + if not (file := request.files.get("file")): + return error("Файл не выбран") + try: + to = int(request.form.get("to", 0)) + gid = int(request.form.get("group", 0)) + except ValueError: + return error("Некорректный получатель") + + data = file.read() + if not data: + return error("Файл пуст") + if len(data) > FILE_MAX: + return error("Файл больше 25 МБ") + name = os.path.basename(file.filename or "файл")[:128] or "файл" + extra = {"name": name, "size": len(data), + "mime": file.mimetype or "application/octet-stream"} + if preview := file_preview(name, data): + extra["preview"] = preview + + # Дедупликация: имя файла — хэш содержимого; одинаковые файлы не дублируются. + fname = hashlib.sha256(data).hexdigest() + ".bin" + path = os.path.join(FILE_DIR, fname) + if not os.path.exists(path): + with open(path, "wb") as f: + f.write(fernet.encrypt(data)) # шифрование на диске + + _, err = (deliver_group_message(user, gid, "file", fname, extra) if gid + else deliver_message(user, to, "file", fname, extra)) + if err: + # Удалится, только если файл не используется другими сообщениями. + remove_attachment_file("file", fname) + return error(err) + touch_seen(user["id"]) + return jsonify(ok=True) + + +@app.get("/api/files/") +@auth +def get_file(user, fname): + """Файл из сообщения — только участникам переписки; расшифровка на лету. + ?download=1 — скачать с исходным именем, иначе — открыть (предпросмотр).""" + # 32 символа — старые случайные имена, 64 — SHA-256 (дедупликация). + if not re.fullmatch(r"[0-9a-f]{32,64}\.bin", fname): + return error("Не найдено", 404) + m = db().execute( + """SELECT m.extra FROM messages m WHERE m.kind = 'file' AND m.content = :f + AND (m.sender_id = :u OR m.recipient_id = :u + OR (m.group_id IS NOT NULL AND EXISTS ( + SELECT 1 FROM group_members gm + WHERE gm.group_id = m.group_id AND gm.user_id = :u)))""", + {"f": fname, "u": user["id"]}, + ).fetchone() + if not m: + return error("Не найдено", 404) + extra = json.loads(m["extra"]) if m["extra"] else {} + with open(os.path.join(FILE_DIR, fname), "rb") as f: + raw = fernet.decrypt(f.read()) + return send_file(io.BytesIO(raw), + mimetype=extra.get("mime") or "application/octet-stream", + download_name=extra.get("name", "файл"), + as_attachment=request.args.get("download") == "1", + max_age=31536000) + + +# ---------- GIF: общая библиотека (бот @addgif) ---------- + +GIF_NAME_RE = re.compile(r"[0-9a-f]{64}\.gif") + + +@app.post("/api/gifs") +@auth +def add_gif(user): + """Добавление GIF в общую библиотеку через бота @addgif. + Лимит — GIF_DAILY_LIMIT штук в сутки на пользователя. Поиск идёт по ключевым словам.""" + if not (file := request.files.get("gif")): + return error("Файл не выбран") + keywords = " ".join(str(request.form.get("keywords", "")).lower().split()) + if not 1 <= len(keywords) <= KEYWORDS_MAX: + return error(f"Укажите ключевые слова (до {KEYWORDS_MAX} символов)") + used = db().execute( + "SELECT COUNT(*) FROM gifs WHERE uploader_id = ? " + "AND created_at >= datetime('now', '-1 day')", (user["id"],)).fetchone()[0] + if used >= GIF_DAILY_LIMIT: + return error(f"Лимит {GIF_DAILY_LIMIT} GIF в сутки исчерпан — попробуйте позже") + + raw = file.read() + if len(raw) > GIF_MAX: + return error("GIF больше 15 МБ") + try: + if Image.open(io.BytesIO(raw)).format != "GIF": + return error("Это должен быть файл GIF") + except Exception: + return error("Не удалось прочитать GIF") + + fname = hashlib.sha256(raw).hexdigest() + ".gif" # дедупликация по содержимому + path = os.path.join(GIF_DIR, fname) + if not os.path.exists(path): + with open(path, "wb") as f: + f.write(fernet.encrypt(raw)) # шифрование на диске + conn = db() + conn.execute("INSERT INTO gifs (file, keywords, uploader_id) VALUES (?, ?, ?)", + (fname, keywords, user["id"])) + conn.commit() + # Эхо-подтверждение в чат с @addgif: сам GIF + сообщение об остатке лимита. + if addgif := by_username(ADDGIF_USERNAME): + store_message(addgif["id"], user["id"], "gif", fname) + store_message(addgif["id"], user["id"], "text", + f"GIF добавлен в общую библиотеку. Поиск по словам: {keywords}. " + f"Сегодня можно добавить ещё {GIF_DAILY_LIMIT - used - 1}.") + return jsonify(ok=True) + + +@app.get("/api/gifs") +@auth +def search_gifs(_user): + """Поиск по общей библиотеке GIF (по ключевым словам). Без q — недавние.""" + q = request.args.get("q", "").strip().lower() + if q: + pat = "%" + q.replace("\\", "\\\\").replace("%", "\\%").replace("_", "\\_") + "%" + rows = db().execute( + """SELECT file FROM gifs WHERE keywords LIKE ? ESCAPE '\\' + GROUP BY file ORDER BY MAX(id) DESC LIMIT 60""", (pat,)).fetchall() + else: + rows = db().execute( + "SELECT file FROM gifs GROUP BY file ORDER BY MAX(id) DESC LIMIT 60").fetchall() + return jsonify([{"file": r["file"]} for r in rows]) + + +@app.get("/api/gifs/") +@auth +def get_gif(user, fname): + """Отдаёт GIF. Библиотечный — всем; ad-hoc GIF из переписки — только её участникам.""" + if not GIF_NAME_RE.fullmatch(fname): + return error("Не найдено", 404) + if not db().execute("SELECT 1 FROM gifs WHERE file = ?", (fname,)).fetchone(): + allowed = db().execute( + """SELECT 1 FROM messages m WHERE m.kind = 'gif' AND m.content = :f + AND (m.sender_id = :u OR m.recipient_id = :u + OR (m.group_id IS NOT NULL AND EXISTS ( + SELECT 1 FROM group_members gm + WHERE gm.group_id = m.group_id AND gm.user_id = :u)))""", + {"f": fname, "u": user["id"]}).fetchone() + if not allowed: + return error("Не найдено", 404) + with open(os.path.join(GIF_DIR, fname), "rb") as f: + raw = f.read() + try: + raw = fernet.decrypt(raw) + except InvalidToken: + pass # файл из БД до включения шифрования + return send_file(io.BytesIO(raw), mimetype="image/gif", max_age=31536000) + + +@app.post("/api/messages/gif") +@auth +def send_gif(user): + """Отправка GIF из общей библиотеки в личный чат или группу.""" + d = request.get_json(silent=True) or {} + fname = str(d.get("file", "")) + if not GIF_NAME_RE.fullmatch(fname) \ + or not db().execute("SELECT 1 FROM gifs WHERE file = ?", (fname,)).fetchone(): + return error("GIF не найден", 404) + gid, to = int(d.get("group", 0)), int(d.get("to", 0)) + _, err = (deliver_group_message(user, gid, "gif", fname) if gid + else deliver_message(user, to, "gif", fname)) + if err: + return error(err) + touch_seen(user["id"]) + return jsonify(ok=True) + + +# ---------- Аудио и голосовые сообщения ---------- + +@app.post("/api/messages/audio") +@auth +def send_audio(user): + """Аудио-сообщение (голосовое с микрофона или аудиофайл, ≤25 МБ): + хранится зашифрованным; в extra — длительность, флаг голосового и имя.""" + if not (file := request.files.get("audio")): + return error("Файл не выбран") + try: + to = int(request.form.get("to", 0)) + gid = int(request.form.get("group", 0)) + except ValueError: + return error("Некорректный получатель") + + data = file.read() + if not data: + return error("Файл пуст") + if len(data) > AUDIO_MAX: + return error("Аудио больше 25 МБ") + mime = (file.mimetype or "").split(";")[0].strip().lower() + ext = AUDIO_EXT.get(mime) + if not ext: # запасной вариант — расширение из имени файла + guess = (os.path.splitext(file.filename or "")[1][1:] or "").lower() + ext = guess if guess in {"webm", "ogg", "mp3", "m4a", "aac", "wav", "flac"} else None + if not ext: + return error("Неподдерживаемый формат аудио") + + voice = request.form.get("voice") == "1" + try: + duration = round(float(request.form.get("duration", 0)), 1) + except ValueError: + duration = 0 + extra = {"mime": mime or f"audio/{ext}", "size": len(data), + "duration": max(0.0, duration), "voice": voice} + if not voice: + extra["name"] = os.path.basename(file.filename or "audio")[:128] or f"audio.{ext}" + + # Дедупликация по содержимому: одинаковое аудио не дублируется на диске. + fname = hashlib.sha256(data).hexdigest() + "." + ext + path = os.path.join(AUDIO_DIR, fname) + if not os.path.exists(path): + with open(path, "wb") as f: + f.write(fernet.encrypt(data)) # шифрование на диске + + _, err = (deliver_group_message(user, gid, "audio", fname, extra) if gid + else deliver_message(user, to, "audio", fname, extra)) + if err: + remove_attachment_file("audio", fname) + return error(err) + touch_seen(user["id"]) + return jsonify(ok=True) + + +@app.get("/api/audio/") +@auth +def get_audio(user, fname): + """Аудио из сообщения — только участникам переписки; расшифровка на лету.""" + if not AUDIO_NAME_RE.fullmatch(fname): + return error("Не найдено", 404) + m = db().execute( + """SELECT m.extra FROM messages m WHERE m.kind = 'audio' AND m.content = :f + AND (m.sender_id = :u OR m.recipient_id = :u + OR (m.group_id IS NOT NULL AND EXISTS ( + SELECT 1 FROM group_members gm + WHERE gm.group_id = m.group_id AND gm.user_id = :u)))""", + {"f": fname, "u": user["id"]}, + ).fetchone() + if not m: + return error("Не найдено", 404) + extra = json.loads(m["extra"]) if m["extra"] else {} + with open(os.path.join(AUDIO_DIR, fname), "rb") as f: + raw = f.read() + try: + raw = fernet.decrypt(raw) + except InvalidToken: + pass # файл из БД до включения шифрования + ext = fname.rsplit(".", 1)[-1] + mime = extra.get("mime") or {"webm": "audio/webm", "ogg": "audio/ogg", + "mp3": "audio/mpeg", "m4a": "audio/mp4", "aac": "audio/aac", + "wav": "audio/wav", "flac": "audio/flac"}.get(ext, "application/octet-stream") + return send_file(io.BytesIO(raw), mimetype=mime, max_age=31536000) + + +# ---------- Push-уведомления (подписка браузера) ---------- + +@app.get("/api/push/key") +@auth +def push_key(_user): + """Публичный VAPID-ключ (applicationServerKey) для подписки в браузере. + enabled=false — push на сервере недоступен (нет библиотеки/ключей).""" + return jsonify(enabled=bool(_WEBPUSH_OK and VAPID_PUBLIC_B64), key=VAPID_PUBLIC_B64) + + +@app.post("/api/push/subscribe") +@auth +def push_subscribe(user): + """Сохраняет PushSubscription браузера (привязка к текущему аккаунту).""" + sub = request.get_json(silent=True) or {} + endpoint = str(sub.get("endpoint", "")) + if not endpoint.startswith("http") or not sub.get("keys"): + return error("Некорректная подписка") + conn = db() + conn.execute( + "INSERT INTO push_subscriptions (endpoint, user_id, data) VALUES (?, ?, ?) " + "ON CONFLICT(endpoint) DO UPDATE SET user_id = excluded.user_id, data = excluded.data", + (endpoint, user["id"], json.dumps(sub, ensure_ascii=False))) + conn.commit() + return jsonify(ok=True) + + +@app.post("/api/push/unsubscribe") +@auth +def push_unsubscribe(_user): + """Удаляет подписку (например, при выходе или отключении уведомлений).""" + endpoint = str((request.get_json(silent=True) or {}).get("endpoint", "")) + if endpoint: + conn = db() + conn.execute("DELETE FROM push_subscriptions WHERE endpoint = ?", (endpoint,)) + conn.commit() + return jsonify(ok=True) + + +# ---------- WebSocket ---------- + +def ws_send(ws, payload: dict): + ws.send(json.dumps(payload, ensure_ascii=False)) + + +def handle_ws(uid: int, msg: dict, ws): + if msg.get("type") == "ping": # keepalive: прокси закрывают «молчащие» соединения + return + sender = by_id(uid) + if not sender or sender["deleted"]: + return + # Забаненный аккаунт — только чтение: доступна лишь история переписки. + if sender["banned"] and msg.get("type") != "history": + return ws_send(ws, {"type": "error", "error": "Аккаунт заблокирован"}) + match msg.get("type"): + case "send": + text = str(msg.get("text", "")).strip() + if not 1 <= len(text) <= MSG_MAX: + return ws_send(ws, {"type": "error", "error": f"Сообщение: 1–{MSG_MAX} символов"}) + fmt = clean_fmt(msg.get("fmt"), len(text)) # пометки форматирования текста + extra = {"fmt": fmt} if fmt else None + if msg.get("group"): + gid = int(msg["group"]) + reply_to, thread_root = None, None + if msg.get("reply_to"): # ответ/комментарий: на пост-зеркало или другое сообщение + parent = db().execute( + "SELECT * FROM messages WHERE id = ? AND group_id = ?", + (int(msg["reply_to"]), gid)).fetchone() + if parent: + reply_to = parent["id"] + thread_root = resolve_thread_root(parent) + _, err = deliver_group_message(sender, gid, "text", text, extra, + reply_to, thread_root) + else: + _, err = deliver_message(sender, int(msg.get("to", 0)), "text", text, + extra, reply_to=msg.get("reply_to")) + if err: + ws_send(ws, {"type": "error", "error": err}) + case "typing": + # Статус «печатает»: не пишем в БД, просто транслируем собеседникам. + if msg.get("group"): + gid = int(msg["group"]) + if group_member(gid, uid): + for r in db().execute( + "SELECT user_id FROM group_members WHERE group_id = ? AND user_id != ?", + (gid, uid)): + push(r["user_id"], {"type": "typing", "group": gid, + "from": uid, "from_name": sender["display_name"]}) + else: + to = int(msg.get("to", 0)) + target = by_id(to) + # Не выдаём «печатает», если кто-то кого-то заблокировал. + if (target and not target["deleted"] and to != uid + and not is_blocked(to, uid) and not is_blocked(uid, to)): + push(to, {"type": "typing", "from": uid, + "from_name": sender["display_name"]}) + case "read": + # Отметка прочтения личной переписки: запоминаем макс. id сообщения + # собеседника и сообщаем ему, чтобы он увидел галочки «прочитано». + peer = int(msg.get("with", 0)) + if peer and peer != uid: + last = db().execute( + "SELECT MAX(id) FROM messages WHERE sender_id = ? AND recipient_id = ? " + "AND group_id IS NULL", (peer, uid)).fetchone()[0] + if last: + conn = db() + conn.execute( + "INSERT INTO reads (reader_id, peer_id, last_read_id) VALUES (?, ?, ?) " + "ON CONFLICT(reader_id, peer_id) DO UPDATE SET " + "last_read_id = MAX(last_read_id, excluded.last_read_id)", + (uid, peer, last)) + conn.commit() + push(peer, {"type": "read", "with": uid, "upto": last}) + case "edit": + mid, text = int(msg.get("id", 0)), str(msg.get("text", "")).strip() + if not 1 <= len(text) <= MSG_MAX: + return ws_send(ws, {"type": "error", "error": f"Сообщение: 1–{MSG_MAX} символов"}) + m = db().execute("SELECT * FROM messages WHERE id = ?", (mid,)).fetchone() + if not m or m["sender_id"] != uid or m["kind"] != "text": + return ws_send(ws, {"type": "error", "error": "Это сообщение нельзя изменить"}) + extra = json.loads(m["extra"]) if m["extra"] else {} + extra["edited"] = True + fmt = clean_fmt(msg.get("fmt"), len(text)) # пометки форматирования правки + if fmt: + extra["fmt"] = fmt + else: + extra.pop("fmt", None) + conn = db() + conn.execute("UPDATE messages SET content = ?, extra = ? WHERE id = ?", + (enc(text), json.dumps(extra, ensure_ascii=False), mid)) + conn.commit() + row = conn.execute("SELECT * FROM messages WHERE id = ?", (mid,)).fetchone() + push_message_event(m, {"type": "message_edited", "message": message_payload(row)}) + case "delete": + mid = int(msg.get("id", 0)) + m = db().execute("SELECT * FROM messages WHERE id = ?", (mid,)).fetchone() + # Удалить может автор; в группе — ещё и её владелец (модерация). + group_owner = m and m["group_id"] and db().execute( + "SELECT 1 FROM groups WHERE id = ? AND owner_id = ?", + (m["group_id"], uid)).fetchone() + if not m or (m["sender_id"] != uid and not group_owner): + return ws_send(ws, {"type": "error", "error": "Это сообщение нельзя удалить"}) + conn = db() + conn.execute("DELETE FROM reactions WHERE message_id = ?", (mid,)) + conn.execute("DELETE FROM messages WHERE id = ?", (mid,)) + conn.commit() + remove_attachment_file(m["kind"], m["content"]) + push_message_event(m, {"type": "message_deleted", "id": mid, + "from": m["sender_id"], "to": m["recipient_id"], + "group": m["group_id"]}) + case "pin": + mid, want = int(msg.get("id", 0)), bool(msg.get("pinned", True)) + m = db().execute("SELECT * FROM messages WHERE id = ?", (mid,)).fetchone() + if not m or not message_visible_to(m, uid): + return ws_send(ws, {"type": "error", "error": "Сообщение не найдено"}) + extra = json.loads(m["extra"]) if m["extra"] else {} + if want: + extra["pinned"] = True + else: + extra.pop("pinned", None) + conn = db() + conn.execute("UPDATE messages SET extra = ? WHERE id = ?", + (json.dumps(extra, ensure_ascii=False) if extra else None, mid)) + conn.commit() + row = conn.execute("SELECT * FROM messages WHERE id = ?", (mid,)).fetchone() + push_message_event(m, {"type": "message_pinned", "message": message_payload(row)}) + case "react": + mid, emoji = int(msg.get("id", 0)), str(msg.get("emoji", "")) + if emoji not in REACTION_OK: + return ws_send(ws, {"type": "error", "error": "Недопустимая реакция"}) + m = db().execute("SELECT * FROM messages WHERE id = ?", (mid,)).fetchone() + if not m or not message_visible_to(m, uid): + return ws_send(ws, {"type": "error", "error": "Сообщение не найдено"}) + conn = db() + cur = conn.execute("SELECT emoji FROM reactions WHERE message_id = ? AND user_id = ?", + (mid, uid)).fetchone() + if cur and cur["emoji"] == emoji: # повторный клик — снять реакцию + conn.execute("DELETE FROM reactions WHERE message_id = ? AND user_id = ?", + (mid, uid)) + else: + conn.execute("INSERT OR REPLACE INTO reactions (message_id, user_id, emoji) " + "VALUES (?, ?, ?)", (mid, uid, emoji)) + conn.commit() + push_message_event(m, {"type": "reaction", "id": mid, "from": m["sender_id"], + "to": m["recipient_id"], "group": m["group_id"], + "reactions": message_reactions(mid)}) + case "history": + before = msg.get("before") + if msg.get("thread"): # тред комментариев под постом обсуждения + root_id = int(msg["thread"]) + root = db().execute("SELECT * FROM messages WHERE id = ?", (root_id,)).fetchone() + if not root or not root["group_id"] or not group_member(root["group_id"], uid): + return ws_send(ws, {"type": "error", "error": "Комментарии недоступны"}) + rows = db().execute( + """SELECT * FROM messages WHERE thread_root = :r + AND (:before IS NULL OR id < :before) + ORDER BY id DESC LIMIT 50""", + {"r": root_id, "before": before}).fetchall() + return ws_send(ws, {"type": "history", "thread": root_id, + "group": root["group_id"], + "root": message_payload(root), + "messages": [message_payload(r) for r in reversed(rows)]}) + if msg.get("group"): # история группового чата — только участникам + gid = int(msg["group"]) + if not group_member(gid, uid): + return ws_send(ws, {"type": "error", + "error": "Вы не состоите в этой группе"}) + rows = db().execute( + """SELECT * FROM messages WHERE group_id = :g + AND (:before IS NULL OR id < :before) + ORDER BY id DESC LIMIT 50""", + {"g": gid, "before": before}).fetchall() + pinned = db().execute( + """SELECT * FROM messages WHERE group_id = :g + AND extra IS NOT NULL AND json_extract(extra, '$.pinned') + ORDER BY id""", {"g": gid}).fetchall() + return ws_send(ws, {"type": "history", "group": gid, + "messages": [message_payload(r) for r in reversed(rows)], + "pinned": [message_payload(r) for r in pinned]}) + peer = int(msg.get("with", 0)) + rows = db().execute( + """SELECT * FROM messages + WHERE ((sender_id = :a AND recipient_id = :b) + OR (sender_id = :b AND recipient_id = :a)) + AND group_id IS NULL + AND (:before IS NULL OR id < :before) + ORDER BY id DESC LIMIT 50""", + {"a": uid, "b": peer, "before": before}, + ).fetchall() + pinned = db().execute( + """SELECT * FROM messages + WHERE ((sender_id = :a AND recipient_id = :b) + OR (sender_id = :b AND recipient_id = :a)) + AND group_id IS NULL + AND extra IS NOT NULL AND json_extract(extra, '$.pinned') + ORDER BY id""", + {"a": uid, "b": peer}, + ).fetchall() + # Докуда собеседник прочитал мои сообщения (для галочек «прочитано»). + peer_read = db().execute( + "SELECT last_read_id FROM reads WHERE reader_id = ? AND peer_id = ?", + (peer, uid)).fetchone() + ws_send(ws, {"type": "history", "with": peer, + "messages": [message_payload(r) for r in reversed(rows)], + "pinned": [message_payload(r) for r in pinned], + "peer_read": peer_read[0] if peer_read else 0}) + + +@sock.route("/ws") +def ws_route(ws): + user = current_user() + if not user: + return + uid = user["id"] + connections[uid].add(ws) + touch_seen(uid) # зашёл в сеть + try: + while True: + raw = ws.receive() + if raw is None: + break + try: + handle_ws(uid, json.loads(raw), ws) + except (ValueError, TypeError): + pass + except Exception: + pass # соединение оборвано + finally: + connections[uid].discard(ws) + touch_seen(uid) # вышел из сети + + +init_db() +_load_vapid() +hf_backup.start(DATA_DIR) # периодический бэкап в приватный репо (no-op, если не настроено) + + +def run_production(port: int): + """Боевой сервер: gunicorn с потоками (рекомендованный flask-sock режим). + + Запускается программно, чтобы на хостинге хватало `scriptName: app.py`. + Один воркер обязателен: WebSocket-рассылка (push) хранит соединения + в памяти процесса. Каждое WS-соединение занимает поток, поэтому потоков + много; при росте аудитории увеличьте threads или переходите на gevent. + """ + from gunicorn.app.base import BaseApplication + + class Server(BaseApplication): + def load_config(self): + for key, value in {"bind": f"0.0.0.0:{port}", "workers": 1, + "threads": 100, "timeout": 120}.items(): + self.cfg.set(key, value) + + def load(self): + return app + + Server().run() + + +if __name__ == "__main__": + port = int(os.environ.get("PORT", 8000)) + if DATA_DIR == "/data": # прод (Amvera): /data существует только на хостинге + run_production(port) + else: + # Локальная разработка: многопоточный werkzeug (WebSocket через flask-sock). + app.run(host="0.0.0.0", port=port, threaded=True)