import datetime as dt import logging import os import re import secrets import threading import uuid import time import contextvars from typing import Optional from fastapi import BackgroundTasks, Depends, FastAPI, File, Form, HTTPException, Query, Request, UploadFile, status from fastapi.responses import HTMLResponse, JSONResponse, RedirectResponse from fastapi.staticfiles import StaticFiles from fastapi.templating import Jinja2Templates from markupsafe import Markup, escape from sqlalchemy import delete, select, text, update from sqlalchemy.orm import Session from starlette.responses import HTMLResponse as _HR from starlette.exceptions import HTTPException as _StarletteHTTPException import urllib.request as _urllib_request import urllib.parse as _urllib_parse import json as _json from config import ( COOKIE_MAX_AGE, COOKIE_NAME, CSRF_COOKIE, GO_POOL_LOCK_TIMEOUT_SECONDS, GO_USER_LOCK_TIMEOUT_SECONDS, LOG_LEVEL, LOG_SLOW_REQUEST_MS, MAX_ACTIVE_SERVICES_PER_USER, PUBLIC_HOST, SESSION_IDLE_SECONDS, WEB_POOL_BUFFER, WEB_POOL_SIZE, TELEGRAM_BOT_TOKEN, TELEGRAM_CHAT_ID, TELEGRAM_API_URL, SMTP_HOST, SMTP_PORT, SMTP_USERNAME, SMTP_PASSWORD, SMTP_FROM_EMAIL, SMTP_FROM_NAME, PORTAL_URL, ADMIN_NOTIFY_EMAIL, SIGNING_KEY, ) from database import get_db from models import ( AuditLog, Category, RdpSlot, Service, ServiceCategory, ServiceType, PendingAccessRequest, SessionModel, SessionStatus, User, UserServiceAccess, ) from utils import ( audit, ensure_icons_dir, format_service_comment, format_seo_description, log_event, normalize_web_target, now_utc, parse_rdp_target, remove_icon_file, request_id_ctx, set_service_categories, session_closed_reason, store_service_icon, ) from auth import ( get_current_user, has_access, issue_auth_cookie, issue_csrf_cookie, require_admin, require_user, user_is_valid, validate_csrf, verify_password, hash_password, check_login_rate_limit, record_login_failure, record_login_success, serializer, ) from runtime import ( acquire_universal_slot, acquire_web_pool_slot, allocator_lock, container_running, create_runtime_container, desired_pool_size, dispatch_universal_target, dispatch_web_pool_target, docker_client, ensure_universal_pool, ensure_warm_pool, ensure_web_pool, find_active_session_for_service, find_active_session_for_user_service, get_active_sessions_count, get_pool_detailed_status, get_pool_status_for_service, get_universal_pool_status, get_web_pool_status, LockTimeoutError, open_warm_web_url, _rdp_slot_container_name, route_ready, sanitize_client_resolution, service_uses_universal_pool, session_redirect_url, connect_rdp_slot, start_rdp_slot_container, stop_rdp_slot_container, stop_runtime_container, terminate_active_slot_sessions, terminate_session_record, wait_for_session_route, ) from maintenance import on_startup logging.basicConfig( level=LOG_LEVEL, format="%(asctime)s %(levelname)s %(name)s %(message)s", ) logger = logging.getLogger("portal") templates = Jinja2Templates(directory="templates") def _get_real_ip(request) -> str: """Real client IP from X-Forwarded-For (Traefik trusts NPM via trustedIPs).""" forwarded_for = request.headers.get("x-forwarded-for", "") if forwarded_for: return forwarded_for.split(",")[0].strip() return request.client.host if request.client else "unknown" def _get_geo(ip: str) -> str: """Lookup city/country for IP via ip-api.com. Returns formatted string or empty.""" try: if ip in ("unknown", "127.0.0.1", "::1") or ip.startswith("10.") or ip.startswith("192.168."): return "" url = f"http://ip-api.com/json/{ip}?lang=ru&fields=status,country,regionName,city,query" req = _urllib_request.Request(url, headers={"User-Agent": "Mozilla/5.0"}) with _urllib_request.urlopen(req, timeout=5) as resp: data = _json.loads(resp.read()) if data.get("status") == "success": parts = [data.get("country", ""), data.get("regionName", ""), data.get("city", "")] return ", ".join(p for p in parts if p) except Exception: pass return "" import secrets as _secrets import string as _string import smtplib as _smtplib import ssl as _ssl import json as _json2 from email.mime.multipart import MIMEMultipart as _MIMEMultipart from email.mime.text import MIMEText as _MIMEText from email.mime.image import MIMEImage as _MIMEImage from email.header import Header as _Header from email.utils import formataddr as _formataddr # MONT logo for HTML emails, attached inline (Content-ID) rather than # embedded as a base64 data: URI in . Outlook's desktop/Word # renderer does not support data: URIs for images at all (silently # drops them), so a cid: reference to an attached image is the only # way to get the logo to show up there. _LOGO_EMAIL_PATH = os.path.join(os.path.dirname(__file__), "static", "logo-email.png") try: with open(_LOGO_EMAIL_PATH, "rb") as _f: _LOGO_EMAIL_BYTES = _f.read() except OSError: _LOGO_EMAIL_BYTES = b"" def _generate_password(length: int = 10) -> str: alphabet = _string.ascii_letters + _string.digits while True: pwd = ''.join(_secrets.choice(alphabet) for _ in range(length)) if (any(c.isupper() for c in pwd) and any(c.islower() for c in pwd) and any(c.isdigit() for c in pwd)): return pwd _EMAIL_RE = re.compile(r"^[^\s@]+@[^\s@]+\.[^\s@]+$") def _send_email(to: str, subject: str, html_body: str) -> None: if not _EMAIL_RE.match(to or ""): # Some accounts (older/manually-created ones) have a plain username # instead of an email address as their login - PendingAccessRequest.email # mirrors that username for the in-app "request more access" flow, so # this is a normal, expected case here, not a real delivery failure. # Fail fast with a clear message instead of letting smtplib open a # connection just to have the SMTP server reject the recipient a few # seconds later with a cryptic 501. raise ValueError(f"'{to}' не похож на email-адрес - уведомление не отправлено") # "related" wrapper so the inline logo (Content-ID) sits alongside the # alternative/html part rather than being just another attachment - # needed for cid: references in html_body to resolve in mail clients. has_logo = _LOGO_EMAIL_BYTES and "cid:mont_logo" in html_body msg = _MIMEMultipart("related") if has_logo else _MIMEMultipart("alternative") msg["Subject"] = subject msg["From"] = _formataddr((str(_Header(SMTP_FROM_NAME, "utf-8")), SMTP_FROM_EMAIL)) msg["To"] = to if has_logo: alt = _MIMEMultipart("alternative") alt.attach(_MIMEText(html_body, "html", "utf-8")) msg.attach(alt) logo = _MIMEImage(_LOGO_EMAIL_BYTES, "png") logo.add_header("Content-ID", "") logo.add_header("Content-Disposition", "inline", filename="logo.png") msg.attach(logo) else: msg.attach(_MIMEText(html_body, "html", "utf-8")) ctx = _ssl.create_default_context() with _smtplib.SMTP_SSL(SMTP_HOST, SMTP_PORT, context=ctx) as srv: srv.login(SMTP_USERNAME, SMTP_PASSWORD) srv.sendmail(SMTP_FROM_EMAIL, to, msg.as_string()) def _build_light_email(title: str, subhead: str, body_html: str, cta_url: str | None = None, cta_label: str | None = None) -> str: """Shared light-theme card shell (logo, heading, body, optional button, footer) for one-off transactional emails that don't need the full credentials-table layout the access-approval templates have inline. Those keep their own copies (changing one must not risk the others); this is for newer, simpler notifications - welcome, password reset, pilot-granted - so they don't each re-type the whole shell.""" cta_html = "" if cta_url: cta_html = ( f'
' f'{cta_label or "Войти в полигон"}' f'
' ) return f"""
MONT

{title}

{subhead}


{body_html} {cta_html}
Если у вас возникли вопросы, свяжитесь с вашим менеджером MONT или напишите на RGalyaviev@mont.ru
""" def _credentials_table(rows: list[tuple[str, str, bool]]) -> str: """rows: (label, value, monospace) - renders the same light-theme label/value table used across access emails.""" trs = [] for i, (label, value, mono) in enumerate(rows): border = "border-top:1px solid #e7ebf0;" if i > 0 else "" font = "font-family:monospace;" if mono else "" trs.append( f'{label}' f'{value}' ) return ( '' + "".join(trs) + "
" ) def _send_welcome_email(user, password: str) -> str: """Sent once, when an admin manually creates a user on the Users tab - the only way that person learns their login/password, since the admin never sees the plaintext (create_user() generates it).""" body = ( f'

' f'Здравствуйте{", " + escape(user.first_name) + "" if user.first_name else ""}!
' f'Для вас создан аккаунт на Инфраструктурном полигоне MONT.

' + _credentials_table([ ("Адрес портала", f'{PORTAL_URL}', False), ("Логин", escape(user.username), True), ("Пароль", escape(password), True), ]) ) html = _build_light_email("Доступ к Инфраструктурному полигону MONT", "Для вас создан аккаунт", body, PORTAL_URL) try: _send_email(user.username, "Доступ к Инфраструктурному полигону MONT", html) return "Email отправлен" except Exception as ex: log_event("email_send_error", error=str(ex), channel="welcome") return f"Ошибка email: {ex}" def _send_password_reset_email(user, password: str) -> str: """Sent when an admin resets a user's password from the Users tab - same reasoning as the welcome email: the plaintext only ever exists in this message, never in the admin UI.""" body = ( f'

' f'Здравствуйте{", " + escape(user.first_name) + "" if user.first_name else ""}!
' f'Пароль от вашего аккаунта на полигоне MONT был сброшен администратором.

' + _credentials_table([ ("Адрес портала", f'{PORTAL_URL}', False), ("Логин", escape(user.username), True), ("Новый пароль", escape(password), True), ]) ) html = _build_light_email("Пароль обновлён", "Инфраструктурный полигон MONT", body, PORTAL_URL) try: _send_email(user.username, "Пароль обновлён — полигон MONT", html) return "Email отправлен" except Exception as ex: log_event("email_send_error", error=str(ex), channel="password_reset") return f"Ошибка email: {ex}" def _send_pilot_granted_email(user, service, expires_at) -> str: """Sent once, the moment a pilot's RDP slot is first assigned to this user (not on every later tweak of that same assignment - see the "new grant" check in assign_rdp_slot()). No credentials here: this user already has an account and password, this just tells them a new product showed up.""" expiry_row = ( f'

Доступ действует до {expires_at.strftime("%d.%m.%Y")}

' if expires_at else "" ) body = ( f'

' f'Здравствуйте{", " + escape(user.first_name) + "" if user.first_name else ""}!
' f'Вам открыт доступ к пилотному проекту {escape(service.name)} на Инфраструктурном полигоне MONT.

' + expiry_row ) html = _build_light_email(f"Доступ к пилоту: {escape(service.name)}", "Новый продукт на полигоне", body, PORTAL_URL, "Открыть полигон") try: _send_email(user.username, f"Доступ к пилоту «{service.name}» — полигон MONT", html) return "Email отправлен" except Exception as ex: log_event("email_send_error", error=str(ex), channel="pilot_granted") return f"Ошибка email: {ex}" def _tg_api(method: str, payload: dict) -> dict: import urllib.request as _ur import json as _j url = f"{TELEGRAM_API_URL}{TELEGRAM_BOT_TOKEN}/{method}" data = _j.dumps(payload).encode() req = _ur.Request(url, data=data, headers={"Content-Type": "application/json"}) with _ur.urlopen(req, timeout=10) as r: return _j.loads(r.read()) def _make_approval_keyboard(req_id: str) -> dict: return { "inline_keyboard": [ [ {"text": "7 дней", "callback_data": f"a7_{req_id}"}, {"text": "14 дней", "callback_data": f"a14_{req_id}"}, {"text": "30 дней", "callback_data": f"a30_{req_id}"}, {"text": "90 дней", "callback_data": f"a90_{req_id}"}, ], [{"text": "Отказать", "callback_data": f"r_{req_id}"}], ] } # --- Second approval channel: email ---------------------------------------- # Alongside the Telegram message, a notification with click-to-confirm links # goes to ADMIN_NOTIFY_EMAIL. Each link only opens a confirmation page (GET) - # it does NOT grant/reject anything by itself, because corporate mail # security (e.g. Microsoft Safe Links) pre-fetches every link in a message # with a plain GET before a human ever opens the email. The actual decision # is only applied on the POST triggered by the "Подтвердить" button on that # page. Whichever channel (Telegram or email) is acted on first wins - the # other just shows "уже обработано" afterwards. _ACTION_DAYS = {"a7": 7, "a14": 14, "a30": 30, "a90": 90, "r": None} _ACTION_LABELS = {"a7": "7 дней", "a14": "14 дней", "a30": "30 дней", "a90": "90 дней", "r": "Отказать"} from itsdangerous import BadSignature, SignatureExpired, URLSafeTimedSerializer as _URLSafeTimedSerializer _email_decision_serializer = _URLSafeTimedSerializer(SIGNING_KEY, salt="email-access-decision") _EMAIL_DECISION_MAX_AGE_SECONDS = 30 * 24 * 3600 # link stays clickable for 30 days def _make_decision_token(req_id: str, action: str) -> str: return _email_decision_serializer.dumps({"req_id": req_id, "action": action}) def _verify_decision_token(token: str, req_id: str, action: str) -> bool: try: data = _email_decision_serializer.loads(token, max_age=_EMAIL_DECISION_MAX_AGE_SECONDS) except (BadSignature, SignatureExpired): return False return data.get("req_id") == req_id and data.get("action") == action def _apply_access_decision(db, pending, days): """Approve (days=int) or reject (days=None) a pending access request: creates/renews the user, assigns requested services, commits, and emails the requester. Used by the email-approval channel below. Mirrors the Telegram approve/reject logic in telegram_webhook() / _process_callback_query() - keep the email templates and grant logic in sync if you change those.""" import json as _jg products = _jg.loads(pending.products_json or "[]") portal_url = pending.portal_url or PORTAL_URL if days is not None: username = pending.email expires = now_utc() + dt.timedelta(days=days) parts = pending.name.strip().split(None, 1) existing_user = db.scalar(select(User).where(User.username == username)) is_renewal = existing_user is not None if is_renewal: password = None existing_user.expires_at = expires existing_user.active = True target_user = existing_user db.flush() else: password = _generate_password() new_user = User( username=username, password_hash=hash_password(password), expires_at=expires, active=True, is_admin=False, first_name=parts[0] if parts else "", last_name=parts[1] if len(parts) > 1 else "", ) db.add(new_user) db.flush() target_user = new_user if products: # Python-side .strip() (not SQL trim()) so this is robust to any # whitespace - Postgres trim()/btrim() only strip plain spaces by # default, not tabs, which is exactly the byte that silently broke # this match for one product that had a leading tab in its name. _active_by_name = { svc.name.strip().lower(): svc for svc in db.scalars(select(Service).where(Service.active == True)).all() } # Dedup by the *matched service*, not just the raw product # string - a product tagged under 2+ categories (e.g. Deckhouse # Platform in both "Платформы виртуализации" and "Кубернетес") shows up as two # separate checkboxes in the request modal, so "Выбрать все" submits its # name twice - without this, that duplicate (user_id, service_id) # pair violates uq_user_service and rolls back the *entire* # approval (no user, no grants, status stays "pending") with only # a silent tg_callback_error in the logs to show for it. _wanted_names = {p.strip().lower() for p in products} matched = [svc for name, svc in _active_by_name.items() if name in _wanted_names] existing_svc_ids = ( {a.service_id for a in db.scalars(select(UserServiceAccess).where(UserServiceAccess.user_id == target_user.id)).all()} if is_renewal else set() ) for svc in matched: if svc.id not in existing_svc_ids: db.add(UserServiceAccess(user_id=target_user.id, service_id=svc.id)) pending.status = "approved" db.commit() products_html = "" if products: items = "".join(f"
  • {escape(p)}
  • " for p in products) products_html = ( '

    Предоставлен доступ к продуктам:

    ' f'' ) day_word = "день" if days == 1 else ("дня" if days < 5 else "дней") email_action = "продлён" if is_renewal else "предоставлен" email_subject = ("Продление доступа к Инфраструктурному полигону MONT" if is_renewal else "Доступ к Инфраструктурному полигону MONT") email_subhead = "Ваш доступ продлён" if is_renewal else "Ваш запрос одобрен" cred_row = ( f'Пароль' f'{password}' ) if not is_renewal else "" access_text = f"Вам {email_action} доступ к полигону на {days} {day_word}." html_email = f"""
    MONT

    {email_subject}

    {email_subhead}


    Здравствуйте, {escape(pending.name)}!
    {access_text}

    {cred_row}
    Адрес портала {portal_url}
    Логин {username}
    Доступ до {expires.strftime("%d.%m.%Y")}
    {products_html}
    Если у вас возникли вопросы, свяжитесь с вашим менеджером MONT или напишите на RGalyaviev@mont.ru
    """ try: _send_email(pending.email, email_subject, html_email) email_status = "Email отправлен" except Exception as ex: log_event("email_send_error", error=str(ex), channel="email_approval") email_status = f"Ошибка email: {ex}" return { "kind": "approved", "is_renewal": is_renewal, "days": days, "username": username, "password": password, "email_status": email_status, } else: manager_contact = pending.manager if pending.manager else "менеджера MONT" html_email = f"""
    MONT

    Запрос на доступ к полигону MONT


    Здравствуйте, {escape(pending.name)}!

    К сожалению, на данный момент мы не можем предоставить доступ к полигону.

    Для уточнения деталей свяжитесь с {escape(manager_contact)}.
    Если не знаете кто ваш менеджер — напишите на RGalyaviev@mont.ru.

    С уважением, команда MONT
    """ pending.status = "rejected" db.commit() try: _send_email(pending.email, "Запрос на доступ к полигону MONT", html_email) email_status = "Email отправлен" except Exception as ex: log_event("email_send_error", error=str(ex), channel="email_approval") email_status = f"Ошибка email: {ex}" return {"kind": "rejected", "email_status": email_status} def _send_admin_decision_email(pending) -> None: """Send the email-approval notification for a freshly created pending request, alongside the Telegram message. Best-effort: failures are logged, not raised (a broken email channel must never block the Telegram channel from working).""" if not ADMIN_NOTIFY_EMAIL: return decide_base = f"{PORTAL_URL}/admin/access-request/{pending.id}/decide" buttons_html = "" for action in ("a7", "a14", "a30", "a90"): token = _make_decision_token(pending.id, action) url = f"{decide_base}?action={action}&token={_urllib_parse.quote(token)}" buttons_html += ( f'{_ACTION_LABELS[action]}' ) reject_token = _make_decision_token(pending.id, "r") reject_url = f"{decide_base}?action=r&token={_urllib_parse.quote(reject_token)}" buttons_html += ( f'Отказать' ) products_line = "" try: import json as _jg2 items = _jg2.loads(pending.products_json or "[]") if items: products_line = ('

    Продукты: ' + ", ".join(escape(str(p)) for p in items) + '

    ') except Exception: pass html = f"""

    Новый запрос доступа к полигону MONT

    Второй канал согласования — дублирует запрос, отправленный в Telegram


    Имя{escape(pending.name)}
    Компания{escape(pending.company)}
    Email{escape(pending.email)}
    Телефон{escape(pending.phone)}
    Менеджер{escape(pending.manager) if pending.manager else "—"}
    {products_line}

    Выдать доступ на:

    {buttons_html}

    Ссылка открывает страницу подтверждения — доступ выдаётся только по нажатию кнопки на ней, а не по самому переходу (безопасно для антифишинг-сканеров почты, которые сами открывают ссылки).

    Заявка #{pending.id}
    """ try: _send_email(ADMIN_NOTIFY_EMAIL, "Новый запрос доступа к полигону MONT", html) except Exception as ex: log_event("email_send_error", error=str(ex), channel="admin_notify") def _decision_page(title: str, body: str): html = f""" {title}
    {body}
    """ return _HR(html) _tg_poll_offset: int = 0 _tg_poll_lock_file = None async def _telegram_poll_loop(): """Poll Telegram for callback_query updates every 5 minutes.""" import asyncio as _asyncio global _tg_poll_offset await _asyncio.sleep(10) # wait for app to fully start while True: try: result = _tg_api("getUpdates", { "offset": _tg_poll_offset, "timeout": 0, "allowed_updates": ["callback_query"], }) updates = result.get("result", []) for upd in updates: _tg_poll_offset = upd["update_id"] + 1 cq = upd.get("callback_query") if cq: try: await _process_callback_query(cq) except Exception as ex: log_event("tg_callback_error", error=str(ex)) except Exception as ex: log_event("tg_poll_error", error=str(ex)) await _asyncio.sleep(30) # 30 seconds async def _telegram_notify_retry_loop(): """Every 5 minutes, resend the initial 'new request' Telegram notification for any pending request whose first attempt failed (e.g. tel.4mont.ru was down/timed out) - see the 2026-08-18 incident where a request's only notification attempt failed and nobody found out until the requester followed up separately. Stops retrying once the request is no longer 'pending' (decided via Telegram or the email channel) or once a send finally succeeds.""" import asyncio as _asyncio3 from database import SessionLocal as _SL3 await _asyncio3.sleep(60) # let the app finish starting first while True: try: db = _SL3() try: pending_rows = db.scalars( select(PendingAccessRequest).where( PendingAccessRequest.status == "pending", PendingAccessRequest.telegram_notified == False, PendingAccessRequest.telegram_message != "", ) ).all() for pending in pending_rows: try: _tg_api("sendMessage", { "chat_id": TELEGRAM_CHAT_ID, "text": pending.telegram_message, "parse_mode": "HTML", "reply_markup": _make_approval_keyboard(pending.id), }) pending.telegram_notified = True db.commit() log_event("telegram_notify_retry_success", req_id=pending.id) except Exception as ex: log_event("telegram_notify_retry_error", req_id=pending.id, error=str(ex)) finally: db.close() except Exception as ex: log_event("telegram_notify_retry_loop_error", error=str(ex)) await _asyncio3.sleep(300) # 5 минут async def _process_callback_query(cq: dict): import json as _jc import datetime as _dtc import re as _rec from database import SessionLocal as _SL cq_id = cq["id"] cb_data = cq.get("data", "") chat_id = cq["message"]["chat"]["id"] msg_id = cq["message"]["message_id"] try: _tg_api("answerCallbackQuery", {"callback_query_id": cq_id}) except Exception: pass approve_match = _rec.match(r'^a(\d+)_(.+)$', cb_data) reject_match = _rec.match(r'^r_(.+)$', cb_data) if not approve_match and not reject_match: return req_id = approve_match.group(2) if approve_match else reject_match.group(1) db = _SL() try: pending = db.get(PendingAccessRequest, req_id) if not pending: _tg_api("editMessageText", { "chat_id": chat_id, "message_id": msg_id, "text": "Запрос не найден (возможно уже обработан).", }) return if pending.status != "pending": _tg_api("editMessageText", { "chat_id": chat_id, "message_id": msg_id, "text": f"Запрос уже обработан: {pending.status}.", }) return products = _jc.loads(pending.products_json or "[]") portal_url = pending.portal_url or PORTAL_URL if approve_match: days = int(approve_match.group(1)) username = pending.email expires = _dtc.datetime.now(_dtc.timezone.utc) + _dtc.timedelta(days=days) parts = pending.name.strip().split(None, 1) existing_user = db.scalar(select(User).where(User.username == username)) is_renewal = existing_user is not None if is_renewal: password = None existing_user.expires_at = expires existing_user.active = True target_user = existing_user db.flush() else: password = _generate_password() new_user = User( username=username, password_hash=hash_password(password), expires_at=expires, active=True, is_admin=False, first_name=parts[0] if parts else "", last_name=parts[1] if len(parts) > 1 else "", ) db.add(new_user) db.flush() target_user = new_user if products: _active_by_name = { svc.name.strip().lower(): svc for svc in db.scalars(select(Service).where(Service.active == True)).all() } # Dedup by the *matched service*, not just the raw product # string - a product tagged under 2+ categories (e.g. Deckhouse # Platform in both "Платформы виртуализации" and "Кубернетес") shows up as two # separate checkboxes in the request modal, so "Выбрать все" submits its # name twice - without this, that duplicate (user_id, service_id) # pair violates uq_user_service and rolls back the *entire* # approval (no user, no grants, status stays "pending") with only # a silent tg_callback_error in the logs to show for it. _wanted_names = {p.strip().lower() for p in products} matched = [svc for name, svc in _active_by_name.items() if name in _wanted_names] existing_svc_ids = ( {a.service_id for a in db.scalars(select(UserServiceAccess).where(UserServiceAccess.user_id == target_user.id)).all()} if is_renewal else set() ) for svc in matched: if svc.id not in existing_svc_ids: db.add(UserServiceAccess(user_id=target_user.id, service_id=svc.id)) db.commit() products_html = "" if products: items = "".join(f"
  • {escape(p)}
  • " for p in products) products_html = ( '

    Предоставлен доступ к продуктам:

    ' f'' ) day_word = "день" if days == 1 else ("дня" if days < 5 else "дней") email_action = "продлён" if is_renewal else "предоставлен" email_subject = ("Продление доступа к Инфраструктурному полигону MONT" if is_renewal else "Доступ к Инфраструктурному полигону MONT") email_subhead = "Ваш доступ продлён" if is_renewal else "Ваш запрос одобрен" cred_row = ( f'Пароль' f'{password}' ) if not is_renewal else "" access_text = f"Вам {email_action} доступ к полигону на {days} {day_word}." html_email = f"""
    MONT

    {email_subject}

    {email_subhead}


    Здравствуйте, {escape(pending.name)}!
    {access_text}

    {cred_row}
    Адрес портала {portal_url}
    Логин {username}
    Доступ до {expires.strftime("%d.%m.%Y")}
    {products_html}
    Если у вас возникли вопросы, свяжитесь с вашим менеджером MONT или напишите на RGalyaviev@mont.ru
    """ pending.status = "approved" db.commit() try: _send_email(pending.email, email_subject, html_email) email_status = "Email отправлен" except Exception as ex: log_event("email_send_error", error=str(ex)) email_status = f"Ошибка email: {ex}" _tg_api("editMessageText", { "chat_id": chat_id, "message_id": msg_id, "text": ( (f"🔄 Продлено на {days} {day_word}\n" if is_renewal else f"✅ Одобрено на {days} {day_word}\n") + f"👤 Логин: `{username}`\n" + (f"🔑 Пароль: `{password}`\n" if not is_renewal else "") + f"📧 {email_status}" ), "parse_mode": "Markdown", }) elif reject_match: manager_contact = pending.manager if pending.manager else "менеджера MONT" html_email = f"""
    MONT

    Запрос на доступ к полигону MONT


    Здравствуйте, {escape(pending.name)}!

    К сожалению, на данный момент мы не можем предоставить доступ к полигону.

    Для уточнения деталей свяжитесь с {escape(manager_contact)}.
    Если не знаете кто ваш менеджер — напишите на RGalyaviev@mont.ru.

    С уважением, команда MONT
    """ pending.status = "rejected" db.commit() try: _send_email(pending.email, "Запрос на доступ к полигону MONT", html_email) email_status = "Email отправлен" except Exception as ex: log_event("email_send_error", error=str(ex)) email_status = f"Ошибка email: {ex}" _tg_api("editMessageText", { "chat_id": chat_id, "message_id": msg_id, "text": f"❌ Отклонено. {email_status}", }) finally: db.close() app = FastAPI(title="MONT - инфрастуктурный полигон", docs_url=None, redoc_url=None, openapi_url=None) @app.exception_handler(_StarletteHTTPException) async def _auth_redirect_handler(request: Request, exc: _StarletteHTTPException): """A session that expired while a page tab sat open used to surface as a raw {"detail": "Unauthorized"} JSON blob on reload - require_user/ require_admin raise before the page route body ever runs, so there was nowhere in those routes to catch it and show something friendlier. Redirect browser page loads (never /api/* calls, whose JS callers parse the JSON body and show their own error) back to the login page instead.""" if ( exc.status_code in (401, 403) and request.method == "GET" and not request.url.path.startswith("/api/") ): return RedirectResponse(url="/", status_code=303) return JSONResponse(status_code=exc.status_code, content={"detail": exc.detail}, headers=exc.headers) @app.get("/admin/access-request/{req_id}/decide") def access_request_decide_confirm(req_id: str, action: str, token: str, db: Session = Depends(get_db)): """GET only renders a confirmation page - it must never mutate state, because corporate mail scanners (e.g. Microsoft Safe Links) pre-fetch every link in an email with a GET before a human reads it.""" if action not in _ACTION_DAYS or not _verify_decision_token(token, req_id, action): return _decision_page("Ссылка недействительна", "

    Ссылка недействительна или устарела.

    ") pending = db.get(PendingAccessRequest, req_id) if not pending: return _decision_page("Запрос не найден", "

    Запрос не найден (возможно, уже обработан).

    ") if pending.status != "pending": return _decision_page("Уже обработано", f"

    Запрос уже обработан: {pending.status}.

    ") decision_label = "Отказать" if action == "r" else f"Одобрить на {_ACTION_LABELS[action]}" body = f"""

    Подтвердите решение

    {escape(pending.name)} ({escape(pending.company)})
    {escape(pending.email)} · {escape(pending.phone)}

    Решение: {decision_label}

    """ return _decision_page("Подтверждение решения", body) @app.post("/admin/access-request/{req_id}/decide") def access_request_decide_apply(req_id: str, action: str = Form(...), token: str = Form(...), db: Session = Depends(get_db)): if action not in _ACTION_DAYS or not _verify_decision_token(token, req_id, action): return _decision_page("Ссылка недействительна", "

    Ссылка недействительна или устарела.

    ") pending = db.get(PendingAccessRequest, req_id) if not pending: return _decision_page("Запрос не найден", "

    Запрос не найден (возможно, уже обработан).

    ") if pending.status != "pending": return _decision_page("Уже обработано", f"

    Запрос уже обработан: {pending.status}.

    ") days = _ACTION_DAYS[action] result = _apply_access_decision(db, pending, days) if result["kind"] == "approved": day_word = "день" if days == 1 else ("дня" if days < 5 else "дней") cred = f"

    Логин: {escape(result['username'])}

    " if not result["is_renewal"]: cred += f"

    Пароль: {escape(result['password'])}

    " title_word = "Продлено" if result["is_renewal"] else "Одобрено" body = f"""

    {title_word} на {days} {day_word}

    {cred}

    {escape(result['email_status'])}

    """ else: body = f"""

    Отклонено

    {escape(result['email_status'])}

    """ return _decision_page("Готово", body) app.mount("/static", StaticFiles(directory="static"), name="static") @app.middleware("http") async def request_logging_middleware(request: Request, call_next): req_id = request.headers.get("X-Request-ID", str(uuid.uuid4())[:8]) token = request_id_ctx.set(req_id) started = time.time() client_ip = request.client.host if request.client else "-" user_agent = request.headers.get("user-agent", "-") try: response = await call_next(request) except Exception: log_event( "request_failed", level=logging.ERROR, method=request.method, path=request.url.path, client_ip=client_ip, user_agent=user_agent, ) request_id_ctx.reset(token) raise duration_ms = int((time.time() - started) * 1000) level = logging.INFO if response.status_code >= 500: level = logging.ERROR elif response.status_code >= 400: level = logging.WARNING log_event( "request", level=level, method=request.method, path=request.url.path, query=str(request.url.query or ""), status=response.status_code, duration_ms=duration_ms, client_ip=client_ip, user_agent=user_agent, ) if duration_ms >= LOG_SLOW_REQUEST_MS: log_event( "slow_request", level=logging.WARNING, method=request.method, path=request.url.path, duration_ms=duration_ms, threshold_ms=LOG_SLOW_REQUEST_MS, ) response.headers["X-Request-ID"] = req_id request_id_ctx.reset(token) return response _MOBILE_UA_RE = re.compile( r"(Mobile|Android|iPhone|iPad|iPod|BlackBerry|IEMobile|Opera Mini|webOS)", re.IGNORECASE, ) _MOBILE_PAGE = ( "" '' "" '' '' "MONT - инфрастуктурный полигон" "" "" "" '' '
    🖥
    ' "

    Ресурс доступен
    только с ПК

    " "

    Пожалуйста, откройте эту страницу на компьютере или ноутбуке.

    " "" "" ) # Public marketing/auth surface that has to work from a phone: the login # page itself, submitting the login form, the "request access" flow (public # form + the products list it loads), and static legal/SEO pages. Everything # else (the actual stand: /go/, /svc/, /s/, /u/, /w/, /rdp/, /admin, ...) # stays desktop-only - a phone can't usefully drive a remote desktop/browser # session anyway, and the dashboard has its own width-based #mobile-wall on # top of this for narrow desktop windows. _MOBILE_ALLOWED_PATHS = { "/", "/login", "/privacy", "/robots.txt", "/sitemap.xml", "/favicon.ico", "/api/public/services-by-category", "/api/request-access", } _MOBILE_ALLOWED_PREFIXES = ("/static/", "/admin/access-request/", "/product/") def _looks_authenticated(request: Request) -> bool: """Best-effort check for a plausibly-valid session cookie, without a DB round trip. Good enough to decide the mobile gate on "/" - the actual routes still enforce real auth via get_current_user/require_user.""" raw = request.cookies.get(COOKIE_NAME) if not raw: return False try: serializer.loads(raw, max_age=COOKIE_MAX_AGE) except Exception: return False return True @app.middleware("http") async def mobile_block_middleware(request: Request, call_next): path = request.url.path if path.startswith(_MOBILE_ALLOWED_PREFIXES): # email-approval links are meant to work from a phone too, same as # approving from the Telegram app. return await call_next(request) if path in _MOBILE_ALLOWED_PATHS: # "/" is shared by the logged-out login page (fine on a phone) and # the authenticated dashboard (not fine - a phone can't drive a # remote desktop/browser session), so it needs the real auth-cookie # check instead of a blanket allow. Every other allowed path here # (login POST, privacy, robots/sitemap, the public request-access # API) never renders the dashboard, so they're always fine. if path != "/" or not _looks_authenticated(request): return await call_next(request) ua = request.headers.get("user-agent", "") if _MOBILE_UA_RE.search(ua): return _HR(content=_MOBILE_PAGE, status_code=200) return await call_next(request) @app.on_event("startup") async def startup_event(): global _tg_poll_lock_file on_startup() import asyncio as _aio, fcntl as _fcntl _lf = open("/tmp/portal-tg-poll.lock", "w") try: _fcntl.flock(_lf.fileno(), _fcntl.LOCK_EX | _fcntl.LOCK_NB) _tg_poll_lock_file = _lf _aio.create_task(_telegram_poll_loop()) _aio.create_task(_telegram_notify_retry_loop()) except BlockingIOError: _lf.close() def _login_wall_services(db: Session): """Active, non-pilot services shown as a logo wall on the public login page - the self-service "pick a product" catalog. Pilots are invite-only and get their own separate, non-clickable teaser (_login_wall_pilots).""" return db.scalars( select(Service).where(Service.active == True, Service.is_pilot == False).order_by(Service.name) ).all() def _login_wall_pilots(db: Session): """Active pilot services shown as a purely informational logo row on the public login page - awareness only, not part of the self-service catalog (no request-access entry point; admin grants these by hand).""" return db.scalars( select(Service).where(Service.active == True, Service.is_pilot == True).order_by(Service.name) ).all() @app.get("/", response_class=HTMLResponse) def index(request: Request, user: Optional[User] = Depends(get_current_user), db: Session = Depends(get_db)): session_closed = (request.query_params.get("session_closed") or "").strip().lower() launch_error = (request.query_params.get("launch_error") or "").strip().lower() session_notice = "" if session_closed == "idle": session_notice = "Сессия была закрыта из-за простоя. Откройте сервис заново." elif session_closed == "limit": session_notice = ( f"Сессия была закрыта из-за лимита в {MAX_ACTIVE_SERVICES_PER_USER} сервиса(ов). " "Освободите один сервис и попробуйте снова." ) elif launch_error == "max_services": session_notice = ( f"Есть ограничение на {MAX_ACTIVE_SERVICES_PER_USER} сервиса(ов). " "Освободите один сервис и попробуйте снова." ) if not user: csrf = request.cookies.get(CSRF_COOKIE) or secrets.token_urlsafe(24) response = templates.TemplateResponse( "login.html", { "request": request, "csrf_token": csrf, "login_error": "", "session_notice": session_notice, "public_services": _login_wall_services(db), "public_pilots": _login_wall_pilots(db), }, ) response.set_cookie(CSRF_COOKIE, csrf, httponly=False, secure=True, samesite="lax", path="/") return response all_services = db.scalars( select(Service) .where(Service.active == True, Service.type.in_([ServiceType.WEB, ServiceType.RDP])) .order_by(Service.name) ).all() # Per-grant expires_at (used by pilots) can make a grant expired even # though the row still exists and the account itself is still valid - # exclude those rows here so an expired pilot grant falls back into # locked_services instead of staying "granted". granted_ids = set( db.scalars( select(UserServiceAccess.service_id).where( UserServiceAccess.user_id == user.id, (UserServiceAccess.expires_at.is_(None)) | (UserServiceAccess.expires_at > now_utc()), ) ).all() ) # Pilots the user has access to get their own collapsible section above # the regular catalog instead of sitting in the "Ваши сервисы" grid. pilot_services = [svc for svc in all_services if svc.id in granted_ids and svc.is_pilot] services = [svc for svc in all_services if svc.id in granted_ids and not svc.is_pilot] locked_services = [svc for svc in all_services if svc.id not in granted_ids] pilot_ids = {svc.id for svc in all_services if svc.is_pilot} # Categories are computed across the whole catalog (granted + locked) so # the nav lets a user browse into a category they don't have access to # yet and request it from there. service_categories = {svc.id: [] for svc in all_services} categories = [] if all_services: service_ids = [svc.id for svc in all_services] rows = db.execute( select(ServiceCategory.service_id, Category.id, Category.name, Category.slug) .join(Category, Category.id == ServiceCategory.category_id) .where(ServiceCategory.service_id.in_(service_ids)) .order_by(Category.name) ).all() category_map = {} for service_id, category_id, category_name, category_slug in rows: service_categories.setdefault(service_id, []).append( { "id": category_id, "name": category_name, "slug": category_slug, } ) if category_id not in category_map: category_map[category_id] = {"id": category_id, "name": category_name, "slug": category_slug} categories = sorted(category_map.values(), key=lambda x: x["name"].lower()) # Stable per-category counts (granted + locked), computed before any # ?category= filtering is applied - the rail always shows the full # catalog's numbers, not just what's currently on screen. category_counts = { cat["slug"]: sum( 1 for svc in all_services if any(c["slug"] == cat["slug"] for c in service_categories.get(svc.id, [])) ) for cat in categories } selected_category_slug = (request.query_params.get("category") or "").strip().lower() if selected_category_slug: services = [ svc for svc in services if any(cat["slug"] == selected_category_slug for cat in service_categories.get(svc.id, [])) ] locked_services = [ svc for svc in locked_services if any(cat["slug"] == selected_category_slug for cat in service_categories.get(svc.id, [])) ] service_comment_html = {svc.id: format_service_comment(svc.comment) for svc in all_services} return templates.TemplateResponse( "dashboard.html", { "request": request, "user": user, "services": services, "locked_services": locked_services, "pilot_services": pilot_services, "pilot_ids": pilot_ids, "categories": categories, "category_counts": category_counts, "total_catalog_count": len(all_services), "selected_category_slug": selected_category_slug, "service_categories": service_categories, "service_comment_html": service_comment_html, "csrf_token": request.cookies.get(CSRF_COOKIE, ""), "session_notice": session_notice, "max_active_services": MAX_ACTIVE_SERVICES_PER_USER, "idle_timeout_min": SESSION_IDLE_SECONDS // 60, }, ) @app.get("/bundles", response_class=HTMLResponse) def bundles_page(request: Request, user: User = Depends(require_user)): """Stub landing page for the upcoming "готовые связки" feature - curated multi-product infrastructure scenarios (e.g. ALDpro + workstations, or RuBackup + Alt PVE) granted and launched as one bundle instead of picking products one by one. Not wired to real data yet - the scenario list here is static copy, kept in sync by hand with the brochure draft until the Bundle/BundleService models + request-access integration are built.""" return templates.TemplateResponse( "bundles.html", { "request": request, "user": user, "max_active_services": MAX_ACTIVE_SERVICES_PER_USER, "idle_timeout_min": SESSION_IDLE_SECONDS // 60, }, ) @app.get("/admin", response_class=HTMLResponse) def admin_page(request: Request, admin: User = Depends(require_admin), db: Session = Depends(get_db)): users = db.scalars(select(User).order_by(User.id)).all() categories = db.scalars(select(Category).order_by(Category.name)).all() services = db.scalars(select(Service).where(Service.type.in_([ServiceType.WEB, ServiceType.RDP])).order_by(Service.id)).all() web_services = [s for s in services if s.type == ServiceType.WEB] rdp_services = [s for s in services if s.type == ServiceType.RDP and not s.is_pilot] pilot_services = [s for s in services if s.is_pilot] pilot_service_ids = {s.id for s in pilot_services} service_category_map = {s.id: [] for s in services} if services: service_rows = db.execute( select(ServiceCategory.service_id, ServiceCategory.category_id).where( ServiceCategory.service_id.in_([s.id for s in services]) ) ).all() for service_id, category_id in service_rows: service_category_map.setdefault(service_id, []).append(category_id) acl_rows = db.scalars(select(UserServiceAccess)).all() acl = {} for row in acl_rows: acl.setdefault(row.user_id, []).append(row.service_id) for user_id in acl: acl[user_id] = sorted(acl[user_id]) pool_status = {s.id: get_pool_status_for_service(s) for s in services} service_health = {} for sid, st in pool_status.items(): service_health[sid] = { "health": st["health"], "running": st["running"], "desired": st["desired"], "active_sessions": get_active_sessions_count(db, sid), } web_pool = get_web_pool_status() web_totals = { "services": len(web_services), "running": web_pool["running"], "desired": web_pool["desired"], "active_sessions": sum(service_health[s.id]["active_sessions"] for s in web_services), } recent_sessions = db.execute( text( """ SELECT s.id, u.username, sv.name AS service_name, sv.slug AS service_slug, s.status, s.created_at, s.last_access_at FROM sessions s JOIN users u ON u.id = s.user_id JOIN services sv ON sv.id = s.service_id WHERE sv.type IN ('WEB','RDP') ORDER BY s.created_at DESC LIMIT 200 """ ) ).mappings().all() open_stats = db.execute( text( """ SELECT u.username, sv.name AS service_name, sv.slug AS service_slug, COUNT(*) AS opens FROM sessions s JOIN users u ON u.id = s.user_id JOIN services sv ON sv.id = s.service_id WHERE sv.type IN ('WEB','RDP') GROUP BY u.username, sv.name, sv.slug ORDER BY opens DESC, u.username ASC LIMIT 200 """ ) ).mappings().all() cutoff = now_utc() - dt.timedelta(seconds=SESSION_IDLE_SECONDS) online_sessions = db.execute( text( """ SELECT s.id, u.username, sv.name AS service_name, sv.slug AS service_slug, sv.type AS service_type, s.container_id, s.created_at, s.last_access_at FROM sessions s JOIN users u ON u.id = s.user_id JOIN services sv ON sv.id = s.service_id WHERE s.status = 'ACTIVE' AND s.last_access_at >= :cutoff AND sv.type IN ('WEB','RDP') ORDER BY s.last_access_at DESC, s.created_at DESC LIMIT 500 """ ), {"cutoff": cutoff}, ).mappings().all() rdp_slots: dict[int, list] = {} for svc in rdp_services + pilot_services: slots = db.scalars(select(RdpSlot).where(RdpSlot.service_id == svc.id).order_by(RdpSlot.id)).all() slot_list = [] for slot in slots: active_sess = db.scalar( select(SessionModel).where( SessionModel.container_id == f"RDPSLOT:{slot.id}", SessionModel.status == SessionStatus.ACTIVE, ) ) running = False if slot.container_name: try: c = docker_client().containers.get(slot.container_name) running = c.status == "running" except Exception: pass occupied_username = None if active_sess: u = db.get(User, active_sess.user_id) occupied_username = u.username if u else f"id={active_sess.user_id}" assigned_username = None assigned_expires_at = None if slot.assigned_user_id: au = db.get(User, slot.assigned_user_id) assigned_username = au.username if au else f"id={slot.assigned_user_id}" access_row = db.scalar( select(UserServiceAccess).where( UserServiceAccess.user_id == slot.assigned_user_id, UserServiceAccess.service_id == svc.id, ) ) if access_row and access_row.expires_at: assigned_expires_at = access_row.expires_at.isoformat() slot_list.append({ "id": slot.id, "rdp_username": slot.rdp_username, "container_name": slot.container_name or "", "running": running, "occupied_username": occupied_username, "assigned_user_id": slot.assigned_user_id, "assigned_username": assigned_username, "assigned_expires_at": assigned_expires_at, }) rdp_slots[svc.id] = slot_list return templates.TemplateResponse( "admin.html", { "request": request, "admin": admin, "users": users, "users_json": [ {"id": u.id, "label": (f"{u.first_name} {u.last_name}".strip() or u.username) + f" ({u.username})"} for u in users ], "web_services": web_services, "rdp_services": rdp_services, "pilot_services": pilot_services, "pilot_service_ids": pilot_service_ids, "services": services, "categories": categories, "service_category_map": service_category_map, "acl": acl, "pool_status": pool_status, "service_health": service_health, "web_totals": web_totals, "web_pool_size": WEB_POOL_SIZE, "web_pool_buffer": WEB_POOL_BUFFER, "recent_sessions": recent_sessions, "open_stats": open_stats, "online_sessions": online_sessions, "csrf_token": request.cookies.get(CSRF_COOKIE, ""), "max_active_services_per_user": MAX_ACTIVE_SERVICES_PER_USER, "rdp_slots": rdp_slots, }, ) @app.get("/yandex_b847b9b35f967fcc.html", include_in_schema=False) def yandex_verify(): from fastapi.responses import HTMLResponse return HTMLResponse(''' Verification: b847b9b35f967fcc ''') @app.get("/privacy", include_in_schema=False) def privacy_page(): from fastapi.responses import RedirectResponse return RedirectResponse(url="https://www.mont.ru/ru-ru/privacy", status_code=301) @app.post("/api/telegram-webhook", include_in_schema=False) async def telegram_webhook(request: Request, db: Session = Depends(get_db)): import json as _jw import datetime as _dt2 try: data = await request.json() except Exception: return {"ok": True} cq = data.get("callback_query") if not cq: return {"ok": True} cq_id = cq["id"] cb_data = cq.get("data", "") chat_id = cq["message"]["chat"]["id"] msg_id = cq["message"]["message_id"] try: _tg_api("answerCallbackQuery", {"callback_query_id": cq_id}) except Exception: pass # parse callback_data: a7_ID, a14_ID, a30_ID, a90_ID, r_ID import re as _rew approve_match = _rew.match(r'^a(\d+)_(.+)$', cb_data) reject_match = _rew.match(r'^r_(.+)$', cb_data) if not approve_match and not reject_match: return {"ok": True} req_id = approve_match.group(2) if approve_match else reject_match.group(1) pending = db.get(PendingAccessRequest, req_id) if not pending: try: _tg_api("editMessageText", { "chat_id": chat_id, "message_id": msg_id, "text": "Запрос не найден (возможно уже обработан).", }) except Exception: pass return {"ok": True} if pending.status != "pending": try: _tg_api("editMessageText", { "chat_id": chat_id, "message_id": msg_id, "text": f"Запрос уже обработан: {pending.status}.", }) except Exception: pass return {"ok": True} products = _jw.loads(pending.products_json or "[]") portal_url = pending.portal_url or PORTAL_URL if approve_match: days = int(approve_match.group(1)) password = _generate_password() username = pending.email expires = _dt2.datetime.now(_dt2.timezone.utc) + _dt2.timedelta(days=days) parts = pending.name.strip().split(None, 1) existing_user = db.scalar(select(User).where(User.username == username)) is_renewal = existing_user is not None if is_renewal: password = None existing_user.expires_at = expires existing_user.active = True target_user = existing_user db.flush() else: new_user = User( username=username, password_hash=hash_password(password), expires_at=expires, active=True, is_admin=False, first_name=parts[0] if parts else "", last_name=parts[1] if len(parts) > 1 else "", ) db.add(new_user) db.flush() target_user = new_user # assign requested services if products: from sqlalchemy import func as _func matched = db.scalars( select(Service).where( func.lower(Service.name).in_([p.lower() for p in products]), Service.active == True, ) ).all() existing_svc_ids = ( {a.service_id for a in db.scalars(select(UserServiceAccess).where(UserServiceAccess.user_id == target_user.id)).all()} if is_renewal else set() ) for svc in matched: if svc.id not in existing_svc_ids: db.add(UserServiceAccess(user_id=target_user.id, service_id=svc.id)) db.commit() # send approval email products_html = "" if products: items = "".join(f"
  • {escape(p)}
  • " for p in products) products_html = f"

    Предоставлен доступ к продуктам:

    " html_email = f"""
    MONT

    Доступ к Инфраструктурному полигону MONT

    Ваш запрос одобрен


    Здравствуйте, {escape(pending.name)}!
    Вам предоставлен доступ к полигону на {days} {'день' if days==1 else 'дня' if days<5 else 'дней'}.

    Адрес портала {portal_url}
    Логин {username}
    Пароль {password}
    Доступ до {expires.strftime('%d.%m.%Y')}
    {products_html}
    Если у вас возникли вопросы, свяжитесь с вашим менеджером MONT или напишите на RGalyaviev@mont.ru
    """ pending.status = "approved" db.commit() try: _send_email(pending.email, "Доступ к Инфраструктурному полигону MONT", html_email) email_status = "Email отправлен" except Exception as ex: log_event("email_send_error", error=str(ex)) email_status = f"Ошибка отправки email: {ex}" try: _tg_api("editMessageText", { "chat_id": chat_id, "message_id": msg_id, "text": ( f"Одобрено на {days} дней\n" f"Логин: `{username}`\n" f"Пароль: `{password}`\n" f"{email_status}" ), "parse_mode": "Markdown", }) except Exception: pass elif reject_match: manager_contact = pending.manager if pending.manager else "менеджера MONT" html_email = f"""
    MONT

    Запрос на доступ к полигону MONT


    Здравствуйте, {escape(pending.name)}!

    К сожалению, на данный момент мы не можем предоставить доступ к полигону.

    Для уточнения деталей, пожалуйста, свяжитесь с {escape(manager_contact)}.
    Если вы не знаете, кто ваш менеджер, напишите нам на RGalyaviev@mont.ru — мы поможем.

    С уважением, команда MONT
    """ pending.status = "rejected" db.commit() try: _send_email(pending.email, "Запрос на доступ к полигону MONT", html_email) email_status = "Email отправлен" except Exception as ex: log_event("email_send_error", error=str(ex)) email_status = f"Ошибка отправки email: {ex}" try: _tg_api("editMessageText", { "chat_id": chat_id, "message_id": msg_id, "text": f"Отклонено. {email_status}", }) except Exception: pass return {"ok": True} @app.api_route("/favicon.ico", methods=["GET", "HEAD"], include_in_schema=False) def favicon(): from fastapi.responses import FileResponse return FileResponse("static/favicon.ico", media_type="image/x-icon") @app.get("/product/{slug}", response_class=HTMLResponse) def product_page(slug: str, request: Request, db: Session = Depends(get_db)): """Public, unauthenticated SEO landing page for a single catalog product - one per active Service row, keyed by slug. There is no separate "generation" step: as soon as a service is saved active in the admin panel its page is reachable here, and sitemap_xml() below picks it up on its next request too. Mirrors the login page's public "request access" flow (same /api/request-access endpoint + modal), just pre-scoped to this one product. """ service = db.scalar( select(Service).where(Service.slug == slug, Service.active == True) ) if not service: raise HTTPException(status_code=404, detail="Продукт не найден") other_services = db.scalars( select(Service) .where(Service.active == True, Service.id != service.id) .order_by(Service.name) ).all() svc_categories = db.scalars( select(Category) .join(ServiceCategory, ServiceCategory.category_id == Category.id) .where(ServiceCategory.service_id == service.id) .order_by(Category.name) ).all() canonical_url = f"{PORTAL_URL}/product/{service.slug}" meta_source = service.seo_description or service.comment or "" # Drop a leading '# Title' markdown heading line (redundant with the # page's own /<h1>) before flattening to plain text, same as # format_seo_description() does for the on-page HTML. meta_lines = meta_source.strip().split("\n") if meta_lines and re.match(r"^#\s+", meta_lines[0]): meta_lines = meta_lines[1:] meta_plain = re.sub(r"[#*`>\-]+", " ", "\n".join(meta_lines)) meta_plain = re.sub(r"\s+", " ", meta_plain).strip() meta_description = (meta_plain[:157] + "…") if len(meta_plain) > 160 else meta_plain if not meta_description: meta_description = f"{service.name} — протестируйте продукт бесплатно на инфраструктурном полигоне MONT." description_html = format_seo_description(service.seo_description or service.comment) return templates.TemplateResponse( "product.html", { "request": request, "service": service, "categories": svc_categories, "description_html": description_html, "other_services": other_services, "canonical_url": canonical_url, "meta_description": meta_description, "portal_url": PORTAL_URL, }, ) @app.get("/robots.txt", include_in_schema=False) def robots_txt(): from fastapi.responses import FileResponse return FileResponse("static/robots.txt", media_type="text/plain") @app.get("/sitemap.xml", include_in_schema=False) def sitemap_xml(db: Session = Depends(get_db)): # Built from the live catalog (not a static file) so a new /product/{slug} # page shows up here the moment its Service row is saved active, with no # separate publish step - see product_page() above. from fastapi.responses import Response as _XmlResponse slugs = db.scalars( select(Service.slug).where(Service.active == True).order_by(Service.name) ).all() entries = [(PORTAL_URL + "/", "1.0")] entries += [(f"{PORTAL_URL}/product/{slug}", "0.8") for slug in slugs] url_blocks = [ f" <url>\n <loc>{loc}</loc>\n <changefreq>weekly</changefreq>\n <priority>{priority}</priority>\n </url>" for loc, priority in entries ] xml = ( '<?xml version="1.0" encoding="UTF-8"?>\n' '<urlset xmlns="http://www.sitemaps.org/schemas/sitemap/0.9">\n' + "\n".join(url_blocks) + "\n</urlset>\n" ) return _XmlResponse(content=xml, media_type="application/xml") @app.get("/api/public/services-by-category") def public_services_by_category(db: Session = Depends(get_db)): # Pilots are invite-only - admin grants them by hand, so they must not # be selectable from the self-service "request access" form. services = db.execute( select(Service).where(Service.active == True, Service.is_pilot == False).order_by(Service.name) ).scalars().all() categories = db.execute(select(Category).order_by(Category.name)).scalars().all() cat_map = {c.id: c.name for c in categories} svc_cats: dict[int, list[str]] = {} links = db.execute(select(ServiceCategory)).scalars().all() for lnk in links: svc_cats.setdefault(lnk.service_id, []).append(cat_map.get(lnk.category_id, "")) result: dict[str, list[dict]] = {"Без категории": []} for svc in services: cats = svc_cats.get(svc.id, []) entry = {"id": svc.id, "name": svc.name} if cats: for cat in cats: result.setdefault(cat, []).append(entry) else: result["Без категории"].append(entry) if not result["Без категории"]: del result["Без категории"] return result _public_form_attempts: dict = {} _public_form_lock = threading.Lock() _PUBLIC_FORM_MAX = 5 # заявок/сообщений _PUBLIC_FORM_WINDOW = 600 # за 10 минут с одного IP _PUBLIC_FORM_BLOCK = 1800 # затем блок на 30 минут def _public_form_rate_limited(bucket: str, ip: str) -> bool: """Простой лимитер на публичные формы портала (заявка на доступ, форма обратной связи) - защита от скриптов, которые заваливают Telegram/почту мусорными/атакующими payload'ами (form-спам, XSS/SQLi-пробы и т.п.). `bucket` разделяет счётчики между формами, чтобы спам по одной не блокировал другую.""" key = f"{bucket}:{ip}" now = time.monotonic() with _public_form_lock: entry = _public_form_attempts.get(key) if entry and entry["blocked_until"] > now: return True if entry and now - entry["first"] > _PUBLIC_FORM_WINDOW: entry = None if not entry: _public_form_attempts[key] = {"count": 1, "first": now, "blocked_until": 0.0} return False entry["count"] += 1 if entry["count"] > _PUBLIC_FORM_MAX: entry["blocked_until"] = now + _PUBLIC_FORM_BLOCK return True return False @app.post("/api/request-access") async def request_access(request: Request, db: Session = Depends(get_db)): ip_for_limit = _get_real_ip(request) if _public_form_rate_limited("request-access", ip_for_limit): raise HTTPException(status_code=429, detail="Слишком много заявок с вашего адреса. Попробуйте позже.") try: data = await request.json() except Exception: raise HTTPException(status_code=400, detail="Invalid JSON") name = str(data.get("name", "")).strip()[:200] company = str(data.get("company", "")).strip()[:200] email = str(data.get("email", "")).strip()[:254] phone = str(data.get("phone", "")).strip()[:32] manager = str(data.get("manager", "")).strip()[:200] products_raw = data.get("products", []) products = [str(p).strip()[:100] for p in products_raw[:200]] if isinstance(products_raw, list) else [] import re as _re if not name or not company or not email or not phone: raise HTTPException(status_code=422, detail="Заполните все обязательные поля") if not _re.match(r'^[^\s@]+@[^\s@]+\.[^\s@]+$', email): raise HTTPException(status_code=422, detail="Некорректный email") if not _re.match(r'^[\+\d][\d\s\-\(\)]{6,18}$', phone): raise HTTPException(status_code=422, detail="Некорректный номер телефона") import html as _html def _e(s): return _html.escape(str(s)) products_text = "" if products: items = "\n".join(f" • {_e(p)}" for p in products) products_text = f"\n\n🖥 <b>Интересующие продукты:</b>\n{items}" divider = "━━━━━━━━━━━━━━━━━━━━━━" manager_text = f"\n🤝 <b>Менеджер MONT:</b> {_e(manager)}" if manager else "" text = ( f"🔔 <b>Новый запрос доступа к полигону MONT</b>\n" f"{divider}\n\n" f"👤 <b>Имя:</b> {_e(name)}\n" f"🏢 <b>Компания:</b> {_e(company)}\n" f"📧 <b>Email:</b> {_e(email)}\n" f"📱 <b>Телефон:</b> {_e(phone)}" f"{manager_text}" f"{products_text}" ) log_event("ip_headers", xff=request.headers.get("x-forwarded-for","–"), xri=request.headers.get("x-real-ip","–"), client=str(request.client.host if request.client else "–")) ip = _get_real_ip(request) geo = _get_geo(ip) geo_text = "" if geo: geo_text += "\n📍 <b>Местоположение:</b> " + _e(geo) geo_text += "\n🖥 <b>IP:</b> " + _e(ip) text += geo_text if not TELEGRAM_BOT_TOKEN or not TELEGRAM_CHAT_ID: log_event("telegram_not_configured") return {"ok": True} # save pending request import json as _j2 req_id = _secrets.token_urlsafe(8)[:12] _origin = request.headers.get('origin', '') if 'stand.mont.ru' in _origin: _req_portal_url = 'https://stand.mont.ru' else: _req_portal_url = PORTAL_URL pending = PendingAccessRequest( id=req_id, name=name, company=company, email=email, phone=phone, manager=manager, products_json=_j2.dumps(products, ensure_ascii=False), portal_url=_req_portal_url, telegram_message=text, ) db.add(pending) db.commit() try: _tg_api("sendMessage", { "chat_id": TELEGRAM_CHAT_ID, "text": text, "parse_mode": "HTML", "reply_markup": _make_approval_keyboard(req_id), }) pending.telegram_notified = True db.commit() except Exception as e: log_event("telegram_send_error", error=str(e)) # not fatal here - _telegram_notify_retry_loop() will keep retrying # every 5 minutes (telegram_notified stays False) until it goes through # second, independent approval channel - failures here must never affect # the Telegram channel above (already sent) or the API response below _send_admin_decision_email(pending) return {"ok": True} # /api/contact ("написать Руслану") was removed 2026-08-10: the UI entry # point (#btn-contact-ruslan button in login.html) was already taken off the # page earlier, but the route itself stayed registered and reachable by # anyone calling it directly - which is exactly what an attacker was doing. # Not registering the route at all means FastAPI returns a plain 404 for it # without running any app code. If this form is ever brought back, restore # it from git history / main.py.bak.* on the server instead of re-adding a # half-remembered version here. def _notify_pending_request(pending_id: str, text: str) -> None: """Send the Telegram + email approval notifications for a pending request in the background, after the HTTP response has already gone out - both calls are external network I/O (Telegram API, SMTP) and have nothing to do with whether the request was accepted.""" from database import SessionLocal db2 = SessionLocal() try: pending = db2.get(PendingAccessRequest, pending_id) if not pending: return try: _tg_api("sendMessage", { "chat_id": TELEGRAM_CHAT_ID, "text": text, "parse_mode": "HTML", "reply_markup": _make_approval_keyboard(pending_id), }) pending.telegram_notified = True db2.commit() except Exception as e: log_event("telegram_send_error", error=str(e)) # not fatal - _telegram_notify_retry_loop() keeps retrying every 5 min _send_admin_decision_email(pending) finally: db2.close() @app.post("/api/request-more-access") async def request_more_access( request: Request, background_tasks: BackgroundTasks, user: User = Depends(require_user), db: Session = Depends(get_db), ): """A logged-in user asking for additional products beyond what they already have. Deliberately reuses the exact same pending-request / Telegram-approval / _apply_access_decision pipeline as the public "Запросить доступ" flow on the login page - approving it renews the user's expires_at *and* grants the newly-requested services without touching what they already have (see _apply_access_decision). We skip company/phone here since the account already identifies the requester; PendingAccessRequest.company/phone just get an explanatory placeholder so it reads clearly in Telegram/email, not a real company/phone value.""" validate_csrf(request) if _public_form_rate_limited("request-more-access", f"user:{user.id}"): raise HTTPException(status_code=429, detail="Слишком много заявок подряд. Попробуйте позже.") try: data = await request.json() except Exception: raise HTTPException(status_code=400, detail="Invalid JSON") products_raw = data.get("products", []) requested = [str(p).strip()[:200] for p in products_raw[:200]] if isinstance(products_raw, list) else [] requested = [p for p in requested if p] note = str(data.get("note", "")).strip()[:500] if not requested: raise HTTPException(status_code=422, detail="Выберите хотя бы один продукт") already_granted = { row[0].lower() for row in db.execute( select(Service.name) .join(UserServiceAccess, UserServiceAccess.service_id == Service.id) .where(UserServiceAccess.user_id == user.id) ).all() } from sqlalchemy import func as _func4 # is_pilot excluded even if the caller bypasses the UI and posts a pilot's # name directly - pilots are admin-granted only, never self-service. matched = db.scalars( select(Service).where( _func4.lower(Service.name).in_([p.lower() for p in requested]), Service.active == True, Service.is_pilot == False, ) ).all() products = [svc.name for svc in matched if svc.name.lower() not in already_granted] if not products: raise HTTPException(status_code=422, detail="Эти продукты уже доступны или не найдены в каталоге") display_name = (f"{user.first_name} {user.last_name}".strip()) or user.username def _e(s): import html as _html_local return _html_local.escape(str(s)) items = "\n".join(f" • {_e(p)}" for p in products) note_text = f"\n\n💬 <b>Комментарий:</b> {_e(note)}" if note else "" divider = "━━━━━━━━━━━━━━━━━━━━━━" text = ( f"🔔 <b>Запрос дополнительного доступа</b>\n" f"{divider}\n\n" f"👤 <b>Пользователь:</b> {_e(display_name)} ({_e(user.username)})\n" f"🖥 <b>Запрошенные продукты:</b>\n{items}" f"{note_text}" ) ip = _get_real_ip(request) geo = _get_geo(ip) geo_text = "" if geo: geo_text += "\n📍 <b>Местоположение:</b> " + _e(geo) geo_text += "\n🖥 <b>IP:</b> " + _e(ip) text += geo_text req_id = _secrets.token_urlsafe(8)[:12] origin = request.headers.get("origin", "") portal_url = "https://stand.mont.ru" if "stand.mont.ru" in origin else PORTAL_URL pending = PendingAccessRequest( id=req_id, name=display_name, company="Существующий пользователь портала", email=user.username, phone="", manager="", products_json=__import__("json").dumps(products, ensure_ascii=False), portal_url=portal_url, telegram_message=text, ) db.add(pending) db.commit() if TELEGRAM_BOT_TOKEN and TELEGRAM_CHAT_ID: background_tasks.add_task(_notify_pending_request, req_id, text) else: log_event("telegram_not_configured") return {"ok": True} @app.post("/login") def login( request: Request, username: str = Form(...), password: str = Form(...), csrf_token: str = Form(...), db: Session = Depends(get_db), ): cookie_csrf = request.cookies.get(CSRF_COOKIE) if not cookie_csrf or csrf_token != cookie_csrf: raise HTTPException(status_code=403, detail="CSRF failed") ip = _get_real_ip(request) if check_login_rate_limit(ip): raise HTTPException(status_code=429, detail="Слишком много попыток входа. Попробуйте через 15 минут.") user = db.scalar(select(User).where(User.username == username)) if not user or not verify_password(password, user.password_hash): record_login_failure(ip) csrf = request.cookies.get(CSRF_COOKIE) or secrets.token_urlsafe(24) response = templates.TemplateResponse( "login.html", { "request": request, "csrf_token": csrf, "login_error": "Неверный логин или пароль", "session_notice": "", "public_services": _login_wall_services(db), "public_pilots": _login_wall_pilots(db), }, status_code=401, ) response.set_cookie(CSRF_COOKIE, csrf, httponly=False, secure=True, samesite="lax", path="/") return response if not user_is_valid(user): csrf = request.cookies.get(CSRF_COOKIE) or secrets.token_urlsafe(24) response = templates.TemplateResponse( "login.html", { "request": request, "csrf_token": csrf, "login_error": "Доступ к сервису приостоновлен, обратитесь к вашему менеджеру", "session_notice": "", "public_services": _login_wall_services(db), "public_pilots": _login_wall_pilots(db), }, status_code=403, ) response.set_cookie(CSRF_COOKIE, csrf, httponly=False, secure=True, samesite="lax", path="/") return response record_login_success(ip) response = RedirectResponse(url="/", status_code=303) issue_auth_cookie(response, user) issue_csrf_cookie(response) audit(db, "LOGIN", f"login success: {username}", user_id=user.id) return response @app.post("/logout") def logout(request: Request): response = RedirectResponse(url="/", status_code=303) response.delete_cookie(COOKIE_NAME, path="/") response.delete_cookie(CSRF_COOKIE, path="/") return response @app.get("/go/{slug}") def go_service( slug: str, sw: Optional[int] = Query(default=None, ge=320, le=7680), sh: Optional[int] = Query(default=None, ge=240, le=4320), user: User = Depends(require_user), db: Session = Depends(get_db), ): total_started = time.perf_counter() phase_ms = {} def _mark(name: str, started: float) -> None: phase_ms[name] = int((time.perf_counter() - started) * 1000) def _emit(result: str, **extra) -> None: payload = { "user_id": user.id, "service_slug": slug, "result": result, "total_ms": int((time.perf_counter() - total_started) * 1000), } payload.update(phase_ms) payload.update(extra) log_event("go_service_timing", **payload) log_event("session_open_requested", user_id=user.id, service_slug=slug, sw=sw, sh=sh) service = db.scalar(select(Service).where(Service.slug == slug, Service.active == True)) if not service: raise HTTPException(status_code=404, detail="Service not found") if service.type == ServiceType.VNC: raise HTTPException(status_code=410, detail="VNC services are deprecated") if not has_access(db, user.id, service.id): raise HTTPException(status_code=403, detail="ACL denied") client_width, client_height = sanitize_client_resolution(sw, sh) log_event( "session_open_resolution", user_id=user.id, service_slug=slug, sw=sw, sh=sh, client_width=client_width, client_height=client_height, ) user_lock_started = time.perf_counter() try: with allocator_lock(db, 92000 + int(user.id), timeout_seconds=GO_USER_LOCK_TIMEOUT_SECONDS): _mark("wait_user_lock_ms", user_lock_started) t_existing = time.perf_counter() existing_user_session = find_active_session_for_user_service(db, user.id, service.id) _mark("check_existing_ms", t_existing) if existing_user_session: _emit("reuse_session", session_id=existing_user_session.id) if existing_user_session.container_id and existing_user_session.container_id.startswith("RDPSLOT:"): try: _rdp_slot_id = int(existing_user_session.container_id.split(":", 1)[1]) threading.Thread(target=connect_rdp_slot, args=(_rdp_slot_id,), daemon=True).start() except Exception: pass return RedirectResponse(url=session_redirect_url(existing_user_session), status_code=303) t_limit = time.perf_counter() cutoff = now_utc() - dt.timedelta(seconds=SESSION_IDLE_SECONDS) active_rows = db.scalars( select(SessionModel).where( SessionModel.user_id == user.id, SessionModel.status == SessionStatus.ACTIVE, SessionModel.last_access_at >= cutoff, ) ).all() active_rows = sorted(active_rows, key=lambda row: row.created_at) active_service_ids = {row.service_id for row in active_rows} _mark("check_limit_ms", t_limit) if service.id not in active_service_ids and len(active_service_ids) >= MAX_ACTIVE_SERVICES_PER_USER: oldest = next((row for row in active_rows if row.service_id != service.id), None) if oldest: t_rotate = time.perf_counter() terminate_session_record(db, oldest, SessionStatus.ROTATED, stop_container=True) db.commit() _mark("rotate_oldest_ms", t_rotate) log_event( "session_rotated", user_id=user.id, closed_session_id=oldest.id, closed_service_id=oldest.service_id, new_service_id=service.id, ) else: _emit("max_services_redirect") return RedirectResponse(url="/?launch_error=max_services", status_code=303) if service.type == ServiceType.RDP: t_rdp_slots = time.perf_counter() slots = db.scalars(select(RdpSlot).where(RdpSlot.service_id == service.id)).all() _mark("check_rdp_slots_ms", t_rdp_slots) if slots: session_id = str(uuid.uuid4()) try: with allocator_lock(db, 91003, timeout_seconds=GO_POOL_LOCK_TIMEOUT_SECONDS): busy_slot_ids: set[int] = set() for row in db.scalars( select(SessionModel).where( SessionModel.status == SessionStatus.ACTIVE, SessionModel.service_id == service.id, SessionModel.container_id.like("RDPSLOT:%"), ) ).all(): try: busy_slot_ids.add(int(row.container_id.split(":", 1)[1])) except Exception: pass if service.is_pilot: # Pilot slots are reserved per person (admin # assigns the container in advance), never # picked from a shared pool - the earlier # existing_user_session check above already # resumes an active session on this slot, so # reaching here with it "busy" would mean two # concurrent launches; treat that the same as # "no slot available" rather than silently # handing the user a different pilot's machine. free_slot = next( (s for s in slots if s.assigned_user_id == user.id and s.id not in busy_slot_ids), None, ) if not free_slot: _emit("pilot_slot_not_assigned") raise HTTPException( status_code=403, detail="Вам не назначен слот для этого пилота. Обратитесь к администратору.", ) else: free_slot = next((s for s in slots if s.id not in busy_slot_ids), None) if not free_slot: _emit("rdp_all_slots_busy") raise HTTPException( status_code=503, detail="Все слоты этого RDP сервиса заняты. Попробуйте позже.", ) session_obj = SessionModel( id=session_id, user_id=user.id, service_id=service.id, container_id=f"RDPSLOT:{free_slot.id}", status=SessionStatus.ACTIVE, created_at=now_utc(), last_access_at=now_utc(), ) db.add(session_obj) db.commit() except LockTimeoutError: _emit("rdp_slot_lock_timeout") raise HTTPException(status_code=503, detail="Пул RDP занят. Повторите через несколько секунд.") log_event("session_created", user_id=user.id, service_slug=service.slug, session_id=session_id, mode="rdp_slot", slot_id=free_slot.id) audit(db, "SESSION_CREATE_RDP_SLOT", f"service={service.slug} session={session_id} slot={free_slot.id}", user_id=user.id) _emit("session_created_rdp_slot", session_id=session_id, slot_id=free_slot.id) threading.Thread(target=connect_rdp_slot, args=(free_slot.id,), daemon=True).start() return RedirectResponse(url=f"/s/{session_id}/", status_code=303) else: # Legacy: no slots configured — exclusive single-session behaviour active_owner = find_active_session_for_service(db, service.id) if active_owner: if active_owner.user_id != user.id: _emit("rdp_busy_legacy") raise HTTPException(status_code=503, detail="RDP сервис занят. Попробуйте позже.") _emit("reuse_rdp_session", session_id=active_owner.id) return RedirectResponse(url=session_redirect_url(active_owner), status_code=303) session_id = str(uuid.uuid4()) if service.type == ServiceType.WEB and WEB_POOL_SIZE > 0: try: t_pool_lock = time.perf_counter() with allocator_lock(db, 91001, timeout_seconds=GO_POOL_LOCK_TIMEOUT_SECONDS): _mark("wait_web_pool_lock_ms", t_pool_lock) t_ensure = time.perf_counter() ensure_web_pool() _mark("ensure_web_pool_ms", t_ensure) t_acquire = time.perf_counter() slot = acquire_web_pool_slot(db) _mark("acquire_web_slot_ms", t_acquire) slot_cid = f"WEBPOOLIDX:{slot}" t_dispatch = time.perf_counter() terminate_active_slot_sessions(db, slot_cid) dispatch_web_pool_target(slot, service, width=client_width, height=client_height) _mark("dispatch_web_target_ms", t_dispatch) t_commit = time.perf_counter() session_obj = SessionModel( id=session_id, user_id=user.id, service_id=service.id, container_id=slot_cid, status=SessionStatus.ACTIVE, created_at=now_utc(), last_access_at=now_utc(), ) db.add(session_obj) db.commit() _mark("db_commit_ms", t_commit) except LockTimeoutError: _emit("web_pool_lock_timeout") raise HTTPException(status_code=503, detail="Пул WEB занят. Повторите через несколько секунд.") except Exception as exc: logger.exception("web_pool_dispatch_failed slug=%s user_id=%s", slug, user.id) log_event("session_create_failed", level=logging.ERROR, user_id=user.id, service_slug=slug, mode="web_pool", error=str(exc)) audit(db, "SESSION_CREATE_FAILED", f"slug={slug} err={str(exc)}", user_id=user.id) _emit("web_pool_create_failed", error=str(exc)) raise HTTPException(status_code=502, detail="WEB runtime failed to switch target") log_event("session_created", user_id=user.id, service_slug=service.slug, session_id=session_id, mode="web_pool", slot=slot) audit(db, "SESSION_CREATE_WEB_POOL", f"service={service.slug} session={session_id} slot={slot}", user_id=user.id) _emit("session_created_web_pool", session_id=session_id, slot=slot) return RedirectResponse(url=f"/s/{session_id}/", status_code=303) if service_uses_universal_pool(service): try: t_pool_lock = time.perf_counter() with allocator_lock(db, 91002, timeout_seconds=GO_POOL_LOCK_TIMEOUT_SECONDS): _mark("wait_universal_pool_lock_ms", t_pool_lock) t_ensure = time.perf_counter() ensure_universal_pool() _mark("ensure_universal_pool_ms", t_ensure) t_acquire = time.perf_counter() slot = acquire_universal_slot(db) _mark("acquire_universal_slot_ms", t_acquire) slot_cid = f"POOLIDX:{slot}" t_dispatch = time.perf_counter() terminate_active_slot_sessions(db, slot_cid) dispatch_universal_target(slot, service, width=client_width, height=client_height) _mark("dispatch_universal_target_ms", t_dispatch) t_commit = time.perf_counter() session_obj = SessionModel( id=session_id, user_id=user.id, service_id=service.id, container_id=slot_cid, status=SessionStatus.ACTIVE, created_at=now_utc(), last_access_at=now_utc(), ) db.add(session_obj) db.commit() _mark("db_commit_ms", t_commit) except LockTimeoutError: _emit("universal_pool_lock_timeout") raise HTTPException(status_code=503, detail="Пул RDP занят. Повторите через несколько секунд.") except Exception as exc: logger.exception("universal_pool_dispatch_failed slug=%s user_id=%s", slug, user.id) log_event("session_create_failed", level=logging.ERROR, user_id=user.id, service_slug=slug, mode="universal_pool", error=str(exc)) audit(db, "SESSION_CREATE_FAILED", f"slug={slug} err={str(exc)}", user_id=user.id) _emit("universal_pool_create_failed", error=str(exc)) raise HTTPException(status_code=502, detail="Universal runtime failed to switch target") log_event("session_created", user_id=user.id, service_slug=service.slug, session_id=session_id, mode="universal_pool", slot=slot) audit(db, "SESSION_CREATE_POOL", f"service={service.slug} session={session_id} slot={slot}", user_id=user.id) _emit("session_created_universal_pool", session_id=session_id, slot=slot) return RedirectResponse(url=f"/s/{session_id}/", status_code=303) if service.type == ServiceType.WEB and desired_pool_size(service) > 0: t_warm = time.perf_counter() ensure_warm_pool(service) open_warm_web_url(service, service.target) _mark("warm_pool_prepare_ms", t_warm) t_commit = time.perf_counter() session_obj = SessionModel( id=session_id, user_id=user.id, service_id=service.id, container_id=f"POOL:{service.slug}", status=SessionStatus.ACTIVE, created_at=now_utc(), last_access_at=now_utc(), ) db.add(session_obj) db.commit() _mark("db_commit_ms", t_commit) log_event("session_created", user_id=user.id, service_slug=service.slug, session_id=session_id, mode="warm_pool") audit(db, "SESSION_CREATE_POOL", f"service={service.slug} session={session_id}", user_id=user.id) _emit("session_created_warm_pool", session_id=session_id) return RedirectResponse(url=f"/s/{session_id}/", status_code=303) try: t_create = time.perf_counter() container_id = create_runtime_container(service, session_id) _mark("create_runtime_container_ms", t_create) except Exception as exc: logger.exception("session_container_create_failed slug=%s user_id=%s", slug, user.id) log_event("session_create_failed", level=logging.ERROR, user_id=user.id, service_slug=slug, mode="single_runtime", error=str(exc)) audit(db, "SESSION_CREATE_FAILED", f"slug={slug} err={str(exc)}", user_id=user.id) _emit("single_runtime_create_failed", error=str(exc)) raise HTTPException(status_code=502, detail="Session runtime failed to start") t_commit = time.perf_counter() session_obj = SessionModel( id=session_id, user_id=user.id, service_id=service.id, container_id=container_id, status=SessionStatus.ACTIVE, created_at=now_utc(), last_access_at=now_utc(), ) db.add(session_obj) db.commit() _mark("db_commit_ms", t_commit) log_event("session_created", user_id=user.id, service_slug=service.slug, session_id=session_id, mode="single_runtime", container_id=container_id) audit(db, "SESSION_CREATE", f"service={service.slug} session={session_id}", user_id=user.id) t_wait = time.perf_counter() ready = wait_for_session_route(session_id) _mark("wait_session_route_ms", t_wait) log_event("session_route_ready", session_id=session_id, ready=ready) _emit("session_created_single_runtime", session_id=session_id, ready=ready) return RedirectResponse(url=f"/s/{session_id}/", status_code=303) except LockTimeoutError: _emit("user_lock_timeout") raise HTTPException(status_code=429, detail="Слишком много параллельных запусков. Повторите через несколько секунд.") @app.get("/svc/{slug}/", response_class=HTMLResponse) def service_wait_page(slug: str, request: Request, user: User = Depends(require_user), db: Session = Depends(get_db)): service = db.scalar(select(Service).where(Service.slug == slug, Service.active == True)) if not service: raise HTTPException(status_code=404, detail="Service not found") if not has_access(db, user.id, service.id): raise HTTPException(status_code=403, detail="ACL denied") return HTMLResponse( content=""" <!doctype html> <html> <head> <meta charset='utf-8'> <title>Service Starting
    Сервис запускается
    Проверка...
    """.strip(), status_code=200, ) @app.get("/s/{session_id}/", response_class=HTMLResponse) def session_wait_page(session_id: str, request: Request, user: User = Depends(require_user), db: Session = Depends(get_db)): sess = db.get(SessionModel, session_id) if not sess or sess.user_id != user.id: raise HTTPException(status_code=404, detail="Session not found") if sess.status != SessionStatus.ACTIVE: raise HTTPException(status_code=410, detail="Session is not active") service = db.get(Service, sess.service_id) service_title = service.name if service else "Сервис" is_rdp = service and service.type == ServiceType.RDP redirect_target = session_redirect_url(sess) return HTMLResponse( content=f""" {service_title} """, status_code=200, ) @app.get("/s/{session_id}/view", response_class=HTMLResponse) def session_view_page(session_id: str, request: Request, user: User = Depends(require_user), db: Session = Depends(get_db)): sess = db.get(SessionModel, session_id) if not sess or sess.user_id != user.id: raise HTTPException(status_code=404, detail="Session not found") if sess.status != SessionStatus.ACTIVE: raise HTTPException(status_code=410, detail="Session is not active") service = db.get(Service, sess.service_id) if not service: raise HTTPException(status_code=404, detail="Service not found") iframe_src = None if sess.container_id and sess.container_id.startswith("POOL:"): iframe_src = f"/svc/{service.slug}/?sid={session_id}" elif sess.container_id and sess.container_id.startswith("WEBPOOLIDX:"): try: slot = int(sess.container_id.split(":", 1)[1]) iframe_src = f"/w/{slot}/?sid={session_id}" except Exception: iframe_src = None elif sess.container_id and sess.container_id.startswith("POOLIDX:"): try: slot = int(sess.container_id.split(":", 1)[1]) iframe_src = f"/u/{slot}/?sid={session_id}" except Exception: iframe_src = None elif sess.container_id and sess.container_id.startswith("RDPSLOT:"): try: slot = int(sess.container_id.split(":", 1)[1]) iframe_src = f"/rdp/{slot}/?sid={session_id}" except Exception: iframe_src = None if iframe_src: creds_html = "" return HTMLResponse( content=f""" {service.name} {creds_html} """.strip() ) return RedirectResponse(url=f"/s/{session_id}/", status_code=303) @app.post("/api/sessions/{session_id}/touch") def touch_session(session_id: str, user: User = Depends(require_user), db: Session = Depends(get_db)): sess = db.get(SessionModel, session_id) if not sess or sess.user_id != user.id: raise HTTPException(status_code=404, detail="Session not found") if sess.status != SessionStatus.ACTIVE: reason = session_closed_reason(sess, db) log_event( "session_touch_rejected", level=logging.WARNING, session_id=session_id, user_id=user.id, status=sess.status.value, reason=reason, ) return JSONResponse( status_code=410, content={ "ok": False, "reason": reason, "status": sess.status.value, }, ) sess.last_access_at = now_utc() db.commit() return {"ok": True} @app.post("/api/sessions/{session_id}/close") def close_session(session_id: str, user: User = Depends(require_user), db: Session = Depends(get_db)): sess = db.get(SessionModel, session_id) if not sess or sess.user_id != user.id: raise HTTPException(status_code=404, detail="Session not found") if sess.status != SessionStatus.ACTIVE: log_event( "session_close_already_closed", session_id=session_id, user_id=user.id, status=sess.status.value, reason=session_closed_reason(sess, db), ) return {"ok": True, "status": sess.status.value} terminate_session_record(db, sess, SessionStatus.TERMINATED, stop_container=True) db.commit() log_event("session_closed_by_user", session_id=session_id, user_id=user.id) return {"ok": True, "status": "TERMINATED"} @app.get("/api/services/{slug}/status") def service_status(slug: str, user: User = Depends(require_user), db: Session = Depends(get_db)): service = db.scalar(select(Service).where(Service.slug == slug, Service.active == True)) if not service: raise HTTPException(status_code=404, detail="Service not found") if service.type == ServiceType.VNC: raise HTTPException(status_code=410, detail="VNC services are deprecated") if not has_access(db, user.id, service.id): raise HTTPException(status_code=403, detail="ACL denied") pool = get_pool_status_for_service(service) route_ok = route_ready(f"/svc/{slug}/") ready = route_ok and (pool["running"] > 0 if desired_pool_size(service) > 0 else True) steps = [ f"ACL: OK ({user.username})", f"Пул: {pool['running']} / {pool['desired']}", f"Маршрут /svc/{slug}/: {'OK' if route_ok else 'ожидание'}", ] return { "ready": ready, "message": "Готово, открываем..." if ready else "Поднимаем контейнер и маршрут...", "steps": steps, } @app.get("/api/sessions/{session_id}/status") def session_status(session_id: str, user: User = Depends(require_user), db: Session = Depends(get_db)): sess = db.get(SessionModel, session_id) if not sess or sess.user_id != user.id: raise HTTPException(status_code=404, detail="Session not found") if sess.status != SessionStatus.ACTIVE: raise HTTPException(status_code=410, detail="Session is not active") service = db.get(Service, sess.service_id) pooled_web = bool(sess.container_id and sess.container_id.startswith("POOL:") and service and service.type == ServiceType.WEB) web_pool_idx = None universal_pool_idx = None if sess.container_id and sess.container_id.startswith("WEBPOOLIDX:"): try: web_pool_idx = int(sess.container_id.split(":", 1)[1]) except Exception: web_pool_idx = None if sess.container_id and sess.container_id.startswith("POOLIDX:"): try: universal_pool_idx = int(sess.container_id.split(":", 1)[1]) except Exception: universal_pool_idx = None pooled_rdp = bool(sess.container_id and sess.container_id.startswith("POOL:") and service and service.type == ServiceType.RDP) rdp_slot_idx = None if sess.container_id and sess.container_id.startswith("RDPSLOT:"): try: rdp_slot_idx = int(sess.container_id.split(":", 1)[1]) except Exception: rdp_slot_idx = None if pooled_web and service: route_path = f"/svc/{service.slug}/" elif pooled_rdp and service: route_path = f"/svc/{service.slug}/" elif rdp_slot_idx is not None: route_path = f"/rdp/{rdp_slot_idx}/" else: route_path = f"/s/{session_id}/" if web_pool_idx is not None: route_path = f"/w/{web_pool_idx}/" if universal_pool_idx is not None: route_path = f"/u/{universal_pool_idx}/" route_ok = route_ready(route_path) running = container_running(sess.container_id) ready = running and route_ok steps = [ f"Контейнер: {'running' if running else 'starting'}", f"Маршрут {route_path}: {'OK' if route_ok else 'ожидание'}", ] payload = { "ready": ready, "message": "Готово, открываем..." if ready else "Запуск сессии...", "steps": steps, } if pooled_web or pooled_rdp: payload["redirect_url"] = f"/s/{session_id}/view" if web_pool_idx is not None: payload["redirect_url"] = f"/s/{session_id}/view" if universal_pool_idx is not None: payload["redirect_url"] = f"/s/{session_id}/view" if rdp_slot_idx is not None: payload["redirect_url"] = f"/s/{session_id}/view" return payload @app.post("/api/admin/services") def create_service(payload: dict, request: Request, _: User = Depends(require_admin), db: Session = Depends(get_db)): validate_csrf(request) service_type = ServiceType(payload["type"]) if service_type == ServiceType.VNC: raise HTTPException(status_code=400, detail="VNC services are no longer supported") target = payload["target"] if service_type == ServiceType.WEB: target = normalize_web_target(target) elif service_type == ServiceType.RDP: parse_rdp_target(target) service = Service( name=payload["name"], slug=payload["slug"], type=service_type, target=target, comment=payload.get("comment", ""), svc_login=payload.get("svc_login", ""), svc_password=payload.get("svc_password", ""), svc_cred_hint=payload.get("svc_cred_hint", ""), active=payload.get("active", True), warm_pool_size=max(0, int(payload.get("warm_pool_size", 0))), is_pilot=bool(payload.get("is_pilot", False)), ) db.add(service) db.flush() set_service_categories(db, service.id, payload.get("category_ids", [])) db.commit() if service.type == ServiceType.WEB and WEB_POOL_SIZE <= 0: ensure_warm_pool(service) elif service_uses_universal_pool(service): ensure_universal_pool() return {"id": service.id} @app.get("/api/admin/services/{service_id}/containers/status") def service_containers_status(service_id: int, _: User = Depends(require_admin), db: Session = Depends(get_db)): service = db.get(Service, service_id) if not service: raise HTTPException(status_code=404, detail="Service not found") out = get_pool_detailed_status(service) out["active_sessions"] = get_active_sessions_count(db, service.id) return out @app.post("/api/admin/services/{service_id}/icon") async def upload_service_icon( service_id: int, request: Request, file: UploadFile = File(...), _: User = Depends(require_admin), db: Session = Depends(get_db), ): validate_csrf(request) service = db.get(Service, service_id) if not service: raise HTTPException(status_code=404, detail="Service not found") new_path = await store_service_icon(service, file) old_path = service.icon_path service.icon_path = new_path db.commit() if old_path and old_path != new_path: remove_icon_file(old_path) return {"ok": True, "icon_path": new_path} @app.delete("/api/admin/services/{service_id}/icon") def delete_service_icon(service_id: int, request: Request, _: User = Depends(require_admin), db: Session = Depends(get_db)): validate_csrf(request) service = db.get(Service, service_id) if not service: raise HTTPException(status_code=404, detail="Service not found") old_path = service.icon_path service.icon_path = "" db.commit() remove_icon_file(old_path) return {"ok": True} @app.put("/api/admin/services/{service_id}") def edit_service(service_id: int, payload: dict, request: Request, _: User = Depends(require_admin), db: Session = Depends(get_db)): validate_csrf(request) service = db.get(Service, service_id) if not service: raise HTTPException(status_code=404, detail="Service not found") for key in ["name", "slug", "target", "active", "comment", "svc_login", "svc_password", "svc_cred_hint", "is_pilot"]: if key in payload: setattr(service, key, payload[key]) if "type" in payload: service.type = ServiceType(payload["type"]) if service.type == ServiceType.VNC: raise HTTPException(status_code=400, detail="VNC services are no longer supported") if service.type == ServiceType.WEB: service.target = normalize_web_target(service.target) elif service.type == ServiceType.RDP: parse_rdp_target(service.target) if "warm_pool_size" in payload: service.warm_pool_size = max(0, int(payload["warm_pool_size"])) if "category_ids" in payload: set_service_categories(db, service.id, payload.get("category_ids", [])) db.commit() if service.type == ServiceType.WEB: if WEB_POOL_SIZE <= 0: ensure_warm_pool(service) open_warm_web_url(service, service.target) elif service_uses_universal_pool(service): ensure_universal_pool() return {"ok": True} @app.delete("/api/admin/services/{service_id}") def delete_service(service_id: int, request: Request, _: User = Depends(require_admin), db: Session = Depends(get_db)): validate_csrf(request) service = db.get(Service, service_id) if not service: raise HTTPException(status_code=404, detail="Service not found") if service.type == ServiceType.WEB and WEB_POOL_SIZE <= 0: ensure_warm_pool(service, 0) remove_icon_file(service.icon_path) db.delete(service) db.commit() return {"ok": True} @app.post("/api/admin/services/{service_id}/prewarm") def prewarm_now(service_id: int, request: Request, _: User = Depends(require_admin), db: Session = Depends(get_db)): validate_csrf(request) service = db.get(Service, service_id) if not service: raise HTTPException(status_code=404, detail="Service not found") if service.type == ServiceType.WEB: ensure_web_pool() return {"ok": True, "pool": get_web_pool_status()} if service_uses_universal_pool(service): ensure_universal_pool() return {"ok": True, "pool": get_universal_pool_status()} if service.type == ServiceType.RDP: return {"ok": True, "pool": get_pool_status_for_service(service), "message": "RDP запускается on-demand"} ensure_warm_pool(service) return {"ok": True, "pool": get_pool_status_for_service(service)} @app.post("/api/admin/services/{service_id}/rdp-slots") def create_rdp_slot(service_id: int, payload: dict, request: Request, _: User = Depends(require_admin), db: Session = Depends(get_db)): validate_csrf(request) service = db.get(Service, service_id) if not service or service.type != ServiceType.RDP: raise HTTPException(status_code=404, detail="RDP service not found") rdp_username = (payload.get("rdp_username") or "").strip() rdp_password = (payload.get("rdp_password") or "").strip() if not rdp_username: raise HTTPException(status_code=400, detail="rdp_username is required") slot = RdpSlot(service_id=service_id, rdp_username=rdp_username, rdp_password=rdp_password) db.add(slot) db.flush() try: container_name = start_rdp_slot_container(slot, service) slot.container_name = container_name except Exception as exc: logger.exception("rdp_slot_container_start_failed service_id=%s", service_id) raise HTTPException(status_code=502, detail=f"Контейнер не запустился: {exc}") db.commit() audit(db, "RDP_SLOT_CREATE", f"service={service.slug} slot={slot.id} user={rdp_username}", user_id=None) return {"ok": True, "slot_id": slot.id, "container_name": slot.container_name} @app.delete("/api/admin/rdp-slots/{slot_id}") def delete_rdp_slot(slot_id: int, request: Request, _: User = Depends(require_admin), db: Session = Depends(get_db)): validate_csrf(request) slot = db.get(RdpSlot, slot_id) if not slot: raise HTTPException(status_code=404, detail="Slot not found") container_name = slot.container_name db.delete(slot) db.commit() if container_name: threading.Thread(target=stop_rdp_slot_container, args=(container_name,), daemon=True).start() return {"ok": True} @app.put("/api/admin/rdp-slots/{slot_id}/assign") def assign_rdp_slot(slot_id: int, payload: dict, request: Request, _: User = Depends(require_admin), db: Session = Depends(get_db)): """Bind (or unbind, with user_id null) a pilot's RDP slot to one specific person. This is the single action that gives someone a pilot: it both points the slot at them AND grants/revokes their UserServiceAccess row for the pilot service, so assigning a slot here is enough to make the pilot usable for them - no separate ACL step on the Users tab. Optional "expires_at" (ISO string or null) sets a per-grant expiry shorter than the account's own. Only meaningful for slots that belong to a pilot service - a slot on a regular pooled RDP service doesn't need this, since any free slot in the pool already works for anyone with access.""" validate_csrf(request) slot = db.get(RdpSlot, slot_id) if not slot: raise HTTPException(status_code=404, detail="Slot not found") service = db.get(Service, slot.service_id) if not service or not service.is_pilot: raise HTTPException(status_code=400, detail="Слот принадлежит не пилотному сервису") prev_user_id = slot.assigned_user_id raw_user_id = payload.get("user_id") if raw_user_id in (None, ""): slot.assigned_user_id = None if prev_user_id: prev_access = db.scalar( select(UserServiceAccess).where( UserServiceAccess.user_id == prev_user_id, UserServiceAccess.service_id == service.id, ) ) if prev_access: db.delete(prev_access) db.commit() audit(db, "RDP_SLOT_UNASSIGN", f"service={service.slug} slot={slot.id}", user_id=None) return {"ok": True, "assigned_user_id": None} target_user = db.get(User, int(raw_user_id)) if not target_user: raise HTTPException(status_code=404, detail="User not found") other = db.scalar( select(RdpSlot).where( RdpSlot.service_id == service.id, RdpSlot.assigned_user_id == target_user.id, RdpSlot.id != slot.id, ) ) if other: raise HTTPException( status_code=409, detail=f"У пользователя уже есть слот №{other.id} на этом пилоте", ) raw_expires = payload.get("expires_at") expires_at = dt.datetime.fromisoformat(raw_expires) if raw_expires else None if prev_user_id and prev_user_id != target_user.id: prev_access = db.scalar( select(UserServiceAccess).where( UserServiceAccess.user_id == prev_user_id, UserServiceAccess.service_id == service.id, ) ) if prev_access: db.delete(prev_access) access = db.scalar( select(UserServiceAccess).where( UserServiceAccess.user_id == target_user.id, UserServiceAccess.service_id == service.id, ) ) is_new_grant = access is None if access: access.expires_at = expires_at else: db.add(UserServiceAccess(user_id=target_user.id, service_id=service.id, expires_at=expires_at)) slot.assigned_user_id = target_user.id db.commit() audit(db, "RDP_SLOT_ASSIGN", f"service={service.slug} slot={slot.id} user={target_user.username}", user_id=None) email_status = None if is_new_grant: # Only on a genuinely new grant - not every time the admin tweaks # the expiry date for someone who already has this pilot, which # also calls this endpoint. email_status = _send_pilot_granted_email(target_user, service, expires_at) return {"ok": True, "assigned_user_id": target_user.id, "email_status": email_status} @app.post("/api/admin/categories") def create_category(payload: dict, request: Request, _: User = Depends(require_admin), db: Session = Depends(get_db)): validate_csrf(request) name = (payload.get("name") or "").strip() slug = (payload.get("slug") or "").strip().lower().replace(" ", "-") if not name: raise HTTPException(status_code=400, detail="Category name is required") if not slug: raise HTTPException(status_code=400, detail="Category slug is required") exists = db.scalar(select(Category).where((Category.name == name) | (Category.slug == slug))) if exists: raise HTTPException(status_code=409, detail="Category already exists") category = Category(name=name, slug=slug) db.add(category) db.commit() return {"id": category.id} @app.delete("/api/admin/categories/{category_id}") def delete_category(category_id: int, request: Request, _: User = Depends(require_admin), db: Session = Depends(get_db)): validate_csrf(request) category = db.get(Category, category_id) if not category: raise HTTPException(status_code=404, detail="Category not found") db.delete(category) db.commit() return {"ok": True} @app.put("/api/admin/web-pool-size") def update_web_pool_size(payload: dict, request: Request, _: User = Depends(require_admin)): validate_csrf(request) global WEB_POOL_SIZE value = max(0, int(payload.get("size", WEB_POOL_SIZE))) WEB_POOL_SIZE = value ensure_web_pool() return {"ok": True, "size": WEB_POOL_SIZE, "pool": get_web_pool_status()} @app.post("/api/admin/users") def create_user(payload: dict, request: Request, _: User = Depends(require_admin), db: Session = Depends(get_db)): """The password is never typed by the admin - it's generated here and mailed to the user (username is their email, same convention as the self-service approval flow). If the email fails to send, the plaintext comes back in the response as a one-time fallback so the admin isn't locked out of handing it over; the frontend only surfaces it then.""" validate_csrf(request) expires_at = dt.datetime.fromisoformat(payload["expires_at"]) password = _generate_password() user = User( username=payload["username"], password_hash=hash_password(password), expires_at=expires_at, active=payload.get("active", True), is_admin=payload.get("is_admin", False), first_name=payload.get("first_name", ""), last_name=payload.get("last_name", ""), ) db.add(user) db.commit() email_status = _send_welcome_email(user, password) result = {"id": user.id, "email_status": email_status} if not email_status.startswith("Email отправлен"): result["password"] = password return result @app.put("/api/admin/users/{user_id}") def edit_user(user_id: int, payload: dict, request: Request, _: User = Depends(require_admin), db: Session = Depends(get_db)): validate_csrf(request) user = db.get(User, user_id) if not user: raise HTTPException(status_code=404, detail="User not found") for key in ["username", "active", "is_admin", "first_name", "last_name"]: if key in payload: setattr(user, key, payload[key]) if "expires_at" in payload: user.expires_at = dt.datetime.fromisoformat(payload["expires_at"]) db.commit() return {"ok": True} @app.post("/api/admin/users/{user_id}/reset-password") def reset_user_password(user_id: int, request: Request, _: User = Depends(require_admin), db: Session = Depends(get_db)): """Generates a new password and emails it - the only way a user's password changes now, mirroring create_user(). Same fallback: if the email bounces, the plaintext comes back so the admin can pass it on.""" validate_csrf(request) user = db.get(User, user_id) if not user: raise HTTPException(status_code=404, detail="User not found") password = _generate_password() user.password_hash = hash_password(password) db.commit() email_status = _send_password_reset_email(user, password) result = {"ok": True, "email_status": email_status} if not email_status.startswith("Email отправлен"): result["password"] = password return result @app.delete("/api/admin/users/{user_id}") def delete_user(user_id: int, request: Request, admin: User = Depends(require_admin), db: Session = Depends(get_db)): validate_csrf(request) user = db.get(User, user_id) if not user: raise HTTPException(status_code=404, detail="User not found") if user.id == admin.id: raise HTTPException(status_code=400, detail="Cannot delete current admin") db.delete(user) db.commit() return {"ok": True} @app.put("/api/admin/users/{user_id}/acl") def set_acl(user_id: int, payload: dict, request: Request, _: User = Depends(require_admin), db: Session = Depends(get_db)): """Replace a user's non-pilot product grants with the posted set. Pilots are deliberately out of scope here - they're granted/revoked only via the per-slot assign endpoint on the Pilots tab, which is also what points a specific container at the person. The Users tab ACL grid doesn't render pilot checkboxes at all, so service_ids posted from there never includes them; if this endpoint treated "pilot id missing from service_ids" as "revoke it", saving any other product's ACL would silently strip every pilot grant the user has.""" validate_csrf(request) user = db.get(User, user_id) if not user: raise HTTPException(status_code=404, detail="User not found") service_ids = set(payload.get("service_ids", [])) pilot_ids = set(db.scalars(select(Service.id).where(Service.is_pilot == True)).all()) existing = db.scalars(select(UserServiceAccess).where(UserServiceAccess.user_id == user_id)).all() existing_map = {x.service_id: x for x in existing} for sid in service_ids - pilot_ids: if sid not in existing_map: db.add(UserServiceAccess(user_id=user_id, service_id=sid)) removed_ids = [ sid for sid, row in existing_map.items() if sid not in service_ids and sid not in pilot_ids ] for sid in removed_ids: db.delete(existing_map[sid]) db.commit() return {"ok": True}