cluster

Infrastructure files for Nordgedanken and Midnightthoughts.
git clone git://archive.git.mtrnord.blog/MTRNord/cluster.git
Log | Files | Refs | README

commit 44402263cdb0dd8433674a9b6daffc639e65a1fc
parent 94a351a2b7bcb8e0d837d83baee2ecb0d762efd0
Author: MTRNord <MTRNord@users.noreply.github.com>
Date:   Sun, 29 Mar 2026 21:19:56 +0200

Do better matrix backups

Diffstat:
M.gitleaksignore | 1+
Mapps/talos_cluster/kustomization.yaml | 1+
Aapps/talos_cluster/matrix-backup/backup-cronjob.yaml | 129+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aapps/talos_cluster/matrix-backup/backup-script-configmap.yaml | 718+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aapps/talos_cluster/matrix-backup/kustomization.yaml | 8++++++++
Aapps/talos_cluster/matrix-backup/namespace.yaml | 8++++++++
Aapps/talos_cluster/matrix-backup/secret.yaml | 34++++++++++++++++++++++++++++++++++
7 files changed, 899 insertions(+), 0 deletions(-)

diff --git a/.gitleaksignore b/.gitleaksignore @@ -1,3 +1,4 @@ apps/talos_cluster/monitoring-stack/dashboards/connectivity-tester-dashboard.json:generic-api-key:328 apps/talos_cluster/blog/docker/wp-cache-config.php:generic-api-key:9 apps/talos_cluster/cgit/deployment.yaml:generic-api-key:51 +apps/talos_cluster/matrix-backup/backup-script-configmap.yaml:generic-api-key:272 diff --git a/apps/talos_cluster/kustomization.yaml b/apps/talos_cluster/kustomization.yaml @@ -42,4 +42,5 @@ resources: - ./zot - ./image-builder - ./image-policies + - ./matrix-backup #- ./proxmox-ccm diff --git a/apps/talos_cluster/matrix-backup/backup-cronjob.yaml b/apps/talos_cluster/matrix-backup/backup-cronjob.yaml @@ -0,0 +1,129 @@ +apiVersion: batch/v1 +kind: CronJob +metadata: + name: matrix-backup + namespace: matrix-backup +spec: + schedule: "0 * * * *" + concurrencyPolicy: Forbid + activeDeadlineSeconds: 3600 + successfulJobsHistoryLimit: 3 + failedJobsHistoryLimit: 3 + jobTemplate: + spec: + template: + spec: + automountServiceAccountToken: false + enableServiceLinks: false + restartPolicy: OnFailure + securityContext: + runAsNonRoot: true + runAsUser: 1000 + runAsGroup: 1000 + fsGroup: 1000 + seccompProfile: + type: RuntimeDefault + initContainers: + - name: install-deps + image: python:3.12-slim + command: + - /bin/sh + - -c + - pip install --quiet --target /deps 'matrix-nio[e2e]' boto3 cryptography + securityContext: + allowPrivilegeEscalation: false + readOnlyRootFilesystem: false + capabilities: + drop: + - ALL + resources: + requests: + cpu: 200m + memory: 256Mi + limits: + memory: 512Mi + volumeMounts: + - name: deps + mountPath: /deps + containers: + - name: backup + image: python:3.12-slim + command: + - /bin/sh + - -c + - PYTHONPATH=/deps python /scripts/backup.py + securityContext: + allowPrivilegeEscalation: false + readOnlyRootFilesystem: true + capabilities: + drop: + - ALL + resources: + requests: + cpu: 100m + memory: 256Mi + limits: + memory: 512Mi + env: + - name: HOMESERVER + value: "https://matrix.mtrnord.blog" + - name: S3_ENDPOINT + value: "hel1.your-objectstorage.com" + - name: S3_BUCKET + value: "midnightthoughts-matrix-backup" + - name: MTRNORD_PASSWORD + valueFrom: + secretKeyRef: + name: matrix-backup-secret + key: mtrnord_password + - name: MTRNORD_SSSS_KEY + valueFrom: + secretKeyRef: + name: matrix-backup-secret + key: mtrnord_ssss_recovery_key + - name: LEXI_PASSWORD + valueFrom: + secretKeyRef: + name: matrix-backup-secret + key: lexi_password + - name: LEXI_SSSS_KEY + valueFrom: + secretKeyRef: + name: matrix-backup-secret + key: lexi_ssss_recovery_key + - name: S3_ACCESS_KEY + valueFrom: + secretKeyRef: + name: matrix-backup-secret + key: s3_access_key + - name: S3_SECRET_KEY + valueFrom: + secretKeyRef: + name: matrix-backup-secret + key: s3_secret_key + - name: KEY_EXPORT_PASSPHRASE + valueFrom: + secretKeyRef: + name: matrix-backup-secret + key: key_export_passphrase + volumeMounts: + - name: scripts + mountPath: /scripts + readOnly: true + - name: crypto-store + mountPath: /data/crypto + - name: deps + mountPath: /deps + readOnly: true + - name: tmp + mountPath: /tmp + volumes: + - name: scripts + configMap: + name: matrix-backup-script + - name: crypto-store + emptyDir: {} + - name: deps + emptyDir: {} + - name: tmp + emptyDir: {} diff --git a/apps/talos_cluster/matrix-backup/backup-script-configmap.yaml b/apps/talos_cluster/matrix-backup/backup-script-configmap.yaml @@ -0,0 +1,718 @@ +apiVersion: v1 +kind: ConfigMap +metadata: + name: matrix-backup-script + namespace: matrix-backup +data: + backup.py: | + #!/usr/bin/env python3 + """ + Matrix backup script. + + Auth lifecycle: + - First run: login with password → creates a dedicated "matrix-backup" device. + session.json (access_token + device_id) is saved into the crypto store. + - Subsequent runs: restore_login from session.json. + - The password in the secret is ONLY used on first run. + + SSSS bootstrap (first run only): + - Decrypts the megolm key backup private key from Secure Secret Storage using + the provided recovery key. + - Fetches all backed-up room keys from the server. + - Decrypts them with the backup private key (olm PK) and imports them into + the local crypto store so the bot can decrypt E2EE messages. + + Per-run backup: + - Crypto store round-tripped via S3 tarball (download at start, upload at end). + - Room list with IDs, aliases, server hints, names, types. + - E2EE exported key file (encrypted with KEY_EXPORT_PASSPHRASE). + - Incremental message history as JSONL. + - Images sent by *:mtrnord.blog users (downloaded and stored to S3). + - Display-name and avatar changes for *:mtrnord.blog users. + """ + import asyncio + import base64 + import hashlib + import hmac as hmac_lib + import io + import json + import mimetypes + import os + import re + import struct + import sys + import tarfile + import tempfile + from collections import Counter + from datetime import datetime, timezone + + import boto3 + from botocore.exceptions import ClientError + from cryptography.hazmat.backends import default_backend + from cryptography.hazmat.primitives import hashes + from cryptography.hazmat.primitives.ciphers import Cipher, algorithms, modes + from cryptography.hazmat.primitives.kdf.hkdf import HKDF + from cryptography.hazmat.primitives.kdf.pbkdf2 import PBKDF2HMAC + from nio import ( + AsyncClient, + AsyncClientConfig, + DownloadError, + LoginError, + MessageDirection, + RoomMessagesError, + ) + + # Environment + HOMESERVER = os.environ["HOMESERVER"] + S3_ENDPOINT = os.environ["S3_ENDPOINT"] + S3_BUCKET = os.environ["S3_BUCKET"] + S3_ACCESS_KEY = os.environ["S3_ACCESS_KEY"] + S3_SECRET_KEY = os.environ["S3_SECRET_KEY"] + KEY_EXPORT_PASSPHRASE = os.environ["KEY_EXPORT_PASSPHRASE"] + MTRNORD_PASSWORD = os.environ["MTRNORD_PASSWORD"] + MTRNORD_SSSS_KEY = os.environ["MTRNORD_SSSS_KEY"] + LEXI_PASSWORD = os.environ["LEXI_PASSWORD"] + LEXI_SSSS_KEY = os.environ["LEXI_SSSS_KEY"] + + OUR_HOMESERVER_SUFFIX = ":mtrnord.blog" + DATE_STR = datetime.now(timezone.utc).strftime("%Y%m%d") + + ACCOUNTS = [ + { + "user_id": "@mtrnord:mtrnord.blog", + "password": MTRNORD_PASSWORD, + "ssss_key": MTRNORD_SSSS_KEY, + "store": "/data/crypto/mtrnord", + "prefix": "mtrnord", + }, + { + "user_id": "@lexi:mtrnord.blog", + "password": LEXI_PASSWORD, + "ssss_key": LEXI_SSSS_KEY, + "store": "/data/crypto/lexi", + "prefix": "lexi", + }, + ] + + # S3 helpers + s3 = boto3.client( + "s3", + endpoint_url=f"https://{S3_ENDPOINT}", + aws_access_key_id=S3_ACCESS_KEY, + aws_secret_access_key=S3_SECRET_KEY, + region_name="hel1", + ) + + + def s3_get_json(key): + try: + obj = s3.get_object(Bucket=S3_BUCKET, Key=key) + return json.loads(obj["Body"].read()) + except ClientError as e: + if e.response["Error"]["Code"] in ("NoSuchKey", "404"): + return None + raise + + + def s3_get_bytes(key): + try: + obj = s3.get_object(Bucket=S3_BUCKET, Key=key) + return obj["Body"].read() + except ClientError as e: + if e.response["Error"]["Code"] in ("NoSuchKey", "404"): + return None + raise + + + def s3_put(key, data: bytes, content_type: str = "application/octet-stream"): + s3.put_object(Bucket=S3_BUCKET, Key=key, Body=data, ContentType=content_type) + + + def s3_exists(key): + try: + s3.head_object(Bucket=S3_BUCKET, Key=key) + return True + except ClientError as e: + if e.response["Error"]["Code"] in ("404", "NoSuchKey"): + return False + raise + + + # Crypto store S3 round-trip + def download_crypto_store(store_path: str, s3_key: str): + data = s3_get_bytes(s3_key) + if not data: + print(f" No existing crypto store at {s3_key} — starting fresh.") + return + os.makedirs(store_path, exist_ok=True) + with tarfile.open(fileobj=io.BytesIO(data), mode="r:gz") as tar: + tar.extractall(store_path) + print(f" Crypto store restored from S3 ({len(data)} bytes).") + + + def upload_crypto_store(store_path: str, s3_key: str): + buf = io.BytesIO() + with tarfile.open(fileobj=buf, mode="w:gz") as tar: + tar.add(store_path, arcname=".") + data = buf.getvalue() + s3_put(s3_key, data) + print(f" Crypto store saved to S3 ({len(data)} bytes).") + + + # Session management (password login, once only) + def load_session(store_path: str): + """Return stored session dict or None if not yet created.""" + session_file = os.path.join(store_path, "session.json") + if not os.path.exists(session_file): + return None + try: + with open(session_file) as f: + return json.load(f) + except Exception: + return None + + + def save_session(store_path: str, session: dict): + os.makedirs(store_path, exist_ok=True) + with open(os.path.join(store_path, "session.json"), "w") as f: + json.dump(session, f) + + + async def ensure_session(account: dict) -> dict: + """ + Return a session dict with access_token, device_id, ssss_bootstrapped. + Performs a password login on first run. + """ + session = load_session(account["store"]) + if session: + return session + + print(f" First run: logging in with password for {account['user_id']} ...") + tmp_client = AsyncClient(HOMESERVER, account["user_id"]) + try: + resp = await tmp_client.login( + password=account["password"], + device_name="matrix-backup", + ) + if isinstance(resp, LoginError): + raise RuntimeError(f"Login failed: {resp}") + finally: + await tmp_client.close() + + session = { + "access_token": resp.access_token, + "device_id": resp.device_id, + "ssss_bootstrapped": False, + } + save_session(account["store"], session) + print(f" Logged in. Device ID: {resp.device_id}") + return session + + + # SSSS recovery key decoding + _BASE58 = "123456789ABCDEFGHJKLMNPQRSTUVWXYZabcdefghijkmnopqrstuvwxyz" + + + def decode_recovery_key(key_str: str) -> bytes: + """ + Decode a Matrix recovery/security key to its 32 raw key bytes. + Accepts keys with or without spaces/dashes. + """ + cleaned = key_str.replace(" ", "").replace("-", "") + n = 0 + for ch in cleaned: + n = n * 58 + _BASE58.index(ch) + raw = n.to_bytes(35, "big") + if raw[0] != 0x8B or raw[1] != 0x01: + raise ValueError("Invalid recovery key prefix") + parity = 0 + for b in raw[:34]: + parity ^= b + if parity != raw[34]: + raise ValueError("Invalid recovery key parity byte") + return raw[2:34] # 32-byte key + + + # SSSS AES-HMAC-SHA2 decryption + def _ssss_derive(raw_key: bytes, secret_name: str): + """Derive (aes_key, mac_key) for a given secret name via HKDF-SHA256.""" + derived = HKDF( + algorithm=hashes.SHA256(), + length=64, + salt=b"\x00" * 32, + info=secret_name.encode(), + backend=default_backend(), + ).derive(raw_key) + return derived[:32], derived[32:] + + + def ssss_decrypt(encrypted: dict, raw_key: bytes, secret_name: str) -> bytes: + """Decrypt one entry from the 'encrypted' map of an SSSS account-data event.""" + aes_key, mac_key = _ssss_derive(raw_key, secret_name) + ciphertext = base64.b64decode(encrypted["ciphertext"]) + iv = base64.b64decode(encrypted["iv"]) + stored_mac = base64.b64decode(encrypted["mac"]) + + expected_mac = hmac_lib.new(mac_key, ciphertext, hashlib.sha256).digest() + if not hmac_lib.compare_digest(expected_mac, stored_mac): + raise ValueError("SSSS MAC mismatch — wrong recovery key?") + + cipher = Cipher(algorithms.AES(aes_key), modes.CTR(iv), backend=default_backend()) + dec = cipher.decryptor() + return dec.update(ciphertext) + dec.finalize() + + + # Megolm export file creation (for client.import_keys) + def build_megolm_export(sessions: list, passphrase: str = "_backup_import_") -> bytes: + """ + Encode a list of session dicts into the standard Megolm key export format + so they can be loaded with AsyncClient.import_keys(). + + sessions: [{"algorithm","room_id","sender_key","session_id","session_key", + "sender_claimed_keys","forwarding_curve25519_key_chain"}] # gitleaks:allow + """ + salt = os.urandom(16) + iv_bytes = bytearray(os.urandom(16)) + iv_bytes[0] &= 0x7F # ensure highest bit is 0 as spec requires + iv = bytes(iv_bytes) + iterations = 100_000 + + kdf = PBKDF2HMAC( + algorithm=hashes.SHA512(), + length=64, + salt=salt, + iterations=iterations, + backend=default_backend(), + ) + key_material = kdf.derive(passphrase.encode()) + aes_key = key_material[:32] + hmac_key = key_material[32:] + + plaintext = json.dumps(sessions).encode() + cipher = Cipher(algorithms.AES(aes_key), modes.CTR(iv), backend=default_backend()) + encryptor = cipher.encryptor() + ciphertext = encryptor.update(plaintext) + encryptor.finalize() + + payload = b"\x01" + salt + iv + struct.pack(">I", iterations) + ciphertext + mac = hmac_lib.new(hmac_key, payload, hashlib.sha256).digest() + + encoded = base64.b64encode(payload + mac).decode() + lines = "\n".join(encoded[i : i + 76] for i in range(0, len(encoded), 76)) + return ( + "-----BEGIN MEGOLM SESSION DATA-----\n" + + lines + + "\n-----END MEGOLM SESSION DATA-----\n" + ).encode() + + + # SSSS bootstrap: fetch and import key backup + async def bootstrap_crypto_from_ssss(client: AsyncClient, recovery_key_str: str): + """ + Decrypt the megolm backup private key from SSSS, fetch all backed-up room + keys, decrypt them, and import them into the local crypto store. + Called once per account on first login. + """ + print(" Bootstrapping crypto from SSSS key backup...") + + raw_key = decode_recovery_key(recovery_key_str) + + # 1. Find the default SSSS key ID + default_key_event = client.account_data.get("m.secret_storage.default_key") + if not default_key_event: + print(" Warning: no m.secret_storage.default_key in account data — skipping SSSS bootstrap") + return + key_id = getattr(default_key_event, "content", {}).get("key") + if not key_id: + print(" Warning: m.secret_storage.default_key has no 'key' field — skipping") + return + print(f" SSSS key ID: {key_id}") + + # 2. Decrypt the megolm backup private key from SSSS + backup_secret_event = client.account_data.get("m.megolm_backup.v1") + if not backup_secret_event: + print(" Warning: no m.megolm_backup.v1 account data — skipping SSSS bootstrap") + return + encrypted_map = getattr(backup_secret_event, "content", {}).get("encrypted", {}) + if key_id not in encrypted_map: + print(f" Warning: key_id {key_id} not in backup secret encrypted map — skipping") + return + + backup_key_bytes = ssss_decrypt(encrypted_map[key_id], raw_key, "m.megolm_backup.v1") + backup_private_key_b64 = backup_key_bytes.decode().strip() + print(" Decrypted megolm backup private key from SSSS.") + + # 3. Fetch all backed-up sessions from the server + backup_version_resp = await client.get_backup_keys_version() + if isinstance(backup_version_resp, Exception) or not hasattr(backup_version_resp, "version"): + # Fall back to direct HTTP call + import urllib.request + url = f"{HOMESERVER}/_matrix/client/v3/room_keys/version" + req = urllib.request.Request( + url, headers={"Authorization": f"Bearer {client.access_token}"} + ) + with urllib.request.urlopen(req, timeout=15) as r: + backup_info = json.loads(r.read()) + backup_version = backup_info["version"] + else: + backup_version = backup_version_resp.version + print(f" Key backup version: {backup_version}") + + import urllib.request + url = f"{HOMESERVER}/_matrix/client/v3/room_keys/keys?version={backup_version}" + req = urllib.request.Request( + url, headers={"Authorization": f"Bearer {client.access_token}"} + ) + with urllib.request.urlopen(req, timeout=60) as r: + backup_data = json.loads(r.read()) + rooms = backup_data.get("rooms", {}) + print(f" Fetched key backup: {len(rooms)} rooms") + + # 4. Decrypt each session using olm PK decryption + from olm.pk import PkDecryption, PkMessage + + private_key_bytes = base64.b64decode(backup_private_key_b64) + try: + pk_dec = PkDecryption.from_private_key(private_key_bytes) + except AttributeError: + # older python-olm: construct object and inject private key + pk_dec = PkDecryption() + pk_dec._private_key = private_key_bytes + + sessions_to_import = [] + for room_id, room_data in rooms.items(): + for session_id, session_info in room_data.get("sessions", {}).items(): + sd = session_info.get("session_data", {}) + try: + msg = PkMessage( + ephemeral_key=sd["ephemeral"], + mac=sd["mac"], + ciphertext=sd["ciphertext"], + ) + plaintext = pk_dec.decrypt(msg) + session_obj = json.loads(plaintext) + sessions_to_import.append({ + "algorithm": "m.megolm.v1.aes-sha2", + "room_id": room_id, + "session_id": session_id, + "session_key": session_obj["session_key"], + "sender_key": session_obj.get("sender_key", ""), + "sender_claimed_keys": session_obj.get("sender_claimed_keys", {}), + "forwarding_curve25519_key_chain": + session_obj.get("forwarding_curve25519_key_chain", []), + }) + except Exception as e: + print(f" Warning: failed to decrypt session {session_id} in {room_id}: {e}") + + print(f" Decrypted {len(sessions_to_import)} sessions; importing...") + + if sessions_to_import: + tmp = tempfile.mktemp(suffix=".megolm") + try: + with open(tmp, "wb") as f: + f.write(build_megolm_export(sessions_to_import)) + await client.import_keys(tmp, "_backup_import_") + finally: + if os.path.exists(tmp): + os.unlink(tmp) + print(f" Imported {len(sessions_to_import)} room sessions from key backup.") + + + # Room helpers + def safe_room_key(room_id: str) -> str: + return room_id.replace("/", "_").replace(":", "_") + + + def derive_server_hints(room) -> list: + counter = Counter() + for user_id in room.joined_members: + if ":" in user_id: + counter[user_id.split(":", 1)[1]] += 1 + return [s for s, _ in counter.most_common(5)] + + + def get_aliases(room) -> list: + aliases = [] + if room.canonical_alias: + aliases.append(room.canonical_alias) + if hasattr(room, "alt_aliases"): + for a in room.alt_aliases: + if a not in aliases: + aliases.append(a) + if hasattr(room, "aliases") and isinstance(room.aliases, dict): + for server_aliases in room.aliases.values(): + for a in server_aliases: + if a not in aliases: + aliases.append(a) + return aliases + + + async def get_room_type(client, room_id: str, room, dm_rooms: set) -> str: + if room_id in dm_rooms: + return "dm" + try: + resp = await client.room_get_state_event(room_id, "m.room.create", "") + if hasattr(resp, "content") and resp.content.get("type") == "m.space": + return "space" + except Exception: + pass + return "normal" + + + # Media helpers + _MXC_RE = re.compile(r"^mxc://([^/]+)/(.+)$") + + + def parse_mxc(mxc_url: str): + m = _MXC_RE.match(mxc_url or "") + return (m.group(1), m.group(2)) if m else (None, None) + + + def ext_from_content_type(ct: str) -> str: + ext = mimetypes.guess_extension(ct.split(";")[0].strip()) or "" + return {".jpe": ".jpg", ".jpeg": ".jpg"}.get(ext, ext) + + + async def store_media(client, mxc_url: str, prefix: str, label: str = "media"): + server, media_id = parse_mxc(mxc_url) + if not server or not media_id: + return + base_key = f"{prefix}/{label}/{server}/{media_id}" + if s3_exists(base_key): + return + resp = await client.download(server_name=server, media_id=media_id) + if isinstance(resp, DownloadError): + print(f" Warning: failed to download {mxc_url}") + return + ext = ext_from_content_type(resp.content_type or "") + key = base_key + ext if ext else base_key + s3_put(key, resp.body, resp.content_type or "application/octet-stream") + print(f" Stored {key} ({len(resp.body)} bytes)") + + + # Profile change tracking + async def record_profile_update(client, src: dict, prefix: str): + state_key = src.get("state_key", "") + if not state_key.endswith(OUR_HOMESERVER_SUFFIX): + return + content = src.get("content", {}) + prev_content = src.get("unsigned", {}).get("prev_content", {}) + new_name = content.get("displayname") + prev_name = prev_content.get("displayname") + new_avatar = content.get("avatar_url") + prev_avatar = prev_content.get("avatar_url") + if new_name == prev_name and new_avatar == prev_avatar: + return + record = { + "ts": src.get("origin_server_ts"), + "event_id": src.get("event_id"), + "user_id": state_key, + "room_id": src.get("room_id"), + } + if new_name != prev_name: + record["displayname_old"] = prev_name + record["displayname_new"] = new_name + if new_avatar != prev_avatar: + record["avatar_old"] = prev_avatar + record["avatar_new"] = new_avatar + if new_avatar: + await store_media(client, new_avatar, prefix, label="avatars") + log_key = f"{prefix}/profile-updates.jsonl" + existing = s3_get_bytes(log_key) or b"" + s3_put(log_key, existing + json.dumps(record).encode() + b"\n", "application/x-ndjson") + + + # History pagination + async def paginate_history(client, room_id: str, prefix: str): + room_key = safe_room_key(room_id) + cursor_s3 = f"{prefix}/history-cursor/{room_key}.json" + history_s3 = f"{prefix}/history/{room_key}/{DATE_STR}.jsonl" + + cursor_data = s3_get_json(cursor_s3) + + if not cursor_data: + room = client.rooms.get(room_id) + token = getattr(room, "prev_batch", None) + if token: + s3_put(cursor_s3, json.dumps({"token": token}).encode(), "application/json") + print(f" History cursor initialised for {room_id}") + return + + start_token = cursor_data["token"] + messages = [] + next_token = start_token + + for _ in range(100): + resp = await client.room_messages( + room_id, + start=next_token, + limit=100, + direction=MessageDirection.forward, + ) + if isinstance(resp, RoomMessagesError): + print(f" Warning: room_messages error for {room_id}: {resp}") + break + if not resp.chunk: + break + for event in resp.chunk: + src = event.source if hasattr(event, "source") else {} + content = src.get("content", {}) + messages.append({ + "event_id": src.get("event_id"), + "sender": src.get("sender"), + "type": src.get("type"), + "origin_server_ts": src.get("origin_server_ts"), + "content": content, + }) + # Images sent by our users + if ( + src.get("type") == "m.room.message" + and content.get("msgtype") == "m.image" + and str(src.get("sender", "")).endswith(OUR_HOMESERVER_SUFFIX) + ): + mxc = content.get("url", "") + if mxc: + await store_media(client, mxc, prefix, label="media") + # Profile changes + if src.get("type") == "m.room.member": + src["room_id"] = room_id + await record_profile_update(client, src, prefix) + new_end = getattr(resp, "end", None) + if not new_end or new_end == next_token: + next_token = None + break + next_token = new_end + + if messages: + existing = s3_get_bytes(history_s3) or b"" + chunk = b"\n".join(json.dumps(m).encode() for m in messages) + b"\n" + s3_put(history_s3, existing + chunk, "application/x-ndjson") + print(f" History: +{len(messages)} events for {room_id}") + + if next_token and next_token != start_token: + s3_put(cursor_s3, json.dumps({"token": next_token}).encode(), "application/json") + + + # Per-account backup + async def backup_account(account: dict): + prefix = account["prefix"] + print(f"\n=== Backing up {account['user_id']} ===") + + # 1. Restore crypto store from S3 (must happen before opening the client) + store_s3_key = f"{prefix}/crypto-store.tar.gz" + download_crypto_store(account["store"], store_s3_key) + + # 2. Ensure we have a session (password-login on first run, persisted inside the store) + session = await ensure_session(account) + + # 3. Open the client (only after store is ready and session is known) + config = AsyncClientConfig(store_sync_tokens=True) + client = AsyncClient( + HOMESERVER, + account["user_id"], + store_path=account["store"], + config=config, + ) + client.restore_login( + user_id=account["user_id"], + device_id=session["device_id"], + access_token=session["access_token"], + ) + + try: + # Load sync token for incremental syncs + token_file = os.path.join(account["store"], "sync_token.json") + since = None + if os.path.exists(token_file): + try: + with open(token_file) as f: + since = json.load(f).get("next_batch") + except Exception: + pass + + print(f" Syncing (full_state=True, since={since!r})...") + sync_resp = await client.sync(full_state=True, since=since, timeout=60000) + with open(token_file, "w") as f: + json.dump({"next_batch": sync_resp.next_batch}, f) + print(f" Sync done. Rooms: {len(client.rooms)}") + + # 4. SSSS bootstrap on first login (after sync so account_data is populated) + if not session.get("ssss_bootstrapped"): + try: + await bootstrap_crypto_from_ssss(client, account["ssss_key"]) + session["ssss_bootstrapped"] = True + save_session(account["store"], session) + except Exception as e: + print(f" Warning: SSSS bootstrap failed: {e}", file=sys.stderr) + + # 5. DM room set + dm_map = {} + ad_event = client.account_data.get("m.direct") + if ad_event: + dm_map = getattr(ad_event, "content", {}) or {} + dm_rooms = set() + for room_ids in dm_map.values(): + if isinstance(room_ids, list): + dm_rooms.update(room_ids) + + # 6. Room list + print(" Building room list...") + rooms_data = [] + for room_id, room in client.rooms.items(): + room_type = await get_room_type(client, room_id, room, dm_rooms) + aliases = get_aliases(room) + via = derive_server_hints(room) + rooms_data.append({ + "room_id": room_id, + "name": room.display_name or room.name or room_id, + "type": room_type, + "aliases": aliases, + "canonical_alias": room.canonical_alias, + "via": via, + "member_count": room.member_count, + "encrypted": room.encrypted, + }) + print(f" [{room_type:6s}] {room.display_name or room_id}") + rooms_json = json.dumps(rooms_data, indent=2).encode() + s3_put(f"{prefix}/rooms-{DATE_STR}.json", rooms_json, "application/json") + s3_put(f"{prefix}/rooms-latest.json", rooms_json, "application/json") + print(f" Uploaded rooms ({len(rooms_data)} rooms)") + + # 7. E2EE key export + print(" Exporting E2EE keys...") + try: + key_file = tempfile.mktemp(suffix=".bin") + await client.export_keys(key_file, KEY_EXPORT_PASSPHRASE) + with open(key_file, "rb") as f: + key_data = f.read() + os.unlink(key_file) + s3_put(f"{prefix}/crypto-keys-{DATE_STR}.bin", key_data) + s3_put(f"{prefix}/crypto-keys-latest.bin", key_data) + print(f" Uploaded crypto-keys-latest.bin ({len(key_data)} bytes)") + except Exception as e: + print(f" Warning: key export failed: {e}", file=sys.stderr) + + # 8. Incremental history + media + profile updates + print(" Backing up history / media / profile updates...") + for room_id in list(client.rooms.keys()): + try: + await paginate_history(client, room_id, prefix) + except Exception as e: + print(f" Warning: history error for {room_id}: {e}", file=sys.stderr) + + finally: + # Close the client BEFORE uploading the store so nothing is mid-write + await client.close() + upload_crypto_store(account["store"], store_s3_key) + + print(f" Done: {account['user_id']}") + + + # Entry point + async def main(): + for account in ACCOUNTS: + await backup_account(account) + print("\nAll accounts backed up successfully.") + + + if __name__ == "__main__": + asyncio.run(main()) diff --git a/apps/talos_cluster/matrix-backup/kustomization.yaml b/apps/talos_cluster/matrix-backup/kustomization.yaml @@ -0,0 +1,8 @@ +apiVersion: kustomize.config.k8s.io/v1beta1 +kind: Kustomization +namespace: matrix-backup +resources: + - namespace.yaml + - secret.yaml + - backup-script-configmap.yaml + - backup-cronjob.yaml diff --git a/apps/talos_cluster/matrix-backup/namespace.yaml b/apps/talos_cluster/matrix-backup/namespace.yaml @@ -0,0 +1,8 @@ +apiVersion: v1 +kind: Namespace +metadata: + name: matrix-backup + labels: + pod-security.kubernetes.io/enforce: restricted + pod-security.kubernetes.io/audit: restricted + pod-security.kubernetes.io/warn: restricted diff --git a/apps/talos_cluster/matrix-backup/secret.yaml b/apps/talos_cluster/matrix-backup/secret.yaml @@ -0,0 +1,34 @@ +apiVersion: v1 +kind: Secret +metadata: + name: matrix-backup-secret + namespace: matrix-backup +type: Opaque +stringData: + mtrnord_password: ENC[AES256_GCM,data:rLhj2nEU49Vi,iv:spylZf2eUmJAz1sxYYs3FH7ogyvL/Uh0IRn3gIxEOh4=,tag:TdiNa8DLfSUThCoKqCw4tA==,type:str] + mtrnord_ssss_recovery_key: ENC[AES256_GCM,data:ETn7dgqZpDCaLk8JZN/8ALsoL59USF/QWhGpZcsyqDM6bkSjBNPPiXr7Sbs/5rnFQH4NY0MRlpdG4Yc=,iv:VUXCLS4mwMmmviI0ZrttaKFRpFGNDnF9dzCJ4RREfTc=,tag:yBmhc/aX5+U/Y4q5bmp3ng==,type:str] + lexi_password: ENC[AES256_GCM,data:QVU6xUkfHcys,iv:poksZTGryhSU6S7M5K4uUjhxflgT3AwOiThgG9rf9nk=,tag:vnEr1NtWa3PSR+cz+7WFwQ==,type:str] + lexi_ssss_recovery_key: ENC[AES256_GCM,data:/uK2cf5d65VrU4/rentCPHLMdTkvc5gVL1yOJmi4LUq/L4RReO2tgAq+ZBl9fPaJYVdsQRlSlOojWjU=,iv:Q7gurpu4xJmO8NNM7yXpHFYzu60IMYCugIRmJTXrmRI=,tag:K8+hXLX5yv4ZmRr0iOXxcg==,type:str] + s3_access_key: ENC[AES256_GCM,data:idq3pSR7TekY9Pcr4DLnLIK0Awk=,iv:Y3awxY8Jau+niPqhdYST/ekm75g6nES2HY8NOghr1gM=,tag:ymp64EtXsHv5hxmJJ2QMcQ==,type:str] + s3_secret_key: ENC[AES256_GCM,data:7U2jWXr87Rj6f3b80C5DZM9K1NSkIbzCXPN9ukecz2owhajeyuQ8hQ==,iv:r4XpfSjn2cAvOxbAcUbL8gB7+2CA7geioJhpdETcqac=,tag:b0qx8wJEQhFT41+JPiBm4w==,type:str] + key_export_passphrase: ENC[AES256_GCM,data:yB1+XS1YoiAAL23Mu2a/V0sI/fGOGN9vaKyH9ePw/yvA66J1QoRBXKhxr4rJCTkYhfzaCn4IoTLrFPxIvvzZ,iv:RV9KjtobYjsz7jVSvCaD7SOJW5Wxt9na9SPjArT67xg=,tag:WcAkwSgCzc5Mdb96OrkBfA==,type:str] +sops: + kms: [] + gcp_kms: [] + azure_kv: [] + hc_vault: [] + age: + - recipient: age1esjyg2qfy49awv0ptkzvpk425adczjr38m37w2mmcahzc4p8n54sll2nzh + enc: | + -----BEGIN AGE ENCRYPTED FILE----- + YWdlLWVuY3J5cHRpb24ub3JnL3YxCi0+IFgyNTUxOSBxUktNQmlTUk1RN2JicTla + RjVFb0V5amtIczVuRVhWTjFIQXdEUVJ3TVdzClJadDNEYStwQW02L3VkK0NuNmxr + TVRCdHRObXBBOUxPaHZ3ZnozbFIwVk0KLS0tIFZBVFRWOUtZaG9rTmZQbDBKVThr + RmxKbkZkanZkSTRyeWM4Rm9jcWhUWkUKzmGRfGP9LQT1ypQYb70+Sfls29m3B/Va + w2WG8BfC2L3Pt/eKalgcPdrqafZ0KWwn/YK6TFw2UDFaQhlq6CkAYw== + -----END AGE ENCRYPTED FILE----- + lastmodified: "2026-03-29T19:15:46Z" + mac: ENC[AES256_GCM,data:B0jvpT7VB0V9rkm1HEUHOLV9EcrVYDaP2Wzwjxg90QvsB5xkhPem6Wd3+jqNcWHnBsKG9XvE0skrzbGdK0ESBMeRUVBrmV9jCxWDPNRhkcjuRqvaxbyrr2C79tsEqWMkJmlUZy+1DsWQ/pyDourKHyt2nw1dZwOJUDexLXfiOsQ=,iv:fvQ+s3OpUrDitaOzJepQT0DoeZh2G8GZ5iOAi56Rz5U=,tag:TngETHbq43LZGX9oHiz4kg==,type:str] + pgp: [] + encrypted_regex: ^(apiKey|appUserPassword|otelUserPassword|harborAdminPassword|totpVaultKey|kimaiAppSecret|kimaiAdminPassword|GITHUB_CLIENT_ID|GITHUB_CLIENT_SECRET|GITHUB_PRIVATE_KEY|woosh|root_password|rspamd_password|pgdb_password|matrix_access_token|pgdb_remote_url|hmac_secret_key|adminPassword|adminEmail|jenkinsAdminEmail|securityRealm|gerrit.config|routing_key|DATABASE_URL|SMTP_PASSWORD|SECRET_KEY_BASE|admin_password|extraCommands|key|clickhouseDatabaseURL|databaseURL|client_id|client_secret|secret_key_base|otp_secret|private_key|public_key|primaryKey|deterministicKey|keyDerivationSalt|token|clientId|secretKey|installationId|installationKey|uriOverride|adminToken.value|password.value|sql_password|erlangCookie|AUTHENTICATION_PASSWORD|ROOM_API_SECRET_KEY|adminPassword|configPassword|adminUser|configUser|MAIL_PASSWORD|APP_KEY|api_key|api_secret|keys|livekit_key|livekit_secret|secret_key|admin_pass|admin_email|mariadbPassword|mariadbRootPassword|privateKey|data|stringData|PASSWD|password|pass|postgresPassword|smtp_auth_password|addresses|smtp_auth_username|authorization_credentials|postgresqlPassword|redminePassword|smtpPassword|registration_shared_secret|shared_secret|secret|admin_token|integrationKey|integration_key|rootPassword|adminPassword|adminUser|adminEmail|emailPassword|secretKey|appId|clientSecret|webhookSecret)$ + version: 3.9.1