")
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 _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"""
{email_subject}
{email_subhead}
|
|
Здравствуйте, {escape(pending.name)}!
{access_text}
| Адрес портала |
{portal_url} |
| Логин |
{username} |
{cred_row}
| Доступ до |
{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
|
|
Здравствуйте, {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}
"""
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"""
{email_subject}
{email_subhead}
|
|
Здравствуйте, {escape(pending.name)}!
{access_text}
| Адрес портала |
{portal_url} |
| Логин |
{username} |
{cred_row}
| Доступ до |
{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
|
|
Здравствуйте, {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.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()
)
services = [svc for svc in all_services if svc.id in granted_ids]
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_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]
pilot_service_ids = {s.id for s in services if s.is_pilot}
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 = {}
acl_expires = {}
for row in acl_rows:
acl.setdefault(row.user_id, []).append(row.service_id)
if row.expires_at is not None:
acl_expires.setdefault(row.user_id, {})[row.service_id] = row.expires_at.isoformat()
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:
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
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}"
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,
})
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_service_ids": pilot_service_ids,
"services": services,
"categories": categories,
"service_category_map": service_category_map,
"acl": acl,
"acl_expires": acl_expires,
"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
Ваш запрос одобрен
|
|
Здравствуйте, {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
|
|
Здравствуйте, {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 /) 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" \n {loc}\n weekly\n {priority}\n "
for loc, priority in entries
]
xml = (
'\n'
'\n'
+ "\n".join(url_blocks) + "\n\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🖥 Интересующие продукты:\n{items}"
divider = "━━━━━━━━━━━━━━━━━━━━━━"
manager_text = f"\n🤝 Менеджер MONT: {_e(manager)}" if manager else ""
text = (
f"🔔 Новый запрос доступа к полигону MONT\n"
f"{divider}\n\n"
f"👤 Имя: {_e(name)}\n"
f"🏢 Компания: {_e(company)}\n"
f"📧 Email: {_e(email)}\n"
f"📱 Телефон: {_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📍 Местоположение: " + _e(geo)
geo_text += "\n🖥 IP: " + _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💬 Комментарий: {_e(note)}" if note else ""
divider = "━━━━━━━━━━━━━━━━━━━━━━"
text = (
f"🔔 Запрос дополнительного доступа\n"
f"{divider}\n\n"
f"👤 Пользователь: {_e(display_name)} ({_e(user.username)})\n"
f"🖥 Запрошенные продукты:\n{items}"
f"{note_text}"
)
ip = _get_real_ip(request)
geo = _get_geo(ip)
geo_text = ""
if geo:
geo_text += "\n📍 Местоположение: " + _e(geo)
geo_text += "\n🖥 IP: " + _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="""
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. 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="Слот принадлежит не пилотному сервису")
raw_user_id = payload.get("user_id")
if raw_user_id in (None, ""):
slot.assigned_user_id = None
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} на этом пилоте",
)
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)
return {"ok": True, "assigned_user_id": target_user.id}
@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)):
validate_csrf(request)
expires_at = dt.datetime.fromisoformat(payload["expires_at"])
user = User(
username=payload["username"],
password_hash=hash_password(payload["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()
return {"id": user.id}
@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 "password" in payload and payload["password"]:
user.password_hash = hash_password(payload["password"])
if "expires_at" in payload:
user.expires_at = dt.datetime.fromisoformat(payload["expires_at"])
db.commit()
return {"ok": True}
@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)):
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", []))
# Optional per-service expiry override, e.g. {"12": "2026-11-01T00:00:00+00:00"}
# or {"12": null} to clear it back to "follow the account expiry".
# Used for pilots, whose access window can be shorter than the account's.
expires_by_service = payload.get("expires_at_by_service") or {}
existing = db.scalars(select(UserServiceAccess).where(UserServiceAccess.user_id == user_id)).all()
existing_map = {x.service_id: x for x in existing}
rows_by_service = dict(existing_map)
for sid in service_ids:
if sid not in existing_map:
new_row = UserServiceAccess(user_id=user_id, service_id=sid)
db.add(new_row)
rows_by_service[sid] = new_row
removed_ids = [sid for sid, row in existing_map.items() if sid not in service_ids]
for sid in removed_ids:
db.delete(existing_map[sid])
for sid_str, iso_value in expires_by_service.items():
try:
sid = int(sid_str)
except (TypeError, ValueError):
continue
row = rows_by_service.get(sid)
if row is None:
continue
row.expires_at = dt.datetime.fromisoformat(iso_value) if iso_value else None
if removed_ids:
# Revoking a pilot immediately frees any slot reserved for this user
# on it, instead of waiting for the next cleanup_loop sweep.
for slot in db.scalars(
select(RdpSlot).where(
RdpSlot.assigned_user_id == user_id,
RdpSlot.service_id.in_(removed_ids),
)
).all():
slot.assigned_user_id = None
db.commit()
return {"ok": True}