feat(userbot): add Kubernetes control panel

Manage Telegram instances through Kubernetes with legacy adoption for forust and anna. Build and deploy the panel image alongside the runtime.
This commit is contained in:
forust committed 2026-09-06 20:39:12 +02:00
1 parent 30e6f05584
commit 71cddd6a91
41 files changed
+8268 -21

No files matched your search

+1
View File
@@ -0,0 +1 @@
"""Userbot Kubernetes control panel backend."""
+202
View File
@@ -0,0 +1,202 @@
from __future__ import annotations
import asyncio
import secrets
from contextlib import suppress
from dataclasses import dataclass
from datetime import UTC, datetime, timedelta
from pyrogram import Client
from pyrogram.errors import (
ApiIdInvalid,
FloodWait,
PasswordHashInvalid,
PhoneCodeExpired,
PhoneCodeInvalid,
PhoneNumberInvalid,
RPCError,
SessionPasswordNeeded,
Unauthorized,
)
from .errors import PanelError
from .models import AccountBase, PhoneStart, StringSessionStart
@dataclass
class AuthFlow:
flow_id: str
account: PhoneStart
client: Client
phone_code_hash: str
expires_at: datetime
@dataclass
class AuthorizedAccount:
account: AccountBase
session_string: str
class TelegramAuthService:
def __init__(self, ttl_seconds: int = 600) -> None:
self.ttl = timedelta(seconds=ttl_seconds)
self.flows: dict[str, AuthFlow] = {}
self._lock = asyncio.Lock()
async def start_phone(self, account: PhoneStart) -> str:
await self._cleanup_expired()
flow_id = secrets.token_urlsafe(24)
telegram = Client(
f"auth-{flow_id}",
api_id=account.api_id,
api_hash=account.api_hash,
in_memory=True,
no_updates=True,
)
try:
await telegram.connect()
sent = await telegram.send_code(account.phone)
except Exception as exc:
await self._disconnect(telegram)
raise self._translate(exc) from exc
flow = AuthFlow(
flow_id=flow_id,
account=account,
client=telegram,
phone_code_hash=sent.phone_code_hash,
expires_at=datetime.now(UTC) + self.ttl,
)
async with self._lock:
self.flows[flow_id] = flow
return flow_id
async def submit_code(
self,
flow_id: str,
code: str,
) -> AuthorizedAccount | None:
flow = await self._get(flow_id)
try:
await flow.client.sign_in(
flow.account.phone,
flow.phone_code_hash,
code,
)
except SessionPasswordNeeded:
return None
except PhoneCodeExpired as exc:
await self._discard(flow)
raise self._translate(exc) from exc
except Exception as exc:
raise self._translate(exc) from exc
return await self._finish(flow)
async def submit_password(
self,
flow_id: str,
password: str,
) -> AuthorizedAccount:
flow = await self._get(flow_id)
try:
await flow.client.check_password(password)
except Exception as exc:
raise self._translate(exc) from exc
return await self._finish(flow)
async def validate_string_session(
self,
account: StringSessionStart,
) -> AuthorizedAccount:
telegram = Client(
f"validate-{secrets.token_urlsafe(12)}",
api_id=account.api_id,
api_hash=account.api_hash,
session_string=account.session_string,
in_memory=True,
no_updates=True,
)
try:
await telegram.connect()
await telegram.get_me()
session_string = await telegram.export_session_string()
except Exception as exc:
raise self._translate(exc, session=True) from exc
finally:
await self._disconnect(telegram)
return AuthorizedAccount(account=account, session_string=session_string)
async def close(self) -> None:
async with self._lock:
flows = list(self.flows.values())
self.flows.clear()
await asyncio.gather(
*(self._disconnect(flow.client) for flow in flows),
return_exceptions=True,
)
async def _get(self, flow_id: str) -> AuthFlow:
await self._cleanup_expired()
async with self._lock:
flow = self.flows.get(flow_id)
if flow is None:
raise PanelError(410, "Authorization flow expired; start again")
return flow
async def _finish(self, flow: AuthFlow) -> AuthorizedAccount:
try:
await flow.client.get_me()
session_string = await flow.client.export_session_string()
except Exception as exc:
raise self._translate(exc) from exc
finally:
async with self._lock:
self.flows.pop(flow.flow_id, None)
await self._disconnect(flow.client)
return AuthorizedAccount(account=flow.account, session_string=session_string)
async def _discard(self, flow: AuthFlow) -> None:
async with self._lock:
self.flows.pop(flow.flow_id, None)
await self._disconnect(flow.client)
async def _cleanup_expired(self) -> None:
now = datetime.now(UTC)
async with self._lock:
expired = [
self.flows.pop(flow_id)
for flow_id, flow in list(self.flows.items())
if flow.expires_at <= now
]
if expired:
await asyncio.gather(
*(self._disconnect(flow.client) for flow in expired),
return_exceptions=True,
)
@staticmethod
async def _disconnect(telegram: Client) -> None:
with suppress(Exception):
if telegram.is_connected:
await telegram.disconnect()
@staticmethod
def _translate(exc: Exception, *, session: bool = False) -> PanelError:
if isinstance(exc, FloodWait):
return PanelError(429, f"Telegram rate limit; retry in {exc.value} seconds")
if isinstance(exc, ApiIdInvalid):
return PanelError(422, "Telegram API ID or API Hash is invalid")
if isinstance(exc, PhoneNumberInvalid):
return PanelError(422, "Phone number is invalid")
if isinstance(exc, PhoneCodeInvalid):
return PanelError(422, "Telegram code is invalid")
if isinstance(exc, PhoneCodeExpired):
return PanelError(410, "Telegram code expired; start again")
if isinstance(exc, PasswordHashInvalid):
return PanelError(422, "2FA password is invalid")
if session and isinstance(exc, (Unauthorized, RPCError)):
return PanelError(422, "StringSession is invalid or expired")
if isinstance(exc, RPCError):
return PanelError(422, "Telegram rejected the authorization request")
return PanelError(503, "Telegram authorization is unavailable")
+46
View File
@@ -0,0 +1,46 @@
from __future__ import annotations
import os
from dataclasses import dataclass
from pathlib import Path
@dataclass(frozen=True)
class Settings:
namespace: str = os.environ.get("USERBOT_NAMESPACE", "userbot")
legacy_namespaces: tuple[str, ...] = tuple(
value.strip()
for value in os.environ.get("USERBOT_LEGACY_NAMESPACES", "default").split(",")
if value.strip()
)
image: str = os.environ.get(
"USERBOT_IMAGE",
"gcr.forust.xyz/forust/userbot:latest",
)
common_secret: str = os.environ.get(
"USERBOT_COMMON_SECRET",
"userbot-common-secrets",
)
common_config: str = os.environ.get(
"USERBOT_COMMON_CONFIG",
"userbot-common-config",
)
storage_class: str = os.environ.get(
"USERBOT_STORAGE_CLASS",
"local-path-retain",
)
downloads_host_path: str = os.environ.get(
"USERBOT_DOWNLOADS_HOST_PATH",
"/srv/homelab/userbot/Downloads",
)
static_dir: Path = Path(os.environ.get("PANEL_STATIC_DIR", "/app/static"))
auth_ttl_seconds: int = int(os.environ.get("PANEL_AUTH_TTL_SECONDS", "600"))
default_storage: str = os.environ.get("USERBOT_DEFAULT_STORAGE", "1Gi")
default_cpu_limit: str = os.environ.get("USERBOT_DEFAULT_CPU_LIMIT", "300m")
default_memory_limit: str = os.environ.get(
"USERBOT_DEFAULT_MEMORY_LIMIT",
"1536Mi",
)
settings = Settings()
+5
View File
@@ -0,0 +1,5 @@
class PanelError(Exception):
def __init__(self, status_code: int, detail: str) -> None:
super().__init__(detail)
self.status_code = status_code
self.detail = detail
@@ -0,0 +1,522 @@
from __future__ import annotations
from datetime import UTC, datetime
from typing import Any
from kubernetes import client, config
from kubernetes.client.exceptions import ApiException
from .config import Settings
from .errors import PanelError
from .models import AccountBase, InstanceSummary
MANAGED_LABEL = "app.kubernetes.io/name=userbot"
INSTANCE_LABEL = "app.kubernetes.io/instance"
MANAGED_BY_LABEL = "app.kubernetes.io/managed-by"
DISPLAY_ANNOTATION = "userbot.forust.xyz/display-name"
LEGACY_ANNOTATION = "userbot.forust.xyz/legacy"
CREDENTIALS_ANNOTATION = "userbot.forust.xyz/credentials-secret"
PVC_ANNOTATION = "userbot.forust.xyz/pvc"
RESTART_ANNOTATION = "userbot.forust.xyz/restarted-at"
def _selector(labels: dict[str, str] | None) -> str:
return ",".join(f"{key}={value}" for key, value in (labels or {}).items())
def _as_datetime(value: Any) -> datetime | None:
if value is None:
return None
if isinstance(value, datetime):
return value
return getattr(value, "replace", lambda **_: None)(tzinfo=UTC)
class KubernetesService:
def __init__(
self,
settings: Settings,
*,
core: client.CoreV1Api | None = None,
apps: client.AppsV1Api | None = None,
custom: client.CustomObjectsApi | None = None,
) -> None:
self.settings = settings
if core is None or apps is None:
try:
config.load_incluster_config()
except config.ConfigException:
config.load_kube_config()
self.core = core or client.CoreV1Api()
self.apps = apps or client.AppsV1Api()
self.custom = custom or client.CustomObjectsApi()
def ensure_prerequisites(self) -> None:
try:
self.core.read_namespaced_secret(
self.settings.common_secret,
self.settings.namespace,
)
self.core.read_namespaced_config_map(
self.settings.common_config,
self.settings.namespace,
)
except ApiException as exc:
if exc.status == 404:
raise PanelError(
503,
"Userbot common Secret or ConfigMap is missing in the userbot namespace",
) from exc
raise self._api_error(exc, "Could not verify userbot prerequisites") from exc
def list_instances(
self,
query: str = "",
status: str = "",
) -> list[InstanceSummary]:
instances: list[InstanceSummary] = []
for namespace in (self.settings.namespace, *self.settings.legacy_namespaces):
try:
deployments = self.apps.list_namespaced_deployment(
namespace,
label_selector=MANAGED_LABEL,
).items
except ApiException as exc:
raise self._api_error(exc, f"Could not list Deployments in {namespace}") from exc
instances.extend(self._summarize(namespace, deployment) for deployment in deployments)
query = query.strip().lower()
if query:
instances = [
item
for item in instances
if query in item.instance_id.lower()
or query in item.display_name.lower()
or query in (item.pod or "").lower()
]
if status:
instances = [item for item in instances if item.status == status]
return sorted(instances, key=lambda item: (item.legacy, item.instance_id))
def get_instance(self, instance_id: str) -> InstanceSummary:
namespace, deployment = self._find_deployment(instance_id)
return self._summarize(namespace, deployment)
def assert_available(self, instance_id: str) -> None:
names = self._resource_names(instance_id)
checks = (
(self.apps.read_namespaced_deployment, names["deployment"], "Deployment"),
(self.core.read_namespaced_secret, names["secret"], "Secret"),
(self.core.read_namespaced_persistent_volume_claim, names["pvc"], "PVC"),
)
for read, name, kind in checks:
try:
read(name, self.settings.namespace)
except ApiException as exc:
if exc.status == 404:
continue
raise self._api_error(exc, f"Could not check {kind} {name}") from exc
raise PanelError(409, f"{kind} {name} already exists")
def provision(self, account: AccountBase, session_string: str) -> InstanceSummary:
self.ensure_prerequisites()
self.assert_available(account.instance_id)
names = self._resource_names(account.instance_id)
created: list[tuple[str, str]] = []
try:
self.core.create_namespaced_secret(
self.settings.namespace,
self._secret(account, session_string, names),
)
created.append(("secret", names["secret"]))
self.core.create_namespaced_persistent_volume_claim(
self.settings.namespace,
self._pvc(account, names),
)
created.append(("pvc", names["pvc"]))
self.apps.create_namespaced_deployment(
self.settings.namespace,
self._deployment(account, names),
)
created.append(("deployment", names["deployment"]))
except ApiException as exc:
self._rollback(created)
raise self._api_error(exc, "Could not create userbot instance") from exc
return self.get_instance(account.instance_id)
def scale(self, instance_id: str, replicas: int) -> InstanceSummary:
namespace, deployment = self._find_deployment(instance_id)
try:
self.apps.patch_namespaced_deployment_scale(
deployment.metadata.name,
namespace,
{"spec": {"replicas": replicas}},
)
except ApiException as exc:
raise self._api_error(exc, "Could not scale userbot instance") from exc
return self.get_instance(instance_id)
def restart(self, instance_id: str) -> InstanceSummary:
namespace, deployment = self._find_deployment(instance_id)
if (deployment.spec.replicas or 0) == 0:
raise PanelError(409, "Stopped instance cannot be restarted")
timestamp = datetime.now(UTC).isoformat()
try:
self.apps.patch_namespaced_deployment(
deployment.metadata.name,
namespace,
{
"spec": {
"template": {
"metadata": {
"annotations": {RESTART_ANNOTATION: timestamp},
}
}
}
},
)
except ApiException as exc:
raise self._api_error(exc, "Could not restart userbot instance") from exc
return self.get_instance(instance_id)
def logs(self, instance_id: str, tail: int = 250) -> str:
namespace, deployment = self._find_deployment(instance_id)
pods = self._pods_for_deployment(namespace, deployment)
if not pods:
raise PanelError(409, "Userbot Pod is not running")
pod = pods[0]
container = deployment.spec.template.spec.containers[0].name
try:
return self.core.read_namespaced_pod_log(
pod.metadata.name,
namespace,
container=container,
tail_lines=max(1, min(tail, 1000)),
timestamps=True,
)
except ApiException as exc:
raise self._api_error(exc, "Could not read userbot logs") from exc
def delete(self, instance_id: str, *, delete_data: bool) -> None:
namespace, deployment = self._find_deployment(instance_id)
annotations = deployment.metadata.annotations or {}
if annotations.get(LEGACY_ANNOTATION) == "true" or namespace != self.settings.namespace:
raise PanelError(409, "Legacy instances cannot be deleted from the panel")
names = self._resource_names(instance_id)
secret_name = annotations.get(CREDENTIALS_ANNOTATION, names["secret"])
pvc_name = annotations.get(PVC_ANNOTATION, names["pvc"])
operations = [
(
self.apps.delete_namespaced_deployment,
(deployment.metadata.name, namespace),
{"propagation_policy": "Foreground"},
),
(self.core.delete_namespaced_secret, (secret_name, namespace), {}),
]
if delete_data:
operations.append(
(
self.core.delete_namespaced_persistent_volume_claim,
(pvc_name, namespace),
{},
)
)
for delete_resource, args, kwargs in operations:
try:
delete_resource(*args, **kwargs)
except ApiException as exc:
if exc.status != 404:
raise self._api_error(exc, "Could not delete userbot instance") from exc
def _find_deployment(self, instance_id: str) -> tuple[str, Any]:
selector = f"{MANAGED_LABEL},{INSTANCE_LABEL}={instance_id}"
for namespace in (self.settings.namespace, *self.settings.legacy_namespaces):
try:
items = self.apps.list_namespaced_deployment(
namespace,
label_selector=selector,
).items
except ApiException as exc:
raise self._api_error(exc, "Could not find userbot instance") from exc
if items:
return namespace, items[0]
raise PanelError(404, f"Instance {instance_id} does not exist")
def _summarize(self, namespace: str, deployment: Any) -> InstanceSummary:
labels = deployment.metadata.labels or {}
annotations = deployment.metadata.annotations or {}
instance_id = labels.get(INSTANCE_LABEL, deployment.metadata.name)
legacy = annotations.get(LEGACY_ANNOTATION) == "true"
pods = self._pods_for_deployment(namespace, deployment)
pod = pods[0] if pods else None
desired = deployment.spec.replicas or 0
ready = False
reason: str | None = None
restarts = 0
updated_at = deployment.metadata.creation_timestamp
if pod is not None:
statuses = pod.status.container_statuses or []
ready = bool(statuses) and all(item.ready for item in statuses)
restarts = sum(item.restart_count for item in statuses)
for item in statuses:
state = item.state
if state and state.waiting and state.waiting.reason:
reason = state.waiting.reason
break
if state and state.terminated and state.terminated.reason:
reason = state.terminated.reason
break
updated_at = pod.status.start_time or pod.metadata.creation_timestamp
if desired == 0:
status = "stopped"
elif reason in {
"CrashLoopBackOff",
"Error",
"ImagePullBackOff",
"ErrImagePull",
"CreateContainerConfigError",
"RunContainerError",
} or (pod is not None and pod.status.phase == "Failed"):
status = "error"
elif ready and (deployment.status.available_replicas or 0) > 0:
status = "running"
else:
status = "pending"
reason = reason or (pod.status.phase if pod is not None else "Scheduling")
container_spec = deployment.spec.template.spec.containers[0]
limits = (container_spec.resources.limits or {}) if container_spec.resources else {}
pvc_name = annotations.get(PVC_ANNOTATION) or self._deployment_pvc(deployment)
storage = self._pvc_storage(namespace, pvc_name)
cpu_usage, memory_usage = self._pod_metrics(namespace, pod.metadata.name if pod else None)
return InstanceSummary(
instance_id=instance_id,
display_name=annotations.get(DISPLAY_ANNOTATION, instance_id),
namespace=namespace,
deployment=deployment.metadata.name,
pod=pod.metadata.name if pod else None,
status=status,
reason=reason,
ready=ready,
restarts=restarts,
pvc=pvc_name,
storage=storage,
image=container_spec.image,
cpu_limit=limits.get("cpu"),
memory_limit=limits.get("memory"),
cpu_usage=cpu_usage,
memory_usage=memory_usage,
updated_at=_as_datetime(updated_at),
legacy=legacy,
deletable=not legacy and namespace == self.settings.namespace,
)
def _pods_for_deployment(self, namespace: str, deployment: Any) -> list[Any]:
selector = _selector(deployment.spec.selector.match_labels)
try:
pods = self.core.list_namespaced_pod(
namespace,
label_selector=selector,
).items
except ApiException as exc:
raise self._api_error(exc, "Could not list userbot Pods") from exc
return sorted(
pods,
key=lambda pod: pod.metadata.creation_timestamp or datetime.min.replace(tzinfo=UTC),
reverse=True,
)
def _pod_metrics(self, namespace: str, pod_name: str | None) -> tuple[str | None, str | None]:
if not pod_name:
return None, None
try:
metrics = self.custom.get_namespaced_custom_object(
"metrics.k8s.io",
"v1beta1",
namespace,
"pods",
pod_name,
)
except (ApiException, AttributeError):
return None, None
containers = metrics.get("containers", [])
cpu = containers[0].get("usage", {}).get("cpu") if containers else None
memory = containers[0].get("usage", {}).get("memory") if containers else None
return cpu, memory
def _pvc_storage(self, namespace: str, pvc_name: str | None) -> str | None:
if not pvc_name:
return None
try:
pvc = self.core.read_namespaced_persistent_volume_claim(pvc_name, namespace)
except ApiException:
return None
requests = pvc.spec.resources.requests or {}
return requests.get("storage")
@staticmethod
def _deployment_pvc(deployment: Any) -> str | None:
for volume in deployment.spec.template.spec.volumes or []:
if volume.persistent_volume_claim:
return volume.persistent_volume_claim.claim_name
return None
@staticmethod
def _resource_names(instance_id: str) -> dict[str, str]:
return {
"deployment": f"userbot-{instance_id}",
"secret": f"userbot-{instance_id}-credentials",
"pvc": f"userbot-{instance_id}-data",
}
def _metadata(
self,
account: AccountBase,
names: dict[str, str],
resource_name: str,
) -> client.V1ObjectMeta:
return client.V1ObjectMeta(
name=resource_name,
namespace=self.settings.namespace,
labels={
"app.kubernetes.io/name": "userbot",
INSTANCE_LABEL: account.instance_id,
MANAGED_BY_LABEL: "userbot-panel",
},
annotations={
DISPLAY_ANNOTATION: account.display_name,
CREDENTIALS_ANNOTATION: names["secret"],
PVC_ANNOTATION: names["pvc"],
},
)
def _secret(
self,
account: AccountBase,
session_string: str,
names: dict[str, str],
) -> client.V1Secret:
return client.V1Secret(
metadata=self._metadata(account, names, names["secret"]),
type="Opaque",
string_data={
"API_ID": str(account.api_id),
"API_HASH": account.api_hash,
"STRINGSESSION": session_string,
},
)
def _pvc(
self,
account: AccountBase,
names: dict[str, str],
) -> client.V1PersistentVolumeClaim:
return client.V1PersistentVolumeClaim(
metadata=self._metadata(account, names, names["pvc"]),
spec=client.V1PersistentVolumeClaimSpec(
access_modes=["ReadWriteOnce"],
storage_class_name=self.settings.storage_class,
resources=client.V1VolumeResourceRequirements(
requests={"storage": account.resources.storage},
),
),
)
def _deployment(
self,
account: AccountBase,
names: dict[str, str],
) -> client.V1Deployment:
pod_labels = {
"app.kubernetes.io/name": "userbot",
INSTANCE_LABEL: account.instance_id,
MANAGED_BY_LABEL: "userbot-panel",
}
container = client.V1Container(
name="userbot",
image=self.settings.image,
image_pull_policy="Always",
env_from=[
client.V1EnvFromSource(
secret_ref=client.V1SecretEnvSource(name=self.settings.common_secret)
),
client.V1EnvFromSource(
config_map_ref=client.V1ConfigMapEnvSource(name=self.settings.common_config)
),
client.V1EnvFromSource(
secret_ref=client.V1SecretEnvSource(name=names["secret"])
),
],
resources=client.V1ResourceRequirements(
requests={"cpu": "80m", "memory": "512Mi"},
limits={
"cpu": account.resources.cpu_limit,
"memory": account.resources.memory_limit,
},
),
volume_mounts=[
client.V1VolumeMount(name="data", mount_path="/app/data"),
client.V1VolumeMount(name="downloads", mount_path="/app/downloads"),
],
)
pod_spec = client.V1PodSpec(
service_account_name="userbot-runtime",
automount_service_account_token=False,
containers=[container],
termination_grace_period_seconds=30,
volumes=[
client.V1Volume(
name="data",
persistent_volume_claim=client.V1PersistentVolumeClaimVolumeSource(
claim_name=names["pvc"]
),
),
client.V1Volume(
name="downloads",
host_path=client.V1HostPathVolumeSource(
path=self.settings.downloads_host_path,
type="DirectoryOrCreate",
),
),
],
)
return client.V1Deployment(
metadata=self._metadata(account, names, names["deployment"]),
spec=client.V1DeploymentSpec(
replicas=1,
strategy=client.V1DeploymentStrategy(type="Recreate"),
selector=client.V1LabelSelector(match_labels=pod_labels),
template=client.V1PodTemplateSpec(
metadata=client.V1ObjectMeta(labels=pod_labels),
spec=pod_spec,
),
),
)
def _rollback(self, created: list[tuple[str, str]]) -> None:
for kind, name in reversed(created):
try:
if kind == "deployment":
self.apps.delete_namespaced_deployment(name, self.settings.namespace)
elif kind == "pvc":
self.core.delete_namespaced_persistent_volume_claim(
name,
self.settings.namespace,
)
else:
self.core.delete_namespaced_secret(name, self.settings.namespace)
except ApiException:
pass
@staticmethod
def _api_error(exc: ApiException, detail: str) -> PanelError:
if exc.status == 403:
return PanelError(503, f"{detail}: Kubernetes RBAC denied the operation")
if exc.status == 409:
return PanelError(409, f"{detail}: resource conflict")
if exc.status == 404:
return PanelError(404, f"{detail}: resource not found")
return PanelError(503, detail)
+198
View File
@@ -0,0 +1,198 @@
from __future__ import annotations
from contextlib import asynccontextmanager
from typing import Annotated
from fastapi import FastAPI, Query, Request
from fastapi.exceptions import RequestValidationError
from fastapi.responses import FileResponse, JSONResponse, Response
from fastapi.staticfiles import StaticFiles
from .auth_service import TelegramAuthService
from .config import settings
from .errors import PanelError
from .kubernetes_service import KubernetesService
from .models import (
AuthResult,
CodeSubmit,
DeleteRequest,
InstanceSummary,
PasswordSubmit,
PhoneStart,
StringSessionStart,
)
@asynccontextmanager
async def lifespan(app: FastAPI):
app.state.kubernetes = KubernetesService(settings)
app.state.telegram = TelegramAuthService(settings.auth_ttl_seconds)
yield
await app.state.telegram.close()
app = FastAPI(
title="Userbot Kubernetes Control",
version="1.0.0",
lifespan=lifespan,
)
@app.exception_handler(PanelError)
async def panel_error_handler(_request: Request, exc: PanelError) -> JSONResponse:
return JSONResponse(status_code=exc.status_code, content={"detail": exc.detail})
@app.exception_handler(RequestValidationError)
async def validation_error_handler(
_request: Request,
exc: RequestValidationError,
) -> JSONResponse:
errors = [
{
"loc": error.get("loc", ()),
"msg": error.get("msg", "Invalid value"),
"type": error.get("type", "value_error"),
}
for error in exc.errors()
]
return JSONResponse(status_code=422, content={"detail": errors})
def kube(request: Request) -> KubernetesService:
return request.app.state.kubernetes
def telegram(request: Request) -> TelegramAuthService:
return request.app.state.telegram
@app.get("/api/health")
def health(request: Request) -> dict[str, object]:
service = kube(request)
try:
service.apps.list_namespaced_deployment(
settings.namespace,
limit=1,
)
except Exception:
return {"ok": False, "kubernetes": False}
return {"ok": True, "kubernetes": True}
@app.get("/api/instances", response_model=list[InstanceSummary])
def list_instances(
request: Request,
query: str = "",
status: str = "",
) -> list[InstanceSummary]:
return kube(request).list_instances(query=query, status=status)
@app.get("/api/instances/{instance_id}", response_model=InstanceSummary)
def get_instance(instance_id: str, request: Request) -> InstanceSummary:
return kube(request).get_instance(instance_id)
@app.get("/api/instances/{instance_id}/logs")
def get_logs(
instance_id: str,
request: Request,
tail: Annotated[int, Query(ge=1, le=1000)] = 250,
) -> dict[str, str]:
return {"logs": kube(request).logs(instance_id, tail)}
@app.post("/api/instances/{instance_id}/start", response_model=InstanceSummary)
def start_instance(instance_id: str, request: Request) -> InstanceSummary:
return kube(request).scale(instance_id, 1)
@app.post("/api/instances/{instance_id}/stop", response_model=InstanceSummary)
def stop_instance(instance_id: str, request: Request) -> InstanceSummary:
return kube(request).scale(instance_id, 0)
@app.post("/api/instances/{instance_id}/restart", response_model=InstanceSummary)
def restart_instance(instance_id: str, request: Request) -> InstanceSummary:
return kube(request).restart(instance_id)
@app.post("/api/instances/{instance_id}/delete", status_code=204)
def delete_instance(
instance_id: str,
payload: DeleteRequest,
request: Request,
) -> Response:
if payload.confirmation != instance_id:
raise PanelError(422, "Type the instance id exactly to confirm deletion")
kube(request).delete(instance_id, delete_data=payload.delete_data)
return Response(status_code=204)
@app.post("/api/auth/phone/start", response_model=AuthResult)
async def auth_phone_start(
payload: PhoneStart,
request: Request,
) -> AuthResult:
service = kube(request)
service.ensure_prerequisites()
service.assert_available(payload.instance_id)
flow_id = await telegram(request).start_phone(payload)
return AuthResult(status="code_required", flow_id=flow_id)
@app.post("/api/auth/phone/{flow_id}/code", response_model=AuthResult)
async def auth_phone_code(
flow_id: str,
payload: CodeSubmit,
request: Request,
) -> AuthResult:
authorized = await telegram(request).submit_code(flow_id, payload.code)
if authorized is None:
return AuthResult(status="password_required", flow_id=flow_id)
instance = kube(request).provision(authorized.account, authorized.session_string)
return AuthResult(status="ready", instance=instance)
@app.post("/api/auth/phone/{flow_id}/password", response_model=AuthResult)
async def auth_phone_password(
flow_id: str,
payload: PasswordSubmit,
request: Request,
) -> AuthResult:
authorized = await telegram(request).submit_password(flow_id, payload.password)
instance = kube(request).provision(authorized.account, authorized.session_string)
return AuthResult(status="ready", instance=instance)
@app.post("/api/auth/string-session", response_model=AuthResult)
async def auth_string_session(
payload: StringSessionStart,
request: Request,
) -> AuthResult:
service = kube(request)
service.ensure_prerequisites()
service.assert_available(payload.instance_id)
authorized = await telegram(request).validate_string_session(payload)
instance = service.provision(authorized.account, authorized.session_string)
return AuthResult(status="ready", instance=instance)
static_dir = settings.static_dir
assets_dir = static_dir / "assets"
if assets_dir.exists():
app.mount("/assets", StaticFiles(directory=assets_dir), name="assets")
@app.get("/", include_in_schema=False)
def index() -> FileResponse:
return FileResponse(static_dir / "index.html")
@app.get("/{path:path}", include_in_schema=False)
def spa_fallback(path: str) -> FileResponse:
candidate = (static_dir / path).resolve()
if candidate.is_file() and static_dir.resolve() in candidate.parents:
return FileResponse(candidate)
return FileResponse(static_dir / "index.html")
+89
View File
@@ -0,0 +1,89 @@
from __future__ import annotations
from datetime import datetime
from typing import Literal
from pydantic import BaseModel, Field, field_validator
INSTANCE_PATTERN = r"^[a-z0-9](?:[a-z0-9-]{0,38}[a-z0-9])?$"
class Resources(BaseModel):
storage: str = "1Gi"
cpu_limit: str = "300m"
memory_limit: str = "1536Mi"
class AccountBase(BaseModel):
instance_id: str = Field(pattern=INSTANCE_PATTERN, max_length=40)
display_name: str = Field(min_length=1, max_length=80)
api_id: int = Field(gt=0)
api_hash: str = Field(min_length=16, max_length=128)
resources: Resources = Field(default_factory=Resources)
@field_validator("display_name", "api_hash")
@classmethod
def strip_text(cls, value: str) -> str:
return value.strip()
class PhoneStart(AccountBase):
phone: str = Field(min_length=7, max_length=32)
@field_validator("phone")
@classmethod
def normalize_phone(cls, value: str) -> str:
value = value.strip()
if not value.startswith("+"):
raise ValueError("phone must use international format")
return value
class StringSessionStart(AccountBase):
session_string: str = Field(min_length=32)
class CodeSubmit(BaseModel):
code: str = Field(min_length=3, max_length=12)
@field_validator("code")
@classmethod
def normalize_code(cls, value: str) -> str:
return "".join(value.split())
class PasswordSubmit(BaseModel):
password: str = Field(min_length=1, max_length=256)
class DeleteRequest(BaseModel):
confirmation: str
delete_data: bool = False
class InstanceSummary(BaseModel):
instance_id: str
display_name: str
namespace: str
deployment: str
pod: str | None = None
status: Literal["running", "stopped", "pending", "error"]
reason: str | None = None
ready: bool = False
restarts: int = 0
pvc: str | None = None
storage: str | None = None
image: str
cpu_limit: str | None = None
memory_limit: str | None = None
cpu_usage: str | None = None
memory_usage: str | None = None
updated_at: datetime | None = None
legacy: bool = False
deletable: bool = True
class AuthResult(BaseModel):
status: Literal["code_required", "password_required", "ready"]
flow_id: str | None = None
instance: InstanceSummary | None = None
@@ -0,0 +1,4 @@
-r requirements.txt
httpx==0.28.1
pytest==8.4.1
pytest-asyncio==1.1.0
+5
View File
@@ -0,0 +1,5 @@
fastapi==0.115.12
kubernetes==33.1.0
pyrofork==2.3.68
tgcrypto==1.2.5
uvicorn[standard]==0.34.2
@@ -0,0 +1,84 @@
from types import SimpleNamespace
import pytest
from app.auth_service import TelegramAuthService
from app.models import PhoneStart, StringSessionStart
class FakeClient:
instances = []
requires_password = False
def __init__(self, *_args, **kwargs):
self.kwargs = kwargs
self.is_connected = False
self.disconnected = False
self.__class__.instances.append(self)
async def connect(self):
self.is_connected = True
async def disconnect(self):
self.is_connected = False
self.disconnected = True
async def send_code(self, _phone):
return SimpleNamespace(phone_code_hash="hash")
async def sign_in(self, *_args):
if self.requires_password:
from pyrogram.errors import SessionPasswordNeeded
raise SessionPasswordNeeded()
async def check_password(self, _password):
return None
async def get_me(self):
return SimpleNamespace(id=1)
async def export_session_string(self):
return "exported-session"
def phone_payload() -> PhoneStart:
return PhoneStart(
instance_id="test",
display_name="Test",
api_id=123,
api_hash="0123456789abcdef",
phone="+421900000000",
)
@pytest.mark.asyncio
async def test_phone_code_success_closes_and_forgets_client(monkeypatch) -> None:
monkeypatch.setattr("app.auth_service.Client", FakeClient)
auth = TelegramAuthService()
flow_id = await auth.start_phone(phone_payload())
result = await auth.submit_code(flow_id, "12345")
assert result.session_string == "exported-session"
assert flow_id not in auth.flows
assert FakeClient.instances[-1].disconnected is True
assert FakeClient.instances[-1].kwargs["in_memory"] is True
assert FakeClient.instances[-1].kwargs["no_updates"] is True
@pytest.mark.asyncio
async def test_string_session_is_validated_and_closed(monkeypatch) -> None:
monkeypatch.setattr("app.auth_service.Client", FakeClient)
auth = TelegramAuthService()
payload = StringSessionStart(
instance_id="test",
display_name="Test",
api_id=123,
api_hash="0123456789abcdef",
session_string="x" * 64,
)
result = await auth.validate_string_session(payload)
assert result.session_string == "exported-session"
assert FakeClient.instances[-1].disconnected is True
@@ -0,0 +1,179 @@
from datetime import UTC, datetime
from types import SimpleNamespace
from unittest.mock import Mock
import pytest
from app.config import Settings
from app.errors import PanelError
from app.kubernetes_service import KubernetesService
from app.models import AccountBase
from kubernetes.client.exceptions import ApiException
def service() -> KubernetesService:
core = Mock()
apps = Mock()
custom = Mock()
custom.get_namespaced_custom_object.side_effect = ApiException(status=404)
return KubernetesService(Settings(), core=core, apps=apps, custom=custom)
def account() -> AccountBase:
return AccountBase(
instance_id="test-account",
display_name="Test Account",
api_id=12345,
api_hash="0123456789abcdef0123456789abcdef",
)
def test_renders_managed_resources_without_leaking_credentials() -> None:
kube = service()
names = kube._resource_names("test-account")
secret = kube._secret(account(), "SESSION", names)
pvc = kube._pvc(account(), names)
deployment = kube._deployment(account(), names)
assert secret.string_data == {
"API_ID": "12345",
"API_HASH": "0123456789abcdef0123456789abcdef",
"STRINGSESSION": "SESSION",
}
assert pvc.spec.storage_class_name == "local-path-retain"
assert pvc.spec.resources.requests["storage"] == "1Gi"
assert deployment.spec.strategy.type == "Recreate"
assert deployment.spec.replicas == 1
assert deployment.spec.template.spec.automount_service_account_token is False
assert deployment.spec.template.spec.service_account_name == "userbot-runtime"
assert deployment.spec.template.spec.volumes[1].host_path.path.endswith("/Downloads")
assert deployment.spec.template.spec.containers[0].resources.requests == {
"cpu": "80m",
"memory": "512Mi",
}
def test_collision_is_reported_before_auth() -> None:
kube = service()
kube.apps.read_namespaced_deployment.return_value = object()
with pytest.raises(PanelError) as error:
kube.assert_available("test-account")
assert error.value.status_code == 409
assert "Deployment" in error.value.detail
def test_partial_provision_rolls_back_only_created_resources() -> None:
kube = service()
kube.ensure_prerequisites = Mock()
kube.assert_available = Mock()
kube.core.create_namespaced_persistent_volume_claim.side_effect = ApiException(status=500)
with pytest.raises(PanelError):
kube.provision(account(), "SESSION")
kube.core.delete_namespaced_secret.assert_called_once_with(
"userbot-test-account-credentials",
"userbot",
)
kube.core.delete_namespaced_persistent_volume_claim.assert_not_called()
kube.apps.delete_namespaced_deployment.assert_not_called()
def test_delete_retains_pvc_unless_explicitly_requested() -> None:
kube = service()
deployment = SimpleNamespace(
metadata=SimpleNamespace(
name="userbot-test-account",
annotations={
"userbot.forust.xyz/credentials-secret": "credentials",
"userbot.forust.xyz/pvc": "data",
},
)
)
kube._find_deployment = Mock(return_value=("userbot", deployment))
kube.delete("test-account", delete_data=False)
kube.apps.delete_namespaced_deployment.assert_called_once()
kube.core.delete_namespaced_secret.assert_called_once_with("credentials", "userbot")
kube.core.delete_namespaced_persistent_volume_claim.assert_not_called()
kube.delete("test-account", delete_data=True)
kube.core.delete_namespaced_persistent_volume_claim.assert_called_once_with(
"data",
"userbot",
)
def test_legacy_delete_is_blocked() -> None:
kube = service()
deployment = SimpleNamespace(
metadata=SimpleNamespace(
name="forust-userbot-deployment",
annotations={"userbot.forust.xyz/legacy": "true"},
)
)
kube._find_deployment = Mock(return_value=("default", deployment))
with pytest.raises(PanelError) as error:
kube.delete("forust", delete_data=False)
assert error.value.status_code == 409
def pod(*, ready: bool, phase: str = "Running", reason: str | None = None):
waiting = SimpleNamespace(reason=reason) if reason else None
state = SimpleNamespace(waiting=waiting, terminated=None)
status = SimpleNamespace(
container_statuses=[
SimpleNamespace(ready=ready, restart_count=2, state=state),
],
phase=phase,
start_time=datetime.now(UTC),
)
return SimpleNamespace(
metadata=SimpleNamespace(name="pod-1", creation_timestamp=datetime.now(UTC)),
status=status,
)
def deployment(replicas: int = 1):
resources = SimpleNamespace(limits={"cpu": "300m", "memory": "1536Mi"})
container = SimpleNamespace(image="userbot:latest", resources=resources)
template_spec = SimpleNamespace(containers=[container], volumes=[])
return SimpleNamespace(
metadata=SimpleNamespace(
name="userbot-test-account",
labels={"app.kubernetes.io/instance": "test-account"},
annotations={},
creation_timestamp=datetime.now(UTC),
),
spec=SimpleNamespace(
replicas=replicas,
selector=SimpleNamespace(match_labels={"app": "test"}),
template=SimpleNamespace(spec=template_spec),
),
status=SimpleNamespace(available_replicas=1 if replicas else 0),
)
@pytest.mark.parametrize(
("replicas", "pod_value", "expected"),
[
(0, None, "stopped"),
(1, pod(ready=True), "running"),
(1, pod(ready=False, phase="Pending"), "pending"),
(1, pod(ready=False, reason="CrashLoopBackOff"), "error"),
(1, pod(ready=False, phase="Failed"), "error"),
],
)
def test_status_classification(replicas, pod_value, expected) -> None:
kube = service()
kube._pods_for_deployment = Mock(return_value=[pod_value] if pod_value else [])
kube._pvc_storage = Mock(return_value=None)
kube._pod_metrics = Mock(return_value=(None, None))
result = kube._summarize("userbot", deployment(replicas))
assert result.status == expected
@@ -0,0 +1,44 @@
import pytest
from app.main import validation_error_handler
from app.models import AccountBase, DeleteRequest
from fastapi.exceptions import RequestValidationError
from pydantic import ValidationError
@pytest.mark.parametrize(
"instance_id",
["Upper", "has_space", "-leading", "trailing-", "x" * 41],
)
def test_invalid_instance_ids(instance_id: str) -> None:
with pytest.raises(ValidationError):
AccountBase(
instance_id=instance_id,
display_name="Test",
api_id=1,
api_hash="0123456789abcdef",
)
def test_delete_data_defaults_to_false() -> None:
request = DeleteRequest(confirmation="test")
assert request.delete_data is False
@pytest.mark.asyncio
async def test_validation_response_does_not_echo_secret_input() -> None:
sensitive_value = "-".join(("very", "private", "string", "session"))
error = RequestValidationError(
[
{
"type": "string_too_short",
"loc": ("body", "session_string"),
"msg": "String should have at least 32 characters",
"input": sensitive_value,
}
]
)
response = await validation_error_handler(None, error)
assert sensitive_value.encode() not in response.body
assert b'"input"' not in response.body