forgejo/lambda/backup-verification/app.py

214 lines
9.2 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_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"})
db_entries = [n for n in names if n.endswith(".sqlite3") or n.endswith(".sql")]
if db_entries:
import sqlite3 as sqlite_mod
tf.extract(db_entries[0], path=tmpdir, filter="data")
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": False, "msg": "Restore test: No SQLite DB 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:
raise ValueError(f"Unknown mode: {mode!r} (expected 'daily' or 'restore-test')")
_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,
}),
}