forgejo/lambda/backup-verification/app.py

247 lines
11 KiB
Python
Raw Normal View History

Add 3-2-1 backup strategy (#2) * Add 3-2-1 backup strategy with cross-region replication and GCS offsite Implements a fully compliant 3-2-1 backup architecture: - Copy 1 (live): Harden existing EBS snapshots to 30-day retention - Copy 2 (near-site): S3 cross-region replication to us-west-2 with Object Lock (governance 90d) and versioning - Copy 3 (offsite): GCS bucket in dedicated seahaven-backups GCP project with 2-year irreversible retention lock Also adds a verification Lambda that checks all 3 locations daily and runs monthly restore tests with SQLite integrity checks. * Enable QEMU in CI for arm64 Lambda Docker builds * Commit cdk.context.json for CI synth without AWS credentials Vpc.fromLookup requires cached context to synthesize without AWS credentials. Required for CI which runs cdk synth without an OIDC role. * Fix GCP project ID to sea-haven-backups * Address code review findings for backup verification Fix 4 critical issues: - Add filter/priority/deleteMarkerReplication to S3 CRR rule (deploy would fail without) - Add stack dependency so replica deploys before main stack - Fix DB file extension matching (.sqlite3/.sql instead of .db) - Replace nonexistent `forgejo restore` command with actual restore steps in README Fix 4 moderate issues: - Add timeout=10 to Slack webhook urlopen call - Add filter='data' to tarfile.extract for PEP 706 compliance - Add explicit ValueError for unknown handler mode - Use date-scoped S3/GCS prefix instead of unbounded listing * Fix backup strategy bug findings * Handle SQL text dumps separately from binary SQLite in restore test Forgejo dump produces gitea-db.sql as a text SQL dump (XORM export), not a binary SQLite file. Opening it directly with sqlite3.connect() throws DatabaseError. Now imports the SQL dump into a temp DB first. * Fix GCS backup check: align staleness cutoff and add size validation GCS check used a 72h cutoff but only listed 2 days of prefixes (~48h), making the staleness check unreachable. Also added 1MB minimum file size validation to match the S3 check. * Rename SECRET_ARN env vars to SECRET_NAME to match actual values * Fix EBS snapshot state check, drop unused GCS write grant and dead lifecycle rule * Fix restore runbook, DLM snapshot tagging, README cleanup, and gsutil prompt * Fix restore runbook: trailing-dot cp idiom and Glacier restore step * Rename GCS service account to match read-only permissions * Add 4 GiB ephemeral storage to verification Lambda Monthly restore-test downloads and extracts the full dump tarball in /tmp. As the dump grows with LFS data, the default 512 MB will eventually cause ENOSPC failures. --------- Co-authored-by: Cursor Agent <cursoragent@cursor.com>
2026-05-14 18:08:06 -04:00
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_NAME = os.environ["GCS_SA_SECRET_NAME"]
SLACK_WEBHOOK_SECRET_NAME = os.environ["SLACK_WEBHOOK_SECRET_NAME"]
_gcs_client = None
def _get_gcs_client():
global _gcs_client
if _gcs_client is None:
raw = secrets.get_secret_value(SecretId=GCS_SA_SECRET_NAME)["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:
today = now.strftime("%Y-%m-%d")
yesterday = (now - timedelta(days=1)).strftime("%Y-%m-%d")
contents = []
for date_prefix in [today, yesterday]:
resp = client.list_objects_v2(Bucket=bucket, Prefix=f"archive/{date_prefix}/")
contents.extend(resp.get("Contents", []))
if not contents:
return False, f"{label}: No objects found under archive/ for last 2 days"
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)
now = datetime.now(timezone.utc)
today = now.strftime("%Y-%m-%d")
yesterday = (now - timedelta(days=1)).strftime("%Y-%m-%d")
blobs = []
for date_prefix in [today, yesterday]:
blobs.extend(list(bucket.list_blobs(prefix=f"archive/{date_prefix}/")))
if not blobs:
return False, "GCS Offsite: No objects found under archive/ for last 2 days"
cutoff = now - timedelta(hours=48)
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})"
if latest.size < 1_000_000:
return False, f"GCS Offsite: Latest dump suspiciously small ({latest.size} bytes)"
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 and s.get("State") == "completed"]
if not recent:
pending = sum(1 for s in snapshots if s["StartTime"] >= cutoff and s.get("State") == "pending")
errored = sum(1 for s in snapshots if s["StartTime"] >= cutoff and s.get("State") == "error")
latest = max(snapshots, key=lambda s: s["StartTime"])
age = (now - latest["StartTime"]).total_seconds() / 3600
return False, (
f"EBS Snapshots: No completed snapshot in last 48h "
f"(latest {age:.0f}h old, state={latest.get('State')}; "
f"pending={pending}, error={errored})"
)
return True, f"EBS Snapshots: OK — {len(recent)} completed in last 48h"
except Exception as e:
return False, f"EBS Snapshots: Error — {e}"
def _restore_test():
results = []
try:
now = datetime.now(timezone.utc)
contents = []
for days_ago in range(7):
date_prefix = (now - timedelta(days=days_ago)).strftime("%Y-%m-%d")
resp = s3.list_objects_v2(Bucket=SOURCE_BUCKET, Prefix=f"archive/{date_prefix}/")
contents.extend(resp.get("Contents", []))
if not contents:
return [{"pass": False, "msg": "Restore test: No dumps found in source bucket (last 7 days)"}]
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"})
sqlite_entries = [n for n in names if n.endswith(".sqlite3")]
sql_entries = [n for n in names if n.endswith(".sql")]
if sqlite_entries:
import sqlite3 as sqlite_mod
tf.extract(sqlite_entries[0], path=tmpdir, filter="data")
db_path = os.path.join(tmpdir, sqlite_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]}"})
elif sql_entries:
import sqlite3 as sqlite_mod
tf.extract(sql_entries[0], path=tmpdir, filter="data")
sql_path = os.path.join(tmpdir, sql_entries[0])
with open(sql_path, "r") as f:
sql_text = f.read()
if len(sql_text) < 100:
results.append({"pass": False, "msg": f"Restore test: SQL dump suspiciously small ({len(sql_text)} bytes)"})
else:
db_path = os.path.join(tmpdir, "restore-test.db")
conn = sqlite_mod.connect(db_path)
conn.executescript(sql_text)
result = conn.execute("PRAGMA integrity_check").fetchone()
conn.close()
if result[0] == "ok":
results.append({"pass": True, "msg": "Restore test: SQL dump import + integrity OK"})
else:
results.append({"pass": False, "msg": f"Restore test: Integrity FAILED after SQL import — {result[0]}"})
else:
results.append({"pass": False, "msg": "Restore test: No database file found in archive"})
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_NAME)["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, timeout=10)
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}"}})
else:
return {
"statusCode": 400,
"body": json.dumps({
"mode": mode,
"error": f"Unsupported backup verification mode: {mode}",
}),
}
_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,
}),
}