forgejo/lambda/backup-verification/app.py
Adam Moussa 0ab9cc9f30 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.
2026-05-13 18:54:04 -04:00

238 lines
11 KiB
Python

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:
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=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:
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_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, 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,
}),
}