"""Deployment orchestrator — generic bot hosting. Deploys *any* bot to its own HF Space by cloning its source repository. The panel itself ships no bot code — every bot's runtime comes from its ``source_repo`` field. Supported source types: * ``hf`` — another HuggingFace repo (model/dataset/space). We ``snapshot_download`` it and ``upload_folder`` to the Space. * ``github`` — any git URL. We shallow-clone it and upload the tree. After the code is in place, we set ``HF Space secrets`` for whatever the bot has configured (Telegram token, NVIDIA key, etc.) and ``HF Space variables`` for ``env_vars_json``. """ from __future__ import annotations import contextlib import json import logging import os import shutil import subprocess import tempfile from pathlib import Path from typing import Any, Optional from huggingface_hub import HfApi from auth import VAULT from bot_runner import bot_workdir, start_bot, stop_bot as runner_stop, set_record_quota from bot_store import materialize_source, push_source, push_run_log from config import PANEL_CONFIG from models import ( BotInstance, BotStatus, DeploymentLog, DeploymentMode, SPACE_STATE_ON_DATASET, SPACE_STATE_ON_SPACE, SPACE_STATE_ON_SPACE_KEPT, SPACE_STATE_PENDING_MATERIALIZE, ) log = logging.getLogger("hosting-panel.deploy") class DeployError(Exception): pass class DeployOrchestrator: """Manages a bot's deployment lifecycle on HuggingFace Spaces.""" def __init__(self, hf_token: str = "") -> None: self.hf_token = hf_token self._api: Optional[HfApi] = None def _get_api(self) -> HfApi: if self._api is None: self._api = HfApi(token=self.hf_token) return self._api @contextlib.contextmanager def _use_token(self, token: str): """Temporarily swap ``self.hf_token`` for the duration of the ``with`` block, then restore it. Used by legacy-space deploys: those operate on the bot's own HF token (because the bot owns its Space), so we need the bot's token in scope without permanently overwriting the orchestrator's panel-level token. Restoring on every exit path (including exceptions) prevents cross-bot token confusion on subsequent deploys. If ``token`` is empty, ``self.hf_token`` is left untouched — we never want to clobber the panel's token with ``""``, because the next HF API call would fail unauthenticated. """ if not token: yield return saved = self.hf_token self.hf_token = token # Also reset the cached HfApi — it was bound to the old token. self._api = None try: yield finally: self.hf_token = saved self._api = None # ------------------------------------------------------------------ # Source resolution — pull the bot's repo into a local temp dir. # ------------------------------------------------------------------ def _resolve_owner(self, bot: BotInstance) -> str: """Get the HF owner/namespace from the bot's HF dataset repo or space name. Required because HF Spaces need an explicit owner/name target. """ # Prefer the dataset repo's owner — it's usually the same # username that owns the bot's Space. for ref in (bot.hf_dataset_repo, bot.hf_space_name): if ref and "/" in ref: return ref.split("/", 1)[0] # Fall back to the HF token's namespace (best-effort). try: info = self._get_api().whoami(token=self.hf_token) return info.get("name") or info.get("fullname") or "" except Exception: return "" def _resolve_space_name(self, bot: BotInstance) -> str: return bot.hf_space_name or bot.slug def _resolve_repo_id(self, bot: BotInstance) -> str: owner = self._resolve_owner(bot) if not owner: raise DeployError("Cannot determine HF owner. Set hf_dataset_repo or hf_space_name.") return f"{owner}/{self._resolve_space_name(bot)}" async def _fetch_source(self, bot: BotInstance, dest: Path) -> None: """Materialise the bot's source into ``dest``. Supports HF Hub repos and GitHub (any git URL). """ repo = bot.source_repo if not repo: raise DeployError( "Bot has no source_repo configured. Set it in the bot's edit form " "to the HuggingFace repo or GitHub URL containing the bot code." ) if bot.source_type == "github": await self._fetch_github(bot, repo, bot.source_branch, dest) else: await self._fetch_hf(repo, dest) async def _fetch_hf(self, repo_id: str, dest: Path) -> None: from huggingface_hub import snapshot_download def _do(): return snapshot_download( repo_id=repo_id, repo_type="space", # spaces have Dockerfile + app config local_dir=str(dest), token=self.hf_token, ) await self._run(_do) async def _fetch_github( self, bot: BotInstance, url: str, branch: str, dest: Path, ) -> None: """Shallow-clone a git URL into ``dest``. Token authentication (for private repos): if the URL is plain ``https://github.com/owner/repo.git`` we inject an embedded token. For other providers the user must include the token themselves in the URL. ``bot`` is required so we can read ``bot.env_vars_json`` for the optional ``GITHUB_TOKEN`` / ``GH_TOKEN`` env var. The previous version of this method referenced an undefined ``bot`` closure variable and would raise ``NameError`` on the first private-repo clone. """ clone_url = url if url.startswith("https://github.com/"): # HuggingFace tokens can't auth GitHub directly, but we # accept a separate ``GITHUB_TOKEN``-style token in # env_vars if the user wants private-repo access. gh_token = "" if bot.env_vars_json: try: extra = json.loads(bot.env_vars_json) gh_token = str( extra.get("GITHUB_TOKEN") or extra.get("GH_TOKEN") or "" ) except Exception: gh_token = "" if gh_token: clone_url = url.replace( "https://github.com/", f"https://oauth2:{gh_token}@github.com/", ) def _do(): subprocess.run( ["git", "clone", "--depth", "1", "--branch", branch, clone_url, str(dest)], check=True, capture_output=True, ) await self._run(_do) # ------------------------------------------------------------------ # Async runner — huggingface_hub + git are sync, so we offload. # ------------------------------------------------------------------ async def _run(self, fn, *args, **kwargs): import asyncio loop = asyncio.get_event_loop() def _call(): try: return fn(*args, **kwargs) except subprocess.CalledProcessError as exc: stderr = (exc.stderr or b"").decode("utf-8", "replace") raise DeployError(f"git failed: {stderr.strip() or exc}") from exc except Exception as exc: raise DeployError(str(exc)) from exc return await loop.run_in_executor(None, _call) # ------------------------------------------------------------------ # HF Space lifecycle # ------------------------------------------------------------------ async def create_hf_space( self, owner: str, space_name: str, sdk: str = "docker", hardware: str = "cpu-basic", private: bool = True, ) -> Any: def _do(): return self._get_api().create_repo( repo_id=f"{owner}/{space_name}", repo_type="space", space_sdk=sdk, space_hardware=hardware, private=private, exist_ok=True, ) return await self._run(_do) # ------------------------------------------------------------------ # Push the bot's source tree to the Space # ------------------------------------------------------------------ async def push_source(self, bot: BotInstance, source_dir: Path, message: str = "") -> None: """Upload the local ``source_dir`` to the bot's HF Space.""" owner = self._resolve_owner(bot) if not owner: raise DeployError("Cannot determine HF owner.") repo_id = f"{owner}/{self._resolve_space_name(bot)}" def _do(): self._get_api().upload_folder( folder_path=str(source_dir), repo_id=repo_id, repo_type="space", commit_message=message or f"Deploy {bot.slug}", token=self.hf_token, ) await self._run(_do) # ------------------------------------------------------------------ # Secret / env push # ------------------------------------------------------------------ async def set_secrets(self, bot: BotInstance, secrets: dict[str, str]) -> None: """Push encrypted secrets to the HF Space.""" if not secrets: return owner = self._resolve_owner(bot) repo_id = f"{owner}/{self._resolve_space_name(bot)}" def _do(): for k, v in secrets.items(): if v is None or v == "": continue try: self._get_api().add_space_secret( repo_id=repo_id, key=k, value=v, token=self.hf_token, ) except Exception as exc: log.warning("Failed to set HF secret %s: %s", k, exc) await self._run(_do) # ------------------------------------------------------------------ # End-to-end deploy # ------------------------------------------------------------------ async def deploy_bot(self, bot: BotInstance, db_session: Any) -> None: """Full deploy. Dispatches on ``bot.deployment_mode``: - ``MULTITENANT``: pull source, push to shared dataset, write env file, spawn subprocess on the bot's allocated port. - ``LEGACY_SPACE``: create dedicated HF Space (legacy path). """ if bot.deployment_mode == DeploymentMode.MULTITENANT: await self._deploy_multitenant(bot, db_session) else: await self._deploy_legacy_space(bot, db_session) # ------------------------------------------------------------------ # Multi-tenant deploy path # ------------------------------------------------------------------ async def _deploy_multitenant( self, bot: BotInstance, db_session: Any, ) -> None: """Pull source в†’ push to shared dataset в†’ write .env в†’ spawn subprocess.""" if not bot.port or not bot.space_id: raise DeployError( "Bot has no space/port allocated; cannot deploy multi-tenant." ) # 1. Pull source into the bot's local workdir. bot.status = BotStatus.STARTING # Track where the bot's source lives relative to its Space. # 'on_dataset' is the resting state after a stop; once we # start copying to the local workdir we move through # 'pending_materialize' and finally land on 'on_space'. bot.space_state = SPACE_STATE_PENDING_MATERIALIZE await db_session.commit() wd = bot_workdir(bot.id) # Fetch source: either from the configured source_repo or, if # the bot has no source_repo, we re-materialise from the # shared dataset (zip upload path). if bot.source_repo: try: await self._fetch_source_to(bot, wd) except DeployError as exc: # Source might already be in shared dataset from a # previous upload -- still try to materialise. log.warning("source_repo fetch failed (%s); falling back to dataset", exc) materialize_source(bot) else: materialize_source(bot) if not wd.exists() or not any(wd.iterdir()): raise DeployError( "Bot has no source files. Upload a zip or set source_repo." ) # Materialisation succeeded \u2014 files now live on the Space's # local filesystem in addition to the central HF Dataset. bot.space_state = SPACE_STATE_ON_SPACE # 2. Push source to the shared dataset (so other Spaces can # materialise the same bot if/when this Space is replaced). try: push_source(bot, wd, message=f"deploy {bot.slug}") except Exception as exc: log.warning("push_source to dataset failed (non-fatal): %s", exc) # 3. Write the bot's env file (secrets + user env vars) into # the workdir. This is read by the bot's run.sh / app.py at # boot via os.environ or a small loader. secrets = self._collect_secrets(bot) # Add the 4 admin-panel secrets; these are exposed to the bot # as ADMIN_LOGIN / ADMIN_PASSWORD / SECRET_WORD / SECRET_NUMBER. # They're two-tier encrypted in the DB, so we need the owner's # per-user Fernet key to decrypt. from auth import PER_USER_KEYS from models import User as _User owner = await db_session.get(_User, bot.owner_id) user_key = PER_USER_KEYS.get_user_key(owner) if owner else None admin_secret_keys = ( ("ADMIN_LOGIN", bot.admin_login_enc), ("ADMIN_PASSWORD", bot.admin_password_enc), ("SECRET_WORD", bot.secret_word_enc), ("SECRET_NUMBER", bot.secret_number_enc), ) for env_key, enc in admin_secret_keys: if not enc: continue try: secrets[env_key] = VAULT.decrypt(enc, user_key=user_key) if user_key else VAULT.decrypt(enc) except Exception as exc: log.warning("Could not decrypt %s for bot %s: %s", env_key, bot.id, exc) env_lines: list[str] = [] for k, v in secrets.items(): # Escape any embedded double-quotes / newlines. v_esc = v.replace("\\", "\\\\").replace('"', '\\"').replace("\n", "\\n") env_lines.append(f'{k}="{v_esc}"') if bot.env_vars_json: try: for k, v in json.loads(bot.env_vars_json).items(): if v in (None, ""): continue v_esc = str(v).replace("\\", "\\\\").replace('"', '\\"') env_lines.append(f'{k}="{v_esc}"') except json.JSONDecodeError: log.warning("bot %s has invalid env_vars_json; skipping", bot.id) env_lines.append(f'PANEL_BOT_PORT={bot.port}') env_lines.append(f'PANEL_BOT_ID={bot.id}') env_file = wd / ".env" env_file.write_text("\n".join(env_lines) + "\n", encoding="utf-8") try: os.chmod(env_file, 0o600) except Exception: pass # 4. Spawn the subprocess. try: set_record_quota(bot.id, bot.ram_mb or 0) rec = await start_bot(bot) log.info( "Bot %s running locally: pid=%d port=%d", bot.slug, rec.pid, rec.port, ) except Exception as exc: bot.status = BotStatus.ERROR await db_session.commit() db_session.add(DeploymentLog( bot_id=bot.id, action="deploy", status="failed", message=f"start_bot failed: {exc}", )) await db_session.commit() raise DeployError(f"Failed to start subprocess: {exc}") from exc # 5. Persist state. bot.status = BotStatus.RUNNING await db_session.commit() db_session.add(DeploymentLog( bot_id=bot.id, action="deploy", status="success", message=( f"Multi-tenant deploy on space {bot.space_id} port {bot.port} " f"(pid {rec.pid}, {bot.cpu_cores} CPU / {bot.ram_mb} MB)." ), details_json=json.dumps({ "mode": "multitenant", "space_id": bot.space_id, "port": bot.port, "pid": rec.pid, "cpu_cores": bot.cpu_cores, "ram_mb": bot.ram_mb, "source_repo": bot.source_repo, "source_type": bot.source_type, "secrets_count": len(secrets), }), )) await db_session.commit() # 6. Best-effort: flush run log placeholder to dataset. try: push_run_log(bot, "latest", message="deploy start") except Exception: pass async def _fetch_source_to(self, bot: BotInstance, dest: Path) -> None: """Same as _fetch_source but uses the given dest.""" # _fetch_source always goes to a tempdir; we adapt by passing # the workdir directly via a tempdir-shaped call. repo = bot.source_repo if not repo: raise DeployError("Bot has no source_repo.") if bot.source_type == "github": await self._fetch_github(bot, repo, bot.source_branch, dest) else: await self._fetch_hf(repo, dest) # ------------------------------------------------------------------ # Legacy one-bot-per-Space deploy path (unchanged behaviour) # ------------------------------------------------------------------ async def _deploy_legacy_space(self, bot: BotInstance, db_session: Any) -> None: hf_token = VAULT.decrypt(bot.hf_token_enc) if bot.hf_token_enc else "" if not hf_token: raise DeployError("HF token not configured for this bot") # Use the bot own HF token for the duration of the deploy # (bot owns the Space), but restore the panel token on # every exit path. Without _use_token the previous code # permanently clobbered self.hf_token and silently routed # subsequent bot deploys through the wrong HF account. # NOTE: concurrent deploys of different bots on the same # orchestrator instance remain unsupported - serialise at # the route layer. with self._use_token(hf_token): owner = self._resolve_owner(bot) if not owner: raise DeployError( "Cannot resolve HF owner. Set hf_dataset_repo or hf_space_name." ) repo_id = f"{owner}/{self._resolve_space_name(bot)}" # 1. Create the Space if missing try: await self.create_hf_space(owner, self._resolve_space_name(bot), private=True) log.info("HF Space %s created or already exists", repo_id) except DeployError as exc: log.warning("HF Space create: %s", exc) # 2. Fetch the source tree into a temp dir with tempfile.TemporaryDirectory(prefix=f"deploy-{bot.slug}-") as tmp: src = Path(tmp) try: await self._fetch_source(bot, src) except DeployError as exc: raise DeployError(f"Failed to fetch source: {exc}") from exc # Basic sanity: a Docker-based Space needs a Dockerfile. if not (src / "Dockerfile").exists(): log.warning("Source has no Dockerfile — HF Space may not start.") # 3. Push the tree to the Space await self.push_source(bot, src, message=f"Deploy {bot.slug} via panel") # 4. Set encrypted secrets (whatever the bot has configured). secrets = self._collect_secrets(bot) await self.set_secrets(bot, secrets) # 5. Set user env vars (non-secret) as Space variables. if bot.env_vars_json: try: extra = json.loads(bot.env_vars_json) if extra: await self.set_env_vars(bot, extra) except json.JSONDecodeError: log.warning("bot %s has invalid env_vars_json; skipping", bot.id) # 6. Update status and persist the resolved URL. bot.status = BotStatus.BUILDING bot.hf_space_url = f"https://huggingface.co/spaces/{repo_id}" await db_session.commit() log_entry = DeploymentLog( bot_id=bot.id, action="deploy", status="success", message=f"HF Space {repo_id} configured. Source from {bot.source_repo or '?'}.", details_json=json.dumps({ "source_repo": bot.source_repo, "source_type": bot.source_type, "branch": bot.source_branch, "secrets_count": len(secrets), }), ) db_session.add(log_entry) await db_session.commit() async def set_env_vars(self, bot: BotInstance, env: dict[str, str]) -> None: """Push user env vars to the Space as variables.""" owner = self._resolve_owner(bot) repo_id = f"{owner}/{self._resolve_space_name(bot)}" def _do(): for k, v in env.items(): if v is None or v == "": continue try: self._get_api().add_space_variable( repo_id=repo_id, key=k, value=v, token=self.hf_token, ) except Exception as exc: log.warning("Failed to set HF variable %s: %s", k, exc) await self._run(_do) # ------------------------------------------------------------------ # Helpers # ------------------------------------------------------------------ def _collect_secrets(self, bot: BotInstance) -> dict[str, str]: """Decrypt and return every secret the bot has configured. Anything that's blank is omitted — the deploy never pushes empty strings. Only fields that are actually filled in are sent. """ out: dict[str, str] = {} def add(key: str, enc: Optional[str]) -> None: if not enc: return try: value = VAULT.decrypt(enc) except Exception: log.warning("Failed to decrypt %s for bot %s", key, bot.id) return if value: out[key] = value add("BOT_TOKEN", bot.bot_token_enc) add("HF_TOKEN", bot.hf_token_enc) add("NVIDIA_API_KEY", bot.nvidia_api_key_enc) add("LLM_PROXY_TOKEN", bot.llm_proxy_token_enc) # Optional plaintext values (the bot owner can override). if bot.nvidia_model: out["NVIDIA_MODEL"] = bot.nvidia_model if bot.hf_dataset_repo: out["HF_DATASET_REPO"] = bot.hf_dataset_repo if bot.render_url: out["LLM_PROXY_URL"] = bot.render_url out["TELEGRAM_PROXY_URL"] = bot.render_url.rstrip("/") + "/telegram" out["RENDER_URL"] = bot.render_url # Platform-wide Render proxy (overrides bot.render_url if set). try: from proxy_client import bot_env_extras out.update(bot_env_extras()) except Exception as exc: log.warning("proxy_client.bot_env_extras failed: %s", exc) return out async def stop_bot(self, bot: BotInstance, db_session: Any) -> None: """Stop a bot. Multi-tenant kills the local subprocess; legacy pauses the HF Space. """ if bot.deployment_mode == DeploymentMode.MULTITENANT: try: await runner_stop(bot) from bot_runner import clear_record_quota clear_record_quota(bot.id) except Exception as exc: log.warning("runner_stop failed for %s: %s", bot.slug, exc) bot.status = BotStatus.STOPPED # After stop the on-Space files are kept (so the next start # is instant); the dataset remains the source of truth. bot.space_state = SPACE_STATE_ON_SPACE_KEPT await db_session.commit() db_session.add(DeploymentLog( bot_id=bot.id, action="stop", status="success", message="Multi-tenant bot stopped", )) await db_session.commit() return # Legacy HF Space pause. hf_token = VAULT.decrypt(bot.hf_token_enc) if bot.hf_token_enc else "" # Same save/restore pattern as _deploy_legacy_space: operate on # the bot's token without permanently clobbering the panel's. # ``_use_token("")`` is a no-op (we never want to blank the # panel token — keep using it if the bot has none). with self._use_token(hf_token): owner = self._resolve_owner(bot) if not owner: bot.status = BotStatus.STOPPED await db_session.commit() return repo_id = f"{owner}/{self._resolve_space_name(bot)}" def _do(): try: self._get_api().pause_space(repo_id=repo_id, token=self.hf_token) except Exception as exc: log.warning("pause_space failed for %s: %s", repo_id, exc) await self._run(_do) bot.status = BotStatus.STOPPED await db_session.commit() log_entry = DeploymentLog( bot_id=bot.id, action="stop", status="success", message="Bot stopped via panel", ) db_session.add(log_entry) await db_session.commit()