diff --git a/lib/tempest/blobs/s3_storage.ex b/lib/tempest/blobs/s3_storage.ex index 09a7332..27ab5bf 100644 --- a/lib/tempest/blobs/s3_storage.ex +++ b/lib/tempest/blobs/s3_storage.ex @@ -46,13 +46,21 @@ defmodule Tempest.Blobs.S3Storage do {:ok, destination} {:ok, %{status: 404}} -> - {:error, :blob_not_found} + promote_blob_via_read_write(config, did, cid, source, destination) {:ok, %{status: status}} -> - {:error, {:s3_status, status}} + case promote_blob_via_read_write(config, did, cid, source, destination) do + {:ok, promoted} -> {:ok, promoted} + {:error, :blob_not_found} -> {:error, {:s3_status, status}} + {:error, reason} -> {:error, reason} + end {:error, reason} -> - {:error, reason} + case promote_blob_via_read_write(config, did, cid, source, destination) do + {:ok, promoted} -> {:ok, promoted} + {:error, :blob_not_found} -> {:error, reason} + {:error, fallback_reason} -> {:error, fallback_reason} + end end end end @@ -124,6 +132,35 @@ defmodule Tempest.Blobs.S3Storage do e in KeyError -> {:error, {:missing_s3_config, e.key}} end + defp promote_blob_via_read_write(config, did, cid, source, destination) do + with {:ok, bytes} <- read_key(config, source), + :ok <- write_key(config, destination, bytes), + :ok <- delete_temp_blob(config, did, cid) do + {:ok, destination} + end + end + + defp read_key(config, key) do + with {:ok, request} <- request_options(config, key, method: :get) do + case Req.request(request) do + {:ok, %{status: status, body: bytes}} when status in 200..299 and is_binary(bytes) -> {:ok, bytes} + {:ok, %{status: 404}} -> {:error, :blob_not_found} + {:ok, %{status: status}} -> {:error, {:s3_status, status}} + {:error, reason} -> {:error, reason} + end + end + end + + defp write_key(config, key, bytes) when is_binary(bytes) do + with {:ok, request} <- request_options(config, key, method: :put, body: bytes) do + case Req.request(request) do + {:ok, %{status: status}} when status in 200..299 -> :ok + {:ok, %{status: status}} -> {:error, {:s3_status, status}} + {:error, reason} -> {:error, reason} + end + end + end + defp object_url(endpoint_url, bucket, key) do endpoint_url |> String.trim_trailing("/") diff --git a/test/tempest/blobs/s3_integration_test.exs b/test/tempest/blobs/s3_integration_test.exs index 705b2cb..0a3cafe 100644 --- a/test/tempest/blobs/s3_integration_test.exs +++ b/test/tempest/blobs/s3_integration_test.exs @@ -89,7 +89,7 @@ defmodule Tempest.Blobs.S3IntegrationTest do assert response(response, 200) == "s3 integration bytes" end - test "record writes report S3 promotion failures clearly", %{conn: conn} do + test "record writes recover when S3 copy fails but the temp object is readable", %{conn: conn} do account = create_account!(conn) cid = Blobs.cid_for("s3 promotion failure") @@ -112,20 +112,75 @@ defmodule Tempest.Blobs.S3IntegrationTest do Plug.Conn.send_resp(conn, 403, "") end) + Req.Test.expect(__MODULE__, fn conn -> + assert conn.method == "GET" + assert conn.request_path == "/tempest-test/temp/blobs/#{encoded_did(account["did"])}/#{cid}" + Plug.Conn.send_resp(conn, 200, "s3 promotion failure") + end) + + Req.Test.expect(__MODULE__, fn conn -> + assert conn.method == "PUT" + assert conn.request_path == "/tempest-test/blobs/#{encoded_did(account["did"])}/#{cid}" + assert {:ok, "s3 promotion failure", conn} = Plug.Conn.read_body(conn) + Plug.Conn.send_resp(conn, 200, "") + end) + + Req.Test.expect(__MODULE__, fn conn -> + assert conn.method == "DELETE" + assert conn.request_path == "/tempest-test/temp/blobs/#{encoded_did(account["did"])}/#{cid}" + Plug.Conn.send_resp(conn, 204, "") + end) + + conn + |> auth_json(account) + |> post(~p"/xrpc/com.atproto.repo.createRecord", %{ + "repo" => account["did"], + "collection" => "app.tempest.blob", + "rkey" => "s3-failure", + "validate" => false, + "record" => %{"$type" => "app.tempest.blob", "image" => blob} + }) + |> json_response(200) + end + + test "record writes report missing temp object bytes clearly", %{conn: conn} do + account = create_account!(conn) + cid = Blobs.cid_for("s3 missing temp") + + Req.Test.expect(__MODULE__, fn conn -> + assert conn.method == "PUT" + assert conn.request_path == "/tempest-test/temp/blobs/#{encoded_did(account["did"])}/#{cid}" + Plug.Conn.send_resp(conn, 200, "") + end) + + blob = upload_blob!(conn, account, "s3 missing temp")["blob"] + + Req.Test.expect(__MODULE__, fn conn -> + assert conn.method == "PUT" + assert conn.request_path == "/tempest-test/blobs/#{encoded_did(account["did"])}/#{cid}" + Plug.Conn.send_resp(conn, 404, "") + end) + + Req.Test.expect(__MODULE__, fn conn -> + assert conn.method == "GET" + assert conn.request_path == "/tempest-test/temp/blobs/#{encoded_did(account["did"])}/#{cid}" + Plug.Conn.send_resp(conn, 404, "") + end) + response = conn |> auth_json(account) |> post(~p"/xrpc/com.atproto.repo.createRecord", %{ "repo" => account["did"], "collection" => "app.tempest.blob", - "rkey" => "s3-failure", + "rkey" => "s3-missing-temp", "validate" => false, "record" => %{"$type" => "app.tempest.blob", "image" => blob} }) - |> json_response(502) + |> json_response(500) - assert response["error"] == "UpstreamFailure" - assert response["message"] == "blob storage request failed" + assert response["error"] == "InternalServerError" + assert response["message"] == "uploaded blob bytes could not be found" end defp create_account!(conn) do diff --git a/test/tempest/blobs/s3_storage_test.exs b/test/tempest/blobs/s3_storage_test.exs index 686f52b..c8dfb05 100644 --- a/test/tempest/blobs/s3_storage_test.exs +++ b/test/tempest/blobs/s3_storage_test.exs @@ -59,6 +59,60 @@ defmodule Tempest.Blobs.S3StorageTest do assert promoted_path == "blobs/#{did}/#{cid}" end + test "promote_blob falls back to read and write when copy cannot find the source", %{ + config: config, + did: did, + cid: cid + } do + Req.Test.expect(__MODULE__, fn conn -> + assert conn.method == "PUT" + assert conn.request_path == "/tempest-test/blobs/did%3Aplc%3As3storage/#{cid}" + Plug.Conn.send_resp(conn, 404, "") + end) + + Req.Test.expect(__MODULE__, fn conn -> + assert conn.method == "GET" + assert conn.request_path == "/tempest-test/temp/blobs/did%3Aplc%3As3storage/#{cid}" + Plug.Conn.send_resp(conn, 200, "s3 bytes") + end) + + Req.Test.expect(__MODULE__, fn conn -> + assert conn.method == "PUT" + assert conn.request_path == "/tempest-test/blobs/did%3Aplc%3As3storage/#{cid}" + assert {:ok, "s3 bytes", conn} = Plug.Conn.read_body(conn) + Plug.Conn.send_resp(conn, 200, "") + end) + + Req.Test.expect(__MODULE__, fn conn -> + assert conn.method == "DELETE" + assert conn.request_path == "/tempest-test/temp/blobs/did%3Aplc%3As3storage/#{cid}" + Plug.Conn.send_resp(conn, 204, "") + end) + + assert {:ok, promoted_path} = S3Storage.promote_blob(config, did, cid) + assert promoted_path == "blobs/#{did}/#{cid}" + end + + test "promote_blob reports missing only when fallback read cannot find the temp object", %{ + config: config, + did: did, + cid: cid + } do + Req.Test.expect(__MODULE__, fn conn -> + assert conn.method == "PUT" + assert conn.request_path == "/tempest-test/blobs/did%3Aplc%3As3storage/#{cid}" + Plug.Conn.send_resp(conn, 404, "") + end) + + Req.Test.expect(__MODULE__, fn conn -> + assert conn.method == "GET" + assert conn.request_path == "/tempest-test/temp/blobs/did%3Aplc%3As3storage/#{cid}" + Plug.Conn.send_resp(conn, 404, "") + end) + + assert {:error, :blob_not_found} = S3Storage.promote_blob(config, did, cid) + end + test "get_blob reads public object bytes", %{config: config, did: did, cid: cid} do Req.Test.expect(__MODULE__, fn conn -> assert conn.method == "GET"