feat(backup): add GCS offsite backup executor
Add gcs_backup_executor with git_bundle mode for backing up bare git repos to Google Cloud Storage. Includes retention management and bundle verification. Adds google-cloud-storage dependency and GCS_CREDENTIALS_FILE setting. Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
+3
-1
@@ -1,6 +1,6 @@
|
||||
[project]
|
||||
name = "the-scheduler"
|
||||
version = "1.1.3"
|
||||
version = "1.2.0"
|
||||
description = "System-wide maintenance orchestration - backups, doc mirroring, cleanup, task automation"
|
||||
readme = "README.md"
|
||||
requires-python = ">=3.12"
|
||||
@@ -26,6 +26,8 @@ dependencies = [
|
||||
"gitpython~=3.1.43",
|
||||
# Security/Auth
|
||||
"python-jose[cryptography]~=3.3.0",
|
||||
# Cloud storage
|
||||
"google-cloud-storage~=2.18.0",
|
||||
]
|
||||
|
||||
[project.optional-dependencies]
|
||||
|
||||
@@ -26,6 +26,9 @@ gitpython~=3.1.43 # Latest stable
|
||||
# Security/Auth
|
||||
python-jose[cryptography]~=3.3.0 # JWT handling
|
||||
|
||||
# Cloud storage
|
||||
google-cloud-storage~=2.18.0 # GCS offsite backups
|
||||
|
||||
# Testing
|
||||
pytest~=8.3.4 # Test framework
|
||||
pytest-asyncio~=0.25.2 # Async test support
|
||||
|
||||
@@ -66,6 +66,9 @@ class Settings(BaseSettings):
|
||||
backup_retention_weekly: int = Field(default=4, alias="BACKUP_RETENTION_WEEKLY")
|
||||
backup_retention_monthly: int = Field(default=12, alias="BACKUP_RETENTION_MONTHLY")
|
||||
|
||||
# Google Cloud Storage
|
||||
gcs_credentials_file: str = Field(default="", alias="GCS_CREDENTIALS_FILE")
|
||||
|
||||
# Documentation Mirroring
|
||||
docs_mirror_path: str = Field(default="/docs-mirror", alias="DOCS_MIRROR_PATH")
|
||||
docs_check_interval: int = Field(default=21600, alias="DOCS_CHECK_INTERVAL") # 6 hours
|
||||
|
||||
@@ -0,0 +1,199 @@
|
||||
"""
|
||||
Google Cloud Storage Backup Executor
|
||||
Backs up data to GCS buckets. Supports multiple modes:
|
||||
- git_bundle: Creates a git bundle from a bare repo and uploads it
|
||||
"""
|
||||
import asyncio
|
||||
import logging
|
||||
import time
|
||||
from datetime import datetime
|
||||
from pathlib import Path
|
||||
from typing import List, Optional
|
||||
|
||||
from google.cloud import storage
|
||||
|
||||
from src.config import Settings
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
async def execute(config: dict, settings: Settings) -> str:
|
||||
"""
|
||||
Execute GCS backup task.
|
||||
|
||||
Config schema:
|
||||
{
|
||||
"mode": "git_bundle",
|
||||
"bucket": "bucket-name",
|
||||
"prefix": "gitea/settled-reach",
|
||||
"credentials_path": "/secrets/gcs-sa-key.json",
|
||||
"repo_path": "/data/docker-data/gitea/data/git/repositories/user/repo.git",
|
||||
"retention_count": 7,
|
||||
"dry_run": false
|
||||
}
|
||||
|
||||
Args:
|
||||
config: Backup configuration
|
||||
settings: Global scheduler settings
|
||||
|
||||
Returns:
|
||||
Summary of backup operation
|
||||
|
||||
Raises:
|
||||
Exception: On backup failure
|
||||
"""
|
||||
mode = config.get('mode')
|
||||
if not mode:
|
||||
raise ValueError("Missing required config field: mode")
|
||||
|
||||
if mode == 'git_bundle':
|
||||
return await _mode_git_bundle(config, settings)
|
||||
else:
|
||||
raise ValueError(f"Unknown backup mode: {mode}")
|
||||
|
||||
|
||||
async def _mode_git_bundle(config: dict, settings: Settings) -> str:
|
||||
"""Create a git bundle from a bare repo and upload to GCS."""
|
||||
bucket_name = config.get('bucket')
|
||||
prefix = config.get('prefix', '').strip('/')
|
||||
credentials_path = config.get('credentials_path', settings.gcs_credentials_file)
|
||||
repo_path = Path(config.get('repo_path', ''))
|
||||
retention_count = config.get('retention_count', 7)
|
||||
dry_run = config.get('dry_run', False)
|
||||
|
||||
# Validate
|
||||
if not bucket_name:
|
||||
raise ValueError("Missing required config field: bucket")
|
||||
if not repo_path or not str(repo_path).strip():
|
||||
raise ValueError("Missing required config field: repo_path")
|
||||
if not repo_path.exists():
|
||||
raise ValueError(f"Repository path does not exist: {repo_path}")
|
||||
if not (repo_path / 'HEAD').exists():
|
||||
raise ValueError(f"Not a valid git repository: {repo_path}")
|
||||
|
||||
repo_name = repo_path.name.removesuffix('.git')
|
||||
timestamp = datetime.now().strftime('%Y%m%d-%H%M%S')
|
||||
bundle_filename = f"{repo_name}-{timestamp}.bundle"
|
||||
bundle_path = Path(f"/tmp/{bundle_filename}")
|
||||
|
||||
logger.info(f"Starting GCS backup: mode=git_bundle, repo={repo_name}, dry_run={dry_run}")
|
||||
|
||||
try:
|
||||
# Step 1: Create git bundle
|
||||
logger.info(f"Creating git bundle from {repo_path}")
|
||||
start = time.monotonic()
|
||||
|
||||
await _run_command([
|
||||
'git', 'bundle', 'create',
|
||||
str(bundle_path),
|
||||
'--all'
|
||||
], cwd=repo_path)
|
||||
|
||||
bundle_duration = time.monotonic() - start
|
||||
|
||||
if not bundle_path.exists():
|
||||
raise Exception("Git bundle was not created")
|
||||
|
||||
bundle_size = bundle_path.stat().st_size
|
||||
bundle_size_mb = bundle_size / (1024 * 1024)
|
||||
logger.info(f"Bundle created: {bundle_filename} ({bundle_size_mb:.1f} MB) in {bundle_duration:.1f}s")
|
||||
|
||||
# Step 2: Verify bundle
|
||||
await _run_command([
|
||||
'git', 'bundle', 'verify',
|
||||
str(bundle_path)
|
||||
], cwd=repo_path)
|
||||
logger.info("Bundle verified OK")
|
||||
|
||||
if dry_run:
|
||||
return (
|
||||
f"[DRY RUN] Would upload {bundle_filename} ({bundle_size_mb:.1f} MB) "
|
||||
f"to gs://{bucket_name}/{prefix}/{bundle_filename}"
|
||||
)
|
||||
|
||||
# Step 3: Upload to GCS
|
||||
gcs_path = f"{prefix}/{bundle_filename}" if prefix else bundle_filename
|
||||
logger.info(f"Uploading to gs://{bucket_name}/{gcs_path}")
|
||||
start = time.monotonic()
|
||||
|
||||
client = storage.Client.from_service_account_json(credentials_path)
|
||||
bucket = client.bucket(bucket_name)
|
||||
blob = bucket.blob(gcs_path)
|
||||
blob.upload_from_filename(str(bundle_path), timeout=3600)
|
||||
|
||||
upload_duration = time.monotonic() - start
|
||||
logger.info(f"Upload complete in {upload_duration:.1f}s")
|
||||
|
||||
# Step 4: Retention cleanup in GCS
|
||||
deleted_count = await _cleanup_gcs_retention(
|
||||
client, bucket_name, prefix, repo_name, retention_count
|
||||
)
|
||||
|
||||
return (
|
||||
f"Backup completed: {bundle_filename} ({bundle_size_mb:.1f} MB). "
|
||||
f"Bundle: {bundle_duration:.1f}s, Upload: {upload_duration:.1f}s. "
|
||||
f"GCS: gs://{bucket_name}/{gcs_path}. "
|
||||
f"Retention: {deleted_count} old bundle(s) removed."
|
||||
)
|
||||
|
||||
finally:
|
||||
# Always clean up the local temp file
|
||||
if bundle_path.exists():
|
||||
bundle_path.unlink()
|
||||
logger.debug(f"Cleaned up temp file: {bundle_path}")
|
||||
|
||||
|
||||
async def _cleanup_gcs_retention(
|
||||
client: storage.Client,
|
||||
bucket_name: str,
|
||||
prefix: str,
|
||||
repo_name: str,
|
||||
retention_count: int
|
||||
) -> int:
|
||||
"""Delete old bundles from GCS, keeping only the most recent retention_count."""
|
||||
if retention_count <= 0:
|
||||
return 0
|
||||
|
||||
bucket = client.bucket(bucket_name)
|
||||
blob_prefix = f"{prefix}/{repo_name}-" if prefix else f"{repo_name}-"
|
||||
|
||||
blobs = list(bucket.list_blobs(prefix=blob_prefix))
|
||||
bundle_blobs = [b for b in blobs if b.name.endswith('.bundle')]
|
||||
|
||||
if len(bundle_blobs) <= retention_count:
|
||||
logger.info(f"Retention OK: {len(bundle_blobs)} bundles (limit: {retention_count})")
|
||||
return 0
|
||||
|
||||
# Sort by name (timestamp in name ensures chronological order)
|
||||
bundle_blobs.sort(key=lambda b: b.name)
|
||||
to_delete = bundle_blobs[:-retention_count]
|
||||
|
||||
for blob in to_delete:
|
||||
logger.info(f"Deleting old bundle: {blob.name}")
|
||||
blob.delete()
|
||||
|
||||
logger.info(f"Retention cleanup: deleted {len(to_delete)} old bundle(s)")
|
||||
return len(to_delete)
|
||||
|
||||
|
||||
async def _run_command(
|
||||
cmd: List[str],
|
||||
cwd: Optional[Path] = None,
|
||||
) -> str:
|
||||
"""Run shell command asynchronously."""
|
||||
logger.debug(f"Running: {' '.join(cmd)} (cwd: {cwd})")
|
||||
|
||||
proc = await asyncio.create_subprocess_exec(
|
||||
*cmd,
|
||||
cwd=cwd,
|
||||
stdout=asyncio.subprocess.PIPE,
|
||||
stderr=asyncio.subprocess.PIPE
|
||||
)
|
||||
|
||||
stdout, stderr = await proc.communicate()
|
||||
|
||||
if proc.returncode != 0:
|
||||
error_msg = stderr.decode() if stderr else "Unknown error"
|
||||
raise Exception(f"Command failed: {' '.join(cmd)}\n{error_msg}")
|
||||
|
||||
return stdout.decode()
|
||||
Reference in New Issue
Block a user