mirror of
https://github.com/Sea-Haven-Industries/forgejo.git
synced 2026-09-30 13:23:10 +00:00
199 lines
8.5 KiB
Python
199 lines
8.5 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:
|
||
|
|
resp = client.list_objects_v2(Bucket=bucket, Prefix="archive/", MaxKeys=1000)
|
||
|
|
contents = resp.get("Contents", [])
|
||
|
|
if not contents:
|
||
|
|
return False, f"{label}: No objects found under archive/"
|
||
|
|
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)
|
||
|
|
blobs = list(bucket.list_blobs(prefix="archive/", max_results=1000))
|
||
|
|
if not blobs:
|
||
|
|
return False, "GCS Offsite: No objects found under archive/"
|
||
|
|
now = datetime.now(timezone.utc)
|
||
|
|
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:
|
||
|
|
resp = s3.list_objects_v2(Bucket=SOURCE_BUCKET, Prefix="archive/", MaxKeys=1000)
|
||
|
|
contents = resp.get("Contents", [])
|
||
|
|
if not contents:
|
||
|
|
return [{"pass": False, "msg": "Restore test: No dumps found in source bucket"}]
|
||
|
|
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(".db") or n.endswith("forgejo.db")]
|
||
|
|
if db_entries:
|
||
|
|
import sqlite3 as sqlite_mod
|
||
|
|
tf.extract(db_entries[0], path=tmpdir)
|
||
|
|
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": True, "msg": "Restore test: No .db file found in archive (may use different format)"})
|
||
|
|
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)
|
||
|
|
|
||
|
|
|
||
|
|
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}"}})
|
||
|
|
|
||
|
|
_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,
|
||
|
|
}),
|
||
|
|
}
|