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 }