diff --git a/docs/tasks/18-account-management.md b/docs/tasks/18-account-management.md index c8b49a1..84c6d21 100644 --- a/docs/tasks/18-account-management.md +++ b/docs/tasks/18-account-management.md @@ -104,22 +104,22 @@ External personal backups become an admin-only account-management feature. byte size, hash, source PDS, handle, and DID. - [x] T18-30: Reuse repo-core verification to validate commit DID, commit signature, MST completeness, record paths, record CIDs, and CAR integrity. -- [ ] T18-31: Extract blob references from repo records and merge them with +- [x] T18-31: Extract blob references from repo records and merge them with paginated `listBlobs` output. -- [ ] T18-32: Add bounded concurrent blob download with CID verification, +- [x] T18-32: Add bounded concurrent blob download with CID verification, missing-blob recording, and retry handling for transient source failures. -- [ ] T18-33: Add credentialed private preference backup through +- [x] T18-33: Add credentialed private preference backup through `app.bsky.actor.getPreferences`, with auth failures reported separately from public repo/blob backup status. -- [ ] T18-34: Write immutable snapshot manifests and verification reports, first +- [x] T18-34: Write immutable snapshot manifests and verification reports, first to a temporary workspace and then atomically mark snapshots complete. -- [ ] T18-35: Store personal backup snapshots through the existing local and +- [x] T18-35: Store personal backup snapshots through the existing local and S3/R2 backup storage shape. -- [ ] T18-36: Add retention policies: keep all, keep last N, keep for days, and +- [x] T18-36: Add retention policies: keep all, keep last N, keep for days, and pinned snapshots. -- [ ] T18-37: Add portable export bundle creation containing manifest, repo CAR, +- [x] T18-37: Add portable export bundle creation containing manifest, repo CAR, blobs, preferences JSON when present, and verification report. -- [ ] T18-38: Add offline snapshot verification that can run without contacting +- [x] T18-38: Add offline snapshot verification that can run without contacting the source PDS. - [ ] T18-39: Add tests proving a snapshot can be understood from its manifest and files without relying on Tempest database rows. diff --git a/lib/tempest/personal_backups.ex b/lib/tempest/personal_backups.ex index 542961f..3876d41 100644 --- a/lib/tempest/personal_backups.ex +++ b/lib/tempest/personal_backups.ex @@ -16,10 +16,15 @@ defmodule Tempest.PersonalBackups do SecretStore, Snapshot, SourceClient, + Storage, Verifier } alias Tempest.Repo + alias Tempest.RepoCore.{Cid, Drisl} + + @default_blob_concurrency 4 + @transient_retries 2 def list_accounts do Account @@ -114,6 +119,65 @@ defmodule Tempest.PersonalBackups do def verify_account_source(account_id) when is_integer(account_id), do: account_id |> get_account!() |> verify_account_source() + def update_retention_setting(%Account{} = account, attrs) when is_map(attrs) do + setting = account |> Repo.preload(:retention_setting, force: true) |> Map.fetch!(:retention_setting) + + setting + |> RetentionSetting.changeset(Map.put(stringify_keys(attrs), "account_id", account.id)) + |> Repo.update() + end + + def pin_snapshot(%Snapshot{} = snapshot, pinned? \\ true) when is_boolean(pinned?) do + snapshot + |> Snapshot.pin_changeset(%{pinned: pinned?}) + |> Repo.update() + end + + def prune_snapshots(%Account{} = account, opts \\ []) do + config = Keyword.get(opts, :config, Tempest.Config.load!()) + account = Repo.preload(account, [:retention_setting, :snapshots], force: true) + + account.retention_setting + |> snapshots_to_prune(account.snapshots, Keyword.get(opts, :now, DateTime.utc_now(:second))) + |> Enum.reduce_while({:ok, []}, fn snapshot, {:ok, pruned} -> + case delete_snapshot(snapshot, config) do + :ok -> {:cont, {:ok, [snapshot | pruned]}} + {:error, reason} -> {:halt, {:error, reason}} + end + end) + |> case do + {:ok, pruned} -> {:ok, Enum.reverse(pruned)} + {:error, reason} -> {:error, reason} + end + end + + def export_snapshot_bundle(%Snapshot{} = snapshot, opts \\ []) do + config = Keyword.get(opts, :config, Tempest.Config.load!()) + + target_path = + Keyword.get(opts, :path) || Path.join([config.data_dir, "tmp", Path.basename(snapshot.storage_key) <> ".zip"]) + + Storage.archive_snapshot(config, snapshot.storage_key, target_path) + end + + def verify_snapshot_offline(snapshot_or_dir, opts \\ []) + + def verify_snapshot_offline(%Snapshot{} = snapshot, opts) do + config = Keyword.get(opts, :config, Tempest.Config.load!()) + verify_snapshot_offline(Path.join(config.data_dir, snapshot.storage_key), opts) + end + + def verify_snapshot_offline(snapshot_dir, _opts) when is_binary(snapshot_dir) do + with {:ok, manifest} <- read_json_file(Path.join(snapshot_dir, "manifest.json")), + {:ok, report} <- read_json_file(Path.join(snapshot_dir, "verification.json")), + :ok <- verify_report_shape(report), + :ok <- verify_repo_file(snapshot_dir, manifest), + :ok <- verify_blob_files(snapshot_dir, manifest), + :ok <- verify_preferences_file(snapshot_dir, manifest) do + {:ok, %{status: "ok", manifest: manifest, report: report}} + end + end + def create_repo_snapshot(%Account{} = account, opts \\ []) do config = Keyword.get(opts, :config, Tempest.Config.load!()) @@ -123,11 +187,14 @@ defmodule Tempest.PersonalBackups do {:ok, run} <- create_run(verified_account), {:ok, repo_car} <- SourceClient.get_repo(verified_account.source_pds_url, verified_account.did, opts), {:ok, verified} <- Verifier.verify_repo_car(repo_car, did_document, verified_account.did), - {:ok, snapshot_attrs} <- write_repo_snapshot(config, verified_account, verified, repo_car), + {:ok, workspace} <- + prepare_snapshot_workspace(config, verified_account, did_document, verified, repo_car, opts), + {:ok, snapshot_attrs} <- Storage.finalize_snapshot(config, workspace), {:ok, snapshot} <- insert_snapshot(verified_account, run, snapshot_attrs), + {:ok, blobs} <- insert_blob_records(snapshot, workspace.blobs), {:ok, run} <- finish_run(run, snapshot), {:ok, account} <- mark_snapshot_success(verified_account, snapshot) do - %{account: account, run: run, snapshot: snapshot, verification: verified} + %{account: account, run: run, snapshot: %{snapshot | blobs: blobs}, verification: verified} else {:error, reason} -> Repo.rollback(reason) end @@ -292,24 +359,197 @@ defmodule Tempest.PersonalBackups do |> Repo.update() end - defp write_repo_snapshot(config, %Account{} = account, verified, repo_car) do - storage_key = snapshot_storage_key(account.did, verified.rev) - snapshot_dir = Path.join(config.data_dir, storage_key) - repo_car_path = Path.join(snapshot_dir, "repo.car") + defp snapshots_to_prune(%RetentionSetting{policy: "keep_all"}, _snapshots, _now), do: [] + + defp snapshots_to_prune(%RetentionSetting{policy: "keep_last_n", keep_last: keep_last}, snapshots, _now) do + snapshots + |> Enum.reject(& &1.pinned) + |> Enum.sort_by(&snapshot_sort_key/1, :desc) + |> Enum.drop(keep_last) + end + + defp snapshots_to_prune(%RetentionSetting{policy: "keep_for_days", keep_days: keep_days}, snapshots, now) + when is_integer(keep_days) do + cutoff = DateTime.add(now, -keep_days, :day) + + snapshots + |> Enum.reject(& &1.pinned) + |> Enum.filter(fn snapshot -> + case snapshot.completed_at || snapshot.inserted_at do + %DateTime{} = completed_at -> DateTime.compare(completed_at, cutoff) == :lt + _missing -> false + end + end) + end + + defp snapshots_to_prune(_setting, _snapshots, _now), do: [] + + defp snapshot_sort_key(snapshot), do: snapshot.completed_at || snapshot.inserted_at || ~U[1970-01-01 00:00:00Z] + + defp delete_snapshot(%Snapshot{} = snapshot, config) do + Repo.transaction(fn -> + with {:ok, _snapshot} <- Repo.delete(snapshot), + :ok <- Storage.delete_snapshot(config, snapshot.storage_key), + :ok <- maybe_clear_last_snapshot(snapshot) do + :ok + else + {:error, %Ecto.Changeset{} = changeset} -> Repo.rollback(changeset) + {:error, reason} -> Repo.rollback(reason) + end + end) + |> case do + {:ok, :ok} -> :ok + {:error, reason} -> {:error, reason} + end + end + + defp maybe_clear_last_snapshot(%Snapshot{} = snapshot) do + case Repo.get(Account, snapshot.account_id) do + %Account{last_snapshot_id: id} = account when id == snapshot.id -> + account + |> Account.registration_changeset(%{ + label: account.label, + did: account.did, + handle: account.handle, + source_pds_url: account.source_pds_url, + pinned_source_pds_url: account.pinned_source_pds_url, + credential_state: account.credential_state, + last_checked_at: account.last_checked_at, + last_success_at: account.last_success_at, + last_snapshot_id: nil, + status: account.status, + status_reason: account.status_reason + }) + |> Repo.update() + |> case do + {:ok, _account} -> :ok + {:error, %Ecto.Changeset{} = changeset} -> {:error, changeset} + end + + _account -> + :ok + end + end + + defp read_json_file(path) do + with {:ok, bytes} <- File.read(path), + {:ok, json} <- Jason.decode(bytes) do + {:ok, json} + else + {:error, reason} -> {:error, {:invalid_json_file, path, reason}} + end + end + + defp verify_report_shape(%{"status" => status, "checked_at" => checked_at}) + when status in ["ok", "warning"] and is_binary(checked_at), + do: :ok + + defp verify_report_shape(%{"status" => status, "checkedAt" => checked_at}) + when status in ["ok", "warning"] and is_binary(checked_at), + do: :ok + + defp verify_report_shape(_report), do: {:error, :invalid_verification_report} + + defp verify_repo_file(snapshot_dir, manifest) do + with %{"account" => %{"did" => did}, "repo" => repo, "identity" => %{"didDocument" => did_document}} <- manifest, + repo_path = Path.join(snapshot_dir, Map.fetch!(repo, "carPath")), + {:ok, bytes} <- File.read(repo_path), + :ok <- verify_file_hash(bytes, Map.fetch!(repo, "sha256")), + true <- byte_size(bytes) == Map.fetch!(repo, "byteSize"), + {:ok, verified} <- Verifier.verify_repo_car(bytes, did_document, did), + true <- verified.commit_cid_string == Map.fetch!(repo, "commit"), + true <- verified.rev == Map.fetch!(repo, "rev") do + :ok + else + false -> {:error, :repo_manifest_mismatch} + {:error, reason} -> {:error, reason} + _match -> {:error, :invalid_manifest_repo} + end + end - with :ok <- File.mkdir_p(snapshot_dir), - :ok <- File.write(repo_car_path, repo_car) do + defp verify_blob_files(snapshot_dir, %{"blobFiles" => blob_files}) when is_list(blob_files) do + Enum.reduce_while(blob_files, :ok, fn blob, :ok -> + case verify_blob_file(snapshot_dir, blob) do + :ok -> {:cont, :ok} + {:error, reason} -> {:halt, {:error, reason}} + end + end) + end + + defp verify_blob_files(_snapshot_dir, _manifest), do: {:error, :invalid_manifest_blobs} + + defp verify_blob_file(_snapshot_dir, %{"status" => status}) when status in ["missing", "failed"], do: :ok + + defp verify_blob_file(snapshot_dir, %{ + "cid" => cid, + "path" => path, + "byteSize" => byte_size, + "sha256" => hash, + "status" => "stored" + }) do + with {:ok, bytes} <- File.read(Path.join(snapshot_dir, path)), + :ok <- verify_file_hash(bytes, hash), + true <- byte_size(bytes) == byte_size, + {:ok, expected_cid} <- Cid.parse(cid), + true <- Cid.for_raw(bytes) == expected_cid do + :ok + else + false -> {:error, {:blob_manifest_mismatch, cid}} + {:error, reason} -> {:error, {:blob_verify_failed, cid, reason}} + end + end + + defp verify_blob_file(_snapshot_dir, _blob), do: {:error, :invalid_manifest_blob} + + defp verify_preferences_file(_snapshot_dir, %{"preferences" => %{"included" => false}}), do: :ok + + defp verify_preferences_file(snapshot_dir, %{"preferences" => %{"included" => true, "path" => path}}) do + with {:ok, _preferences} <- read_json_file(Path.join(snapshot_dir, path)), do: :ok + end + + defp verify_preferences_file(_snapshot_dir, _manifest), do: {:error, :invalid_manifest_preferences} + + defp verify_file_hash(bytes, expected_hash) when is_binary(expected_hash) do + if sha256(bytes) == expected_hash, do: :ok, else: {:error, :sha256_mismatch} + end + + defp verify_file_hash(_bytes, _hash), do: {:error, :missing_sha256} + + defp prepare_snapshot_workspace(config, %Account{} = account, did_document, verified, repo_car, opts) do + storage_key = snapshot_storage_key(account.did, verified.rev) + temp_dir = snapshot_temp_dir(config) + final_dir = Path.join(config.data_dir, storage_key) + repo_car_path = Path.join(temp_dir, "repo.car") + + with :ok <- File.mkdir_p(temp_dir), + :ok <- File.write(repo_car_path, repo_car), + {:ok, blob_cids} <- snapshot_blob_cids(account, verified, opts), + {:ok, blob_results} <- download_snapshot_blobs(account, temp_dir, blob_cids, opts), + {:ok, preferences} <- maybe_write_preferences(account, temp_dir, opts), + report = verification_report(account, verified, blob_results, preferences), + manifest = snapshot_manifest(account, did_document, verified, repo_car, blob_results, preferences, report), + :ok <- File.write(Path.join(temp_dir, "verification.json"), Jason.encode!(report, pretty: true)), + :ok <- File.write(Path.join(temp_dir, "manifest.json"), Jason.encode!(manifest, pretty: true)) do {:ok, %{ + account: account, storage_key: storage_key, + temp_dir: temp_dir, + final_dir: final_dir, repo_car_path: Path.join(storage_key, "repo.car"), + manifest_path: Path.join(storage_key, "manifest.json"), + verification_report_path: Path.join(storage_key, "verification.json"), commit_cid: verified.commit_cid_string, rev: verified.rev, source_pds_url: account.source_pds_url, handle: account.handle, did: account.did, byte_size: byte_size(repo_car), - sha256: sha256(repo_car) + sha256: sha256(repo_car), + status: snapshot_status(blob_results), + verification_status: snapshot_verification_status(blob_results, preferences), + blobs: blob_results, + preferences: preferences }} end end @@ -320,21 +560,311 @@ defmodule Tempest.PersonalBackups do Map.merge(attrs, %{ account_id: account.id, run_id: run.id, - status: "complete", + status: attrs.status, completed_at: DateTime.utc_now(:second), - verification_status: "ok" + manifest_path: attrs.manifest_path, + verification_status: attrs.verification_status, + verification_report_path: attrs.verification_report_path }) ) |> Repo.insert() end + defp insert_blob_records(%Snapshot{} = snapshot, blobs) do + Enum.reduce_while(blobs, {:ok, []}, fn blob, {:ok, inserted} -> + attrs = + blob + |> Map.take([:cid, :path, :byte_size, :sha256, :status, :error_reason]) + |> Map.put(:snapshot_id, snapshot.id) + + case %Tempest.PersonalBackups.Blob{} |> Tempest.PersonalBackups.Blob.changeset(attrs) |> Repo.insert() do + {:ok, row} -> {:cont, {:ok, [row | inserted]}} + {:error, %Ecto.Changeset{} = changeset} -> {:halt, {:error, changeset}} + end + end) + |> case do + {:ok, rows} -> {:ok, Enum.reverse(rows)} + {:error, reason} -> {:error, reason} + end + end + defp snapshot_storage_key(did, rev) do stamp = DateTime.utc_now(:second) |> DateTime.to_iso8601(:basic) |> String.replace(~r/[^0-9A-Za-z]/, "") safe_did = String.replace(did, ~r/[^A-Za-z0-9._-]/, "_") safe_rev = String.replace(rev, ~r/[^A-Za-z0-9._-]/, "_") - Path.join(["personal-backups", safe_did, "snapshots", stamp <> "-" <> safe_rev]) + suffix = System.unique_integer([:positive]) + Path.join(["personal-backups", safe_did, "snapshots", stamp <> "-" <> safe_rev <> "-" <> Integer.to_string(suffix)]) + end + + defp snapshot_temp_dir(config) do + Path.join([config.data_dir, "tmp", "personal-backups", Integer.to_string(System.unique_integer([:positive]))]) + end + + defp snapshot_blob_cids(account, verified, opts) do + with {:ok, listed_cids} <- list_all_source_blob_cids(account, opts), + referenced_cids = referenced_blob_cids(verified) do + {:ok, Enum.sort(Enum.uniq(referenced_cids ++ listed_cids))} + end + end + + defp list_all_source_blob_cids(account, opts, cursor \\ nil, cids \\ []) do + request_opts = + opts + |> Keyword.put(:limit, Keyword.get(opts, :blob_page_limit, 500)) + |> maybe_put_keyword(:cursor, cursor) + + with {:ok, response} <- SourceClient.list_blobs(account.source_pds_url, account.did, request_opts), + page_cids <- Map.get(response, "cids", Map.get(response, "blobs", [])), + :ok <- validate_cid_list(page_cids) do + case Map.get(response, "cursor") do + next_cursor when is_binary(next_cursor) and next_cursor != "" -> + list_all_source_blob_cids(account, opts, next_cursor, cids ++ page_cids) + + _cursor -> + {:ok, cids ++ page_cids} + end + end + end + + defp validate_cid_list(cids) when is_list(cids) do + if Enum.all?(cids, &is_binary/1), do: :ok, else: {:error, :invalid_blob_list} + end + + defp validate_cid_list(_cids), do: {:error, :invalid_blob_list} + + defp referenced_blob_cids(verified) do + blocks = Map.new(verified.car.blocks, fn %{cid: cid, data: data} -> {Cid.to_string(cid), data} end) + + verified.entries + |> Map.values() + |> Enum.flat_map(fn cid -> + cid + |> Cid.to_string() + |> then(&Map.get(blocks, &1)) + |> decode_record_block() + |> extract_blob_cids() + end) + |> Enum.uniq() + end + + defp decode_record_block(nil), do: nil + + defp decode_record_block(bytes) do + case Drisl.decode(bytes) do + {:ok, value} -> value + {:error, _reason} -> Jason.decode(bytes) |> elem_or_nil() + end + end + + defp elem_or_nil({:ok, value}), do: value + defp elem_or_nil({:error, _reason}), do: nil + + defp extract_blob_cids(%{"$type" => "blob", "ref" => %{"$link" => cid}}) when is_binary(cid), do: [cid] + defp extract_blob_cids(%{"$type" => "blob", "cid" => cid}) when is_binary(cid), do: [cid] + + defp extract_blob_cids(map) when is_map(map) do + map + |> Map.values() + |> Enum.flat_map(&extract_blob_cids/1) + end + + defp extract_blob_cids(list) when is_list(list), do: Enum.flat_map(list, &extract_blob_cids/1) + defp extract_blob_cids(_value), do: [] + + defp download_snapshot_blobs(account, temp_dir, cids, opts) do + blob_dir = Path.join(temp_dir, "blobs") + + with :ok <- File.mkdir_p(blob_dir) do + cids + |> Task.async_stream( + fn cid -> download_snapshot_blob(account, blob_dir, cid, opts) end, + max_concurrency: Keyword.get(opts, :blob_concurrency, @default_blob_concurrency), + timeout: :infinity + ) + |> Enum.reduce_while({:ok, []}, fn + {:ok, blob}, {:ok, blobs} -> {:cont, {:ok, [blob | blobs]}} + {:exit, reason}, {:ok, _blobs} -> {:halt, {:error, {:blob_task_exit, reason}}} + end) + |> case do + {:ok, blobs} -> {:ok, Enum.reverse(blobs)} + {:error, reason} -> {:error, reason} + end + end + end + + defp download_snapshot_blob(account, blob_dir, cid, opts) do + case fetch_blob_with_retry(account, cid, opts, @transient_retries) do + {:ok, bytes} -> + verify_and_write_blob(blob_dir, cid, bytes) + + {:error, {:source_pds_http_error, status}} when status in [404, 410] -> + %{cid: cid, path: nil, byte_size: 0, sha256: nil, status: "missing", error_reason: "source returned #{status}"} + + {:error, reason} -> + %{cid: cid, path: nil, byte_size: 0, sha256: nil, status: "failed", error_reason: inspect(reason)} + end + end + + defp fetch_blob_with_retry(account, cid, opts, retries_left) do + case SourceClient.get_blob(account.source_pds_url, account.did, cid, opts) do + {:error, reason} -> + if retries_left > 0 and transient_blob_error?(reason) do + fetch_blob_with_retry(account, cid, opts, retries_left - 1) + else + {:error, reason} + end + + result -> + result + end + end + + defp transient_blob_error?({:source_pds_http_error, status}) when status in 500..599, do: true + defp transient_blob_error?({:source_pds_request_failed, _reason}), do: true + defp transient_blob_error?(_reason), do: false + + defp verify_and_write_blob(blob_dir, cid, bytes) do + with {:ok, expected_cid} <- Cid.parse(cid), + actual_cid = Cid.for_raw(bytes), + true <- actual_cid == expected_cid, + path = Path.join(blob_dir, cid), + :ok <- File.write(path, bytes) do + %{ + cid: cid, + path: Path.join(["blobs", cid]), + byte_size: byte_size(bytes), + sha256: sha256(bytes), + status: "stored", + error_reason: nil + } + else + false -> + %{cid: cid, path: nil, byte_size: 0, sha256: nil, status: "failed", error_reason: "cid_mismatch"} + + {:error, reason} -> + %{cid: cid, path: nil, byte_size: 0, sha256: nil, status: "failed", error_reason: inspect(reason)} + end + end + + defp maybe_write_preferences(account, temp_dir, opts) do + with {:ok, credential} <- preferences_credential(account), + {:ok, secret} <- SecretStore.decrypt(credential.secret_ciphertext), + {:ok, preferences} <- SourceClient.get_preferences(account.source_pds_url, secret, opts), + :ok <- File.write(Path.join(temp_dir, "preferences.json"), Jason.encode!(preferences, pretty: true)) do + {:ok, %{included: true, path: "preferences.json", warning: nil}} + else + {:skip, reason} when reason in [:no_credentials, :credential_deleted] -> + {:ok, %{included: false, path: nil, warning: nil}} + + {:skip, reason} -> + {:ok, %{included: false, path: nil, warning: Atom.to_string(reason)}} + + {:error, reason} -> + {:ok, %{included: false, path: nil, warning: "preferences_auth_or_fetch_failed: #{inspect(reason)}"}} + end + end + + defp preferences_credential(account) do + credential = account |> Repo.preload(:credential, force: true) |> Map.fetch!(:credential) + + case credential do + %Credential{mode: "none"} -> {:skip, :no_credentials} + %Credential{deleted_at: %DateTime{}} -> {:skip, :credential_deleted} + %Credential{secret_ciphertext: secret} when is_binary(secret) and secret != "" -> {:ok, credential} + %Credential{} -> {:skip, :credential_missing_secret} + end + end + + defp snapshot_manifest(account, did_document, verified, repo_car, blobs, preferences, report) do + missing = blobs |> Enum.filter(&(&1.status == "missing")) |> Enum.map(& &1.cid) + + %{ + "version" => 1, + "account" => %{ + "did" => account.did, + "handle" => account.handle, + "sourcePds" => account.source_pds_url + }, + "identity" => %{ + "didDocument" => did_document + }, + "repo" => %{ + "carPath" => "repo.car", + "commit" => verified.commit_cid_string, + "rev" => verified.rev, + "byteSize" => byte_size(repo_car), + "sha256" => sha256(repo_car) + }, + "blobs" => %{ + "count" => Enum.count(blobs, &(&1.status == "stored")), + "expected" => length(blobs), + "complete" => Enum.all?(blobs, &(&1.status == "stored")), + "missing" => missing + }, + "blobFiles" => Enum.map(blobs, &blob_manifest_entry/1), + "preferences" => Map.take(preferences, [:included, :path]) |> stringify_map_keys(), + "verification" => %{ + "status" => report.status, + "checkedAt" => report.checked_at, + "path" => "verification.json" + } + } + end + + defp blob_manifest_entry(blob) do + %{ + "cid" => blob.cid, + "path" => blob.path, + "byteSize" => blob.byte_size, + "sha256" => blob.sha256, + "status" => blob.status, + "errorReason" => blob.error_reason + } end + defp verification_report(account, verified, blobs, preferences) do + warnings = + [] + |> add_warning(preferences.warning) + |> add_warning_if(Enum.any?(blobs, &(&1.status == "missing")), "missing_blobs") + |> add_warning_if(Enum.any?(blobs, &(&1.status == "failed")), "failed_blobs") + + %{ + status: if(warnings == [], do: "ok", else: "warning"), + checked_at: DateTime.utc_now(:second) |> DateTime.to_iso8601(), + did: account.did, + handle: account.handle, + source_pds_url: account.source_pds_url, + commit_cid: verified.commit_cid_string, + rev: verified.rev, + record_count: verified.record_count, + blob_count: length(blobs), + stored_blob_count: Enum.count(blobs, &(&1.status == "stored")), + warnings: Enum.reverse(warnings) + } + end + + defp snapshot_status(blobs) do + if Enum.all?(blobs, &(&1.status == "stored")), do: "complete", else: "incomplete" + end + + defp snapshot_verification_status(blobs, preferences) do + if snapshot_status(blobs) == "complete" and is_nil(preferences.warning), do: "ok", else: "warning" + end + + defp add_warning(warnings, nil), do: warnings + defp add_warning(warnings, warning), do: [warning | warnings] + + defp add_warning_if(warnings, true, warning), do: [warning | warnings] + defp add_warning_if(warnings, false, _warning), do: warnings + + defp stringify_map_keys(map) do + Map.new(map, fn {key, value} -> {Atom.to_string(key), value} end) + end + + defp maybe_put_keyword(keyword, _key, nil), do: keyword + defp maybe_put_keyword(keyword, key, value), do: Keyword.put(keyword, key, value) + defp sha256(bytes), do: :sha256 |> :crypto.hash(bytes) |> Base.encode16(case: :lower) defp stringify_keys(attrs) do diff --git a/lib/tempest/personal_backups/snapshot.ex b/lib/tempest/personal_backups/snapshot.ex index 4bc9876..b8b4f37 100644 --- a/lib/tempest/personal_backups/snapshot.ex +++ b/lib/tempest/personal_backups/snapshot.ex @@ -63,4 +63,10 @@ defmodule Tempest.PersonalBackups.Snapshot do |> validate_inclusion(:verification_status, @verification_statuses) |> unique_constraint(:storage_key) end + + def pin_changeset(snapshot, attrs) do + snapshot + |> cast(attrs, [:pinned]) + |> validate_required([:pinned]) + end end diff --git a/lib/tempest/personal_backups/storage.ex b/lib/tempest/personal_backups/storage.ex new file mode 100644 index 0000000..697072b --- /dev/null +++ b/lib/tempest/personal_backups/storage.ex @@ -0,0 +1,103 @@ +defmodule Tempest.PersonalBackups.Storage do + @moduledoc """ + Storage boundary for personal backup snapshot directories and archives. + """ + + alias Tempest.Admin.S3BackupStorage + alias Tempest.Config + + def finalize_snapshot(%Config{} = config, workspace) when is_map(workspace) do + if File.exists?(workspace.final_dir) do + {:error, :snapshot_already_exists} + else + with :ok <- File.mkdir_p(Path.dirname(workspace.final_dir)), + :ok <- File.rename(workspace.temp_dir, workspace.final_dir), + {:ok, upload} <- maybe_upload_snapshot(config, workspace.storage_key) do + {:ok, + workspace + |> Map.delete(:temp_dir) + |> Map.put(:storage_upload, upload)} + end + end + end + + def delete_snapshot(%Config{} = config, storage_key) when is_binary(storage_key) do + config.data_dir + |> Path.join(storage_key) + |> File.rm_rf() + |> case do + {:ok, _files} -> :ok + {:error, _file, reason} -> {:error, reason} + end + end + + def archive_snapshot(%Config{} = config, storage_key, target_path) + when is_binary(storage_key) and is_binary(target_path) do + snapshot_dir = Path.join(config.data_dir, storage_key) + + with true <- File.dir?(snapshot_dir), + :ok <- File.mkdir_p(Path.dirname(target_path)), + {:ok, files} <- snapshot_files(snapshot_dir), + {:ok, archive} <- create_zip(target_path, snapshot_dir, files) do + {:ok, %{path: archive, byte_size: File.stat!(archive).size}} + else + false -> {:error, :snapshot_not_found} + {:error, reason} -> {:error, reason} + end + end + + defp maybe_upload_snapshot(config, storage_key) do + storage_config = Application.get_env(:tempest, __MODULE__, []) + + case Keyword.get(storage_config, :store, :local) do + :s3 -> upload_snapshot_archive(config, storage_key, Keyword.fetch!(storage_config, :s3)) + "s3" -> upload_snapshot_archive(config, storage_key, Keyword.fetch!(storage_config, :s3)) + _local -> {:ok, nil} + end + rescue + e in KeyError -> {:error, {:missing_personal_backup_storage_config, e.key}} + end + + defp upload_snapshot_archive(config, storage_key, s3_config) do + archive_path = Path.join([config.data_dir, "tmp", Path.basename(storage_key) <> ".zip"]) + key = Keyword.get(s3_config, :key) || storage_key <> ".zip" + + with {:ok, %{path: path, byte_size: byte_size}} <- archive_snapshot(config, storage_key, archive_path), + {:ok, uploaded} <- S3BackupStorage.upload_file(s3_config, key, path) do + {:ok, Map.merge(uploaded, %{archive_path: path, archive_bytes: byte_size})} + end + end + + defp snapshot_files(snapshot_dir) do + snapshot_dir + |> files_under() + |> case do + [] -> {:error, :snapshot_empty} + files -> {:ok, files} + end + end + + defp files_under(dir) do + dir + |> File.ls!() + |> Enum.flat_map(fn entry -> + path = Path.join(dir, entry) + if File.dir?(path), do: files_under(path), else: [path] + end) + end + + defp create_zip(target_path, snapshot_dir, files) do + archive = String.to_charlist(target_path) + + entries = + Enum.map(files, fn path -> + relative = path |> Path.relative_to(snapshot_dir) |> String.to_charlist() + {relative, File.read!(path)} + end) + + case :zip.create(archive, entries) do + {:ok, archive} -> {:ok, List.to_string(archive)} + {:error, reason} -> {:error, {:zip_create, reason}} + end + end +end diff --git a/test/tempest/personal_backups_test.exs b/test/tempest/personal_backups_test.exs index bb65ce5..e972c25 100644 --- a/test/tempest/personal_backups_test.exs +++ b/test/tempest/personal_backups_test.exs @@ -3,8 +3,8 @@ defmodule Tempest.PersonalBackupsTest do alias Tempest.Config alias Tempest.PersonalBackups - alias Tempest.PersonalBackups.{Account, Credential, RetentionSetting, Run, Snapshot, SourceClient} - alias Tempest.RepoCore.{Car, Cid, Commit, Mst, Tid} + alias Tempest.PersonalBackups.{Account, Credential, RetentionSetting, Run, Snapshot, SourceClient, Storage} + alias Tempest.RepoCore.{Car, Cid, Commit, Drisl, Mst, Tid} alias Tempest.Repo import Ecto.Query @@ -19,6 +19,7 @@ defmodule Tempest.PersonalBackupsTest do old_identity_config = Application.get_env(:tempest, Tempest.Identity, []) old_source_client_config = Application.get_env(:tempest, SourceClient, []) + old_storage_config = Application.get_env(:tempest, Storage, []) Application.put_env(:tempest, Tempest.Identity, plc_directory_url: "https://plc.test", @@ -32,6 +33,7 @@ defmodule Tempest.PersonalBackupsTest do on_exit(fn -> Application.put_env(:tempest, Tempest.Identity, old_identity_config) Application.put_env(:tempest, SourceClient, old_source_client_config) + Application.put_env(:tempest, Storage, old_storage_config) end) :ok @@ -200,12 +202,16 @@ defmodule Tempest.PersonalBackupsTest do expect_did_document(did_document) assert {:ok, account} = PersonalBackups.register_account(%{did: @did, handle: @handle}) - Req.Test.expect(__MODULE__, 3, fn conn -> + Req.Test.expect(__MODULE__, 4, fn conn -> case conn.request_path do "/xrpc/com.atproto.sync.getRepo" -> assert conn.query_params["did"] == @did Plug.Conn.resp(conn, 200, repo_car) + "/xrpc/com.atproto.sync.listBlobs" -> + assert conn.query_params["did"] == @did + Req.Test.json(conn, %{"cids" => []}) + "/" <> encoded_did -> assert URI.decode(encoded_did) == @did Req.Test.json(conn, did_document) @@ -230,6 +236,220 @@ defmodule Tempest.PersonalBackupsTest do assert verification.record_count == 1 assert File.read!(Path.join(config.data_dir, snapshot.repo_car_path)) == repo_car + assert File.exists?(Path.join(config.data_dir, snapshot.manifest_path)) + assert File.exists?(Path.join(config.data_dir, snapshot.verification_report_path)) + end + + test "create_repo_snapshot stores merged referenced and listed blobs plus preferences" do + data_dir = + Path.join(System.tmp_dir!(), "tempest_personal_backup_blob_snapshot_#{System.unique_integer([:positive])}") + + config = snapshot_test_config(data_dir) + on_exit(fn -> File.rm_rf(data_dir) end) + + referenced_blob = "referenced blob" + listed_blob = "listed blob" + referenced_cid = Cid.for_raw(referenced_blob) |> Cid.to_string() + listed_cid = Cid.for_raw(listed_blob) |> Cid.to_string() + {repo_car, _commit, _commit_cid, did_document} = fixture_car(referenced_cid) + + expect_did_document(did_document) + assert {:ok, account} = PersonalBackups.register_account(%{did: @did, handle: @handle}) + assert {:ok, %{account: account}} = PersonalBackups.rotate_credential(account, "access_token", "access-token") + + Req.Test.expect(__MODULE__, 7, fn conn -> + case conn.request_path do + "/xrpc/com.atproto.sync.getRepo" -> + Plug.Conn.resp(conn, 200, repo_car) + + "/xrpc/com.atproto.sync.listBlobs" -> + Req.Test.json(conn, %{"cids" => [listed_cid]}) + + "/xrpc/com.atproto.sync.getBlob" -> + case conn.query_params["cid"] do + ^referenced_cid -> Plug.Conn.resp(conn, 200, referenced_blob) + ^listed_cid -> Plug.Conn.resp(conn, 200, listed_blob) + end + + "/xrpc/app.bsky.actor.getPreferences" -> + assert ["Bearer access-token"] = Plug.Conn.get_req_header(conn, "authorization") + Req.Test.json(conn, %{"preferences" => [%{"$type" => "app.bsky.actor.defs#adultContentPref"}]}) + + "/" <> encoded_did -> + assert URI.decode(encoded_did) == @did + Req.Test.json(conn, did_document) + end + end) + + assert {:ok, %{snapshot: snapshot}} = PersonalBackups.create_repo_snapshot(account, config: config) + + snapshot = Repo.preload(snapshot, :blobs) + assert snapshot.status == "complete" + assert snapshot.verification_status == "ok" + assert Enum.map(snapshot.blobs, & &1.cid) |> Enum.sort() == Enum.sort([referenced_cid, listed_cid]) + assert Enum.all?(snapshot.blobs, &(&1.status == "stored")) + assert File.read!(Path.join([config.data_dir, snapshot.storage_key, "blobs", referenced_cid])) == referenced_blob + assert File.read!(Path.join([config.data_dir, snapshot.storage_key, "preferences.json"])) =~ "adultContentPref" + + manifest = snapshot |> manifest_json(config) + assert manifest["blobs"]["expected"] == 2 + assert manifest["blobs"]["complete"] == true + assert manifest["preferences"]["included"] == true + end + + test "create_repo_snapshot records missing blobs and preference auth warnings separately" do + data_dir = + Path.join(System.tmp_dir!(), "tempest_personal_backup_missing_snapshot_#{System.unique_integer([:positive])}") + + config = snapshot_test_config(data_dir) + on_exit(fn -> File.rm_rf(data_dir) end) + + missing_cid = Cid.for_raw("missing blob") |> Cid.to_string() + {repo_car, _commit, _commit_cid, did_document} = fixture_car(missing_cid) + + expect_did_document(did_document) + assert {:ok, account} = PersonalBackups.register_account(%{did: @did, handle: @handle}) + assert {:ok, %{account: account}} = PersonalBackups.rotate_credential(account, "access_token", "bad-token") + + Req.Test.expect(__MODULE__, 6, fn conn -> + case conn.request_path do + "/xrpc/com.atproto.sync.getRepo" -> + Plug.Conn.resp(conn, 200, repo_car) + + "/xrpc/com.atproto.sync.listBlobs" -> + Req.Test.json(conn, %{"cids" => []}) + + "/xrpc/com.atproto.sync.getBlob" -> + Plug.Conn.resp(conn, 404, "missing") + + "/xrpc/app.bsky.actor.getPreferences" -> + Plug.Conn.resp(conn, 401, "bad auth") + + "/" <> encoded_did -> + assert URI.decode(encoded_did) == @did + Req.Test.json(conn, did_document) + end + end) + + assert {:ok, %{snapshot: snapshot}} = PersonalBackups.create_repo_snapshot(account, config: config) + + snapshot = Repo.preload(snapshot, :blobs) + assert snapshot.status == "incomplete" + assert snapshot.verification_status == "warning" + assert [%{cid: ^missing_cid, status: "missing"}] = snapshot.blobs + + manifest = snapshot |> manifest_json(config) + assert manifest["blobs"]["complete"] == false + assert manifest["blobs"]["missing"] == [missing_cid] + + report = snapshot |> verification_report_json(config) + assert "missing_blobs" in report["warnings"] + assert Enum.any?(report["warnings"], &String.starts_with?(&1, "preferences_auth_or_fetch_failed")) + end + + test "create_repo_snapshot can upload a snapshot archive through S3-compatible storage" do + data_dir = Path.join(System.tmp_dir!(), "tempest_personal_backup_s3_snapshot_#{System.unique_integer([:positive])}") + config = snapshot_test_config(data_dir) + on_exit(fn -> File.rm_rf(data_dir) end) + + Application.put_env(:tempest, Storage, + store: :s3, + s3: [ + endpoint_url: "https://objects.example.test", + bucket: "tempest-backups", + req_options: [plug: {Req.Test, __MODULE__}], + headers: [{"authorization", "Bearer backup-token"}] + ] + ) + + {repo_car, _commit, _commit_cid, did_document} = fixture_car() + expect_did_document(did_document) + assert {:ok, account} = PersonalBackups.register_account(%{did: @did, handle: @handle}) + + Req.Test.expect(__MODULE__, 5, fn conn -> + case conn.request_path do + "/xrpc/com.atproto.sync.getRepo" -> + Plug.Conn.resp(conn, 200, repo_car) + + "/xrpc/com.atproto.sync.listBlobs" -> + Req.Test.json(conn, %{"cids" => []}) + + "/tempest-backups/" <> key -> + assert conn.method == "PUT" + assert String.ends_with?(URI.decode(key), ".zip") + assert Plug.Conn.get_req_header(conn, "authorization") == ["Bearer backup-token"] + assert {:ok, bytes, conn} = Plug.Conn.read_body(conn) + assert byte_size(bytes) > 0 + Plug.Conn.send_resp(conn, 200, "") + + "/" <> encoded_did -> + assert URI.decode(encoded_did) == @did + Req.Test.json(conn, did_document) + end + end) + + assert {:ok, %{snapshot: snapshot}} = PersonalBackups.create_repo_snapshot(account, config: config) + assert File.exists?(Path.join(config.data_dir, snapshot.storage_key)) + end + + test "prune_snapshots applies retention policies without deleting pinned snapshots" do + expect_did_document() + assert {:ok, account} = PersonalBackups.register_account(%{did: @did, handle: @handle}) + + data_dir = Path.join(System.tmp_dir!(), "tempest_personal_backup_retention_#{System.unique_integer([:positive])}") + config = snapshot_test_config(data_dir) + on_exit(fn -> File.rm_rf(data_dir) end) + + old = insert_dummy_snapshot!(account, config, "old", ~U[2026-01-01 00:00:00Z]) + pinned = insert_dummy_snapshot!(account, config, "pinned", ~U[2026-01-02 00:00:00Z], pinned: true) + recent = insert_dummy_snapshot!(account, config, "recent", ~U[2026-01-03 00:00:00Z]) + + assert {:ok, _setting} = PersonalBackups.update_retention_setting(account, %{policy: "keep_last_n", keep_last: 1}) + assert {:ok, pruned} = PersonalBackups.prune_snapshots(account, config: config) + + assert Enum.map(pruned, & &1.id) == [old.id] + refute Repo.get(Snapshot, old.id) + assert Repo.get(Snapshot, pinned.id) + assert Repo.get(Snapshot, recent.id) + refute File.exists?(Path.join(config.data_dir, old.storage_key)) + assert File.exists?(Path.join(config.data_dir, pinned.storage_key)) + end + + test "export_snapshot_bundle creates a portable zip with snapshot files" do + data_dir = Path.join(System.tmp_dir!(), "tempest_personal_backup_export_#{System.unique_integer([:positive])}") + config = snapshot_test_config(data_dir) + on_exit(fn -> File.rm_rf(data_dir) end) + + {snapshot, _did_document} = create_no_blob_snapshot!(config) + export_path = Path.join([data_dir, "exports", "snapshot.zip"]) + + assert {:ok, %{path: ^export_path, byte_size: byte_size}} = + PersonalBackups.export_snapshot_bundle(snapshot, config: config, path: export_path) + + assert byte_size > 0 + assert {:ok, entries} = :zip.list_dir(String.to_charlist(export_path)) + + names = + entries + |> Enum.filter(&match?({:zip_file, _, _, _, _, _}, &1)) + |> Enum.map(fn {:zip_file, name, _, _, _, _} -> List.to_string(name) end) + + assert "manifest.json" in names + assert "repo.car" in names + assert "verification.json" in names + end + + test "verify_snapshot_offline validates manifest files without source PDS access" do + data_dir = Path.join(System.tmp_dir!(), "tempest_personal_backup_offline_#{System.unique_integer([:positive])}") + config = snapshot_test_config(data_dir) + on_exit(fn -> File.rm_rf(data_dir) end) + + {snapshot, _did_document} = create_no_blob_snapshot!(config) + + assert {:ok, %{status: "ok"}} = PersonalBackups.verify_snapshot_offline(snapshot, config: config) + + File.write!(Path.join(config.data_dir, snapshot.repo_car_path), "corrupt") + assert {:error, :sha256_mismatch} = PersonalBackups.verify_snapshot_offline(snapshot, config: config) end test "create_repo_snapshot rejects invalid commit signatures" do @@ -277,6 +497,77 @@ defmodule Tempest.PersonalBackupsTest do refute Repo.exists?(from snapshot in Snapshot, where: snapshot.account_id == ^account.id) end + defp create_no_blob_snapshot!(config) do + {repo_car, _commit, _commit_cid, did_document} = fixture_car() + expect_did_document(did_document) + assert {:ok, account} = PersonalBackups.register_account(%{did: @did, handle: @handle}) + + Req.Test.expect(__MODULE__, 4, fn conn -> + case conn.request_path do + "/xrpc/com.atproto.sync.getRepo" -> + Plug.Conn.resp(conn, 200, repo_car) + + "/xrpc/com.atproto.sync.listBlobs" -> + Req.Test.json(conn, %{"cids" => []}) + + "/" <> encoded_did -> + assert URI.decode(encoded_did) == @did + Req.Test.json(conn, did_document) + end + end) + + assert {:ok, %{snapshot: snapshot}} = PersonalBackups.create_repo_snapshot(account, config: config) + {snapshot, did_document} + end + + defp insert_dummy_snapshot!(account, config, name, completed_at, opts \\ []) do + storage_key = Path.join(["personal-backups", account.did |> String.replace(":", "_"), "snapshots", name]) + snapshot_dir = Path.join(config.data_dir, storage_key) + File.mkdir_p!(snapshot_dir) + File.write!(Path.join(snapshot_dir, "manifest.json"), "{}") + + %Snapshot{} + |> Snapshot.changeset(%{ + account_id: account.id, + status: "complete", + storage_key: storage_key, + source_pds_url: account.source_pds_url, + handle: account.handle, + did: account.did, + completed_at: completed_at, + verification_status: "ok", + pinned: Keyword.get(opts, :pinned, false) + }) + |> Repo.insert!() + end + + defp snapshot_test_config(data_dir) do + Config.validate!( + [ + hostname: "localhost", + public_url: "http://localhost:4002", + data_dir: data_dir, + blob_max_bytes: 10_000_000 + ], + env: :test, + endpoint_config: Application.fetch_env!(:tempest, TempestWeb.Endpoint) + ) + end + + defp manifest_json(snapshot, config) do + config.data_dir + |> Path.join(snapshot.manifest_path) + |> File.read!() + |> Jason.decode!() + end + + defp verification_report_json(snapshot, config) do + config.data_dir + |> Path.join(snapshot.verification_report_path) + |> File.read!() + |> Jason.decode!() + end + defp expect_did_document(overrides \\ %{}) do Req.Test.expect(__MODULE__, fn conn -> assert conn.request_path == "/" <> URI.encode(@did) @@ -302,9 +593,26 @@ defmodule Tempest.PersonalBackupsTest do end) end - defp fixture_car do + defp fixture_car(blob_cid \\ nil) do {public_key, private_key} = :crypto.generate_key(:ecdh, :secp256k1) - record_cid = Cid.for_raw(~s({"$type":"app.bsky.feed.post","text":"fixture"})) + + record = + if blob_cid do + %{ + "$type" => "app.bsky.actor.profile", + "avatar" => %{ + "$type" => "blob", + "ref" => %{"$link" => blob_cid}, + "mimeType" => "image/png", + "size" => 12 + } + } + else + %{"$type" => "app.bsky.feed.post", "text" => "fixture"} + end + + record_bytes = Drisl.encode!(record) + record_cid = Cid.for_drisl(record_bytes) mst = Mst.from_entries!([{"app.bsky.feed.post/3jui7kd54zh2y", record_cid}]) {:ok, %{root: mst_root, blocks: mst_blocks}} = Mst.serialize(mst) @@ -322,7 +630,7 @@ defmodule Tempest.PersonalBackupsTest do blocks = [ {commit_cid, commit_bytes}, - {record_cid, ~s({"$type":"app.bsky.feed.post","text":"fixture"})} | mst_blocks + {record_cid, record_bytes} | mst_blocks ] {:ok, car_bytes} = Car.encode([commit_cid], blocks)