cluster

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

backup.go (5373B)


      1 // backup.go — per-account backup orchestration.
      2 package main
      3 
      4 import (
      5 	"context"
      6 	"encoding/json"
      7 	"log/slog"
      8 
      9 	"maunium.net/go/mautrix"
     10 	"maunium.net/go/mautrix/event"
     11 	"maunium.net/go/mautrix/id"
     12 )
     13 
     14 // backupAccount runs the full backup pipeline for a single Matrix account:
     15 // restore crypto state, sync, export E2EE keys, save room list, paginate
     16 // history, save account data, then persist state back to S3.
     17 func backupAccount(ctx context.Context, acc accountCfg) error {
     18 	slog.Info("=== Backing up account ===", "user_id", acc.UserID)
     19 
     20 	storeS3Key := acc.Prefix + "/crypto-store.tar.gz"
     21 	sessionS3Key := acc.Prefix + "/session.json"
     22 	syncTokenS3Key := acc.Prefix + "/sync-token.json"
     23 
     24 	if err := downloadStore(ctx, acc.StoreDir, storeS3Key); err != nil {
     25 		slog.Warn("Could not restore store", "error", err)
     26 	}
     27 
     28 	client, err := mautrix.NewClient(homeserver, "", "")
     29 	if err != nil {
     30 		return err
     31 	}
     32 
     33 	if _, err := ensureSession(ctx, client, acc, sessionS3Key); err != nil {
     34 		return err
     35 	}
     36 
     37 	var syncToken struct {
     38 		NextBatch string `json:"next_batch"`
     39 	}
     40 	_ = s3GetJSON(ctx, syncTokenS3Key, &syncToken)
     41 
     42 	slog.Info("Syncing", "since", syncToken.NextBatch)
     43 	syncResp, err := client.SyncRequest(ctx, 60000, syncToken.NextBatch, "", true, event.PresenceUnavailable)
     44 	if err != nil {
     45 		return err
     46 	}
     47 	if data, err := json.Marshal(map[string]string{"next_batch": syncResp.NextBatch}); err == nil {
     48 		if err := s3Put(ctx, syncTokenS3Key, data, "application/json"); err != nil {
     49 			slog.Warn("Failed to save sync token", "error", err)
     50 		}
     51 	}
     52 	slog.Info("Sync done", "rooms", len(syncResp.Rooms.Join))
     53 
     54 	var sessions megolmSessions
     55 	if acc.SSSSKey != "" {
     56 		sessions, err = fetchAndExportKeyBackup(ctx, client, acc.SSSSKey, acc.Prefix)
     57 		if err != nil {
     58 			slog.Warn("Key backup export failed", "error", err)
     59 		}
     60 	}
     61 
     62 	if err := saveRoomList(ctx, client, syncResp, acc.Prefix); err != nil {
     63 		slog.Warn("Room list failed", "error", err)
     64 	}
     65 
     66 	// Back up account-data, per-room account data, and profile snapshot.
     67 	backupAccountData(ctx, client, syncResp, acc.Prefix)
     68 	backupRoomAccountData(ctx, client, syncResp, acc.Prefix)
     69 	backupProfileSnapshot(ctx, client, acc.Prefix)
     70 
     71 	// Build the DM room set: m.direct API + 2-member heuristic from stored list.
     72 	dmRooms := getDMRooms(ctx, client, syncResp)
     73 	if raw, err := getDecryptedAgeFromS3(ctx, acc.Prefix+"/rooms-latest.json.age"); err == nil && raw != nil {
     74 		var storedRooms []roomEntry
     75 		if json.Unmarshal(raw, &storedRooms) == nil {
     76 			for _, r := range storedRooms {
     77 				if r.Type == "dm" {
     78 					dmRooms[id.RoomID(r.RoomID)] = true
     79 				}
     80 			}
     81 		}
     82 	}
     83 
     84 	// Paginate history for all rooms present in the sync delta.
     85 	for roomID, joinedRoom := range syncResp.Rooms.Join {
     86 		isDM := dmRooms[roomID]
     87 		if err := paginateRoom(ctx, client, roomID, acc.Prefix, isDM, joinedRoom.Timeline.PrevBatch, syncResp.NextBatch, sessions); err != nil {
     88 			slog.Warn("History error", "room_id", roomID, "error", err)
     89 		}
     90 	}
     91 
     92 	// Catch-up: DM rooms absent from this sync delta (no new events) may still
     93 	// need their history backfilled. Use the stored rooms list as the authoritative
     94 	// DM source, falling back to the API-based dmRooms map.
     95 	catchupDMs := make(map[id.RoomID]bool)
     96 	for roomID := range dmRooms {
     97 		catchupDMs[roomID] = true
     98 	}
     99 	if raw, err := getDecryptedAgeFromS3(ctx, acc.Prefix+"/rooms-latest.json.age"); err != nil {
    100 		slog.Warn("DM catch-up: failed to read rooms list", "error", err)
    101 	} else if raw == nil {
    102 		slog.Warn("DM catch-up: rooms list not found or ageIdentity not set")
    103 	} else {
    104 		var storedRooms []roomEntry
    105 		if err := json.Unmarshal(raw, &storedRooms); err != nil {
    106 			slog.Warn("DM catch-up: failed to parse rooms list", "error", err)
    107 		} else {
    108 			slog.Info("DM catch-up: loaded rooms list", "total_rooms", len(storedRooms))
    109 			for _, r := range storedRooms {
    110 				if r.Type == "dm" {
    111 					catchupDMs[id.RoomID(r.RoomID)] = true
    112 				}
    113 			}
    114 		}
    115 	}
    116 	slog.Info("DM catch-up: checking rooms", "total_dm_rooms", len(catchupDMs), "in_sync", len(syncResp.Rooms.Join))
    117 	for roomID := range catchupDMs {
    118 		if _, inSync := syncResp.Rooms.Join[roomID]; inSync {
    119 			continue // already handled above
    120 		}
    121 		safeKey := roomSafeKey(roomID)
    122 		dmHistoryDoneKey := acc.Prefix + "/history-cursor/" + safeKey + ".dm-done"
    123 		dmDone, dmDoneErr := s3Get(ctx, dmHistoryDoneKey)
    124 		if dmDoneErr != nil {
    125 			slog.Warn("DM catch-up: error checking dm-done marker", "room_id", roomID, "error", dmDoneErr)
    126 		}
    127 		if dmDone != nil {
    128 			continue // already backfilled
    129 		}
    130 		cursorKey := acc.Prefix + "/history-cursor/" + safeKey + ".json"
    131 		var cursor struct {
    132 			Token string `json:"token"`
    133 		}
    134 		if err := s3GetJSON(ctx, cursorKey, &cursor); err != nil {
    135 			slog.Warn("DM catch-up: error reading cursor", "room_id", roomID, "error", err)
    136 			continue
    137 		}
    138 		if cursor.Token == "" {
    139 			slog.Info("DM catch-up: no cursor yet, skipping", "room_id", roomID)
    140 			continue
    141 		}
    142 		slog.Info("DM catch-up: backfilling room", "room_id", roomID, "cursor", cursor.Token)
    143 		if err := paginateRoom(ctx, client, roomID, acc.Prefix, true, cursor.Token, syncResp.NextBatch, sessions); err != nil {
    144 			slog.Warn("DM catch-up error", "room_id", roomID, "error", err)
    145 		}
    146 	}
    147 
    148 	if err := uploadStore(ctx, acc.StoreDir, storeS3Key); err != nil {
    149 		slog.Warn("Could not upload store", "error", err)
    150 	}
    151 
    152 	slog.Info("Done", "user_id", acc.UserID)
    153 	return nil
    154 }