mirror of
https://github.com/Sea-Haven-Industries/forgejo.git
synced 2026-09-30 18:53:14 +00:00
240 lines
11 KiB
Python
240 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 = []
|
|
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]
|
|
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_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,
|
|
}),
|
|
}
|