backend optimizations + redundant commits cleanup
This commit is contained in:
+32
-3
@@ -4,6 +4,7 @@ import os
|
||||
from time import perf_counter
|
||||
|
||||
from fastapi import Depends, FastAPI, Request
|
||||
from fastapi.responses import ORJSONResponse
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
from starlette.middleware.cors import CORSMiddleware
|
||||
from starlette.middleware.gzip import GZipMiddleware
|
||||
@@ -28,6 +29,7 @@ app = FastAPI(
|
||||
docs_url="/api/docs",
|
||||
redoc_url="/api/redoc",
|
||||
openapi_url="/api/openapi.json",
|
||||
default_response_class=ORJSONResponse,
|
||||
)
|
||||
|
||||
_cors_origins = API_CORS_ORIGINS if API_CORS_ORIGINS != ["*"] else API_CORS_ORIGINS
|
||||
@@ -43,6 +45,8 @@ app.add_middleware(
|
||||
|
||||
app.add_middleware(GZipMiddleware, minimum_size=1024, compresslevel=6)
|
||||
|
||||
_ETAG_MAX_BODY_BYTES = 256 * 1024
|
||||
|
||||
|
||||
@app.middleware("http")
|
||||
async def security_and_cache_middleware(request: Request, call_next):
|
||||
@@ -60,9 +64,32 @@ async def security_and_cache_middleware(request: Request, call_next):
|
||||
return response
|
||||
|
||||
if request.method == "GET" and response.status_code == 200 and "application/json" in content_type:
|
||||
body = b""
|
||||
content_length_header = response.headers.get("content-length")
|
||||
try:
|
||||
cl = int(content_length_header) if content_length_header is not None else None
|
||||
except (TypeError, ValueError):
|
||||
cl = None
|
||||
if cl is not None and cl > _ETAG_MAX_BODY_BYTES:
|
||||
response.headers.setdefault("Cache-Control", "no-cache")
|
||||
return response
|
||||
chunks: list[bytes] = []
|
||||
total = 0
|
||||
too_big = False
|
||||
async for chunk in response.body_iterator:
|
||||
body += chunk
|
||||
total += len(chunk)
|
||||
if total > _ETAG_MAX_BODY_BYTES:
|
||||
chunks.append(chunk)
|
||||
too_big = True
|
||||
async for remaining in response.body_iterator:
|
||||
chunks.append(remaining)
|
||||
break
|
||||
chunks.append(chunk)
|
||||
body = b"".join(chunks)
|
||||
if too_big:
|
||||
headers = dict(response.headers)
|
||||
headers["Cache-Control"] = "no-cache"
|
||||
headers.pop("content-length", None)
|
||||
return StarletteResponse(content=body, status_code=200, headers=headers, media_type=response.media_type)
|
||||
etag = '"' + hashlib.md5(body).hexdigest() + '"'
|
||||
if_none_match = request.headers.get("if-none-match", "")
|
||||
client_etags = [t.strip() for t in if_none_match.split(",") if t.strip()]
|
||||
@@ -80,12 +107,13 @@ async def security_and_cache_middleware(request: Request, call_next):
|
||||
@app.middleware("http")
|
||||
async def api_access_log_middleware(request: Request, call_next):
|
||||
context = ensure_api_context(request)
|
||||
started = perf_counter()
|
||||
if not API_LOGGING:
|
||||
response = await call_next(request)
|
||||
response.headers["X-Request-Id"] = context.request_id
|
||||
response.headers["X-Response-Time"] = f"{int((perf_counter() - started) * 1000)}ms"
|
||||
return response
|
||||
|
||||
started = perf_counter()
|
||||
try:
|
||||
response = await call_next(request)
|
||||
except Exception as exc:
|
||||
@@ -110,6 +138,7 @@ async def api_access_log_middleware(request: Request, call_next):
|
||||
|
||||
duration_ms = int((perf_counter() - started) * 1000)
|
||||
response.headers["X-Request-Id"] = context.request_id
|
||||
response.headers["X-Response-Time"] = f"{duration_ms}ms"
|
||||
result = "success" if response.status_code < 400 else "fail"
|
||||
log_api_access(
|
||||
request,
|
||||
|
||||
@@ -155,7 +155,6 @@ def generate_crud_router(
|
||||
data["days"] = None
|
||||
obj = model(**data)
|
||||
session.add(obj)
|
||||
await session.commit()
|
||||
await session.refresh(obj)
|
||||
return to_schema(schema_response, obj)
|
||||
|
||||
@@ -181,7 +180,6 @@ def generate_crud_router(
|
||||
for k, v in validated.model_dump(exclude_unset=True).items():
|
||||
setattr(obj, k, v)
|
||||
|
||||
await session.commit()
|
||||
await session.refresh(obj)
|
||||
return to_schema(schema_response, obj)
|
||||
|
||||
@@ -203,7 +201,6 @@ def generate_crud_router(
|
||||
raise HTTPException(status_code=404, detail=f"{model.__name__} not found")
|
||||
|
||||
await session.delete(obj)
|
||||
await session.commit()
|
||||
return {"detail": f"{model.__name__} deleted"}
|
||||
|
||||
return router
|
||||
|
||||
@@ -65,5 +65,4 @@ async def delete_gift_with_usages(
|
||||
|
||||
await session.execute(delete(GiftUsage).where(GiftUsage.gift_id == gift_id))
|
||||
await session.delete(gift)
|
||||
await session.commit()
|
||||
return {"message": "Подарок и связанные использования удалены"}
|
||||
|
||||
@@ -44,7 +44,6 @@ async def delete_key_by_email(
|
||||
cluster_id=db_key.server_id,
|
||||
)
|
||||
await session.delete(db_key)
|
||||
await session.commit()
|
||||
logger.info(f"[API] Ключ удалён: {db_key.client_id}")
|
||||
return {"message": "Ключ успешно удалён"}
|
||||
|
||||
@@ -108,7 +107,6 @@ async def edit_key_by_email(
|
||||
hwid_device_limit=getattr(db_key, "device_limit", None),
|
||||
reset_traffic=True,
|
||||
)
|
||||
await session.commit()
|
||||
|
||||
logger.info(f"[API] Ключ обновлён: {db_key.client_id}")
|
||||
return db_key
|
||||
|
||||
@@ -195,7 +195,6 @@ async def change_domain(
|
||||
)
|
||||
)
|
||||
result = await session.execute(stmt)
|
||||
await session.commit()
|
||||
|
||||
return {"updated": result.rowcount or 0}
|
||||
|
||||
@@ -214,7 +213,6 @@ async def restore_trials(
|
||||
.values(trial=0)
|
||||
)
|
||||
result = await session.execute(stmt)
|
||||
await session.commit()
|
||||
|
||||
return {"restored": result.rowcount or 0}
|
||||
|
||||
|
||||
+954
-964
File diff suppressed because it is too large
Load Diff
@@ -42,5 +42,4 @@ async def delete_one_referral(
|
||||
if not obj:
|
||||
raise HTTPException(status_code=404, detail="Referral not found")
|
||||
await session.delete(obj)
|
||||
await session.commit()
|
||||
return {"status": "deleted_one"}
|
||||
|
||||
@@ -125,7 +125,6 @@ async def upsert_setting(
|
||||
value=payload.value,
|
||||
description=payload.description,
|
||||
)
|
||||
await session.commit()
|
||||
await session.refresh(obj)
|
||||
settings_cache.update(
|
||||
key,
|
||||
@@ -148,6 +147,5 @@ async def delete_setting(
|
||||
if not obj:
|
||||
raise HTTPException(status_code=404, detail="Setting not found")
|
||||
await session.delete(obj)
|
||||
await session.commit()
|
||||
settings_cache.delete(key)
|
||||
return {"detail": "Setting deleted"}
|
||||
|
||||
@@ -56,7 +56,6 @@ async def apply_coupon(
|
||||
tg_id=tg_id,
|
||||
code=str(body.code or ""),
|
||||
)
|
||||
await session.commit()
|
||||
return CouponApplyResponse(
|
||||
ok=True,
|
||||
message="Купон успешно активирован",
|
||||
|
||||
@@ -76,7 +76,6 @@ async def create_flow(
|
||||
)
|
||||
session.add(flow)
|
||||
await bump_site_revision(session)
|
||||
await session.commit()
|
||||
await session.refresh(flow)
|
||||
return _flow_to_response(flow)
|
||||
|
||||
@@ -101,7 +100,6 @@ async def update_flow(
|
||||
flow.updated_at = datetime.now(UTC)
|
||||
|
||||
await bump_site_revision(session)
|
||||
await session.commit()
|
||||
await session.refresh(flow)
|
||||
return _flow_to_response(flow)
|
||||
|
||||
@@ -117,4 +115,3 @@ async def delete_flow(
|
||||
raise HTTPException(404, "Flow not found")
|
||||
await session.delete(flow)
|
||||
await bump_site_revision(session)
|
||||
await session.commit()
|
||||
|
||||
@@ -244,7 +244,6 @@ async def delete_my_gift(
|
||||
raise HTTPException(status_code=404, detail="Подарок не найден")
|
||||
await session.execute(delete(GiftUsage).where(GiftUsage.gift_id == gift_id))
|
||||
await session.delete(gift)
|
||||
await session.commit()
|
||||
return {"ok": True, "message": "Подарок удалён"}
|
||||
|
||||
|
||||
@@ -330,5 +329,4 @@ async def delete_gift_with_usages(
|
||||
raise HTTPException(status_code=404, detail="Gift not found")
|
||||
await session.execute(delete(GiftUsage).where(GiftUsage.gift_id == gift_id))
|
||||
await session.delete(gift)
|
||||
await session.commit()
|
||||
return {"message": "Подарок и связанные использования удалены"}
|
||||
|
||||
@@ -21,7 +21,6 @@ async def delete_key_by_email(
|
||||
cluster_id=db_key.server_id,
|
||||
)
|
||||
await session.delete(db_key)
|
||||
await session.commit()
|
||||
logger.info(f"[API] Ключ удалён: {db_key.client_id}")
|
||||
return {"message": "Ключ успешно удалён"}
|
||||
except Exception as e:
|
||||
@@ -82,7 +81,6 @@ async def edit_key_by_email(
|
||||
hwid_device_limit=getattr(db_key, "device_limit", None),
|
||||
reset_traffic=True,
|
||||
)
|
||||
await session.commit()
|
||||
logger.info(f"[API] Ключ обновлён: {db_key.client_id}")
|
||||
return db_key
|
||||
except Exception as e:
|
||||
|
||||
@@ -647,7 +647,6 @@ async def user_key_apply_addons(
|
||||
await update_balance(session, int(billing_user_id), -int(final_extra_price_rub))
|
||||
if coupon_id is not None:
|
||||
await mark_coupon_used(session, int(coupon_id), int(billing_user_id))
|
||||
await session.commit()
|
||||
return AccountKeyApplyAddonsResponse(
|
||||
ok=True,
|
||||
message="Доп. опции применены",
|
||||
|
||||
@@ -195,7 +195,6 @@ async def user_key_update_alias(
|
||||
if db_key is None:
|
||||
raise HTTPException(status_code=404, detail="Подписка не найдена")
|
||||
db_key.alias = alias
|
||||
await session.commit()
|
||||
return AccountKeyResponse(
|
||||
email=str(getattr(db_key, "email", "") or ""),
|
||||
alias=getattr(db_key, "alias", None),
|
||||
@@ -237,5 +236,4 @@ async def user_key_delete(
|
||||
session=session,
|
||||
)
|
||||
await session.delete(db_key)
|
||||
await session.commit()
|
||||
return AccountKeyActionResponse(ok=True, message="Подписка удалена")
|
||||
|
||||
@@ -215,7 +215,6 @@ async def user_key_change_location(
|
||||
db_key.client_id = key_client_id
|
||||
db_key.key = public_link if isinstance(public_link, str) and public_link.strip() else None
|
||||
db_key.remnawave_link = remnawave_link
|
||||
await session.commit()
|
||||
return AccountKeyChangeLocationResponse(
|
||||
ok=True,
|
||||
message="Локация успешно изменена",
|
||||
|
||||
@@ -174,7 +174,6 @@ async def user_key_renew(
|
||||
)
|
||||
except ServiceError as e:
|
||||
raise HTTPException(status_code=400, detail=e.message)
|
||||
await session.commit()
|
||||
return AccountKeyRenewResponse(
|
||||
ok=True,
|
||||
message="Подписка продлена",
|
||||
|
||||
@@ -215,7 +215,6 @@ async def change_domain(
|
||||
)
|
||||
)
|
||||
result = await session.execute(stmt)
|
||||
await session.commit()
|
||||
return {"updated": result.rowcount or 0}
|
||||
|
||||
|
||||
@@ -235,7 +234,6 @@ async def restore_trials(
|
||||
.values(trial=0)
|
||||
)
|
||||
result = await session.execute(stmt)
|
||||
await session.commit()
|
||||
return {"restored": result.rowcount or 0}
|
||||
|
||||
|
||||
|
||||
+44
-56
@@ -9,7 +9,7 @@ from urllib.parse import urlsplit
|
||||
import qrcode
|
||||
|
||||
from fastapi import APIRouter, Depends, HTTPException, Path, Query, Request
|
||||
from fastapi.responses import JSONResponse, StreamingResponse
|
||||
from fastapi.responses import ORJSONResponse, StreamingResponse
|
||||
from sqlalchemy import text
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
@@ -227,7 +227,6 @@ async def partner_apply(
|
||||
)
|
||||
except Exception:
|
||||
pass
|
||||
await session.commit()
|
||||
return PartnerApplyResponse(
|
||||
ok=True,
|
||||
message="Партнерский код применен",
|
||||
@@ -516,7 +515,6 @@ async def partner_create_payout_request(
|
||||
text("UPDATE users SET partner_balance = :balance WHERE id = :id"),
|
||||
{"balance": new_balance, "id": int(user_id)},
|
||||
)
|
||||
await session.commit()
|
||||
return PartnerPayoutRequestResponse(
|
||||
ok=True,
|
||||
message="Заявка на вывод создана",
|
||||
@@ -577,7 +575,7 @@ async def get_all_partners(
|
||||
"method": partner[5] or None,
|
||||
"referred_count": int(partner[6] or 0),
|
||||
})
|
||||
return JSONResponse(content={"total": total, "items": partners_list})
|
||||
return ORJSONResponse(content={"total": total, "items": partners_list})
|
||||
|
||||
|
||||
@router.get("/stats/all")
|
||||
@@ -623,7 +621,7 @@ async def get_partners_stats(
|
||||
"top_partner_tg_id": 0,
|
||||
"top_partner_refs": 0,
|
||||
}
|
||||
return JSONResponse(content=stats)
|
||||
return ORJSONResponse(content=stats)
|
||||
|
||||
|
||||
@router.patch("/{tg_id}")
|
||||
@@ -644,15 +642,14 @@ async def update_partner(
|
||||
"""
|
||||
)
|
||||
result = await session.execute(stmt, {"tg_id": tg_id, "balance": balance, "percent": percent})
|
||||
await session.commit()
|
||||
if result.rowcount > 0:
|
||||
return JSONResponse(
|
||||
return ORJSONResponse(
|
||||
content={"success": True, "message": f"Партнёр {tg_id} успешно обновлён"}, status_code=200
|
||||
)
|
||||
return JSONResponse(content={"success": False, "message": "Партнёр не найден"}, status_code=404)
|
||||
return ORJSONResponse(content={"success": False, "message": "Партнёр не найден"}, status_code=404)
|
||||
except Exception as e:
|
||||
await session.rollback()
|
||||
return JSONResponse(content={"success": False, "message": str(e)}, status_code=500)
|
||||
return ORJSONResponse(content={"success": False, "message": str(e)}, status_code=500)
|
||||
|
||||
|
||||
@router.get("/{tg_id}")
|
||||
@@ -706,7 +703,7 @@ async def get_partner_data(
|
||||
for row in invited_rows
|
||||
],
|
||||
}
|
||||
return JSONResponse(content=response)
|
||||
return ORJSONResponse(content=response)
|
||||
|
||||
|
||||
@router.post("/{tg_id}/invited")
|
||||
@@ -718,18 +715,18 @@ async def add_partner_invited(
|
||||
):
|
||||
"""Добавляет приглашённого пользователя партнёру."""
|
||||
if joined_tg_id == tg_id:
|
||||
return JSONResponse(
|
||||
return ORJSONResponse(
|
||||
content={"success": False, "message": "Нельзя привязать пользователя к самому себе"}, status_code=400
|
||||
)
|
||||
try:
|
||||
partner_exists = await session.execute(text("SELECT 1 FROM users WHERE tg_id = :tg_id"), {"tg_id": tg_id})
|
||||
if not partner_exists.scalar():
|
||||
return JSONResponse(content={"success": False, "message": "Партнёр не найден"}, status_code=404)
|
||||
return ORJSONResponse(content={"success": False, "message": "Партнёр не найден"}, status_code=404)
|
||||
invited_exists = await session.execute(
|
||||
text("SELECT 1 FROM users WHERE tg_id = :joined_tg_id"), {"joined_tg_id": joined_tg_id}
|
||||
)
|
||||
if not invited_exists.scalar():
|
||||
return JSONResponse(
|
||||
return ORJSONResponse(
|
||||
content={"success": False, "message": "Приглашённый пользователь не найден"}, status_code=404
|
||||
)
|
||||
existing = await session.execute(
|
||||
@@ -738,7 +735,7 @@ async def add_partner_invited(
|
||||
)
|
||||
existing_partner = existing.scalar()
|
||||
if existing_partner is not None:
|
||||
return JSONResponse(
|
||||
return ORJSONResponse(
|
||||
content={"success": False, "message": f"Пользователь уже привязан к партнёру {existing_partner}"},
|
||||
status_code=409,
|
||||
)
|
||||
@@ -746,8 +743,7 @@ async def add_partner_invited(
|
||||
text("INSERT INTO partners (partner_tg_id, joined_tg_id) VALUES (:partner_tg_id, :joined_tg_id)"),
|
||||
{"partner_tg_id": tg_id, "joined_tg_id": joined_tg_id},
|
||||
)
|
||||
await session.commit()
|
||||
return JSONResponse(
|
||||
return ORJSONResponse(
|
||||
content={
|
||||
"success": True,
|
||||
"message": "Приглашённый добавлен",
|
||||
@@ -758,7 +754,7 @@ async def add_partner_invited(
|
||||
)
|
||||
except Exception as e:
|
||||
await session.rollback()
|
||||
return JSONResponse(content={"success": False, "message": str(e)}, status_code=500)
|
||||
return ORJSONResponse(content={"success": False, "message": str(e)}, status_code=500)
|
||||
|
||||
|
||||
@router.delete("/{tg_id}/invited/{joined_tg_id}")
|
||||
@@ -774,9 +770,8 @@ async def delete_partner_invited(
|
||||
text("DELETE FROM partners WHERE partner_tg_id = :partner_tg_id AND joined_tg_id = :joined_tg_id"),
|
||||
{"partner_tg_id": tg_id, "joined_tg_id": joined_tg_id},
|
||||
)
|
||||
await session.commit()
|
||||
if result.rowcount > 0:
|
||||
return JSONResponse(
|
||||
return ORJSONResponse(
|
||||
content={
|
||||
"success": True,
|
||||
"message": "Приглашённый удалён",
|
||||
@@ -785,12 +780,12 @@ async def delete_partner_invited(
|
||||
},
|
||||
status_code=200,
|
||||
)
|
||||
return JSONResponse(
|
||||
return ORJSONResponse(
|
||||
content={"success": False, "message": "Связка партнёр-приглашённый не найдена"}, status_code=404
|
||||
)
|
||||
except Exception as e:
|
||||
await session.rollback()
|
||||
return JSONResponse(content={"success": False, "message": str(e)}, status_code=500)
|
||||
return ORJSONResponse(content={"success": False, "message": str(e)}, status_code=500)
|
||||
|
||||
|
||||
@router.patch("/{tg_id}/percent")
|
||||
@@ -803,7 +798,7 @@ async def update_partner_percent(
|
||||
"""Обновляет персональный процент партнёра."""
|
||||
normalized = _parse_percent(percent)
|
||||
if normalized is None:
|
||||
return JSONResponse(
|
||||
return ORJSONResponse(
|
||||
content={"success": False, "message": "Неверный процент. Допустимо 0-100 или 0.0-1.0"}, status_code=400
|
||||
)
|
||||
try:
|
||||
@@ -811,15 +806,14 @@ async def update_partner_percent(
|
||||
text("UPDATE users SET partner_percent = :percent, partner_percent_custom = true WHERE tg_id = :tg_id"),
|
||||
{"tg_id": tg_id, "percent": normalized},
|
||||
)
|
||||
await session.commit()
|
||||
if result.rowcount > 0:
|
||||
return JSONResponse(
|
||||
return ORJSONResponse(
|
||||
content={"success": True, "message": "Процент обновлён", "percent": normalized}, status_code=200
|
||||
)
|
||||
return JSONResponse(content={"success": False, "message": "Партнёр не найден"}, status_code=404)
|
||||
return ORJSONResponse(content={"success": False, "message": "Партнёр не найден"}, status_code=404)
|
||||
except Exception as e:
|
||||
await session.rollback()
|
||||
return JSONResponse(content={"success": False, "message": str(e)}, status_code=500)
|
||||
return ORJSONResponse(content={"success": False, "message": str(e)}, status_code=500)
|
||||
|
||||
|
||||
@router.patch("/{tg_id}/balance")
|
||||
@@ -833,22 +827,22 @@ async def update_partner_balance(
|
||||
"""Изменяет баланс партнёрской программы."""
|
||||
mode_normalized = (mode or "set").strip().lower()
|
||||
if mode_normalized not in {"set", "add", "subtract"}:
|
||||
return JSONResponse(
|
||||
return ORJSONResponse(
|
||||
content={"success": False, "message": "Неверный режим. Используйте set, add или subtract"}, status_code=400
|
||||
)
|
||||
try:
|
||||
amount_val = float(amount)
|
||||
except (TypeError, ValueError):
|
||||
return JSONResponse(content={"success": False, "message": "Неверная сумма"}, status_code=400)
|
||||
return ORJSONResponse(content={"success": False, "message": "Неверная сумма"}, status_code=400)
|
||||
if amount_val < 0:
|
||||
return JSONResponse(content={"success": False, "message": "Сумма не может быть отрицательной"}, status_code=400)
|
||||
return ORJSONResponse(content={"success": False, "message": "Сумма не может быть отрицательной"}, status_code=400)
|
||||
try:
|
||||
current_res = await session.execute(
|
||||
text("SELECT partner_balance FROM users WHERE tg_id = :tg_id"), {"tg_id": tg_id}
|
||||
)
|
||||
current_balance = current_res.scalar()
|
||||
if current_balance is None:
|
||||
return JSONResponse(content={"success": False, "message": "Партнёр не найден"}, status_code=404)
|
||||
return ORJSONResponse(content={"success": False, "message": "Партнёр не найден"}, status_code=404)
|
||||
current_balance = float(current_balance or 0.0)
|
||||
if mode_normalized == "set":
|
||||
new_balance = amount_val
|
||||
@@ -856,19 +850,18 @@ async def update_partner_balance(
|
||||
new_balance = current_balance + amount_val
|
||||
else:
|
||||
if current_balance < amount_val:
|
||||
return JSONResponse(content={"success": False, "message": "Недостаточно средств"}, status_code=400)
|
||||
return ORJSONResponse(content={"success": False, "message": "Недостаточно средств"}, status_code=400)
|
||||
new_balance = current_balance - amount_val
|
||||
await session.execute(
|
||||
text("UPDATE users SET partner_balance = :balance WHERE tg_id = :tg_id"),
|
||||
{"tg_id": tg_id, "balance": new_balance},
|
||||
)
|
||||
await session.commit()
|
||||
return JSONResponse(
|
||||
return ORJSONResponse(
|
||||
content={"success": True, "message": "Баланс обновлён", "balance": new_balance}, status_code=200
|
||||
)
|
||||
except Exception as e:
|
||||
await session.rollback()
|
||||
return JSONResponse(content={"success": False, "message": str(e)}, status_code=500)
|
||||
return ORJSONResponse(content={"success": False, "message": str(e)}, status_code=500)
|
||||
|
||||
|
||||
@router.get("/{tg_id}/invited")
|
||||
@@ -901,7 +894,7 @@ async def get_partner_invited(
|
||||
}
|
||||
for row in invited_rows
|
||||
]
|
||||
return JSONResponse(content=invited_list)
|
||||
return ORJSONResponse(content=invited_list)
|
||||
|
||||
|
||||
@router.get("/payouts/pending")
|
||||
@@ -944,7 +937,7 @@ async def get_partner_payouts_pending(
|
||||
}
|
||||
for row in result.fetchall()
|
||||
]
|
||||
return JSONResponse(content={"total": int(total), "items": items})
|
||||
return ORJSONResponse(content={"total": int(total), "items": items})
|
||||
|
||||
|
||||
@router.get("/payouts/history")
|
||||
@@ -987,7 +980,7 @@ async def get_partner_payouts_history(
|
||||
}
|
||||
for row in result.fetchall()
|
||||
]
|
||||
return JSONResponse(content={"total": int(total), "items": items})
|
||||
return ORJSONResponse(content={"total": int(total), "items": items})
|
||||
|
||||
|
||||
@router.post("/payouts/{payout_id}/approve")
|
||||
@@ -1003,7 +996,7 @@ async def approve_partner_payout(
|
||||
)
|
||||
req = req_row.fetchone()
|
||||
if not req:
|
||||
return JSONResponse(
|
||||
return ORJSONResponse(
|
||||
content={"success": False, "message": "Заявка не найдена или уже обработана"}, status_code=404
|
||||
)
|
||||
user_row = await session.execute(
|
||||
@@ -1019,8 +1012,7 @@ async def approve_partner_payout(
|
||||
),
|
||||
{"id": payout_id, "method": payout_method, "destination": destination},
|
||||
)
|
||||
await session.commit()
|
||||
return JSONResponse(content={"success": True, "message": "Заявка одобрена"}, status_code=200)
|
||||
return ORJSONResponse(content={"success": True, "message": "Заявка одобрена"}, status_code=200)
|
||||
|
||||
|
||||
@router.post("/payouts/{payout_id}/reject")
|
||||
@@ -1036,7 +1028,7 @@ async def reject_partner_payout(
|
||||
)
|
||||
req = req_row.fetchone()
|
||||
if not req:
|
||||
return JSONResponse(
|
||||
return ORJSONResponse(
|
||||
content={"success": False, "message": "Заявка не найдена или уже обработана"}, status_code=404
|
||||
)
|
||||
user_row = await session.execute(
|
||||
@@ -1058,8 +1050,7 @@ async def reject_partner_payout(
|
||||
text("UPDATE users SET partner_balance = :balance WHERE tg_id = :tg_id"),
|
||||
{"balance": current_balance + float(req[2] or 0.0), "tg_id": req[1]},
|
||||
)
|
||||
await session.commit()
|
||||
return JSONResponse(content={"success": True, "message": "Заявка отклонена"}, status_code=200)
|
||||
return ORJSONResponse(content={"success": True, "message": "Заявка отклонена"}, status_code=200)
|
||||
|
||||
|
||||
@router.patch("/{tg_id}/percent/reset")
|
||||
@@ -1073,10 +1064,9 @@ async def reset_partner_percent(
|
||||
text("UPDATE users SET partner_percent = NULL, partner_percent_custom = false WHERE tg_id = :tg_id"),
|
||||
{"tg_id": tg_id},
|
||||
)
|
||||
await session.commit()
|
||||
if result.rowcount > 0:
|
||||
return JSONResponse(content={"success": True, "message": "Процент сброшен"}, status_code=200)
|
||||
return JSONResponse(content={"success": False, "message": "Партнёр не найден"}, status_code=404)
|
||||
return ORJSONResponse(content={"success": True, "message": "Процент сброшен"}, status_code=200)
|
||||
return ORJSONResponse(content={"success": False, "message": "Партнёр не найден"}, status_code=404)
|
||||
|
||||
|
||||
@router.patch("/{tg_id}/code")
|
||||
@@ -1089,23 +1079,22 @@ async def update_partner_code(
|
||||
"""Обновляет код партнёрской ссылки."""
|
||||
raw = (code or "").strip().lower()
|
||||
if not raw:
|
||||
return JSONResponse(content={"success": False, "message": "Код не может быть пустым"}, status_code=400)
|
||||
return ORJSONResponse(content={"success": False, "message": "Код не может быть пустым"}, status_code=400)
|
||||
if not re.fullmatch(r"[a-z0-9_]{3,32}", raw):
|
||||
return JSONResponse(
|
||||
return ORJSONResponse(
|
||||
content={"success": False, "message": "Неверный код. Разрешены a-z, 0-9, _ (3-32 символа)"}, status_code=400
|
||||
)
|
||||
exists = await session.execute(
|
||||
text("SELECT 1 FROM users WHERE partner_code = :code AND tg_id != :tg_id"), {"code": raw, "tg_id": tg_id}
|
||||
)
|
||||
if exists.first():
|
||||
return JSONResponse(content={"success": False, "message": "Такой код уже занят"}, status_code=409)
|
||||
return ORJSONResponse(content={"success": False, "message": "Такой код уже занят"}, status_code=409)
|
||||
result = await session.execute(
|
||||
text("UPDATE users SET partner_code = :code WHERE tg_id = :tg_id"), {"code": raw, "tg_id": tg_id}
|
||||
)
|
||||
await session.commit()
|
||||
if result.rowcount > 0:
|
||||
return JSONResponse(content={"success": True, "message": "Код обновлён", "code": raw}, status_code=200)
|
||||
return JSONResponse(content={"success": False, "message": "Партнёр не найден"}, status_code=404)
|
||||
return ORJSONResponse(content={"success": True, "message": "Код обновлён", "code": raw}, status_code=200)
|
||||
return ORJSONResponse(content={"success": False, "message": "Партнёр не найден"}, status_code=404)
|
||||
|
||||
|
||||
@router.post("/reset-disabled-methods")
|
||||
@@ -1138,12 +1127,11 @@ async def reset_disabled_payout_methods(
|
||||
if not ENABLE_PAYOUT_SBP and B:
|
||||
disabled.append(B.METHOD_SBP)
|
||||
if not disabled:
|
||||
return JSONResponse(content={"success": True, "message": "Отключённых методов нет"}, status_code=200)
|
||||
return ORJSONResponse(content={"success": True, "message": "Отключённых методов нет"}, status_code=200)
|
||||
await session.execute(
|
||||
text("UPDATE users SET card_number = NULL WHERE payout_method = ANY(:methods)"), {"methods": disabled}
|
||||
)
|
||||
await session.commit()
|
||||
return JSONResponse(content={"success": True, "message": "Отключённые методы сброшены"}, status_code=200)
|
||||
return ORJSONResponse(content={"success": True, "message": "Отключённые методы сброшены"}, status_code=200)
|
||||
|
||||
|
||||
@router.get("/{tg_id}/export")
|
||||
@@ -1159,7 +1147,7 @@ async def export_partner_invites_csv(
|
||||
)
|
||||
data = rows.fetchall()
|
||||
if not data:
|
||||
return JSONResponse(content={"success": False, "message": "Нет приглашённых"}, status_code=404)
|
||||
return ORJSONResponse(content={"success": False, "message": "Нет приглашённых"}, status_code=404)
|
||||
buffer = StringIO()
|
||||
writer = csv.writer(buffer, delimiter=";")
|
||||
writer.writerow(["joined_tg_id", "created_at"])
|
||||
|
||||
@@ -120,7 +120,6 @@ async def upsert_setting(
|
||||
value=payload.value,
|
||||
description=payload.description,
|
||||
)
|
||||
await session.commit()
|
||||
await session.refresh(obj)
|
||||
settings_cache.update(
|
||||
key,
|
||||
@@ -144,6 +143,5 @@ async def delete_setting(
|
||||
if not obj:
|
||||
raise HTTPException(status_code=404, detail="Setting not found")
|
||||
await session.delete(obj)
|
||||
await session.commit()
|
||||
settings_cache.delete(key)
|
||||
return {"detail": "Setting deleted"}
|
||||
|
||||
@@ -302,7 +302,6 @@ async def purchase_tariff_with_balance(
|
||||
)
|
||||
if coupon_id is not None:
|
||||
await mark_coupon_used(session, int(coupon_id), int(tg_id))
|
||||
await session.commit()
|
||||
except Exception:
|
||||
logger.exception("web tariff purchase failed")
|
||||
raise HTTPException(status_code=500, detail="Не удалось оформить подписку") from None
|
||||
|
||||
@@ -40,6 +40,10 @@ from logger import logger
|
||||
UPLOAD_DIR = Path("static/web_uploads")
|
||||
ALLOWED_EXTENSIONS = frozenset({".png", ".jpg", ".jpeg", ".gif", ".webp", ".svg", ".mp4", ".webm"})
|
||||
MAX_FILE_SIZE = 100 * 1024 * 1024
|
||||
_IMAGE_RESIZE_EXTENSIONS = frozenset({".png", ".jpg", ".jpeg", ".webp"})
|
||||
_IMAGE_MAX_SIDE = 2048
|
||||
_IMAGE_JPEG_QUALITY = 85
|
||||
_IMAGE_WEBP_QUALITY = 85
|
||||
|
||||
_SLUG_RE = re.compile(r"^[a-z0-9][a-z0-9\-]*$")
|
||||
|
||||
@@ -55,6 +59,39 @@ EXTENSION_CONTENT_TYPES: dict[str, frozenset[str]] = {
|
||||
}
|
||||
|
||||
|
||||
def _optimize_image_bytes(data: bytes, ext: str) -> bytes:
|
||||
"""Уменьшает большие картинки до _IMAGE_MAX_SIDE и пережимает с разумным качеством."""
|
||||
try:
|
||||
from io import BytesIO
|
||||
|
||||
from PIL import Image, ImageOps
|
||||
|
||||
with Image.open(BytesIO(data)) as img:
|
||||
img = ImageOps.exif_transpose(img)
|
||||
w, h = img.size
|
||||
if max(w, h) <= _IMAGE_MAX_SIDE and len(data) < 500_000:
|
||||
return data
|
||||
if max(w, h) > _IMAGE_MAX_SIDE:
|
||||
img.thumbnail((_IMAGE_MAX_SIDE, _IMAGE_MAX_SIDE), Image.Resampling.LANCZOS)
|
||||
buffer = BytesIO()
|
||||
save_kwargs: dict = {}
|
||||
if ext in (".jpg", ".jpeg"):
|
||||
if img.mode not in ("RGB", "L"):
|
||||
img = img.convert("RGB")
|
||||
save_kwargs = {"format": "JPEG", "quality": _IMAGE_JPEG_QUALITY, "optimize": True, "progressive": True}
|
||||
elif ext == ".webp":
|
||||
save_kwargs = {"format": "WEBP", "quality": _IMAGE_WEBP_QUALITY, "method": 6}
|
||||
elif ext == ".png":
|
||||
save_kwargs = {"format": "PNG", "optimize": True}
|
||||
else:
|
||||
return data
|
||||
img.save(buffer, **save_kwargs)
|
||||
optimized = buffer.getvalue()
|
||||
return optimized if len(optimized) < len(data) else data
|
||||
except Exception:
|
||||
return data
|
||||
|
||||
|
||||
def _sanitize_svg(data: bytes) -> bytes:
|
||||
import re as _re
|
||||
|
||||
@@ -529,6 +566,10 @@ async def upload_media(
|
||||
file_data = b"".join(chunks)
|
||||
if ext == ".svg":
|
||||
file_data = _sanitize_svg(file_data)
|
||||
elif ext in _IMAGE_RESIZE_EXTENSIONS:
|
||||
from core.executor import run_cpu
|
||||
|
||||
file_data = await run_cpu(_optimize_image_bytes, file_data, ext)
|
||||
with open(path, "wb") as f:
|
||||
f.write(file_data)
|
||||
url = f"/api/web/uploads/{name}"
|
||||
|
||||
Binary file not shown.
@@ -1,6 +1,7 @@
|
||||
import asyncio
|
||||
|
||||
from apscheduler.triggers.cron import CronTrigger
|
||||
from apscheduler.triggers.interval import IntervalTrigger
|
||||
|
||||
from database import async_session_maker, cancel_expired_pending_payments
|
||||
from logger import logger
|
||||
@@ -65,7 +66,29 @@ def cleanup_expired_gifts_process_runner() -> None:
|
||||
asyncio.run(cleanup_expired_gifts_job())
|
||||
|
||||
|
||||
async def log_db_pool_status() -> None:
|
||||
"""Раз в минуту логирует состояние пула соединений: даёт видимость «упираемся ли в лимит»."""
|
||||
try:
|
||||
from database.db import engine
|
||||
|
||||
pool = engine.pool
|
||||
size = pool.size()
|
||||
checked_out = pool.checkedout()
|
||||
checked_in = getattr(pool, "checkedin", lambda: size - checked_out)()
|
||||
overflow = getattr(pool, "overflow", lambda: -1)()
|
||||
logger.info(
|
||||
"[DBPool] size={} in_use={} idle={} overflow={}",
|
||||
size,
|
||||
checked_out,
|
||||
checked_in,
|
||||
overflow,
|
||||
)
|
||||
except Exception as error:
|
||||
logger.debug("[DBPool] не удалось получить статус пула: {}", error)
|
||||
|
||||
|
||||
AUDIT_DRAIN_TRIGGER = CronTrigger(hour=0, minute=0, timezone="Europe/Moscow")
|
||||
DAILY_STATS_REPORT_TRIGGER = CronTrigger(hour=0, minute=1, timezone="Europe/Moscow")
|
||||
STALE_PAYMENTS_SWEEP_TRIGGER = CronTrigger(minute=0, timezone="Europe/Moscow")
|
||||
EXPIRED_GIFTS_CLEANUP_TRIGGER = CronTrigger(hour=3, minute=0, timezone="Europe/Moscow")
|
||||
DB_POOL_STATUS_TRIGGER = IntervalTrigger(minutes=1)
|
||||
|
||||
@@ -1,10 +1,12 @@
|
||||
from core.tasks.cron_tasks import (
|
||||
AUDIT_DRAIN_TRIGGER,
|
||||
DAILY_STATS_REPORT_TRIGGER,
|
||||
DB_POOL_STATUS_TRIGGER,
|
||||
EXPIRED_GIFTS_CLEANUP_TRIGGER,
|
||||
STALE_PAYMENTS_SWEEP_TRIGGER,
|
||||
cleanup_expired_gifts_job,
|
||||
cleanup_expired_gifts_process_runner,
|
||||
log_db_pool_status,
|
||||
scheduled_audit_drain,
|
||||
scheduled_audit_drain_process_runner,
|
||||
scheduled_stats_report,
|
||||
@@ -110,4 +112,10 @@ def register_periodic_tasks() -> None:
|
||||
EXPIRED_GIFTS_CLEANUP_TRIGGER,
|
||||
)
|
||||
|
||||
periodic_task_manager.register_cron_task(
|
||||
"db_pool_status",
|
||||
log_db_pool_status,
|
||||
DB_POOL_STATUS_TRIGGER,
|
||||
)
|
||||
|
||||
_TASKS_REGISTERED = True
|
||||
|
||||
@@ -1225,6 +1225,28 @@ async def _migration_v23_add_identity_onboarding_stage(conn: AsyncConnection) ->
|
||||
await _exec_ignore(conn, "ALTER TABLE identities ADD COLUMN onboarding_stage VARCHAR(32)")
|
||||
|
||||
|
||||
async def _migration_v26_add_keys_indexes(conn: AsyncConnection) -> None:
|
||||
logger.info("[schema_upgrade] v26: индексы keys(expiry_time/server_id/tariff_id)")
|
||||
if not await _table_exists(conn, "keys"):
|
||||
return
|
||||
if not await _index_exists(conn, "keys", "ix_keys_expiry_time"):
|
||||
await _exec_ignore(conn, "CREATE INDEX ix_keys_expiry_time ON keys(expiry_time)")
|
||||
if not await _index_exists(conn, "keys", "ix_keys_server_id"):
|
||||
await _exec_ignore(conn, "CREATE INDEX ix_keys_server_id ON keys(server_id)")
|
||||
if not await _index_exists(conn, "keys", "ix_keys_tariff_id"):
|
||||
await _exec_ignore(conn, "CREATE INDEX ix_keys_tariff_id ON keys(tariff_id)")
|
||||
|
||||
|
||||
async def _migration_v25_add_partners_indexes(conn: AsyncConnection) -> None:
|
||||
logger.info("[schema_upgrade] v25: индексы на partners(partner_tg_id/joined_tg_id)")
|
||||
if not await _table_exists(conn, "partners"):
|
||||
return
|
||||
if not await _index_exists(conn, "partners", "ix_partners_partner_tg_id"):
|
||||
await _exec_ignore(conn, "CREATE INDEX ix_partners_partner_tg_id ON partners(partner_tg_id)")
|
||||
if not await _index_exists(conn, "partners", "ix_partners_joined_tg_id"):
|
||||
await _exec_ignore(conn, "CREATE INDEX ix_partners_joined_tg_id ON partners(joined_tg_id)")
|
||||
|
||||
|
||||
async def _migration_v24_add_identity_sessions(conn: AsyncConnection) -> None:
|
||||
logger.info("[schema_upgrade] v24: таблица identity_sessions + перенос существующих токенов")
|
||||
if not await _table_exists(conn, "identities"):
|
||||
@@ -1302,6 +1324,8 @@ _MIGRATIONS = [
|
||||
(22, "identities.onboarding_completed_at", _migration_v22_add_identity_onboarding_completed_at),
|
||||
(23, "identities.onboarding_stage", _migration_v23_add_identity_onboarding_stage),
|
||||
(24, "таблица identity_sessions (мультидевайс)", _migration_v24_add_identity_sessions),
|
||||
(25, "индексы на partners(partner_tg_id/joined_tg_id)", _migration_v25_add_partners_indexes),
|
||||
(26, "индексы keys(expiry_time/server_id/tariff_id)", _migration_v26_add_keys_indexes),
|
||||
]
|
||||
|
||||
|
||||
|
||||
@@ -13,11 +13,11 @@ class Key(DictLikeMixin, Base):
|
||||
tg_id = Column(BigInteger, ForeignKey("users.tg_id"), nullable=True, index=True)
|
||||
email = Column(String, unique=True)
|
||||
created_at = Column(BigInteger)
|
||||
expiry_time = Column(BigInteger)
|
||||
expiry_time = Column(BigInteger, index=True)
|
||||
key = Column(String)
|
||||
server_id = Column(String)
|
||||
server_id = Column(String, index=True)
|
||||
remnawave_link = Column(String)
|
||||
tariff_id = Column(Integer, ForeignKey("tariffs.id", ondelete="SET NULL"))
|
||||
tariff_id = Column(Integer, ForeignKey("tariffs.id", ondelete="SET NULL"), index=True)
|
||||
is_frozen = Column(Boolean, default=False)
|
||||
alias = Column(String)
|
||||
notified = Column(Boolean, default=False)
|
||||
|
||||
@@ -1,21 +1,30 @@
|
||||
from sqlalchemy import select
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from core.redis_cache import cache_delete, cache_get, cache_set
|
||||
from database.models import Setting
|
||||
|
||||
|
||||
_KEY = "CONTENT_REVISION"
|
||||
_CACHE_KEY = "site_revision:value"
|
||||
_CACHE_TTL_SEC = 30
|
||||
|
||||
|
||||
async def get_site_revision(session: AsyncSession) -> int:
|
||||
cached = await cache_get(_CACHE_KEY)
|
||||
if isinstance(cached, int):
|
||||
return cached
|
||||
result = await session.execute(select(Setting).where(Setting.key == _KEY))
|
||||
setting = result.scalar_one_or_none()
|
||||
if setting is None:
|
||||
await cache_set(_CACHE_KEY, 0, _CACHE_TTL_SEC)
|
||||
return 0
|
||||
try:
|
||||
return int(setting.value or 0)
|
||||
value = int(setting.value or 0)
|
||||
except (TypeError, ValueError):
|
||||
return 0
|
||||
value = 0
|
||||
await cache_set(_CACHE_KEY, value, _CACHE_TTL_SEC)
|
||||
return value
|
||||
|
||||
|
||||
async def bump_site_revision(session: AsyncSession) -> int:
|
||||
@@ -23,10 +32,12 @@ async def bump_site_revision(session: AsyncSession) -> int:
|
||||
setting = result.scalar_one_or_none()
|
||||
if setting is None:
|
||||
session.add(Setting(key=_KEY, value=1))
|
||||
await cache_delete(_CACHE_KEY)
|
||||
return 1
|
||||
try:
|
||||
current = int(setting.value or 0)
|
||||
except (TypeError, ValueError):
|
||||
current = 0
|
||||
setting.value = current + 1
|
||||
await cache_delete(_CACHE_KEY)
|
||||
return current + 1
|
||||
|
||||
+11
-1
@@ -1,16 +1,24 @@
|
||||
from sqlalchemy import select
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from core.redis_cache import cache_delete, cache_get, cache_set
|
||||
from database.models import Setting
|
||||
|
||||
|
||||
_KEY = "SITE_INITIALIZED"
|
||||
_CACHE_KEY = "site_state:initialized"
|
||||
_CACHE_TTL_SEC = 300
|
||||
|
||||
|
||||
async def is_site_initialized(session: AsyncSession) -> bool:
|
||||
cached = await cache_get(_CACHE_KEY)
|
||||
if isinstance(cached, bool):
|
||||
return cached
|
||||
result = await session.execute(select(Setting).where(Setting.key == _KEY))
|
||||
setting = result.scalar_one_or_none()
|
||||
return bool(setting and setting.value is True)
|
||||
value = bool(setting and setting.value is True)
|
||||
await cache_set(_CACHE_KEY, value, _CACHE_TTL_SEC)
|
||||
return value
|
||||
|
||||
|
||||
async def mark_site_initialized(session: AsyncSession) -> None:
|
||||
@@ -21,6 +29,7 @@ async def mark_site_initialized(session: AsyncSession) -> None:
|
||||
session.add(Setting(key=_KEY, value=True, description="Сайт прошёл первую настройку админом"))
|
||||
elif setting.value is not True:
|
||||
setting.value = True
|
||||
await cache_delete(_CACHE_KEY)
|
||||
|
||||
|
||||
async def reset_site_initialized(session: AsyncSession) -> None:
|
||||
@@ -29,3 +38,4 @@ async def reset_site_initialized(session: AsyncSession) -> None:
|
||||
setting = result.scalar_one_or_none()
|
||||
if setting is not None:
|
||||
setting.value = False
|
||||
await cache_delete(_CACHE_KEY)
|
||||
|
||||
@@ -161,7 +161,6 @@ async def handle_ads_delete(
|
||||
try:
|
||||
await session.execute(update(User).where(User.source_code == code).values(source_code=None))
|
||||
await session.execute(delete(TrackingSource).where(TrackingSource.code == code))
|
||||
await session.commit()
|
||||
await cache_delete(cache_key("utm_exists", code))
|
||||
await callback_query.message.edit_text(
|
||||
f"🗑️ Ссылка <code>{code}</code> удалена.",
|
||||
|
||||
@@ -232,7 +232,6 @@ async def handle_clear_blocked_users(callback_query: CallbackQuery, session: Asy
|
||||
return
|
||||
|
||||
await session.execute(delete(BlockedUser))
|
||||
await session.commit()
|
||||
|
||||
await callback_query.message.answer(
|
||||
text=f"🗑️ Очищено {total_count} записей забанивших пользователей из базы данных.",
|
||||
@@ -271,7 +270,6 @@ async def handle_clear_shadow_bans(callback_query: CallbackQuery, session: Async
|
||||
)
|
||||
tg_to_invalidate = [r[0] for r in tg_ids_result.all() if r[0] is not None]
|
||||
await session.execute(delete(ManualBan).where(ManualBan.reason == "shadow"))
|
||||
await session.commit()
|
||||
for tid in tg_to_invalidate:
|
||||
await invalidate_ban_cache(tid)
|
||||
|
||||
@@ -314,7 +312,6 @@ async def handle_clear_manual_bans(callback_query: CallbackQuery, session: Async
|
||||
)
|
||||
tg_to_invalidate = [r[0] for r in tg_ids_result.all() if r[0] is not None]
|
||||
await session.execute(delete(ManualBan).where(or_(ManualBan.reason != "shadow", ManualBan.reason.is_(None))))
|
||||
await session.commit()
|
||||
for tid in tg_to_invalidate:
|
||||
await invalidate_ban_cache(tid)
|
||||
|
||||
@@ -408,7 +405,6 @@ async def handle_preemptive_ids_input(message: Message, state: FSMContext, sessi
|
||||
)
|
||||
|
||||
await session.execute(stmt)
|
||||
await session.commit()
|
||||
for tid in cache_tg_ids:
|
||||
await invalidate_ban_cache(tid)
|
||||
|
||||
|
||||
@@ -193,7 +193,6 @@ async def handle_days_input(message: Message, state: FSMContext, session: AsyncS
|
||||
.values(expiry_time=Key.expiry_time + add_ms)
|
||||
.execution_options(synchronize_session=False)
|
||||
)
|
||||
await session.commit()
|
||||
|
||||
await message.answer(
|
||||
f"✅ Время подписки продлено на <b>{days} дней</b> для <b>{affected}</b> пользователей в кластере <b>{cluster_name}</b>."
|
||||
@@ -327,7 +326,6 @@ async def handle_new_cluster_name_input(message: Message, state: FSMContext, ses
|
||||
update(Key).where(Key.server_id == old_cluster_name).values(server_id=new_cluster_name)
|
||||
)
|
||||
|
||||
await session.commit()
|
||||
|
||||
await message.answer(
|
||||
text=f"✅ Название кластера успешно изменено с '{old_cluster_name}' на '{new_cluster_name}'!",
|
||||
@@ -442,7 +440,6 @@ async def handle_new_server_name_input(message: Message, state: FSMContext, sess
|
||||
if keys_count > 0:
|
||||
await session.execute(update(Key).where(Key.server_id == old_server_name).values(server_id=new_server_name))
|
||||
|
||||
await session.commit()
|
||||
|
||||
await message.answer(
|
||||
text=(
|
||||
|
||||
@@ -370,7 +370,6 @@ async def handle_sync_server(
|
||||
.where(Key.user_id == key["user_id"], Key.client_id == key["client_id"])
|
||||
.values(remnawave_link=new_remnawave_link, key=key_value)
|
||||
)
|
||||
await session.commit()
|
||||
logger.info(f"[Sync] Обновлена ссылка для {key['email']}: {new_remnawave_link}")
|
||||
except Exception as e:
|
||||
logger.warning(f"[Sync] Не удалось получить ссылку для {key['email']}: {e}")
|
||||
@@ -683,7 +682,6 @@ async def handle_sync_cluster(
|
||||
await session.run_sync(
|
||||
lambda sync_session: sync_session.bulk_update_mappings(Key, bulk_updates)
|
||||
)
|
||||
await session.commit()
|
||||
logger.info(f"[Sync] Bulk: обновлено {len(bulk_updates)} ключей")
|
||||
except Exception as bulk_error:
|
||||
logger.warning(f"[Sync] Bulk упал, fallback: {bulk_error}")
|
||||
@@ -696,7 +694,6 @@ async def handle_sync_cluster(
|
||||
.where(Key.client_id == upd["client_id"])
|
||||
.values(remnawave_link=upd["remnawave_link"], key=upd["key"])
|
||||
)
|
||||
await session.commit()
|
||||
except Exception as e:
|
||||
logger.error(f"[Sync] Fallback ошибка {upd['client_id']}: {e}")
|
||||
await session.rollback()
|
||||
@@ -708,7 +705,6 @@ async def handle_sync_cluster(
|
||||
await session.execute(
|
||||
delete(Key).where(Key.user_id == key["user_id"], Key.client_id == key["client_id"])
|
||||
)
|
||||
await session.commit()
|
||||
|
||||
cluster_id_for_recreate = key["server_id"] if use_country_selection else cluster_name
|
||||
await create_key_on_cluster(
|
||||
|
||||
@@ -62,7 +62,6 @@ async def apply_tariff_group(callback: CallbackQuery, callback_data: AdminCluste
|
||||
group_code = row["group_code"]
|
||||
|
||||
await session.execute(update(Server).where(Server.cluster_name == cluster_name).values(tariff_group=group_code))
|
||||
await session.commit()
|
||||
|
||||
servers = await get_servers(session=session, include_enabled=True)
|
||||
cluster_servers = servers.get(cluster_name, [])
|
||||
@@ -294,7 +293,6 @@ async def apply_tariffs(
|
||||
for sid in to_insert
|
||||
])
|
||||
|
||||
await session.commit()
|
||||
|
||||
await state.update_data({
|
||||
f"subgrp_sel:{cluster_name}": [],
|
||||
@@ -343,7 +341,6 @@ async def reset_cluster_subgroups(callback: CallbackQuery, callback_data: AdminC
|
||||
return
|
||||
|
||||
await session.execute(delete(ServerSubgroup).where(ServerSubgroup.server_id.in_(server_ids)))
|
||||
await session.commit()
|
||||
|
||||
servers = await get_servers(session=session, include_enabled=True)
|
||||
cluster_servers = servers.get(cluster_name, [])
|
||||
@@ -587,7 +584,6 @@ async def apply_group_to_servers(
|
||||
|
||||
if to_insert:
|
||||
session.add_all([ServerSpecialgroup(server_id=sid, group_code=group_code) for sid in to_insert])
|
||||
await session.commit()
|
||||
|
||||
logger.debug(f"[apply_group_to_servers] group={group_code} server_ids={server_ids}")
|
||||
|
||||
@@ -616,7 +612,6 @@ async def reset_cluster_groups(callback: CallbackQuery, callback_data: AdminClus
|
||||
await callback.answer("В кластере нет серверов", show_alert=True)
|
||||
return
|
||||
await session.execute(delete(ServerSpecialgroup).where(ServerSpecialgroup.server_id.in_(server_ids)))
|
||||
await session.commit()
|
||||
servers = await get_servers(session=session, include_enabled=True)
|
||||
cluster_servers = servers.get(cluster_name, [])
|
||||
await callback.message.edit_text(
|
||||
|
||||
@@ -33,7 +33,6 @@ async def handle_server_transfer(callback_query: CallbackQuery, state: FSMContex
|
||||
)
|
||||
)
|
||||
|
||||
await session.commit()
|
||||
|
||||
use_country_selection = bool(MODES_CONFIG.get("COUNTRY_SELECTION_ENABLED", USE_COUNTRY_SELECTION))
|
||||
|
||||
@@ -77,7 +76,6 @@ async def handle_cluster_transfer(callback_query: CallbackQuery, state: FSMConte
|
||||
)
|
||||
)
|
||||
|
||||
await session.commit()
|
||||
|
||||
await callback_query.message.edit_text(
|
||||
text=(
|
||||
|
||||
@@ -231,7 +231,6 @@ async def handle_panel_type_selection(
|
||||
)
|
||||
|
||||
session.add(new_server)
|
||||
await session.commit()
|
||||
|
||||
await callback_query.message.edit_text(
|
||||
text=f"✅ Сервер <b>{server_name}</b> с панелью <b>{panel_type}</b> успешно добавлен в кластер <b>{cluster_name}</b>!",
|
||||
|
||||
@@ -243,7 +243,6 @@ async def delete_gift(callback: CallbackQuery, session: AsyncSession):
|
||||
|
||||
await session.execute(delete(GiftUsage).where(GiftUsage.gift_id == gift_id))
|
||||
await session.execute(delete(Gift).where(Gift.gift_id == gift_id))
|
||||
await session.commit()
|
||||
|
||||
builder = InlineKeyboardBuilder()
|
||||
builder.button(text="🔙 Назад к списку", callback_data="admin_gifts_all")
|
||||
|
||||
@@ -55,7 +55,6 @@ async def save_new_admin(message: Message, session: AsyncSession, state: FSMCont
|
||||
await message.answer("⚠️ Такой админ уже существует.")
|
||||
else:
|
||||
session.add(Admin(tg_id=tg_id, role="moderator", description="Добавлен вручную"))
|
||||
await session.commit()
|
||||
await message.answer(f"✅ Админ <code>{tg_id}</code> добавлен.", reply_markup=build_admin_back_kb_to_admins())
|
||||
|
||||
await state.clear()
|
||||
@@ -87,7 +86,6 @@ async def generate_token(callback: CallbackQuery, callback_data: AdminPanelCallb
|
||||
token = Admin.generate_token()
|
||||
token_hash = hashlib.sha256(token.encode()).hexdigest()
|
||||
admin.token = token_hash
|
||||
await session.commit()
|
||||
|
||||
msg = await callback.message.edit_text(
|
||||
f"🎟 <b>Новый токен для</b> <code>{tg_id}</code>:\n\n"
|
||||
@@ -135,7 +133,6 @@ async def set_admin_role(callback: CallbackQuery, callback_data: AdminPanelCallb
|
||||
return
|
||||
|
||||
admin.role = role
|
||||
await session.commit()
|
||||
|
||||
await callback.message.edit_text(
|
||||
f"✅ Роль админа <code>{tg_id}</code> изменена на <b>{role}</b>.", reply_markup=build_single_admin_menu(tg_id)
|
||||
@@ -147,7 +144,6 @@ async def delete_admin(callback: CallbackQuery, callback_data: AdminPanelCallbac
|
||||
tg_id = int(callback_data.action.split("|")[1])
|
||||
|
||||
await session.execute(delete(Admin).where(Admin.tg_id == tg_id))
|
||||
await session.commit()
|
||||
|
||||
await callback.message.edit_text(
|
||||
f"🗑 Админ <code>{tg_id}</code> удалён.", reply_markup=build_admin_back_kb_to_admins()
|
||||
|
||||
@@ -55,7 +55,6 @@ async def process_new_domain(message: Message, state: FSMContext, session: Async
|
||||
)
|
||||
)
|
||||
await session.execute(stmt)
|
||||
await session.commit()
|
||||
logger.info("[DomainChange] Запрос на обновление домена выполнен успешно.")
|
||||
except Exception as e:
|
||||
logger.error(f"[DomainChange] Ошибка при выполнении запроса: {e}")
|
||||
|
||||
@@ -140,7 +140,6 @@ async def import_remnawave_users(session: AsyncSession, users: list[dict]) -> in
|
||||
logger.error(f"[Remnawave Import] Ошибка при добавлении пользователя {tg_id}: {e}")
|
||||
continue
|
||||
|
||||
await session.commit()
|
||||
return added
|
||||
|
||||
|
||||
@@ -192,6 +191,5 @@ async def import_remnawave_keys(session: AsyncSession, users: list[dict], server
|
||||
except Exception as e:
|
||||
logger.error(f"[ERROR] Ошибка при добавлении ключа {client_id}: {e}")
|
||||
|
||||
await session.commit()
|
||||
logger.info(f"[IMPORT] Всего добавлено ключей: {added}")
|
||||
return added
|
||||
|
||||
@@ -229,7 +229,6 @@ async def process_callback_delete_server(
|
||||
(Server.cluster_name == cluster_name) & (Server.server_name == server_name)
|
||||
)
|
||||
await session.execute(stmt_delete)
|
||||
await session.commit()
|
||||
await callback_query.message.edit_text(
|
||||
text=(
|
||||
f"✅ Сервер '{server_name}' удален. "
|
||||
@@ -242,7 +241,6 @@ async def process_callback_delete_server(
|
||||
(Server.cluster_name == cluster_name) & (Server.server_name == server_name)
|
||||
)
|
||||
await session.execute(stmt_delete)
|
||||
await session.commit()
|
||||
await callback_query.message.edit_text(
|
||||
text=f"✅ Сервер '{server_name}' удален.",
|
||||
reply_markup=build_admin_back_kb("clusters"),
|
||||
@@ -261,7 +259,6 @@ async def toggle_server_enabled(
|
||||
new_status = action == "enable"
|
||||
|
||||
await session.execute(update(Server).where(Server.server_name == server_name).values(enabled=new_status))
|
||||
await session.commit()
|
||||
|
||||
servers = await get_servers(session=session, include_enabled=True)
|
||||
|
||||
@@ -314,7 +311,6 @@ async def save_server_limit(message: types.Message, state: FSMContext, session:
|
||||
new_value = limit if limit > 0 else None
|
||||
|
||||
await session.execute(update(Server).where(Server.server_name == server_name).values(max_keys=new_value))
|
||||
await session.commit()
|
||||
|
||||
servers = await get_servers(session=session, include_enabled=True)
|
||||
cluster_name, server = next(
|
||||
|
||||
@@ -46,7 +46,6 @@ async def toggle_button_setting(
|
||||
config[key] = not current
|
||||
|
||||
await update_buttons_config(session, config)
|
||||
await session.commit()
|
||||
|
||||
buttons_state = {k: bool(config.get(k, False)) for k in BUTTON_TITLES.keys()}
|
||||
await callback.message.edit_reply_markup(reply_markup=build_settings_buttons_kb(buttons_state))
|
||||
|
||||
@@ -60,7 +60,6 @@ async def toggle_cashbox_setting(
|
||||
session,
|
||||
config,
|
||||
)
|
||||
await session.commit()
|
||||
|
||||
updated_state = {k: bool(config.get(k, False)) for k in PAYMENT_PROVIDER_TITLES.keys()}
|
||||
await callback.message.edit_reply_markup(
|
||||
|
||||
@@ -67,7 +67,6 @@ async def toggle_notification_setting(
|
||||
config[key] = not current
|
||||
|
||||
await update_notifications_config(session, config)
|
||||
await session.commit()
|
||||
|
||||
notifications_state = await load_notification_settings()
|
||||
await callback.message.edit_reply_markup(
|
||||
@@ -126,7 +125,6 @@ async def notification_interval_value_input(message: Message, state: FSMContext,
|
||||
config[key] = new_value
|
||||
|
||||
await update_notifications_config(session, config)
|
||||
await session.commit()
|
||||
await state.clear()
|
||||
|
||||
notifications_state = await load_notification_settings()
|
||||
|
||||
@@ -61,7 +61,6 @@ async def save_device_step(message: Message, state: FSMContext, session: AsyncSe
|
||||
|
||||
tariff.device_step_rub = price
|
||||
tariff.updated_at = datetime.utcnow()
|
||||
await session.commit()
|
||||
|
||||
await state.set_state(TariffConfigState.choosing_section)
|
||||
|
||||
@@ -151,7 +150,6 @@ async def clear_device_overrides(callback: CallbackQuery, state: FSMContext, ses
|
||||
|
||||
tariff.device_overrides = None
|
||||
tariff.updated_at = datetime.utcnow()
|
||||
await session.commit()
|
||||
|
||||
text, markup = build_device_overrides_screen(tariff)
|
||||
await callback.message.edit_text(text=text, reply_markup=markup)
|
||||
@@ -197,7 +195,6 @@ async def save_device_override_price(message: Message, state: FSMContext, sessio
|
||||
tariff.device_overrides = overrides if overrides else None
|
||||
attributes.flag_modified(tariff, "device_overrides")
|
||||
tariff.updated_at = datetime.utcnow()
|
||||
await session.commit()
|
||||
|
||||
await state.update_data(devices_override=None)
|
||||
|
||||
|
||||
@@ -74,7 +74,6 @@ async def save_devices_config(message: Message, state: FSMContext, session: Asyn
|
||||
return
|
||||
|
||||
tariff.updated_at = datetime.utcnow()
|
||||
await session.commit()
|
||||
|
||||
await state.set_state(TariffConfigState.choosing_section)
|
||||
|
||||
@@ -136,7 +135,6 @@ async def save_traffic_config(message: Message, state: FSMContext, session: Asyn
|
||||
return
|
||||
|
||||
tariff.updated_at = datetime.utcnow()
|
||||
await session.commit()
|
||||
|
||||
await state.set_state(TariffConfigState.choosing_section)
|
||||
|
||||
|
||||
@@ -62,7 +62,6 @@ async def save_traffic_step(message: Message, state: FSMContext, session: AsyncS
|
||||
|
||||
tariff.traffic_step_rub = price
|
||||
tariff.updated_at = datetime.utcnow()
|
||||
await session.commit()
|
||||
|
||||
await state.set_state(TariffConfigState.choosing_section)
|
||||
|
||||
@@ -144,7 +143,6 @@ async def clear_traffic_overrides(callback: CallbackQuery, state: FSMContext, se
|
||||
|
||||
tariff.traffic_overrides = None
|
||||
tariff.updated_at = datetime.utcnow()
|
||||
await session.commit()
|
||||
|
||||
text, markup = build_traffic_overrides_screen(tariff)
|
||||
|
||||
@@ -195,7 +193,6 @@ async def save_traffic_override_price(message: Message, state: FSMContext, sessi
|
||||
tariff.traffic_overrides = overrides if overrides else None
|
||||
attributes.flag_modified(tariff, "traffic_overrides")
|
||||
tariff.updated_at = datetime.utcnow()
|
||||
await session.commit()
|
||||
|
||||
await state.update_data(traffic_override_gb=None)
|
||||
|
||||
|
||||
@@ -406,7 +406,6 @@ async def delete_tariff_with_gift_replacement(callback: CallbackQuery, session:
|
||||
if not remaining_tariffs:
|
||||
await session.execute(update(Server).where(Server.tariff_group == group_code).values(tariff_group=None))
|
||||
|
||||
await session.commit()
|
||||
await callback.message.edit_text("🗑 Тариф удалён. Все подарки обновлены.", reply_markup=build_tariff_menu_kb())
|
||||
|
||||
|
||||
@@ -432,7 +431,6 @@ async def delete_tariff(callback: CallbackQuery, session: AsyncSession):
|
||||
if not remaining_tariffs:
|
||||
await session.execute(update(Server).where(Server.tariff_group == group_code).values(tariff_group=None))
|
||||
|
||||
await session.commit()
|
||||
await callback.message.edit_text("🗑 Тариф успешно удалён.", reply_markup=build_tariff_menu_kb())
|
||||
|
||||
|
||||
@@ -506,7 +504,6 @@ async def set_vless_flag(callback: CallbackQuery, state: FSMContext, session: As
|
||||
|
||||
tariff.vless = vless_flag
|
||||
tariff.updated_at = datetime.utcnow()
|
||||
await session.commit()
|
||||
await state.clear()
|
||||
|
||||
text, markup = render_tariff_card(tariff)
|
||||
@@ -542,7 +539,6 @@ async def apply_edit(message: Message, state: FSMContext, session: AsyncSession)
|
||||
value = None
|
||||
setattr(tariff, field, value)
|
||||
tariff.updated_at = datetime.utcnow()
|
||||
await session.commit()
|
||||
await state.clear()
|
||||
|
||||
text, markup = render_tariff_card(tariff)
|
||||
@@ -565,7 +561,6 @@ async def apply_edit(message: Message, state: FSMContext, session: AsyncSession)
|
||||
setattr(tariff, field, value)
|
||||
tariff.updated_at = datetime.utcnow()
|
||||
|
||||
await session.commit()
|
||||
await state.clear()
|
||||
|
||||
text, markup = render_tariff_card(tariff)
|
||||
@@ -584,7 +579,6 @@ async def toggle_tariff_status(callback: CallbackQuery, session: AsyncSession):
|
||||
return
|
||||
|
||||
tariff.is_active = not tariff.is_active
|
||||
await session.commit()
|
||||
|
||||
text, markup = render_tariff_card(tariff)
|
||||
await callback.message.edit_text(text=text, reply_markup=markup)
|
||||
@@ -618,7 +612,6 @@ async def toggle_tariff_configurable(callback: CallbackQuery, session: AsyncSess
|
||||
tariff.configurable = not current
|
||||
tariff.updated_at = datetime.utcnow()
|
||||
|
||||
await session.commit()
|
||||
|
||||
text, markup = render_tariff_card(tariff)
|
||||
await callback.message.edit_text(text=text, reply_markup=markup)
|
||||
|
||||
@@ -147,7 +147,6 @@ async def apply_subgroup_title(message: Message, state: FSMContext, session: Asy
|
||||
await session.execute(
|
||||
update(Tariff).where(Tariff.id.in_(selected_ids)).values(subgroup_title=title, updated_at=datetime.utcnow())
|
||||
)
|
||||
await session.commit()
|
||||
await state.clear()
|
||||
|
||||
await message.answer(
|
||||
@@ -286,7 +285,6 @@ async def save_new_subgroup_title(message: Message, state: FSMContext, session:
|
||||
)
|
||||
.values(subgroup_title=new_title)
|
||||
)
|
||||
await session.commit()
|
||||
await state.clear()
|
||||
|
||||
create_subgroup_hash(new_title, group_code)
|
||||
@@ -347,7 +345,6 @@ async def perform_subgroup_deletion(callback: CallbackQuery, state: FSMContext,
|
||||
.where(Tariff.group_code == group_code, Tariff.subgroup_title == subgroup_title)
|
||||
.values(subgroup_title=None)
|
||||
)
|
||||
await session.commit()
|
||||
await state.clear()
|
||||
|
||||
await callback.message.edit_text(
|
||||
@@ -496,7 +493,6 @@ async def save_subgroup_tariffs_changes(callback: CallbackQuery, state: FSMConte
|
||||
.values(subgroup_title=subgroup_title, updated_at=datetime.utcnow())
|
||||
)
|
||||
|
||||
await session.commit()
|
||||
await state.clear()
|
||||
|
||||
if not selected_tariff_ids:
|
||||
|
||||
@@ -90,7 +90,6 @@ async def handle_ban_forever_reason_input(message: Message, state: FSMContext, s
|
||||
)
|
||||
|
||||
await session.execute(stmt)
|
||||
await session.commit()
|
||||
if u.tg_id is not None:
|
||||
await invalidate_ban_cache(u.tg_id)
|
||||
await state.clear()
|
||||
@@ -179,7 +178,6 @@ async def handle_ban_duration_input(message: Message, state: FSMContext, session
|
||||
)
|
||||
|
||||
await session.execute(stmt)
|
||||
await session.commit()
|
||||
if u.tg_id is not None:
|
||||
await invalidate_ban_cache(u.tg_id)
|
||||
|
||||
@@ -227,7 +225,6 @@ async def handle_ban_shadow(callback: CallbackQuery, callback_data: AdminUserEdi
|
||||
)
|
||||
)
|
||||
await session.execute(stmt)
|
||||
await session.commit()
|
||||
if u.tg_id is not None:
|
||||
await invalidate_ban_cache(u.tg_id)
|
||||
|
||||
@@ -252,7 +249,6 @@ async def handle_user_unban(
|
||||
return
|
||||
|
||||
await session.execute(delete(ManualBan).where(ManualBan.user_id == u.id))
|
||||
await session.commit()
|
||||
if u.tg_id is not None:
|
||||
await invalidate_ban_cache(u.tg_id)
|
||||
|
||||
|
||||
@@ -145,7 +145,6 @@ async def handle_gift_delete_confirm(
|
||||
|
||||
await session.execute(delete(GiftUsage).where(GiftUsage.gift_id == gift_id))
|
||||
await session.execute(delete(Gift).where(Gift.gift_id == gift_id))
|
||||
await session.commit()
|
||||
|
||||
await callback.answer("✅ Подарок удалён", show_alert=True)
|
||||
await show_gifts_list(callback.message, session, tg_id, page=0)
|
||||
|
||||
@@ -130,7 +130,6 @@ async def handle_admin_freeze_subscription(
|
||||
time_left = 0
|
||||
|
||||
await mark_key_as_frozen(session, record["tg_id"], client_id, time_left)
|
||||
await session.commit()
|
||||
session.expire_all()
|
||||
|
||||
await callback_query.answer("✅ Подписка отключена")
|
||||
@@ -207,7 +206,6 @@ async def handle_admin_unfreeze_subscription(
|
||||
new_expiry_time = now_ms + leftover
|
||||
|
||||
await mark_key_as_unfrozen(session, record["tg_id"], client_id, new_expiry_time)
|
||||
await session.commit()
|
||||
session.expire_all()
|
||||
await release_session_early(session)
|
||||
|
||||
|
||||
@@ -343,7 +343,6 @@ async def restore_trials(callback_query: types.CallbackQuery, session: AsyncSess
|
||||
.values(trial=0)
|
||||
)
|
||||
result = await session.execute(stmt)
|
||||
await session.commit()
|
||||
|
||||
builder = InlineKeyboardBuilder()
|
||||
builder.row(build_admin_back_btn())
|
||||
|
||||
@@ -187,7 +187,6 @@ async def key_cluster_mode(
|
||||
selected_price_rub=price_to_charge,
|
||||
)
|
||||
)
|
||||
await session.commit()
|
||||
|
||||
key_record = await get_key_details(session, email)
|
||||
if not key_record:
|
||||
|
||||
@@ -288,7 +288,6 @@ async def finalize_key_creation(
|
||||
if state:
|
||||
await state.update_data(skip_balance_charge=False)
|
||||
|
||||
await session.commit()
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"[Key Finalize] Ошибка при создании ключа для пользователя {tg_id}: {e}")
|
||||
|
||||
@@ -282,7 +282,6 @@ async def handle_new_alias_input(message: Message, state: FSMContext, session: A
|
||||
await state.clear()
|
||||
return
|
||||
await session.execute(update(Key).where(Key.user_id == u.id, Key.client_id == client_id).values(alias=alias))
|
||||
await session.commit()
|
||||
except Exception as error:
|
||||
await message.answer("❌ Не удалось переименовать подписку.")
|
||||
logger.error(f"Ошибка при обновлении alias: {error}")
|
||||
|
||||
@@ -67,9 +67,6 @@ async def execute_bulk_updates(session: AsyncSession, bulk_updates: dict[str, An
|
||||
if to_add or to_delete:
|
||||
logger.info(f"Bulk: {len(to_add)} добавлений, {len(to_delete)} удалений уведомлений")
|
||||
|
||||
await session.commit()
|
||||
|
||||
except Exception as error:
|
||||
logger.error(f"Ошибка в bulk-обновлениях: {error}")
|
||||
await session.rollback()
|
||||
raise
|
||||
|
||||
@@ -162,4 +162,3 @@ async def _run_cycle(bot: Bot, sessionmaker: async_sessionmaker):
|
||||
|
||||
total_time = (datetime.now() - start_time).total_seconds()
|
||||
logger.info(f"Уведомления завершены за {total_time:.2f}s")
|
||||
await session.commit()
|
||||
|
||||
@@ -110,7 +110,6 @@ async def _process_renew_candidates(
|
||||
)
|
||||
try:
|
||||
result = await try_auto_renew(one_ctx, key)
|
||||
await session.commit()
|
||||
return (key, notification_id, result)
|
||||
except Exception as e:
|
||||
logger.error(f"Ошибка продления {key.tg_id} ({key.email}): {e}")
|
||||
|
||||
@@ -78,7 +78,6 @@ async def process_inactive_trial(
|
||||
|
||||
if users_to_extend:
|
||||
await session.execute(update(User).where(User.tg_id.in_(users_to_extend)).values(trial=-1))
|
||||
await session.commit()
|
||||
logger.info(f"[InactiveTrial] {len(users_to_extend)} пользователей с расширенным триалом")
|
||||
|
||||
if messages:
|
||||
|
||||
@@ -130,7 +130,6 @@ async def process_zero_traffic(
|
||||
if keys_to_mark:
|
||||
try:
|
||||
await session.execute(update(Key).where(Key.client_id.in_(keys_to_mark)).values(notified=True))
|
||||
await session.commit()
|
||||
logger.info(f"[ZeroTraffic] Отмечено {len(keys_to_mark)} ключей как notified")
|
||||
except Exception as e:
|
||||
logger.error(f"[ZeroTraffic] Ошибка обновления notified: {e}")
|
||||
|
||||
Binary file not shown.
Binary file not shown.
@@ -664,7 +664,6 @@ async def handle_addons_downgrade_apply(callback: CallbackQuery, state: FSMConte
|
||||
has_traffic_choice=has_traffic_choice,
|
||||
config_mode="downgrade",
|
||||
)
|
||||
await session.commit()
|
||||
except Exception as error:
|
||||
logger.error(f"[ADDONS] Ошибка при сохранении будущих условий для {email}: {error}")
|
||||
await callback.message.answer("❌ Ошибка при сохранении новых условий. Попробуйте позже.")
|
||||
@@ -896,7 +895,6 @@ async def handle_addons_confirm(callback: CallbackQuery, state: FSMContext, sess
|
||||
)
|
||||
|
||||
await update_balance(session, tg_id, -extra_price)
|
||||
await session.commit()
|
||||
|
||||
logger.info(
|
||||
"[ADDONS] Успешное применение расширения: "
|
||||
|
||||
@@ -131,7 +131,6 @@ async def check_server_key_limit(server_info: dict, session: AsyncSession) -> bo
|
||||
pass
|
||||
if anchor_uid is not None:
|
||||
session.add(Notification(user_id=anchor_uid, notification_type=notif_key))
|
||||
await session.commit()
|
||||
|
||||
return await _svc_check(server_info, session, on_capacity_warning=_notify_admin_capacity)
|
||||
|
||||
|
||||
File diff suppressed because one or more lines are too long
+3
-1
@@ -38,6 +38,8 @@ MarkupSafe==3.0.2
|
||||
mdurl==0.1.2
|
||||
multidict==6.1.0
|
||||
netaddr==1.3.0
|
||||
orjson>=3.10.0
|
||||
httptools>=0.6.0
|
||||
bcrypt>=4.0.0
|
||||
pillow==11.3.0
|
||||
ping3==4.0.8
|
||||
@@ -67,7 +69,7 @@ StrEnum==0.4.15
|
||||
typing_extensions==4.12.2
|
||||
tzlocal==5.3.1
|
||||
urllib3==2.6.3
|
||||
uvicorn==0.35.0
|
||||
uvicorn[standard]==0.35.0
|
||||
wrapt==1.16.0
|
||||
yarl>=1.17.0,<2.0
|
||||
yookassa==3.9.0
|
||||
|
||||
@@ -249,7 +249,6 @@ async def update_subscription(
|
||||
await release_session_early(session)
|
||||
await delete_key_from_cluster(old_cluster_id, email, client_id, session=session)
|
||||
await delete_key_by_user_and_email(session, uid, email)
|
||||
await session.commit()
|
||||
|
||||
if country_override or cluster_override:
|
||||
new_cluster_id = country_override or cluster_override
|
||||
|
||||
Reference in New Issue
Block a user