style: format the last 10 files that ruff format disagreed with
ruff.toml has declared `quote-style = "single"` and line-length 120 since the lint job landed, and 118 of 128 files follow it. The panel backend and the netbox configuration were written in black/prettier style instead, so a `ruff format --check` would have failed on them from the start. Bring them onto the style the repository already declares, which is what makes the check adoptable at all. Formatting only: apart from quote style the diff is multi-line expressions joined where they fit inside 120 columns. Both suites still pass afterwards (25 pytest, and ruff check is clean). Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
1 parent
fddd82704f
commit
0ae0df7473
10 files changed
+320
-334
No files matched your search
@@ -1,48 +1,48 @@
|
||||
import os
|
||||
|
||||
|
||||
def _csv(name, default=""):
|
||||
return [item.strip() for item in os.environ.get(name, default).split(",") if item.strip()]
|
||||
def _csv(name, default=''):
|
||||
return [item.strip() for item in os.environ.get(name, default).split(',') if item.strip()]
|
||||
|
||||
|
||||
ALLOWED_HOSTS = _csv("ALLOWED_HOSTS", "localhost,127.0.0.1,[::1]")
|
||||
CSRF_TRUSTED_ORIGINS = _csv("CSRF_TRUSTED_ORIGINS")
|
||||
ALLOWED_HOSTS = _csv('ALLOWED_HOSTS', 'localhost,127.0.0.1,[::1]')
|
||||
CSRF_TRUSTED_ORIGINS = _csv('CSRF_TRUSTED_ORIGINS')
|
||||
USE_X_FORWARDED_HOST = True
|
||||
SECURE_PROXY_SSL_HEADER = ("HTTP_X_FORWARDED_PROTO", "https")
|
||||
SECURE_PROXY_SSL_HEADER = ('HTTP_X_FORWARDED_PROTO', 'https')
|
||||
|
||||
DATABASES = {
|
||||
"default": {
|
||||
"NAME": os.environ["DB_NAME"],
|
||||
"USER": os.environ["DB_USER"],
|
||||
"PASSWORD": os.environ["DB_PASSWORD"],
|
||||
"HOST": os.environ["DB_HOST"],
|
||||
"PORT": os.environ.get("DB_PORT", "5432"),
|
||||
"OPTIONS": {"sslmode": os.environ.get("DB_SSLMODE", "disable")},
|
||||
"CONN_MAX_AGE": int(os.environ.get("DB_CONN_MAX_AGE", "300")),
|
||||
'default': {
|
||||
'NAME': os.environ['DB_NAME'],
|
||||
'USER': os.environ['DB_USER'],
|
||||
'PASSWORD': os.environ['DB_PASSWORD'],
|
||||
'HOST': os.environ['DB_HOST'],
|
||||
'PORT': os.environ.get('DB_PORT', '5432'),
|
||||
'OPTIONS': {'sslmode': os.environ.get('DB_SSLMODE', 'disable')},
|
||||
'CONN_MAX_AGE': int(os.environ.get('DB_CONN_MAX_AGE', '300')),
|
||||
}
|
||||
}
|
||||
|
||||
REDIS = {
|
||||
"tasks": {
|
||||
"HOST": os.environ["REDIS_HOST"],
|
||||
"PORT": int(os.environ.get("REDIS_PORT", "6379")),
|
||||
"PASSWORD": os.environ["REDIS_PASSWORD"],
|
||||
"DATABASE": int(os.environ.get("REDIS_DATABASE", "0")),
|
||||
"SSL": False,
|
||||
'tasks': {
|
||||
'HOST': os.environ['REDIS_HOST'],
|
||||
'PORT': int(os.environ.get('REDIS_PORT', '6379')),
|
||||
'PASSWORD': os.environ['REDIS_PASSWORD'],
|
||||
'DATABASE': int(os.environ.get('REDIS_DATABASE', '0')),
|
||||
'SSL': False,
|
||||
},
|
||||
"caching": {
|
||||
"HOST": os.environ["REDIS_CACHE_HOST"],
|
||||
"PORT": int(os.environ.get("REDIS_CACHE_PORT", "6379")),
|
||||
"PASSWORD": os.environ["REDIS_CACHE_PASSWORD"],
|
||||
"DATABASE": int(os.environ.get("REDIS_CACHE_DATABASE", "1")),
|
||||
"SSL": False,
|
||||
'caching': {
|
||||
'HOST': os.environ['REDIS_CACHE_HOST'],
|
||||
'PORT': int(os.environ.get('REDIS_CACHE_PORT', '6379')),
|
||||
'PASSWORD': os.environ['REDIS_CACHE_PASSWORD'],
|
||||
'DATABASE': int(os.environ.get('REDIS_CACHE_DATABASE', '1')),
|
||||
'SSL': False,
|
||||
},
|
||||
}
|
||||
|
||||
SECRET_KEY = os.environ["SECRET_KEY"]
|
||||
API_TOKEN_PEPPERS = {1: os.environ["API_TOKEN_PEPPER_1"]}
|
||||
TIME_ZONE = os.environ.get("TIME_ZONE", "UTC")
|
||||
MEDIA_ROOT = "/opt/netbox/netbox/media"
|
||||
REPORTS_ROOT = "/opt/netbox/netbox/reports"
|
||||
SCRIPTS_ROOT = "/opt/netbox/netbox/scripts"
|
||||
SECRET_KEY = os.environ['SECRET_KEY']
|
||||
API_TOKEN_PEPPERS = {1: os.environ['API_TOKEN_PEPPER_1']}
|
||||
TIME_ZONE = os.environ.get('TIME_ZONE', 'UTC')
|
||||
MEDIA_ROOT = '/opt/netbox/netbox/media'
|
||||
REPORTS_ROOT = '/opt/netbox/netbox/reports'
|
||||
SCRIPTS_ROOT = '/opt/netbox/netbox/scripts'
|
||||
CENSUS_REPORTING_ENABLED = False
|
||||
@@ -51,11 +51,11 @@ class TelegramAuthService:
|
||||
if len(self.flows) >= self.max_flows:
|
||||
raise PanelError(
|
||||
429,
|
||||
"Too many pending authorization flows; try again later",
|
||||
'Too many pending authorization flows; try again later',
|
||||
)
|
||||
flow_id = secrets.token_urlsafe(24)
|
||||
telegram = Client(
|
||||
f"auth-{flow_id}",
|
||||
f'auth-{flow_id}',
|
||||
api_id=account.api_id,
|
||||
api_hash=account.api_hash,
|
||||
in_memory=True,
|
||||
@@ -117,7 +117,7 @@ class TelegramAuthService:
|
||||
account: StringSessionStart,
|
||||
) -> AuthorizedAccount:
|
||||
telegram = Client(
|
||||
f"validate-{secrets.token_urlsafe(12)}",
|
||||
f'validate-{secrets.token_urlsafe(12)}',
|
||||
api_id=account.api_id,
|
||||
api_hash=account.api_hash,
|
||||
session_string=account.session_string,
|
||||
@@ -148,7 +148,7 @@ class TelegramAuthService:
|
||||
async with self._lock:
|
||||
flow = self.flows.get(flow_id)
|
||||
if flow is None:
|
||||
raise PanelError(410, "Authorization flow expired; start again")
|
||||
raise PanelError(410, 'Authorization flow expired; start again')
|
||||
return flow
|
||||
|
||||
async def _finish(self, flow: AuthFlow) -> AuthorizedAccount:
|
||||
@@ -171,11 +171,7 @@ class TelegramAuthService:
|
||||
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
|
||||
]
|
||||
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),
|
||||
@@ -191,19 +187,19 @@ class TelegramAuthService:
|
||||
@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")
|
||||
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")
|
||||
return PanelError(422, 'Telegram API ID or API Hash is invalid')
|
||||
if isinstance(exc, PhoneNumberInvalid):
|
||||
return PanelError(422, "Phone number is invalid")
|
||||
return PanelError(422, 'Phone number is invalid')
|
||||
if isinstance(exc, PhoneCodeInvalid):
|
||||
return PanelError(422, "Telegram code is invalid")
|
||||
return PanelError(422, 'Telegram code is invalid')
|
||||
if isinstance(exc, PhoneCodeExpired):
|
||||
return PanelError(410, "Telegram code expired; start again")
|
||||
return PanelError(410, 'Telegram code expired; start again')
|
||||
if isinstance(exc, PasswordHashInvalid):
|
||||
return PanelError(422, "2FA password is invalid")
|
||||
return PanelError(422, '2FA password is invalid')
|
||||
if session and isinstance(exc, (Unauthorized, RPCError)):
|
||||
return PanelError(422, "StringSession is invalid or expired")
|
||||
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")
|
||||
return PanelError(422, 'Telegram rejected the authorization request')
|
||||
return PanelError(503, 'Telegram authorization is unavailable')
|
||||
@@ -7,39 +7,37 @@ from pathlib import Path
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class Settings:
|
||||
namespace: str = os.environ.get("USERBOT_NAMESPACE", "userbot")
|
||||
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()
|
||||
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",
|
||||
'USERBOT_IMAGE',
|
||||
'gcr.forust.xyz/forust/userbot:latest',
|
||||
)
|
||||
common_secret: str = os.environ.get(
|
||||
"USERBOT_COMMON_SECRET",
|
||||
"userbot-common-secrets",
|
||||
'USERBOT_COMMON_SECRET',
|
||||
'userbot-common-secrets',
|
||||
)
|
||||
common_config: str = os.environ.get(
|
||||
"USERBOT_COMMON_CONFIG",
|
||||
"userbot-common-config",
|
||||
'USERBOT_COMMON_CONFIG',
|
||||
'userbot-common-config',
|
||||
)
|
||||
storage_class: str = os.environ.get(
|
||||
"USERBOT_STORAGE_CLASS",
|
||||
"local-path-retain",
|
||||
'USERBOT_STORAGE_CLASS',
|
||||
'local-path-retain',
|
||||
)
|
||||
downloads_host_path: str = os.environ.get(
|
||||
"USERBOT_DOWNLOADS_HOST_PATH",
|
||||
"/srv/homelab/userbot/Downloads",
|
||||
'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")
|
||||
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",
|
||||
'USERBOT_DEFAULT_MEMORY_LIMIT',
|
||||
'1536Mi',
|
||||
)
|
||||
|
||||
|
||||
|
||||
@@ -14,18 +14,18 @@ from .models import AccountBase, InstanceSummary
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
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"
|
||||
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())
|
||||
return ','.join(f'{key}={value}' for key, value in (labels or {}).items())
|
||||
|
||||
|
||||
def _as_datetime(value: Any) -> datetime | None:
|
||||
@@ -33,7 +33,7 @@ def _as_datetime(value: Any) -> datetime | None:
|
||||
return None
|
||||
if isinstance(value, datetime):
|
||||
return value
|
||||
return getattr(value, "replace", lambda **_: None)(tzinfo=UTC)
|
||||
return getattr(value, 'replace', lambda **_: None)(tzinfo=UTC)
|
||||
|
||||
|
||||
class KubernetesService:
|
||||
@@ -55,7 +55,7 @@ class KubernetesService:
|
||||
except Exception as exc:
|
||||
raise PanelError(
|
||||
503,
|
||||
"No in-cluster or kubeconfig configuration is available",
|
||||
'No in-cluster or kubeconfig configuration is available',
|
||||
) from exc
|
||||
self.core = core or client.CoreV1Api()
|
||||
self.apps = apps or client.AppsV1Api()
|
||||
@@ -76,14 +76,14 @@ class KubernetesService:
|
||||
if exc.status == 404:
|
||||
raise PanelError(
|
||||
503,
|
||||
"Userbot common Secret or ConfigMap is missing in the userbot namespace",
|
||||
'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
|
||||
raise self._api_error(exc, 'Could not verify userbot prerequisites') from exc
|
||||
|
||||
def list_instances(
|
||||
self,
|
||||
query: str = "",
|
||||
status: str = "",
|
||||
query: str = '',
|
||||
status: str = '',
|
||||
) -> list[InstanceSummary]:
|
||||
instances: list[InstanceSummary] = []
|
||||
for namespace in (self.settings.namespace, *self.settings.legacy_namespaces):
|
||||
@@ -93,7 +93,7 @@ class KubernetesService:
|
||||
label_selector=MANAGED_LABEL,
|
||||
).items
|
||||
except ApiException as exc:
|
||||
raise self._api_error(exc, f"Could not list Deployments in {namespace}") from 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()
|
||||
@@ -103,7 +103,7 @@ class KubernetesService:
|
||||
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()
|
||||
or query in (item.pod or '').lower()
|
||||
]
|
||||
if status:
|
||||
instances = [item for item in instances if item.status == status]
|
||||
@@ -116,9 +116,9 @@ class KubernetesService:
|
||||
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"),
|
||||
(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:
|
||||
@@ -126,8 +126,8 @@ class KubernetesService:
|
||||
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")
|
||||
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:
|
||||
with self._provision_lock:
|
||||
@@ -139,25 +139,25 @@ class KubernetesService:
|
||||
self.settings.namespace,
|
||||
self._secret(account, session_string, names),
|
||||
)
|
||||
created.append(("secret", names["secret"]))
|
||||
created.append(('secret', names['secret']))
|
||||
self.core.create_namespaced_persistent_volume_claim(
|
||||
self.settings.namespace,
|
||||
self._pvc(account, names),
|
||||
)
|
||||
created.append(("pvc", names["pvc"]))
|
||||
created.append(('pvc', names['pvc']))
|
||||
self.apps.create_namespaced_deployment(
|
||||
self.settings.namespace,
|
||||
self._deployment(account, names),
|
||||
)
|
||||
created.append(("deployment", names["deployment"]))
|
||||
created.append(('deployment', names['deployment']))
|
||||
except ApiException as exc:
|
||||
self._rollback(created)
|
||||
if exc.status == 409:
|
||||
raise PanelError(
|
||||
409,
|
||||
f"Instance {account.instance_id} already exists",
|
||||
f'Instance {account.instance_id} already exists',
|
||||
) from exc
|
||||
raise self._api_error(exc, "Could not create userbot instance") from exc
|
||||
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:
|
||||
@@ -166,40 +166,40 @@ class KubernetesService:
|
||||
self.apps.patch_namespaced_deployment_scale(
|
||||
deployment.metadata.name,
|
||||
namespace,
|
||||
{"spec": {"replicas": replicas}},
|
||||
{'spec': {'replicas': replicas}},
|
||||
)
|
||||
except ApiException as exc:
|
||||
raise self._api_error(exc, "Could not scale userbot instance") from 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")
|
||||
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},
|
||||
'spec': {
|
||||
'template': {
|
||||
'metadata': {
|
||||
'annotations': {RESTART_ANNOTATION: timestamp},
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
)
|
||||
except ApiException as exc:
|
||||
raise self._api_error(exc, "Could not restart userbot instance") from 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")
|
||||
raise PanelError(409, 'Userbot Pod is not running')
|
||||
pod = pods[0]
|
||||
container = deployment.spec.template.spec.containers[0].name
|
||||
try:
|
||||
@@ -211,22 +211,22 @@ class KubernetesService:
|
||||
timestamps=True,
|
||||
)
|
||||
except ApiException as exc:
|
||||
raise self._api_error(exc, "Could not read userbot logs") from 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")
|
||||
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"])
|
||||
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"},
|
||||
{'propagation_policy': 'Foreground'},
|
||||
),
|
||||
(self.core.delete_namespaced_secret, (secret_name, namespace), {}),
|
||||
]
|
||||
@@ -243,10 +243,10 @@ class KubernetesService:
|
||||
delete_resource(*args, **kwargs)
|
||||
except ApiException as exc:
|
||||
if exc.status != 404:
|
||||
raise self._api_error(exc, "Could not delete userbot instance") from exc
|
||||
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}"
|
||||
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(
|
||||
@@ -254,16 +254,16 @@ class KubernetesService:
|
||||
label_selector=selector,
|
||||
).items
|
||||
except ApiException as exc:
|
||||
raise self._api_error(exc, "Could not find userbot instance") from 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")
|
||||
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"
|
||||
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
|
||||
@@ -287,21 +287,21 @@ class KubernetesService:
|
||||
updated_at = pod.status.start_time or pod.metadata.creation_timestamp
|
||||
|
||||
if desired == 0:
|
||||
status = "stopped"
|
||||
status = 'stopped'
|
||||
elif reason in {
|
||||
"CrashLoopBackOff",
|
||||
"Error",
|
||||
"ImagePullBackOff",
|
||||
"ErrImagePull",
|
||||
"CreateContainerConfigError",
|
||||
"RunContainerError",
|
||||
} or (pod is not None and pod.status.phase == "Failed"):
|
||||
status = "error"
|
||||
'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"
|
||||
status = 'running'
|
||||
else:
|
||||
status = "pending"
|
||||
reason = reason or (pod.status.phase if pod is not None else "Scheduling")
|
||||
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 {}
|
||||
@@ -321,8 +321,8 @@ class KubernetesService:
|
||||
pvc=pvc_name,
|
||||
storage=storage,
|
||||
image=container_spec.image,
|
||||
cpu_limit=limits.get("cpu"),
|
||||
memory_limit=limits.get("memory"),
|
||||
cpu_limit=limits.get('cpu'),
|
||||
memory_limit=limits.get('memory'),
|
||||
cpu_usage=cpu_usage,
|
||||
memory_usage=memory_usage,
|
||||
updated_at=_as_datetime(updated_at),
|
||||
@@ -338,7 +338,7 @@ class KubernetesService:
|
||||
label_selector=selector,
|
||||
).items
|
||||
except ApiException as exc:
|
||||
raise self._api_error(exc, "Could not list userbot Pods") from 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),
|
||||
@@ -350,17 +350,17 @@ class KubernetesService:
|
||||
return None, None
|
||||
try:
|
||||
metrics = self.custom.get_namespaced_custom_object(
|
||||
"metrics.k8s.io",
|
||||
"v1beta1",
|
||||
'metrics.k8s.io',
|
||||
'v1beta1',
|
||||
namespace,
|
||||
"pods",
|
||||
'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
|
||||
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:
|
||||
@@ -371,7 +371,7 @@ class KubernetesService:
|
||||
except ApiException:
|
||||
return None
|
||||
requests = pvc.spec.resources.requests or {}
|
||||
return requests.get("storage")
|
||||
return requests.get('storage')
|
||||
|
||||
@staticmethod
|
||||
def _deployment_pvc(deployment: Any) -> str | None:
|
||||
@@ -383,9 +383,9 @@ class KubernetesService:
|
||||
@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",
|
||||
'deployment': f'userbot-{instance_id}',
|
||||
'secret': f'userbot-{instance_id}-credentials',
|
||||
'pvc': f'userbot-{instance_id}-data',
|
||||
}
|
||||
|
||||
def _metadata(
|
||||
@@ -398,14 +398,14 @@ class KubernetesService:
|
||||
name=resource_name,
|
||||
namespace=self.settings.namespace,
|
||||
labels={
|
||||
"app.kubernetes.io/name": "userbot",
|
||||
'app.kubernetes.io/name': 'userbot',
|
||||
INSTANCE_LABEL: account.instance_id,
|
||||
MANAGED_BY_LABEL: "userbot-panel",
|
||||
MANAGED_BY_LABEL: 'userbot-panel',
|
||||
},
|
||||
annotations={
|
||||
DISPLAY_ANNOTATION: account.display_name,
|
||||
CREDENTIALS_ANNOTATION: names["secret"],
|
||||
PVC_ANNOTATION: names["pvc"],
|
||||
CREDENTIALS_ANNOTATION: names['secret'],
|
||||
PVC_ANNOTATION: names['pvc'],
|
||||
},
|
||||
)
|
||||
|
||||
@@ -416,12 +416,12 @@ class KubernetesService:
|
||||
names: dict[str, str],
|
||||
) -> client.V1Secret:
|
||||
return client.V1Secret(
|
||||
metadata=self._metadata(account, names, names["secret"]),
|
||||
type="Opaque",
|
||||
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,
|
||||
'API_ID': str(account.api_id),
|
||||
'API_HASH': account.api_hash,
|
||||
'STRINGSESSION': session_string,
|
||||
},
|
||||
)
|
||||
|
||||
@@ -431,12 +431,12 @@ class KubernetesService:
|
||||
names: dict[str, str],
|
||||
) -> client.V1PersistentVolumeClaim:
|
||||
return client.V1PersistentVolumeClaim(
|
||||
metadata=self._metadata(account, names, names["pvc"]),
|
||||
metadata=self._metadata(account, names, names['pvc']),
|
||||
spec=client.V1PersistentVolumeClaimSpec(
|
||||
access_modes=["ReadWriteOnce"],
|
||||
access_modes=['ReadWriteOnce'],
|
||||
storage_class_name=self.settings.storage_class,
|
||||
resources=client.V1VolumeResourceRequirements(
|
||||
requests={"storage": account.resources.storage},
|
||||
requests={'storage': account.resources.storage},
|
||||
),
|
||||
),
|
||||
)
|
||||
@@ -447,63 +447,55 @@ class KubernetesService:
|
||||
names: dict[str, str],
|
||||
) -> client.V1Deployment:
|
||||
pod_labels = {
|
||||
"app.kubernetes.io/name": "userbot",
|
||||
'app.kubernetes.io/name': 'userbot',
|
||||
INSTANCE_LABEL: account.instance_id,
|
||||
MANAGED_BY_LABEL: "userbot-panel",
|
||||
MANAGED_BY_LABEL: 'userbot-panel',
|
||||
}
|
||||
container = client.V1Container(
|
||||
name="userbot",
|
||||
name='userbot',
|
||||
image=self.settings.image,
|
||||
image_pull_policy="Always",
|
||||
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"])
|
||||
),
|
||||
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"},
|
||||
requests={'cpu': '80m', 'memory': '512Mi'},
|
||||
limits={
|
||||
"cpu": account.resources.cpu_limit,
|
||||
"memory": account.resources.memory_limit,
|
||||
'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"),
|
||||
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",
|
||||
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"]
|
||||
),
|
||||
name='data',
|
||||
persistent_volume_claim=client.V1PersistentVolumeClaimVolumeSource(claim_name=names['pvc']),
|
||||
),
|
||||
client.V1Volume(
|
||||
name="downloads",
|
||||
name='downloads',
|
||||
host_path=client.V1HostPathVolumeSource(
|
||||
path=self.settings.downloads_host_path,
|
||||
type="DirectoryOrCreate",
|
||||
type='DirectoryOrCreate',
|
||||
),
|
||||
),
|
||||
],
|
||||
)
|
||||
return client.V1Deployment(
|
||||
metadata=self._metadata(account, names, names["deployment"]),
|
||||
metadata=self._metadata(account, names, names['deployment']),
|
||||
spec=client.V1DeploymentSpec(
|
||||
replicas=1,
|
||||
strategy=client.V1DeploymentStrategy(type="Recreate"),
|
||||
strategy=client.V1DeploymentStrategy(type='Recreate'),
|
||||
selector=client.V1LabelSelector(match_labels=pod_labels),
|
||||
template=client.V1PodTemplateSpec(
|
||||
metadata=client.V1ObjectMeta(labels=pod_labels),
|
||||
@@ -515,9 +507,9 @@ class KubernetesService:
|
||||
def _rollback(self, created: list[tuple[str, str]]) -> None:
|
||||
for kind, name in reversed(created):
|
||||
try:
|
||||
if kind == "deployment":
|
||||
if kind == 'deployment':
|
||||
self.apps.delete_namespaced_deployment(name, self.settings.namespace)
|
||||
elif kind == "pvc":
|
||||
elif kind == 'pvc':
|
||||
self.core.delete_namespaced_persistent_volume_claim(
|
||||
name,
|
||||
self.settings.namespace,
|
||||
@@ -526,7 +518,7 @@ class KubernetesService:
|
||||
self.core.delete_namespaced_secret(name, self.settings.namespace)
|
||||
except ApiException as exc:
|
||||
logger.warning(
|
||||
"Rollback of %s %s in %s failed: %s",
|
||||
'Rollback of %s %s in %s failed: %s',
|
||||
kind,
|
||||
name,
|
||||
self.settings.namespace,
|
||||
@@ -536,9 +528,9 @@ class KubernetesService:
|
||||
@staticmethod
|
||||
def _api_error(exc: ApiException, detail: str) -> PanelError:
|
||||
if exc.status == 403:
|
||||
return PanelError(503, f"{detail}: Kubernetes RBAC denied the operation")
|
||||
return PanelError(503, f'{detail}: Kubernetes RBAC denied the operation')
|
||||
if exc.status == 409:
|
||||
return PanelError(409, f"{detail}: resource conflict")
|
||||
return PanelError(409, f'{detail}: resource conflict')
|
||||
if exc.status == 404:
|
||||
return PanelError(404, f"{detail}: resource not found")
|
||||
return PanelError(404, f'{detail}: resource not found')
|
||||
return PanelError(503, detail)
|
||||
@@ -30,21 +30,21 @@ async def lifespan(app: FastAPI):
|
||||
try:
|
||||
app.state.kubernetes.ensure_prerequisites()
|
||||
except PanelError as exc:
|
||||
print(f"WARNING: userbot prerequisites check failed at startup: {exc.detail}")
|
||||
print(f'WARNING: userbot prerequisites check failed at startup: {exc.detail}')
|
||||
yield
|
||||
await app.state.telegram.close()
|
||||
|
||||
|
||||
app = FastAPI(
|
||||
title="Userbot Kubernetes Control",
|
||||
version="1.0.0",
|
||||
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})
|
||||
return JSONResponse(status_code=exc.status_code, content={'detail': exc.detail})
|
||||
|
||||
|
||||
@app.exception_handler(RequestValidationError)
|
||||
@@ -54,13 +54,13 @@ async def validation_error_handler(
|
||||
) -> JSONResponse:
|
||||
errors = [
|
||||
{
|
||||
"loc": error.get("loc", ()),
|
||||
"msg": error.get("msg", "Invalid value"),
|
||||
"type": error.get("type", "value_error"),
|
||||
'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})
|
||||
return JSONResponse(status_code=422, content={'detail': errors})
|
||||
|
||||
|
||||
def kube(request: Request) -> KubernetesService:
|
||||
@@ -71,7 +71,7 @@ def telegram(request: Request) -> TelegramAuthService:
|
||||
return request.app.state.telegram
|
||||
|
||||
|
||||
@app.get("/api/health")
|
||||
@app.get('/api/health')
|
||||
def health(request: Request) -> dict[str, object]:
|
||||
service = kube(request)
|
||||
try:
|
||||
@@ -80,61 +80,61 @@ def health(request: Request) -> dict[str, object]:
|
||||
limit=1,
|
||||
)
|
||||
except Exception:
|
||||
return {"ok": False, "kubernetes": False}
|
||||
return {"ok": True, "kubernetes": True}
|
||||
return {'ok': False, 'kubernetes': False}
|
||||
return {'ok': True, 'kubernetes': True}
|
||||
|
||||
|
||||
@app.get("/api/instances", response_model=list[InstanceSummary])
|
||||
@app.get('/api/instances', response_model=list[InstanceSummary])
|
||||
def list_instances(
|
||||
request: Request,
|
||||
query: str = "",
|
||||
status: str = "",
|
||||
query: str = '',
|
||||
status: str = '',
|
||||
) -> list[InstanceSummary]:
|
||||
return kube(request).list_instances(query=query, status=status)
|
||||
|
||||
|
||||
@app.get("/api/instances/{instance_id}", response_model=InstanceSummary)
|
||||
@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")
|
||||
@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)}
|
||||
return {'logs': kube(request).logs(instance_id, tail)}
|
||||
|
||||
|
||||
@app.post("/api/instances/{instance_id}/start", response_model=InstanceSummary)
|
||||
@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)
|
||||
@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)
|
||||
@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)
|
||||
@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")
|
||||
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)
|
||||
@app.post('/api/auth/phone/start', response_model=AuthResult)
|
||||
async def auth_phone_start(
|
||||
payload: PhoneStart,
|
||||
request: Request,
|
||||
@@ -142,10 +142,10 @@ async def auth_phone_start(
|
||||
service = kube(request)
|
||||
service.assert_available(payload.instance_id)
|
||||
flow_id = await telegram(request).start_phone(payload)
|
||||
return AuthResult(status="code_required", flow_id=flow_id)
|
||||
return AuthResult(status='code_required', flow_id=flow_id)
|
||||
|
||||
|
||||
@app.post("/api/auth/phone/{flow_id}/code", response_model=AuthResult)
|
||||
@app.post('/api/auth/phone/{flow_id}/code', response_model=AuthResult)
|
||||
async def auth_phone_code(
|
||||
flow_id: str,
|
||||
payload: CodeSubmit,
|
||||
@@ -153,12 +153,12 @@ async def auth_phone_code(
|
||||
) -> AuthResult:
|
||||
authorized = await telegram(request).submit_code(flow_id, payload.code)
|
||||
if authorized is None:
|
||||
return AuthResult(status="password_required", flow_id=flow_id)
|
||||
return AuthResult(status='password_required', flow_id=flow_id)
|
||||
instance = kube(request).provision(authorized.account, authorized.session_string)
|
||||
return AuthResult(status="ready", instance=instance)
|
||||
return AuthResult(status='ready', instance=instance)
|
||||
|
||||
|
||||
@app.post("/api/auth/phone/{flow_id}/password", response_model=AuthResult)
|
||||
@app.post('/api/auth/phone/{flow_id}/password', response_model=AuthResult)
|
||||
async def auth_phone_password(
|
||||
flow_id: str,
|
||||
payload: PasswordSubmit,
|
||||
@@ -166,10 +166,10 @@ async def auth_phone_password(
|
||||
) -> 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)
|
||||
return AuthResult(status='ready', instance=instance)
|
||||
|
||||
|
||||
@app.post("/api/auth/string-session", response_model=AuthResult)
|
||||
@app.post('/api/auth/string-session', response_model=AuthResult)
|
||||
async def auth_string_session(
|
||||
payload: StringSessionStart,
|
||||
request: Request,
|
||||
@@ -178,28 +178,28 @@ async def auth_string_session(
|
||||
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)
|
||||
return AuthResult(status='ready', instance=instance)
|
||||
|
||||
|
||||
static_dir = settings.static_dir
|
||||
assets_dir = static_dir / "assets"
|
||||
assets_dir = static_dir / 'assets'
|
||||
if assets_dir.exists():
|
||||
app.mount("/assets", StaticFiles(directory=assets_dir), name="assets")
|
||||
app.mount('/assets', StaticFiles(directory=assets_dir), name='assets')
|
||||
|
||||
|
||||
@app.get("/", include_in_schema=False)
|
||||
@app.get('/', include_in_schema=False)
|
||||
def index() -> FileResponse:
|
||||
return FileResponse(static_dir / "index.html")
|
||||
return FileResponse(static_dir / 'index.html')
|
||||
|
||||
|
||||
@app.get("/{path:path}", include_in_schema=False)
|
||||
@app.get('/{path:path}', include_in_schema=False)
|
||||
def spa_fallback(path: str) -> FileResponse:
|
||||
root = static_dir.resolve()
|
||||
candidate = (static_dir / path).resolve()
|
||||
try:
|
||||
candidate.relative_to(root)
|
||||
except ValueError:
|
||||
return FileResponse(static_dir / "index.html")
|
||||
return FileResponse(static_dir / 'index.html')
|
||||
if candidate.is_file():
|
||||
return FileResponse(candidate)
|
||||
return FileResponse(static_dir / "index.html")
|
||||
return FileResponse(static_dir / 'index.html')
|
||||
@@ -5,13 +5,13 @@ 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])?$"
|
||||
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"
|
||||
storage: str = '1Gi'
|
||||
cpu_limit: str = '300m'
|
||||
memory_limit: str = '1536Mi'
|
||||
|
||||
|
||||
class AccountBase(BaseModel):
|
||||
@@ -21,7 +21,7 @@ class AccountBase(BaseModel):
|
||||
api_hash: str = Field(min_length=16, max_length=128)
|
||||
resources: Resources = Field(default_factory=Resources)
|
||||
|
||||
@field_validator("display_name", "api_hash")
|
||||
@field_validator('display_name', 'api_hash')
|
||||
@classmethod
|
||||
def strip_text(cls, value: str) -> str:
|
||||
return value.strip()
|
||||
@@ -30,12 +30,12 @@ class AccountBase(BaseModel):
|
||||
class PhoneStart(AccountBase):
|
||||
phone: str = Field(min_length=7, max_length=32)
|
||||
|
||||
@field_validator("phone")
|
||||
@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")
|
||||
if not value.startswith('+'):
|
||||
raise ValueError('phone must use international format')
|
||||
return value
|
||||
|
||||
|
||||
@@ -46,10 +46,10 @@ class StringSessionStart(AccountBase):
|
||||
class CodeSubmit(BaseModel):
|
||||
code: str = Field(min_length=3, max_length=12)
|
||||
|
||||
@field_validator("code")
|
||||
@field_validator('code')
|
||||
@classmethod
|
||||
def normalize_code(cls, value: str) -> str:
|
||||
return "".join(value.split())
|
||||
return ''.join(value.split())
|
||||
|
||||
|
||||
class PasswordSubmit(BaseModel):
|
||||
@@ -67,7 +67,7 @@ class InstanceSummary(BaseModel):
|
||||
namespace: str
|
||||
deployment: str
|
||||
pod: str | None = None
|
||||
status: Literal["running", "stopped", "pending", "error"]
|
||||
status: Literal['running', 'stopped', 'pending', 'error']
|
||||
reason: str | None = None
|
||||
ready: bool = False
|
||||
restarts: int = 0
|
||||
@@ -84,6 +84,6 @@ class InstanceSummary(BaseModel):
|
||||
|
||||
|
||||
class AuthResult(BaseModel):
|
||||
status: Literal["code_required", "password_required", "ready"]
|
||||
status: Literal['code_required', 'password_required', 'ready']
|
||||
flow_id: str | None = None
|
||||
instance: InstanceSummary | None = None
|
||||
@@ -23,7 +23,7 @@ class FakeClient:
|
||||
self.disconnected = True
|
||||
|
||||
async def send_code(self, _phone):
|
||||
return SimpleNamespace(phone_code_hash="hash")
|
||||
return SimpleNamespace(phone_code_hash='hash')
|
||||
|
||||
async def sign_in(self, *_args):
|
||||
if self.requires_password:
|
||||
@@ -38,49 +38,49 @@ class FakeClient:
|
||||
return SimpleNamespace(id=1)
|
||||
|
||||
async def export_session_string(self):
|
||||
return "exported-session"
|
||||
return 'exported-session'
|
||||
|
||||
|
||||
def phone_payload() -> PhoneStart:
|
||||
return PhoneStart(
|
||||
instance_id="test",
|
||||
display_name="Test",
|
||||
instance_id='test',
|
||||
display_name='Test',
|
||||
api_id=123,
|
||||
api_hash="0123456789abcdef",
|
||||
phone="+421900000000",
|
||||
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)
|
||||
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")
|
||||
result = await auth.submit_code(flow_id, '12345')
|
||||
|
||||
assert result.session_string == "exported-session"
|
||||
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
|
||||
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)
|
||||
monkeypatch.setattr('app.auth_service.Client', FakeClient)
|
||||
auth = TelegramAuthService()
|
||||
payload = StringSessionStart(
|
||||
instance_id="test",
|
||||
display_name="Test",
|
||||
instance_id='test',
|
||||
display_name='Test',
|
||||
api_id=123,
|
||||
api_hash="0123456789abcdef",
|
||||
session_string="x" * 64,
|
||||
api_hash='0123456789abcdef',
|
||||
session_string='x' * 64,
|
||||
)
|
||||
|
||||
result = await auth.validate_string_session(payload)
|
||||
|
||||
assert result.session_string == "exported-session"
|
||||
assert result.session_string == 'exported-session'
|
||||
assert FakeClient.instances[-1].disconnected is True
|
||||
|
||||
|
||||
@@ -88,7 +88,7 @@ async def test_string_session_is_validated_and_closed(monkeypatch) -> None:
|
||||
async def test_start_phone_rejects_overflow(monkeypatch) -> None:
|
||||
from app.errors import PanelError
|
||||
|
||||
monkeypatch.setattr("app.auth_service.Client", FakeClient)
|
||||
monkeypatch.setattr('app.auth_service.Client', FakeClient)
|
||||
auth = TelegramAuthService(ttl_seconds=3600, max_flows=2)
|
||||
|
||||
await auth.start_phone(phone_payload())
|
||||
@@ -98,4 +98,4 @@ async def test_start_phone_rejects_overflow(monkeypatch) -> None:
|
||||
await auth.start_phone(phone_payload())
|
||||
|
||||
assert error.value.status_code == 429
|
||||
assert "Too many pending" in error.value.detail
|
||||
assert 'Too many pending' in error.value.detail
|
||||
@@ -20,35 +20,35 @@ def service() -> KubernetesService:
|
||||
|
||||
def account() -> AccountBase:
|
||||
return AccountBase(
|
||||
instance_id="test-account",
|
||||
display_name="Test Account",
|
||||
instance_id='test-account',
|
||||
display_name='Test Account',
|
||||
api_id=12345,
|
||||
api_hash="0123456789abcdef0123456789abcdef",
|
||||
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)
|
||||
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",
|
||||
'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 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.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",
|
||||
'cpu': '80m',
|
||||
'memory': '512Mi',
|
||||
}
|
||||
|
||||
|
||||
@@ -57,10 +57,10 @@ def test_collision_is_reported_before_auth() -> None:
|
||||
kube.apps.read_namespaced_deployment.return_value = object()
|
||||
|
||||
with pytest.raises(PanelError) as error:
|
||||
kube.assert_available("test-account")
|
||||
kube.assert_available('test-account')
|
||||
|
||||
assert error.value.status_code == 409
|
||||
assert "Deployment" in error.value.detail
|
||||
assert 'Deployment' in error.value.detail
|
||||
|
||||
|
||||
def test_partial_provision_rolls_back_only_created_resources() -> None:
|
||||
@@ -70,11 +70,11 @@ def test_partial_provision_rolls_back_only_created_resources() -> None:
|
||||
kube.core.create_namespaced_persistent_volume_claim.side_effect = ApiException(status=500)
|
||||
|
||||
with pytest.raises(PanelError):
|
||||
kube.provision(account(), "SESSION")
|
||||
kube.provision(account(), 'SESSION')
|
||||
|
||||
kube.core.delete_namespaced_secret.assert_called_once_with(
|
||||
"userbot-test-account-credentials",
|
||||
"userbot",
|
||||
'userbot-test-account-credentials',
|
||||
'userbot',
|
||||
)
|
||||
kube.core.delete_namespaced_persistent_volume_claim.assert_not_called()
|
||||
kube.apps.delete_namespaced_deployment.assert_not_called()
|
||||
@@ -86,18 +86,18 @@ def test_provision_conflict_reports_existing_instance() -> None:
|
||||
kube.apps.create_namespaced_deployment.side_effect = ApiException(status=409)
|
||||
|
||||
with pytest.raises(PanelError) as error:
|
||||
kube.provision(account(), "SESSION")
|
||||
kube.provision(account(), 'SESSION')
|
||||
|
||||
assert error.value.status_code == 409
|
||||
assert "already exists" in error.value.detail
|
||||
assert 'already exists' in error.value.detail
|
||||
# Partial resources created before the 409 must be rolled back.
|
||||
kube.core.delete_namespaced_secret.assert_called_once_with(
|
||||
"userbot-test-account-credentials",
|
||||
"userbot",
|
||||
'userbot-test-account-credentials',
|
||||
'userbot',
|
||||
)
|
||||
kube.core.delete_namespaced_persistent_volume_claim.assert_called_once_with(
|
||||
"userbot-test-account-data",
|
||||
"userbot",
|
||||
'userbot-test-account-data',
|
||||
'userbot',
|
||||
)
|
||||
|
||||
|
||||
@@ -105,25 +105,25 @@ def test_delete_retains_pvc_unless_explicitly_requested() -> None:
|
||||
kube = service()
|
||||
deployment = SimpleNamespace(
|
||||
metadata=SimpleNamespace(
|
||||
name="userbot-test-account",
|
||||
name='userbot-test-account',
|
||||
annotations={
|
||||
"userbot.forust.xyz/credentials-secret": "credentials",
|
||||
"userbot.forust.xyz/pvc": "data",
|
||||
'userbot.forust.xyz/credentials-secret': 'credentials',
|
||||
'userbot.forust.xyz/pvc': 'data',
|
||||
},
|
||||
)
|
||||
)
|
||||
kube._find_deployment = Mock(return_value=("userbot", deployment))
|
||||
kube._find_deployment = Mock(return_value=('userbot', deployment))
|
||||
|
||||
kube.delete("test-account", delete_data=False)
|
||||
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_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.delete('test-account', delete_data=True)
|
||||
kube.core.delete_namespaced_persistent_volume_claim.assert_called_once_with(
|
||||
"data",
|
||||
"userbot",
|
||||
'data',
|
||||
'userbot',
|
||||
)
|
||||
|
||||
|
||||
@@ -131,19 +131,19 @@ def test_legacy_delete_is_blocked() -> None:
|
||||
kube = service()
|
||||
deployment = SimpleNamespace(
|
||||
metadata=SimpleNamespace(
|
||||
name="forust-userbot-deployment",
|
||||
annotations={"userbot.forust.xyz/legacy": "true"},
|
||||
name='forust-userbot-deployment',
|
||||
annotations={'userbot.forust.xyz/legacy': 'true'},
|
||||
)
|
||||
)
|
||||
kube._find_deployment = Mock(return_value=("default", deployment))
|
||||
kube._find_deployment = Mock(return_value=('default', deployment))
|
||||
|
||||
with pytest.raises(PanelError) as error:
|
||||
kube.delete("forust", delete_data=False)
|
||||
kube.delete('forust', delete_data=False)
|
||||
|
||||
assert error.value.status_code == 409
|
||||
|
||||
|
||||
def pod(*, ready: bool, phase: str = "Running", reason: str | None = None):
|
||||
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(
|
||||
@@ -154,25 +154,25 @@ def pod(*, ready: bool, phase: str = "Running", reason: str | None = None):
|
||||
start_time=datetime.now(UTC),
|
||||
)
|
||||
return SimpleNamespace(
|
||||
metadata=SimpleNamespace(name="pod-1", creation_timestamp=datetime.now(UTC)),
|
||||
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)
|
||||
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"},
|
||||
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"}),
|
||||
selector=SimpleNamespace(match_labels={'app': 'test'}),
|
||||
template=SimpleNamespace(spec=template_spec),
|
||||
),
|
||||
status=SimpleNamespace(available_replicas=1 if replicas else 0),
|
||||
@@ -180,13 +180,13 @@ def deployment(replicas: int = 1):
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("replicas", "pod_value", "expected"),
|
||||
('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"),
|
||||
(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:
|
||||
@@ -195,6 +195,6 @@ def test_status_classification(replicas, pod_value, expected) -> None:
|
||||
kube._pvc_storage = Mock(return_value=None)
|
||||
kube._pod_metrics = Mock(return_value=(None, None))
|
||||
|
||||
result = kube._summarize("userbot", deployment(replicas))
|
||||
result = kube._summarize('userbot', deployment(replicas))
|
||||
|
||||
assert result.status == expected
|
||||
@@ -6,34 +6,34 @@ from pydantic import ValidationError
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"instance_id",
|
||||
["Upper", "has_space", "-leading", "trailing-", "x" * 41],
|
||||
'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",
|
||||
display_name='Test',
|
||||
api_id=1,
|
||||
api_hash="0123456789abcdef",
|
||||
api_hash='0123456789abcdef',
|
||||
)
|
||||
|
||||
|
||||
def test_delete_data_defaults_to_false() -> None:
|
||||
request = DeleteRequest(confirmation="test")
|
||||
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"))
|
||||
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,
|
||||
'type': 'string_too_short',
|
||||
'loc': ('body', 'session_string'),
|
||||
'msg': 'String should have at least 32 characters',
|
||||
'input': sensitive_value,
|
||||
}
|
||||
]
|
||||
)
|
||||
|
||||
@@ -6,43 +6,43 @@ from app.main import spa_fallback
|
||||
|
||||
@pytest.fixture
|
||||
def static_dir(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> Path:
|
||||
static = tmp_path / "static"
|
||||
static = tmp_path / 'static'
|
||||
static.mkdir()
|
||||
(static / "index.html").write_text("<html>index</html>", encoding="utf-8")
|
||||
(static / "app.js").write_text("console.log('app')", encoding="utf-8")
|
||||
secret = tmp_path / "secret.txt"
|
||||
secret.write_text("TOP SECRET", encoding="utf-8")
|
||||
monkeypatch.setattr("app.main.static_dir", static)
|
||||
(static / 'index.html').write_text('<html>index</html>', encoding='utf-8')
|
||||
(static / 'app.js').write_text("console.log('app')", encoding='utf-8')
|
||||
secret = tmp_path / 'secret.txt'
|
||||
secret.write_text('TOP SECRET', encoding='utf-8')
|
||||
monkeypatch.setattr('app.main.static_dir', static)
|
||||
return static
|
||||
|
||||
|
||||
def static_index(static_dir: Path) -> Path:
|
||||
return static_dir / "index.html"
|
||||
return static_dir / 'index.html'
|
||||
|
||||
|
||||
def test_returns_existing_file(static_dir: Path) -> None:
|
||||
response = spa_fallback("app.js")
|
||||
assert response.path == static_dir / "app.js"
|
||||
response = spa_fallback('app.js')
|
||||
assert response.path == static_dir / 'app.js'
|
||||
|
||||
|
||||
def test_unknown_path_falls_back_to_index(static_dir: Path) -> None:
|
||||
response = spa_fallback("does/not/exist.js")
|
||||
response = spa_fallback('does/not/exist.js')
|
||||
assert response.path == static_index(static_dir)
|
||||
|
||||
|
||||
def test_traversal_does_not_leak_outside_static(static_dir: Path) -> None:
|
||||
response = spa_fallback("../secret.txt")
|
||||
response = spa_fallback('../secret.txt')
|
||||
assert response.path == static_index(static_dir)
|
||||
|
||||
response = spa_fallback("%2e%2e/secret.txt")
|
||||
response = spa_fallback('%2e%2e/secret.txt')
|
||||
assert response.path == static_index(static_dir)
|
||||
|
||||
|
||||
def test_symlink_outside_static_is_blocked(static_dir: Path, tmp_path: Path) -> None:
|
||||
target = tmp_path / "outside.txt"
|
||||
target.write_text("secret", encoding="utf-8")
|
||||
link = static_dir / "leak.txt"
|
||||
target = tmp_path / 'outside.txt'
|
||||
target.write_text('secret', encoding='utf-8')
|
||||
link = static_dir / 'leak.txt'
|
||||
link.symlink_to(target)
|
||||
|
||||
response = spa_fallback("leak.txt")
|
||||
response = spa_fallback('leak.txt')
|
||||
assert response.path == static_index(static_dir)
|
||||
Reference in new issue
Block a user