303 lines
11 KiB
Python
303 lines
11 KiB
Python
import secrets
|
|
from datetime import datetime, timedelta, timezone
|
|
|
|
import httpx
|
|
from fastapi import status
|
|
from sqlalchemy import delete, func, select
|
|
from sqlalchemy.orm import Session
|
|
|
|
from app.core.config import settings
|
|
from app.core.errors import raise_api_error
|
|
from app.models.managed_service import ManagedService
|
|
from app.models.server import Server
|
|
from app.models.telegram_link import TelegramLink
|
|
from app.models.user import User
|
|
from app.schemas.telegram import (
|
|
TelegramCommandServerStatus,
|
|
TelegramCommandServiceStatus,
|
|
TelegramLinkCodeRead,
|
|
TelegramServicesCommandRead,
|
|
TelegramStatusCommandRead,
|
|
)
|
|
from app.services.server_runtime import ServerRuntimeService
|
|
from app.services.service_control import ServiceControlService
|
|
|
|
|
|
class TelegramService:
|
|
def __init__(
|
|
self,
|
|
runtime_service: ServerRuntimeService | None = None,
|
|
control_service: ServiceControlService | None = None,
|
|
) -> None:
|
|
self.runtime_service = runtime_service or ServerRuntimeService()
|
|
self.control_service = control_service or ServiceControlService()
|
|
|
|
def get_status_overview(self, db: Session) -> tuple[bool, int]:
|
|
linked_users_count = db.scalar(
|
|
select(func.count()).select_from(TelegramLink).where(
|
|
TelegramLink.telegram_user_id.is_not(None)
|
|
)
|
|
)
|
|
return bool(settings.telegram_bot_token), int(linked_users_count or 0)
|
|
|
|
def issue_link_code(self, db: Session, user: User) -> TelegramLinkCodeRead:
|
|
link = db.scalar(select(TelegramLink).where(TelegramLink.user_id == user.id))
|
|
expires_at = datetime.now(timezone.utc) + timedelta(
|
|
minutes=settings.telegram_link_code_expire_minutes
|
|
)
|
|
code = secrets.token_urlsafe(9)
|
|
|
|
if link is None:
|
|
link = TelegramLink(user_id=user.id)
|
|
# Refresh the code every time so old handshakes cannot be replayed.
|
|
link.link_code = code
|
|
link.link_code_expires_at = expires_at
|
|
db.add(link)
|
|
db.commit()
|
|
db.refresh(link)
|
|
return TelegramLinkCodeRead(link_code=code, expires_at=expires_at)
|
|
|
|
def complete_link(
|
|
self,
|
|
db: Session,
|
|
*,
|
|
link_code: str,
|
|
telegram_user_id: str,
|
|
chat_id: str,
|
|
) -> tuple[TelegramLink, User]:
|
|
link = db.scalar(select(TelegramLink).where(TelegramLink.link_code == link_code))
|
|
if link is None or link.link_code_expires_at is None:
|
|
raise_api_error(
|
|
status_code=status.HTTP_404_NOT_FOUND,
|
|
code="telegram_link_code_not_found",
|
|
message="Telegram link code was not found.",
|
|
)
|
|
|
|
expires_at = self._ensure_utc_datetime(link.link_code_expires_at)
|
|
if expires_at < datetime.now(timezone.utc):
|
|
raise_api_error(
|
|
status_code=status.HTTP_410_GONE,
|
|
code="telegram_link_code_expired",
|
|
message="Telegram link code has expired.",
|
|
)
|
|
|
|
existing_link = db.scalar(
|
|
select(TelegramLink).where(TelegramLink.telegram_user_id == telegram_user_id)
|
|
)
|
|
if existing_link is not None and existing_link.user_id != link.user_id:
|
|
raise_api_error(
|
|
status_code=status.HTTP_409_CONFLICT,
|
|
code="telegram_account_already_linked",
|
|
message="This Telegram account is already linked to another user.",
|
|
)
|
|
|
|
user = db.get(User, link.user_id)
|
|
if user is None or not user.is_active:
|
|
raise_api_error(
|
|
status_code=status.HTTP_400_BAD_REQUEST,
|
|
code="telegram_link_user_unavailable",
|
|
message="The target user is unavailable for Telegram linking.",
|
|
)
|
|
|
|
link.telegram_user_id = telegram_user_id
|
|
link.chat_id = chat_id
|
|
link.linked_at = datetime.now(timezone.utc)
|
|
link.link_code = None
|
|
link.link_code_expires_at = None
|
|
db.add(link)
|
|
db.commit()
|
|
db.refresh(link)
|
|
return link, user
|
|
|
|
def get_link_for_user(self, db: Session, user: User) -> TelegramLink | None:
|
|
return db.scalar(select(TelegramLink).where(TelegramLink.user_id == user.id))
|
|
|
|
def unlink(self, db: Session, link_id: int, current_user: User) -> None:
|
|
link = db.get(TelegramLink, link_id)
|
|
if link is None:
|
|
raise_api_error(
|
|
status_code=status.HTTP_404_NOT_FOUND,
|
|
code="telegram_link_not_found",
|
|
message="Telegram link was not found.",
|
|
)
|
|
if current_user.role != "admin" and link.user_id != current_user.id:
|
|
raise_api_error(
|
|
status_code=status.HTTP_403_FORBIDDEN,
|
|
code="forbidden",
|
|
message="Insufficient permissions.",
|
|
)
|
|
|
|
db.execute(delete(TelegramLink).where(TelegramLink.id == link.id))
|
|
db.commit()
|
|
|
|
def get_linked_user(self, db: Session, telegram_user_id: str) -> tuple[TelegramLink, User]:
|
|
link = db.scalar(
|
|
select(TelegramLink).where(TelegramLink.telegram_user_id == telegram_user_id)
|
|
)
|
|
if link is None:
|
|
raise_api_error(
|
|
status_code=status.HTTP_404_NOT_FOUND,
|
|
code="telegram_link_not_found",
|
|
message="Telegram account is not linked.",
|
|
)
|
|
|
|
user = db.get(User, link.user_id)
|
|
if user is None or not user.is_active:
|
|
raise_api_error(
|
|
status_code=status.HTTP_403_FORBIDDEN,
|
|
code="telegram_user_unavailable",
|
|
message="Linked user is unavailable.",
|
|
)
|
|
|
|
return link, user
|
|
|
|
def get_status_for_linked_user(
|
|
self,
|
|
db: Session,
|
|
telegram_user_id: str,
|
|
) -> TelegramStatusCommandRead:
|
|
_, user = self.get_linked_user(db, telegram_user_id)
|
|
servers = list(
|
|
db.scalars(
|
|
select(Server).where(Server.is_active.is_(True)).order_by(Server.name)
|
|
)
|
|
)
|
|
status_rows = [
|
|
TelegramCommandServerStatus(
|
|
name=server.name,
|
|
environment=server.environment,
|
|
host=server.host,
|
|
status=health.status,
|
|
reachable=health.reachable,
|
|
checked_at=health.checked_at,
|
|
detail=health.detail,
|
|
)
|
|
for server in servers
|
|
for health in [self.runtime_service.get_server_health(server)]
|
|
]
|
|
return TelegramStatusCommandRead(
|
|
username=user.username,
|
|
role=user.role,
|
|
servers=status_rows,
|
|
)
|
|
|
|
def get_services_for_linked_user(
|
|
self,
|
|
db: Session,
|
|
telegram_user_id: str,
|
|
server_name: str | None = None,
|
|
) -> TelegramServicesCommandRead:
|
|
_, user = self.get_linked_user(db, telegram_user_id)
|
|
server_query = select(Server).where(Server.is_active.is_(True))
|
|
if server_name:
|
|
server_query = server_query.where(Server.name == server_name)
|
|
servers = list(db.scalars(server_query.order_by(Server.name)))
|
|
|
|
if server_name and not servers:
|
|
raise_api_error(
|
|
status_code=status.HTTP_404_NOT_FOUND,
|
|
code="server_not_found",
|
|
message="Requested server was not found.",
|
|
)
|
|
|
|
service_rows: list[TelegramCommandServiceStatus] = []
|
|
for server in servers:
|
|
for service in self.control_service.list_managed_services(db, server):
|
|
service_rows.append(
|
|
TelegramCommandServiceStatus(
|
|
server_name=server.name,
|
|
service_name=service.name,
|
|
service_type=service.service_type,
|
|
status=service.status,
|
|
detail=service.detail,
|
|
)
|
|
)
|
|
|
|
return TelegramServicesCommandRead(
|
|
username=user.username,
|
|
role=user.role,
|
|
services=service_rows,
|
|
)
|
|
|
|
def restart_service_for_linked_user(
|
|
self,
|
|
db: Session,
|
|
*,
|
|
telegram_user_id: str,
|
|
server_name: str,
|
|
service_name: str,
|
|
):
|
|
_, user = self.get_linked_user(db, telegram_user_id)
|
|
if user.role not in {"admin", "operator"}:
|
|
raise_api_error(
|
|
status_code=status.HTTP_403_FORBIDDEN,
|
|
code="forbidden",
|
|
message="Insufficient permissions.",
|
|
)
|
|
|
|
server = db.scalar(select(Server).where(Server.name == server_name))
|
|
if server is None:
|
|
raise_api_error(
|
|
status_code=status.HTTP_404_NOT_FOUND,
|
|
code="server_not_found",
|
|
message="Requested server was not found.",
|
|
)
|
|
managed_service = db.scalar(
|
|
select(ManagedService).where(
|
|
ManagedService.server_id == server.id,
|
|
ManagedService.name == service_name,
|
|
)
|
|
)
|
|
if managed_service is None:
|
|
raise_api_error(
|
|
status_code=status.HTTP_404_NOT_FOUND,
|
|
code="service_not_found",
|
|
message="Requested managed service was not found.",
|
|
)
|
|
|
|
return self.control_service.run_service_action(
|
|
db,
|
|
server=server,
|
|
managed_service=managed_service,
|
|
action_name="restart",
|
|
actor=user,
|
|
source="telegram",
|
|
)
|
|
|
|
def send_test_message_for_user(self, db: Session, user: User) -> None:
|
|
link = self.get_link_for_user(db, user)
|
|
if link is None or not link.chat_id:
|
|
raise_api_error(
|
|
status_code=status.HTTP_404_NOT_FOUND,
|
|
code="telegram_link_not_found",
|
|
message="Current user does not have a linked Telegram chat.",
|
|
)
|
|
self.send_message(chat_id=link.chat_id, text="Server Panel test message: Telegram integration is working.")
|
|
|
|
def send_message(self, *, chat_id: str, text: str) -> None:
|
|
if not settings.telegram_bot_token:
|
|
raise_api_error(
|
|
status_code=status.HTTP_503_SERVICE_UNAVAILABLE,
|
|
code="telegram_disabled",
|
|
message="Telegram bot token is not configured.",
|
|
)
|
|
|
|
response = httpx.post(
|
|
f"https://api.telegram.org/bot{settings.telegram_bot_token}/sendMessage",
|
|
json={"chat_id": chat_id, "text": text},
|
|
timeout=10,
|
|
)
|
|
response.raise_for_status()
|
|
payload = response.json()
|
|
if not payload.get("ok"):
|
|
raise_api_error(
|
|
status_code=status.HTTP_502_BAD_GATEWAY,
|
|
code="telegram_delivery_failed",
|
|
message="Telegram did not accept the message.",
|
|
)
|
|
|
|
def _ensure_utc_datetime(self, value: datetime) -> datetime:
|
|
if value.tzinfo is None:
|
|
return value.replace(tzinfo=timezone.utc)
|
|
return value.astimezone(timezone.utc)
|