Template
285 lines
8.7 KiB
Python
285 lines
8.7 KiB
Python
"""Бэкап / восстановление PostgreSQL через pg_dump и psql."""
|
||
|
||
from __future__ import annotations
|
||
|
||
import asyncio
|
||
import gzip
|
||
import logging
|
||
import os
|
||
import re
|
||
import shutil
|
||
from datetime import datetime, timezone
|
||
from pathlib import Path
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
BACKUP_DIR = Path(os.getenv("BACKUP_DIR") or "/app/backups")
|
||
_SAFE_NAME = re.compile(r"^vpnbot-[\w.\-]+\.sql\.gz$", re.I)
|
||
_UPLOAD_NAME = re.compile(r"^[\w.\-]+\.sql(\.gz)?$", re.I)
|
||
|
||
|
||
def _pg_env() -> tuple[str, str, str, str, str]:
|
||
user = os.getenv("POSTGRES_USER") or "vpnbot"
|
||
dbname = os.getenv("POSTGRES_DB") or "vpnbot"
|
||
password = os.getenv("POSTGRES_PASSWORD") or "vpnbot"
|
||
host = os.getenv("POSTGRES_HOST") or "db"
|
||
port = os.getenv("POSTGRES_PORT") or "5432"
|
||
return user, dbname, password, host, port
|
||
|
||
|
||
def resolve_backup_path(name: str) -> Path:
|
||
"""Безопасный путь к файлу бэкапа (без path traversal)."""
|
||
base = Path(name).name
|
||
if not _SAFE_NAME.match(base):
|
||
raise ValueError("Некорректное имя файла бэкапа")
|
||
path = (BACKUP_DIR / base).resolve()
|
||
root = BACKUP_DIR.resolve()
|
||
if path.parent != root or not path.is_file():
|
||
raise FileNotFoundError("Бэкап не найден")
|
||
return path
|
||
|
||
|
||
def delete_backup(name: str) -> None:
|
||
path = resolve_backup_path(name)
|
||
path.unlink()
|
||
|
||
|
||
async def _run_pg_dump_to_host() -> bytes:
|
||
"""pg_dump по сети Docker к сервису db."""
|
||
user, dbname, password, host, port = _pg_env()
|
||
pg_dump = shutil.which("pg_dump")
|
||
if not pg_dump:
|
||
raise FileNotFoundError(
|
||
"pg_dump не найден. Пересоберите образ: docker compose up -d --build bot"
|
||
)
|
||
|
||
env = os.environ.copy()
|
||
env["PGPASSWORD"] = password
|
||
proc = await asyncio.create_subprocess_exec(
|
||
pg_dump,
|
||
"-h",
|
||
host,
|
||
"-p",
|
||
port,
|
||
"-U",
|
||
user,
|
||
"-d",
|
||
dbname,
|
||
"--no-owner",
|
||
"--no-acl",
|
||
"--clean",
|
||
"--if-exists",
|
||
stdout=asyncio.subprocess.PIPE,
|
||
stderr=asyncio.subprocess.PIPE,
|
||
env=env,
|
||
)
|
||
stdout, stderr = await proc.communicate()
|
||
if proc.returncode != 0:
|
||
raise RuntimeError((stderr or b"").decode("utf-8", "replace")[:500] or "pg_dump failed")
|
||
if not stdout:
|
||
raise RuntimeError("pg_dump вернул пустой дамп")
|
||
return stdout
|
||
|
||
|
||
async def _run_pg_dump_via_docker(compose_dir: Path | None) -> bytes:
|
||
user, dbname, password, _host, _port = _pg_env()
|
||
if not shutil.which("docker"):
|
||
raise FileNotFoundError("docker не найден")
|
||
|
||
cmd = [
|
||
"docker",
|
||
"compose",
|
||
"exec",
|
||
"-T",
|
||
"-e",
|
||
f"PGPASSWORD={password}",
|
||
"db",
|
||
"pg_dump",
|
||
"-U",
|
||
user,
|
||
"-d",
|
||
dbname,
|
||
"--no-owner",
|
||
"--no-acl",
|
||
"--clean",
|
||
"--if-exists",
|
||
]
|
||
proc = await asyncio.create_subprocess_exec(
|
||
*cmd,
|
||
stdout=asyncio.subprocess.PIPE,
|
||
stderr=asyncio.subprocess.PIPE,
|
||
cwd=str(compose_dir) if compose_dir else None,
|
||
)
|
||
stdout, stderr = await proc.communicate()
|
||
if proc.returncode != 0:
|
||
raise RuntimeError((stderr or b"").decode("utf-8", "replace")[:500] or "docker pg_dump failed")
|
||
if not stdout:
|
||
raise RuntimeError("docker pg_dump вернул пустой дамп")
|
||
return stdout
|
||
|
||
|
||
async def run_pg_dump_backup(
|
||
*,
|
||
compose_dir: Path | None = None,
|
||
keep: int = 14,
|
||
) -> Path:
|
||
BACKUP_DIR.mkdir(parents=True, exist_ok=True)
|
||
stamp = datetime.now(timezone.utc).strftime("%Y%m%d-%H%M%S")
|
||
out = BACKUP_DIR / f"vpnbot-{stamp}.sql.gz"
|
||
|
||
errors: list[str] = []
|
||
stdout: bytes | None = None
|
||
|
||
try:
|
||
stdout = await _run_pg_dump_to_host()
|
||
except Exception as exc: # noqa: BLE001
|
||
errors.append(f"network pg_dump: {exc}")
|
||
logger.warning("backup via host pg_dump failed: %s", exc)
|
||
|
||
if stdout is None:
|
||
try:
|
||
stdout = await _run_pg_dump_via_docker(compose_dir)
|
||
except Exception as exc: # noqa: BLE001
|
||
errors.append(f"docker pg_dump: {exc}")
|
||
logger.warning("backup via docker failed: %s", exc)
|
||
|
||
if stdout is None:
|
||
raise RuntimeError(
|
||
"Не удалось сделать бэкап. "
|
||
+ " | ".join(errors)
|
||
+ ". Убедитесь, что bot пересобран с postgresql-client."
|
||
)
|
||
|
||
with gzip.open(out, "wb") as fh:
|
||
fh.write(stdout)
|
||
|
||
_prune(BACKUP_DIR, keep=keep)
|
||
logger.info("postgres backup written %s (%s bytes)", out, len(stdout))
|
||
return out
|
||
|
||
|
||
def _prune(directory: Path, *, keep: int) -> None:
|
||
files = sorted(
|
||
directory.glob("vpnbot-*.sql.gz"),
|
||
key=lambda p: p.stat().st_mtime,
|
||
reverse=True,
|
||
)
|
||
for old in files[max(0, keep) :]:
|
||
try:
|
||
old.unlink()
|
||
except OSError:
|
||
pass
|
||
|
||
|
||
def list_backups(limit: int = 30) -> list[dict]:
|
||
BACKUP_DIR.mkdir(parents=True, exist_ok=True)
|
||
files = sorted(
|
||
BACKUP_DIR.glob("vpnbot-*.sql.gz"),
|
||
key=lambda p: p.stat().st_mtime,
|
||
reverse=True,
|
||
)
|
||
out = []
|
||
for p in files[:limit]:
|
||
st = p.stat()
|
||
out.append(
|
||
{
|
||
"name": p.name,
|
||
"size": st.st_size,
|
||
"mtime": datetime.fromtimestamp(st.st_mtime, tz=timezone.utc).isoformat(),
|
||
}
|
||
)
|
||
return out
|
||
|
||
|
||
def _gunzip_if_needed(raw: bytes, *, filename: str = "") -> bytes:
|
||
name = (filename or "").lower()
|
||
if name.endswith(".gz") or raw[:2] == b"\x1f\x8b":
|
||
return gzip.decompress(raw)
|
||
return raw
|
||
|
||
|
||
async def _run_psql_sql(sql: bytes) -> None:
|
||
user, dbname, password, host, port = _pg_env()
|
||
psql = shutil.which("psql")
|
||
if not psql:
|
||
raise FileNotFoundError(
|
||
"psql не найден. Пересоберите образ: docker compose up -d --build bot"
|
||
)
|
||
|
||
# Чистый restore: рвём сессии и пересоздаём schema public,
|
||
# иначе старые дампы без DROP падают на «relation already exists».
|
||
safe_user = user.replace('"', "")
|
||
pre = (
|
||
"SELECT pg_terminate_backend(pid) FROM pg_stat_activity "
|
||
f"WHERE datname = '{dbname}' AND pid <> pg_backend_pid();\n"
|
||
"DROP SCHEMA IF EXISTS public CASCADE;\n"
|
||
"CREATE SCHEMA public;\n"
|
||
f'GRANT ALL ON SCHEMA public TO "{safe_user}";\n'
|
||
"GRANT ALL ON SCHEMA public TO public;\n"
|
||
f'ALTER SCHEMA public OWNER TO "{safe_user}";\n'
|
||
).encode("utf-8")
|
||
|
||
env = os.environ.copy()
|
||
env["PGPASSWORD"] = password
|
||
proc = await asyncio.create_subprocess_exec(
|
||
psql,
|
||
"-h",
|
||
host,
|
||
"-p",
|
||
port,
|
||
"-U",
|
||
user,
|
||
"-d",
|
||
dbname,
|
||
"-v",
|
||
"ON_ERROR_STOP=1",
|
||
stdin=asyncio.subprocess.PIPE,
|
||
stdout=asyncio.subprocess.PIPE,
|
||
stderr=asyncio.subprocess.PIPE,
|
||
env=env,
|
||
)
|
||
_stdout, stderr = await proc.communicate(pre + sql)
|
||
if proc.returncode != 0:
|
||
err = (stderr or b"").decode("utf-8", "replace")[:800]
|
||
raise RuntimeError(err or "psql restore failed")
|
||
logger.info("postgres restore ok (%s bytes sql)", len(sql))
|
||
|
||
|
||
async def restore_from_bytes(raw: bytes, *, filename: str = "") -> None:
|
||
if not raw:
|
||
raise ValueError("Пустой файл")
|
||
if len(raw) > 200 * 1024 * 1024:
|
||
raise ValueError("Файл слишком большой (лимит 200 МБ)")
|
||
sql = _gunzip_if_needed(raw, filename=filename)
|
||
if not sql.strip():
|
||
raise ValueError("Дамп пустой после распаковки")
|
||
head = sql[:4000].decode("utf-8", "replace").lower()
|
||
if (
|
||
"create " not in head
|
||
and "drop " not in head
|
||
and "copy " not in head
|
||
and "insert " not in head
|
||
):
|
||
raise ValueError("Файл не похож на SQL-дамп PostgreSQL")
|
||
await _run_psql_sql(sql)
|
||
|
||
|
||
async def restore_from_backup_name(name: str) -> None:
|
||
path = resolve_backup_path(name)
|
||
raw = path.read_bytes()
|
||
await restore_from_bytes(raw, filename=path.name)
|
||
|
||
|
||
async def save_uploaded_backup(raw: bytes, *, original_name: str = "") -> Path:
|
||
"""Сохранить загруженный дамп как vpnbot-*.sql.gz на диск."""
|
||
BACKUP_DIR.mkdir(parents=True, exist_ok=True)
|
||
base = Path(original_name or "upload.sql.gz").name
|
||
if not _UPLOAD_NAME.match(base):
|
||
base = "upload.sql.gz"
|
||
sql = _gunzip_if_needed(raw, filename=base)
|
||
stamp = datetime.now(timezone.utc).strftime("%Y%m%d-%H%M%S")
|
||
out = BACKUP_DIR / f"vpnbot-{stamp}.sql.gz"
|
||
with gzip.open(out, "wb") as fh:
|
||
fh.write(sql)
|
||
return out
|