diff --git a/api/v2/routes/notifications.py b/api/v2/routes/notifications.py
index 099d8c5a..f337b185 100644
--- a/api/v2/routes/notifications.py
+++ b/api/v2/routes/notifications.py
@@ -1,4 +1,4 @@
-from fastapi import APIRouter, Depends, Query
+from fastapi import APIRouter, Depends, HTTPException, Path, Query
from pydantic import BaseModel
from sqlalchemy.ext.asyncio import AsyncSession
@@ -84,3 +84,29 @@ async def read_all_notifications(
):
count = await wn_db.mark_all_read_for_identity(session, identity.id)
return {"ok": True, "updated": count}
+
+
+@router.post("/notifications/{notification_id}/read", tags=["Notifications"])
+async def read_one_notification(
+ notification_id: str = Path(..., min_length=1, max_length=64),
+ session: AsyncSession = Depends(get_session),
+ identity: Identity = Depends(verify_identity_token),
+):
+ """Пометить одно уведомление прочитанным. 404 если не найдено или не принадлежит юзеру."""
+ ok = await wn_db.mark_one_read_for_identity(session, identity.id, notification_id)
+ if not ok:
+ raise HTTPException(status_code=404, detail="Уведомление не найдено")
+ return {"ok": True}
+
+
+@router.delete("/notifications/{notification_id}", tags=["Notifications"])
+async def delete_one_notification(
+ notification_id: str = Path(..., min_length=1, max_length=64),
+ session: AsyncSession = Depends(get_session),
+ identity: Identity = Depends(verify_identity_token),
+):
+ """Удалить одно уведомление. 404 если не найдено или не принадлежит юзеру."""
+ ok = await wn_db.delete_one_for_identity(session, identity.id, notification_id)
+ if not ok:
+ raise HTTPException(status_code=404, detail="Уведомление не найдено")
+ return {"ok": True}
diff --git a/audit/__init__.py b/audit/__init__.py
index 9a68feda..364abb93 100644
--- a/audit/__init__.py
+++ b/audit/__init__.py
@@ -11,6 +11,7 @@ from typing import Any
from aiogram.types import CallbackQuery, InlineQuery, Message, TelegramObject, User
from fastapi import Request
+from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from database.audit import (
@@ -1092,6 +1093,20 @@ async def drain_audit_redis_to_db(session_factory: Any) -> int:
session,
[str(rec.get("request_id")) for rec in batch if rec.get("request_id")],
)
+ identity_ids_in_batch = {
+ str(rec.get("actor_identity_id"))
+ for rec in batch
+ if rec.get("actor_identity_id")
+ }
+ if identity_ids_in_batch:
+ from database.models import Identity
+
+ rows = await session.execute(
+ select(Identity.id).where(Identity.id.in_(identity_ids_in_batch))
+ )
+ existing_identity_ids = {row[0] for row in rows.all()}
+ else:
+ existing_identity_ids = set()
inserted_count = 0
seen_batch_request_ids: set[str] = set()
for rec in batch:
@@ -1100,6 +1115,9 @@ async def drain_audit_redis_to_db(session_factory: Any) -> int:
continue
if request_id:
seen_batch_request_ids.add(request_id)
+ actor_identity_id = rec.get("actor_identity_id")
+ if actor_identity_id and actor_identity_id not in existing_identity_ids:
+ actor_identity_id = None
created = rec.get("created_at")
if isinstance(created, str):
try:
@@ -1115,7 +1133,7 @@ async def drain_audit_redis_to_db(session_factory: Any) -> int:
event_type=rec.get("event_type", "telegram_access"),
channel=rec.get("channel", "telegram"),
path_or_handler=rec.get("path_or_handler") or "telegram",
- actor_identity_id=rec.get("actor_identity_id"),
+ actor_identity_id=actor_identity_id,
actor_tg_id=rec.get("actor_tg_id"),
entity_type=rec.get("entity_type"),
entity_id=rec.get("entity_id"),
diff --git a/core/settings/web_config.py b/core/settings/web_config.py
index 9b287aae..d78ee603 100644
--- a/core/settings/web_config.py
+++ b/core/settings/web_config.py
@@ -76,3 +76,7 @@ def get_site_url() -> str:
def is_web_enabled() -> bool:
return bool(WEB_CONFIG.get("WEB_ENABLED", False))
+
+
+def is_email_binding_enabled() -> bool:
+ return bool(WEB_CONFIG.get("EMAIL_BINDING_ENABLED", False))
diff --git a/database/web_notifications.py b/database/web_notifications.py
index 89586941..fe9f5dae 100644
--- a/database/web_notifications.py
+++ b/database/web_notifications.py
@@ -110,6 +110,41 @@ async def mark_all_read_for_identity(
return result.rowcount
+async def mark_one_read_for_identity(
+ session: AsyncSession,
+ identity_id: str,
+ notification_id: str,
+) -> bool:
+ """Mark single notification as read. Returns True if updated."""
+ result = await session.execute(
+ update(WebNotification)
+ .where(
+ WebNotification.identity_id == identity_id,
+ WebNotification.id == notification_id,
+ )
+ .values(read=True)
+ )
+ return (result.rowcount or 0) > 0
+
+
+async def delete_one_for_identity(
+ session: AsyncSession,
+ identity_id: str,
+ notification_id: str,
+) -> bool:
+ """Delete single notification owned by identity. Returns True if deleted."""
+ from sqlalchemy import delete as sql_delete
+
+ result = await session.execute(
+ sql_delete(WebNotification)
+ .where(
+ WebNotification.identity_id == identity_id,
+ WebNotification.id == notification_id,
+ )
+ )
+ return (result.rowcount or 0) > 0
+
+
async def resolve_identity_id_by_tg_id(
session: AsyncSession,
tg_id: int,
diff --git a/handlers/__init__.py b/handlers/__init__.py
index 6af1510d..d001856d 100644
--- a/handlers/__init__.py
+++ b/handlers/__init__.py
@@ -7,6 +7,7 @@ from .captcha import router as captcha_router
from .chat_member import router as chat_member_router
from .coupons import router as coupons_router
from .donate import router as donate_router
+from .email_binding import router as email_binding_router
from .instructions import router as instructions_router
from .keys import router as keys_router
from .notifications import router as notifications_router
@@ -33,4 +34,5 @@ router.include_routers(
admin_router,
refferal_router,
tariff_router,
+ email_binding_router,
)
diff --git a/handlers/admin/settings/settings_web.py b/handlers/admin/settings/settings_web.py
index 2834bf70..a8aa5e04 100644
--- a/handlers/admin/settings/settings_web.py
+++ b/handlers/admin/settings/settings_web.py
@@ -35,6 +35,13 @@ def build_settings_web_kb() -> InlineKeyboardBuilder:
callback_data=AdminPanelCallback(action="settings_web_url").pack(),
)
)
+ email_binding = bool(WEB_CONFIG.get("EMAIL_BINDING_ENABLED", False))
+ builder.row(
+ InlineKeyboardButton(
+ text=f"{'✅' if email_binding else '❌'} Привязка почты {'включена' if email_binding else 'выключена'}",
+ callback_data=AdminPanelCallback(action="settings_web_email_binding_toggle").pack(),
+ )
+ )
builder.row(
InlineKeyboardButton(
text="🔄 Сбросить сайт к исходнику",
@@ -49,12 +56,16 @@ def build_settings_web_kb() -> InlineKeyboardBuilder:
def _web_settings_text() -> str:
enabled = bool(WEB_CONFIG.get("WEB_ENABLED", False))
url = str(WEB_CONFIG.get("SITE_URL") or "не указан")
+ email_binding = bool(WEB_CONFIG.get("EMAIL_BINDING_ENABLED", False))
return (
"🌐 Настройки веб-сайта\n\n"
f"Статус: {'✅ Включён' if enabled else '❌ Выключен'}\n"
- f"URL: {url}\n\n"
+ f"URL: {url}\n"
+ f"Привязка почты: {'✅ Включена' if email_binding else '❌ Выключена'}\n\n"
"Сайт может работать на отдельном домене и сервере.\n"
- "При выключении кнопка «Личный кабинет» скрывается из бота."
+ "При выключении кнопка «Личный кабинет» скрывается из бота.\n"
+ "Привязка почты — кнопка в боте, через которую пользователь указывает email "
+ "для входа на сайт на случай проблем с Telegram."
)
@@ -84,6 +95,23 @@ async def toggle_web_enabled(callback: CallbackQuery) -> None:
)
+@router.callback_query(AdminPanelCallback.filter(F.action == "settings_web_email_binding_toggle"))
+async def toggle_email_binding(callback: CallbackQuery) -> None:
+ current = bool(WEB_CONFIG.get("EMAIL_BINDING_ENABLED", False))
+ new_config = dict(WEB_CONFIG)
+ new_config["EMAIL_BINDING_ENABLED"] = not current
+
+ async with async_session_maker() as session:
+ await update_web_config(session, new_config)
+
+ status = "✅ Привязка почты включена" if new_config["EMAIL_BINDING_ENABLED"] else "❌ Привязка почты выключена"
+ await callback.answer(status, show_alert=True)
+ await callback.message.edit_text(
+ text=_web_settings_text(),
+ reply_markup=build_settings_web_kb().as_markup(),
+ )
+
+
@router.callback_query(AdminPanelCallback.filter(F.action == "settings_web_url"))
async def prompt_web_url(callback: CallbackQuery, state: FSMContext) -> None:
current = str(WEB_CONFIG.get("SITE_URL") or "")
diff --git a/handlers/email_binding.py b/handlers/email_binding.py
new file mode 100644
index 00000000..8534d697
--- /dev/null
+++ b/handlers/email_binding.py
@@ -0,0 +1,90 @@
+import re
+
+from aiogram import F, Router
+from aiogram.fsm.context import FSMContext
+from aiogram.fsm.state import State, StatesGroup
+from aiogram.types import CallbackQuery, InlineKeyboardButton, Message
+from aiogram.utils.keyboard import InlineKeyboardBuilder
+
+from core.settings.web_config import is_email_binding_enabled
+from database.identities import get_identity_by_email, get_or_create_identity_for_tg
+from handlers.buttons import BACK
+from handlers.utils import edit_or_send_message
+from logger import logger
+
+
+router = Router(name="email_binding")
+
+EMAIL_RE = re.compile(r"^[^\s@]+@[^\s@]+\.[^\s@]+$")
+
+
+class EmailBindingState(StatesGroup):
+ waiting_for_email = State()
+
+
+@router.callback_query(F.data == "bind_email")
+async def prompt_email(callback: CallbackQuery, state: FSMContext, session) -> None:
+ if not is_email_binding_enabled():
+ await callback.answer("Привязка почты отключена", show_alert=True)
+ return
+
+ identity = await get_or_create_identity_for_tg(session, callback.from_user.id)
+ if identity.email:
+ await callback.answer("Почта уже привязана", show_alert=True)
+ return
+
+ builder = InlineKeyboardBuilder()
+ builder.row(InlineKeyboardButton(text=BACK, callback_data="profile"))
+
+ await edit_or_send_message(
+ target_message=callback.message,
+ text=(
+ "📧 Привязка почты\n\n"
+ "Укажите email — он понадобится для входа на сайт, "
+ "если возникнут проблемы с Telegram.\n\n"
+ "Отправьте адрес сообщением."
+ ),
+ reply_markup=builder.as_markup(),
+ )
+ await state.set_state(EmailBindingState.waiting_for_email)
+ await callback.answer()
+
+
+@router.message(EmailBindingState.waiting_for_email)
+async def receive_email(message: Message, state: FSMContext, session) -> None:
+ raw = (message.text or "").strip().lower()
+ if not EMAIL_RE.match(raw) or len(raw) > 255:
+ await message.answer("❌ Неверный формат email. Попробуйте ещё раз.")
+ return
+
+ existing = await get_identity_by_email(session, raw)
+ if existing and existing.tg_id and existing.tg_id != message.from_user.id:
+ await message.answer("❌ Этот email уже занят другим пользователем.")
+ return
+
+ identity = await get_or_create_identity_for_tg(session, message.from_user.id)
+ if identity.email:
+ await state.clear()
+ await message.answer("ℹ️ Почта уже была привязана.")
+ return
+
+ if existing and existing.id != identity.id:
+ logger.warning(
+ "email_binding: identity collision tg_id=%s wants email=%s already on identity_id=%s",
+ message.from_user.id,
+ raw,
+ existing.id,
+ )
+ await message.answer("❌ Этот email уже занят. Используйте другой.")
+ return
+
+ identity.email = raw
+ await session.flush()
+ await state.clear()
+
+ builder = InlineKeyboardBuilder()
+ builder.row(InlineKeyboardButton(text="👤 В кабинет", callback_data="profile"))
+ await message.answer(
+ f"✅ Почта {raw} привязана.",
+ reply_markup=builder.as_markup(),
+ )
diff --git a/handlers/payments/wata/service.py b/handlers/payments/wata/service.py
index 7b2ee23a..e29bf4fb 100644
--- a/handlers/payments/wata/service.py
+++ b/handlers/payments/wata/service.py
@@ -19,7 +19,7 @@ from config import (
WATA_SUCCESS_URL,
)
from core.bootstrap import PAYMENTS_CONFIG
-from database import register_pending_payment
+from database import add_payment, async_session_maker
from database.models import User
from handlers.buttons import BACK, PAY_2, WATA_INT, WATA_RU
from handlers.payments.keyboards import (
@@ -438,15 +438,27 @@ async def generate_wata_payment_link(
logger.error(f"[WATA] В ответе нет поля url: {resp_json}")
return None
- await register_pending_payment(
- payment_id=unique_order_id,
- tg_id=tg_id,
- amount=float(int(amount)),
- payment_system="wata",
- currency="RUB",
- metadata=pending_metadata,
- original_amount=pending_original_amount,
- )
+ async with async_session_maker() as db_session:
+ try:
+ await add_payment(
+ session=db_session,
+ tg_id=tg_id,
+ amount=float(int(amount)),
+ payment_system="wata",
+ status="pending",
+ currency="RUB",
+ payment_id=unique_order_id,
+ metadata=pending_metadata,
+ original_amount=pending_original_amount,
+ )
+ await db_session.commit()
+ except Exception as e:
+ logger.error(
+ f"[WATA] Не удалось записать pending платёж в БД "
+ f"(order_id={unique_order_id}, tg_id={tg_id}): {e}"
+ )
+ await db_session.rollback()
+ return None
logger.info(
f"[WATA] Ссылка создана: tg_id={tg_id}, order_id={unique_order_id}, "
f"rub_amount={amount}, api_amount={api_amount} {currency}"
diff --git a/handlers/profile.py b/handlers/profile.py
index fca2a9c7..e72592d5 100644
--- a/handlers/profile.py
+++ b/handlers/profile.py
@@ -108,7 +108,7 @@ async def process_callback_view_profile(
builder = InlineKeyboardBuilder()
- from core.settings.web_config import get_site_url, is_web_enabled
+ from core.settings.web_config import get_site_url, is_email_binding_enabled, is_web_enabled
if is_web_enabled():
site_url = get_site_url()
@@ -116,6 +116,13 @@ async def process_callback_view_profile(
webapp_url = f"{site_url}/dashboard"
builder.row(InlineKeyboardButton(text="🌐 Личный кабинет", web_app=WebAppInfo(url=webapp_url)))
+ if is_email_binding_enabled():
+ from database.identities import get_identity_by_tg_id
+
+ identity = await get_identity_by_tg_id(session, chat_id)
+ if not (identity and identity.email):
+ builder.row(InlineKeyboardButton(text="📧 Привязать почту", callback_data="bind_email"))
+
trial_time_disabled = bool(MODES_CONFIG.get("TRIAL_TIME_DISABLED", TRIAL_TIME_DISABLE))
if key_count > 0:
diff --git a/middlewares/actor.py b/middlewares/actor.py
index a8f5f3bf..299f5528 100644
--- a/middlewares/actor.py
+++ b/middlewares/actor.py
@@ -26,5 +26,11 @@ class ActorMiddleware(BaseMiddleware):
):
data["actor"] = await resolve_actor_from_legacy_ref(session, int(from_user.id))
except Exception as error:
- logger.error(f"[ActorMiddleware] Ошибка резолва actor: {error}")
+ logger.error(f"[ActorMiddleware] Ошибка резолва actor: {error}", exc_info=True)
+ _session = data.get("session")
+ if _session is not None:
+ try:
+ await _session.rollback()
+ except Exception:
+ pass
return await handler(event, data)
diff --git a/middlewares/admin.py b/middlewares/admin.py
index f0f6c3d0..e11286df 100644
--- a/middlewares/admin.py
+++ b/middlewares/admin.py
@@ -68,4 +68,9 @@ class AdminMiddleware(BaseMiddleware):
await cache_set(cache_key("admin_access", user_id), bool(is_admin), _ADMIN_CACHE_TTL)
return is_admin
except Exception:
+ if session is not None:
+ try:
+ await session.rollback()
+ except Exception:
+ pass
return False
diff --git a/middlewares/user.py b/middlewares/user.py
index 18d2011f..5e099df9 100644
--- a/middlewares/user.py
+++ b/middlewares/user.py
@@ -34,7 +34,13 @@ class UserMiddleware(BaseMiddleware):
if db_user:
data["user"] = db_user
except Exception as e:
- logger.error(f"Ошибка при обработке пользователя: {e}")
+ logger.error(f"Ошибка при обработке пользователя: {e}", exc_info=True)
+ _session = data.get("session")
+ if _session is not None:
+ try:
+ await _session.rollback()
+ except Exception:
+ pass
return await handler(event, data)
async def _process_user(self, user: User, session: AsyncSession) -> dict | None:
diff --git a/utils/versioning.py b/utils/versioning.py
index 1a1487e9..9cfbdaea 100644
--- a/utils/versioning.py
+++ b/utils/versioning.py
@@ -92,4 +92,4 @@ def get_git_commit_number() -> str:
def get_version() -> str:
- return f"v.6-b2704260038 {get_git_commit_number()}"
+ return f"v.6-b0505261115 {get_git_commit_number()}"