forgejo/lambda/backup-verification/app.py

247 lines
11 KiB
Python
Raw Normal View History

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:
2026-05-13 22:31:23 +00:00
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,
}),
}