forgejo/lambda/backup-verification/app.py
Adam Moussa 5ed1db788e
Add 3-2-1 backup strategy with cross-region and GCS offsite (#5)
* 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.

* Replace hardcoded instance ID in README with CloudFormation lookup

The instance ID changes on every instance replacement (version
upgrades, stack updates). Using a dynamic query prevents stale
references and removes a manual update step from the deploy process.

* Read backup S3 prefix from SSM parameter at runtime

Adds /forgejo/backup-s3-prefix SSM parameter (value: archive)
and updates the backup script to fetch it instead of hardcoding
the prefix. Eliminates the manual post-deploy sed step.

* Address cross-review findings for backup verification

- Add size guard before downloading dump in restore test (3.5 GB cap)
- Use paginator for list_objects_v2 in S3 checks and restore test
- Remove unnecessary overrideLogicalId on GcsTransferCredentials secret
- Pass explicit { mode: "daily" } to daily EventBridge rule target
- Add fallback for SSM parameter fetch in backup script
- Export replica bucket ARN/name from replica stack, consume via props

* Add CloudWatch alarm for backup verification Lambda errors

Fires on any Lambda error and on missing data (missed schedule).
Catches silent failures where the Slack notification never fires.

* Add .env to .gitignore

Required by org CI conventions check.

---------

Co-authored-by: Cursor Agent <cursoragent@cursor.com>
2026-05-15 17:59:31 -04:00

251 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_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 = []
paginator = client.get_paginator("list_objects_v2")
for date_prefix in [today, yesterday]:
for page in paginator.paginate(Bucket=bucket, Prefix=f"archive/{date_prefix}/"):
contents.extend(page.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 = []
paginator = s3.get_paginator("list_objects_v2")
for days_ago in range(7):
date_prefix = (now - timedelta(days=days_ago)).strftime("%Y-%m-%d")
for page in paginator.paginate(Bucket=SOURCE_BUCKET, Prefix=f"archive/{date_prefix}/"):
contents.extend(page.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"])
max_bytes = 4 * 1024 * 1024 * 1024 - 512 * 1024 * 1024
if latest["Size"] > max_bytes:
return [{"pass": False, "msg": f"Restore test: Dump too large for ephemeral storage ({latest['Size'] / 1_000_000_000:.1f} GB)"}]
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,
}),
}