Fix Remnawave sync

This commit is contained in:
Capybara-z
2025-10-25 00:57:32 +03:00
parent 12dd8f75fb
commit da44b4e041
3 changed files with 119 additions and 19 deletions
+33 -18
View File
@@ -22,25 +22,40 @@ async def store_key(
):
try:
exists = await session.execute(select(Key).where(Key.tg_id == tg_id, Key.client_id == client_id))
if exists.scalar_one_or_none():
logger.info(f"[Store Key] Ключ уже существует — пропускаем: tg_id={tg_id}, client_id={client_id}")
return
new_key = Key(
tg_id=tg_id,
client_id=client_id,
email=email,
created_at=int(datetime.utcnow().timestamp() * 1000),
expiry_time=expiry_time,
key=key,
server_id=server_id,
remnawave_link=remnawave_link,
tariff_id=tariff_id,
alias=alias,
)
session.add(new_key)
existing_key = exists.scalar_one_or_none()
if existing_key:
await session.execute(
update(Key)
.where(Key.tg_id == tg_id, Key.client_id == client_id)
.values(
email=email,
expiry_time=expiry_time,
key=key,
server_id=server_id,
remnawave_link=remnawave_link,
tariff_id=tariff_id,
alias=alias,
)
)
logger.info(f"[Store Key] Ключ обновлён: tg_id={tg_id}, client_id={client_id}, server_id={server_id}")
else:
new_key = Key(
tg_id=tg_id,
client_id=client_id,
email=email,
created_at=int(datetime.utcnow().timestamp() * 1000),
expiry_time=expiry_time,
key=key,
server_id=server_id,
remnawave_link=remnawave_link,
tariff_id=tariff_id,
alias=alias,
)
session.add(new_key)
logger.info(f"[Store Key] Ключ создан: tg_id={tg_id}, client_id={client_id}, server_id={server_id}")
await session.commit()
logger.info(f"Ключ сохранён: tg_id={tg_id}, client_id={client_id}, server_id={server_id}")
except SQLAlchemyError as e:
logger.error(f"❌ Ошибка при сохранении ключа: {e}")
@@ -14,6 +14,7 @@ from sqlalchemy.ext.asyncio import AsyncSession
from config import (
ADMIN_PASSWORD,
ADMIN_USERNAME,
HAPP_CRYPTOLINK,
REMNAWAVE_LOGIN,
REMNAWAVE_PASSWORD,
USE_COUNTRY_SELECTION,
@@ -27,6 +28,7 @@ from handlers.keys.operations import (
delete_key_from_cluster,
renew_key_in_cluster,
)
from handlers.keys.operations.aggregated_links import make_aggregated_link
from handlers.utils import ALLOWED_GROUP_CODES
from logger import logger
from panels.remnawave import RemnawaveAPI
@@ -512,6 +514,7 @@ async def handle_sync_server(
Key.email,
Key.expiry_time,
Key.tariff_id,
Key.remnawave_link,
)
.join(Key, Server.cluster_name == Key.server_id)
.where(Server.server_name == server_name)
@@ -562,6 +565,48 @@ async def handle_sync_server(
hwid_device_limit=hwid_limit,
)
if success:
try:
sub = await remna.get_subscription_by_username(key["email"])
if sub:
new_remnawave_link = sub.get("subscriptionUrl")
if HAPP_CRYPTOLINK:
happ = sub.get("happ") or {}
new_remnawave_link = happ.get("cryptoLink") or happ.get("link") or new_remnawave_link
if new_remnawave_link:
server_result = await session.execute(
select(Server.cluster_name).where(Server.server_name == server_name)
)
cluster_name = server_result.scalar()
servers = await get_servers(session)
cluster_servers = servers.get(cluster_name, [])
key_value = await make_aggregated_link(
session=session,
cluster_all=cluster_servers,
cluster_id=cluster_name,
email=key["email"],
client_id=key["client_id"],
tg_id=key["tg_id"],
remna_link_override=new_remnawave_link,
plan=key["tariff_id"],
)
await session.execute(
update(Key)
.where(Key.tg_id == key["tg_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}")
if not success:
logger.warning("[Sync] ошибка обновления, пробуем пересоздать")
@@ -575,6 +620,7 @@ async def handle_sync_server(
expiry_timestamp=key["expiry_time"],
plan=key["tariff_id"],
session=session,
remnawave_link=key["remnawave_link"],
)
else:
await create_client_on_server(
@@ -694,6 +740,43 @@ async def handle_sync_cluster(
hwid_device_limit=hwid_limit,
)
if success:
try:
sub = await remna.get_subscription_by_username(key["email"])
if sub:
new_remnawave_link = sub.get("subscriptionUrl")
if HAPP_CRYPTOLINK:
happ = sub.get("happ") or {}
new_remnawave_link = happ.get("cryptoLink") or happ.get("link") or new_remnawave_link
if new_remnawave_link:
servers = await get_servers(session)
cluster_servers = servers.get(cluster_name, [])
key_value = await make_aggregated_link(
session=session,
cluster_all=cluster_servers,
cluster_id=cluster_name,
email=key["email"],
client_id=key["client_id"],
tg_id=key["tg_id"],
remna_link_override=new_remnawave_link,
plan=key["tariff_id"],
)
await session.execute(
update(Key)
.where(Key.tg_id == key["tg_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}")
if not success:
logger.warning("[Sync] ошибка обновления, пробуем пересоздать")
+3 -1
View File
@@ -163,7 +163,9 @@ async def make_aggregated_link(
return f"{base}/{email}/{tg_id}"
best_vless, sub_url = await _try_build_remna_vless(remna, email)
if remna_link_override and (
remna_link_override.lower().startswith("vless://") or remna_link_override.startswith("http")
remna_link_override.lower().startswith("vless://") or
remna_link_override.startswith("http") or
remna_link_override.startswith("happ://")
):
logger.info("[agg_link] choose override Remnawave (non-vless)")
return remna_link_override