import json import os import tarfile import tempfile import urllib.request from datetime import datetime, timedelta, timezone import boto3 from google.cloud import storage as gcs from google.oauth2 import service_account s3 = boto3.client("s3") s3_west = boto3.client("s3", region_name="us-west-2") ec2 = boto3.client("ec2") secrets = boto3.client("secretsmanager") SOURCE_BUCKET = os.environ["SOURCE_BUCKET"] REPLICA_BUCKET = os.environ["REPLICA_BUCKET"] GCS_BUCKET = os.environ["GCS_BUCKET"] GCS_SA_SECRET_ARN = os.environ["GCS_SA_SECRET_ARN"] SLACK_WEBHOOK_SECRET_ARN = os.environ["SLACK_WEBHOOK_SECRET_ARN"] _gcs_client = None def _get_gcs_client(): global _gcs_client if _gcs_client is None: raw = secrets.get_secret_value(SecretId=GCS_SA_SECRET_ARN)["SecretString"] info = json.loads(raw) creds = service_account.Credentials.from_service_account_info(info) _gcs_client = gcs.Client(credentials=creds, project=info.get("project_id")) return _gcs_client def _check_s3_bucket(client, bucket, label): now = datetime.now(timezone.utc) cutoff = now - timedelta(hours=48) try: resp = client.list_objects_v2(Bucket=bucket, Prefix="archive/", MaxKeys=1000) contents = resp.get("Contents", []) if not contents: return False, f"{label}: No objects found under archive/" latest = max(contents, key=lambda o: o["LastModified"]) if latest["LastModified"] < cutoff: age = (now - latest["LastModified"]).total_seconds() / 3600 return False, f"{label}: Latest dump is {age:.0f}h old ({latest['Key']})" if latest["Size"] < 1_000_000: return False, f"{label}: Latest dump suspiciously small ({latest['Size']} bytes)" return True, f"{label}: OK — {latest['Key']} ({latest['Size'] / 1_000_000:.1f} MB)" except Exception as e: return False, f"{label}: Error — {e}" def _check_gcs(): try: client = _get_gcs_client() bucket = client.bucket(GCS_BUCKET) blobs = list(bucket.list_blobs(prefix="archive/", max_results=1000)) if not blobs: return False, "GCS Offsite: No objects found under archive/" now = datetime.now(timezone.utc) cutoff = now - timedelta(hours=72) latest = max(blobs, key=lambda b: b.updated) if latest.updated < cutoff: age = (now - latest.updated).total_seconds() / 3600 return False, f"GCS Offsite: Latest object is {age:.0f}h old ({latest.name})" return True, f"GCS Offsite: OK — {latest.name} ({latest.size / 1_000_000:.1f} MB)" except Exception as e: return False, f"GCS Offsite: Error — {e}" def _check_ebs_snapshots(): try: now = datetime.now(timezone.utc) cutoff = now - timedelta(hours=48) resp = ec2.describe_snapshots( Filters=[{"Name": "tag:forgejo-backup", "Values": ["true"]}], OwnerIds=["self"], ) snapshots = resp.get("Snapshots", []) if not snapshots: return False, "EBS Snapshots: No snapshots found with forgejo-backup tag" recent = [s for s in snapshots if s["StartTime"] >= cutoff] if not recent: latest = max(snapshots, key=lambda s: s["StartTime"]) age = (now - latest["StartTime"]).total_seconds() / 3600 return False, f"EBS Snapshots: Latest is {age:.0f}h old ({latest['SnapshotId']})" return True, f"EBS Snapshots: OK — {len(snapshots)} total, {len(recent)} in last 48h" except Exception as e: return False, f"EBS Snapshots: Error — {e}" def _restore_test(): results = [] try: resp = s3.list_objects_v2(Bucket=SOURCE_BUCKET, Prefix="archive/", MaxKeys=1000) contents = resp.get("Contents", []) if not contents: return [{"pass": False, "msg": "Restore test: No dumps found in source bucket"}] latest = max(contents, key=lambda o: o["LastModified"]) results.append({"pass": True, "msg": f"Restore test: Using {latest['Key']} ({latest['Size'] / 1_000_000:.1f} MB)"}) with tempfile.TemporaryDirectory() as tmpdir: local_path = os.path.join(tmpdir, "dump.tar.gz") s3.download_file(SOURCE_BUCKET, latest["Key"], local_path) results.append({"pass": True, "msg": "Restore test: Download OK"}) try: with tarfile.open(local_path, "r:gz") as tf: names = tf.getnames() results.append({"pass": True, "msg": f"Restore test: Archive OK — {len(names)} entries"}) db_entries = [n for n in names if n.endswith(".db") or n.endswith("forgejo.db")] if db_entries: import sqlite3 as sqlite_mod tf.extract(db_entries[0], path=tmpdir) db_path = os.path.join(tmpdir, db_entries[0]) conn = sqlite_mod.connect(db_path) result = conn.execute("PRAGMA integrity_check").fetchone() conn.close() if result[0] == "ok": results.append({"pass": True, "msg": "Restore test: SQLite integrity OK"}) else: results.append({"pass": False, "msg": f"Restore test: SQLite integrity FAILED — {result[0]}"}) else: results.append({"pass": True, "msg": "Restore test: No .db file found in archive (may use different format)"}) except tarfile.TarError as e: results.append({"pass": False, "msg": f"Restore test: Archive extraction FAILED — {e}"}) except Exception as e: results.append({"pass": False, "msg": f"Restore test: Error — {e}"}) return results def _post_slack(blocks): raw = secrets.get_secret_value(SecretId=SLACK_WEBHOOK_SECRET_ARN)["SecretString"] webhook_url = raw.strip() payload = json.dumps({"blocks": blocks}).encode() req = urllib.request.Request( webhook_url, data=payload, headers={"Content-Type": "application/json"}, method="POST", ) urllib.request.urlopen(req) def handler(event, context): mode = event.get("mode", "daily") results = [] if mode == "daily": results.append(_check_s3_bucket(s3, SOURCE_BUCKET, "S3 Source (us-east-1)")) results.append(_check_s3_bucket(s3_west, REPLICA_BUCKET, "S3 Replica (us-west-2)")) results.append(_check_gcs()) results.append(_check_ebs_snapshots()) all_pass = all(r[0] for r in results) header = "Forgejo Backup Verification" blocks = [ {"type": "header", "text": {"type": "plain_text", "text": header}}, {"type": "section", "text": {"type": "mrkdwn", "text": f"*Date:* {datetime.now(timezone.utc).strftime('%Y-%m-%d %H:%M UTC')}"}}, {"type": "divider"}, ] for passed, msg in results: emoji = ":white_check_mark:" if passed else ":x:" blocks.append({"type": "section", "text": {"type": "mrkdwn", "text": f"{emoji} {msg}"}}) blocks.append({"type": "divider"}) overall = ":white_check_mark: All checks passed" if all_pass else ":rotating_light: One or more checks failed" blocks.append({"type": "section", "text": {"type": "mrkdwn", "text": f"*Overall:* {overall}"}}) elif mode == "restore-test": test_results = _restore_test() all_pass = all(r["pass"] for r in test_results) header = "Forgejo Monthly Restore Test" blocks = [ {"type": "header", "text": {"type": "plain_text", "text": header}}, {"type": "section", "text": {"type": "mrkdwn", "text": f"*Date:* {datetime.now(timezone.utc).strftime('%Y-%m-%d %H:%M UTC')}"}}, {"type": "divider"}, ] for r in test_results: emoji = ":white_check_mark:" if r["pass"] else ":x:" blocks.append({"type": "section", "text": {"type": "mrkdwn", "text": f"{emoji} {r['msg']}"}}) blocks.append({"type": "divider"}) overall = ":white_check_mark: Restore test passed" if all_pass else ":rotating_light: Restore test failed" blocks.append({"type": "section", "text": {"type": "mrkdwn", "text": f"*Overall:* {overall}"}}) _post_slack(blocks) return { "statusCode": 200, "body": json.dumps({ "mode": mode, "all_pass": all_pass, "results": [{"pass": r[0], "msg": r[1]} for r in results] if mode == "daily" else test_results, }), }