bot_host / deploy.py
ItsBounvy's picture
deploy: english UI + admin creds + dataset wiring
a4517f0 verified
Raw History Blame Contribute Delete
27.2 kB
"""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()