diff --git a/lib/durable_server/backends/object_store.ex b/lib/durable_server/backends/object_store.ex index e6bd9ba..5a562ed 100644 --- a/lib/durable_server/backends/object_store.ex +++ b/lib/durable_server/backends/object_store.ex @@ -84,36 +84,35 @@ defmodule DurableServer.Backends.ObjectStore do defp resolve_ambiguous_conditional_put( %ObjectStore{} = store, key, - %StoredState{meta: %Meta{} = attempted_meta} = data, + data, encoded, opts ) do - if Keyword.has_key?(opts, :etag) do - case ObjectStore.get_object(store, key, consistent: true) do - {:ok, %{body: ^encoded, etag: etag}} -> - with {:ok, %StoredState{meta: %Meta{} = persisted_meta}} <- decode_body(encoded), - true <- same_boot_owner?(attempted_meta, persisted_meta) do - {:ok, %{body: data, etag: etag}} - else - _other -> {:error, :conflict} - end - - _other -> - {:error, :conflict} - end + with true <- Keyword.has_key?(opts, :etag), + {:ok, etag} <- read_matching_owned_state(store, key, data, encoded) do + {:ok, %{body: data, etag: etag}} else - {:error, :conflict} + _other -> {:error, :conflict} end end - defp resolve_ambiguous_conditional_put( - %ObjectStore{}, - _key, - _data, - _encoded, - _opts - ), - do: {:error, :conflict} + defp read_matching_owned_state( + store, + key, + %StoredState{meta: %Meta{} = attempted_meta}, + encoded + ) do + with {:ok, %{body: ^encoded, etag: etag}} <- + ObjectStore.get_object(store, key, consistent: true), + {:ok, %StoredState{meta: %Meta{} = persisted_meta}} <- decode_body(encoded), + true <- same_boot_owner?(attempted_meta, persisted_meta) do + {:ok, etag} + else + _other -> :error + end + end + + defp read_matching_owned_state(_store, _key, _data, _encoded), do: :error defp same_boot_owner?( %Meta{pid: pid, node_ref: node_ref, node_str: node_str}, @@ -137,7 +136,24 @@ defmodule DurableServer.Backends.ObjectStore do @impl true def try_claim(%ObjectStore{} = store, key, body) do with {:ok, encoded} <- encode_body(body) do - ObjectStore.try_claim(store, key, encoded) + case ObjectStore.try_claim(store, key, encoded) do + {:error, :already_claimed} -> + resolve_ambiguous_claim(store, key, body, encoded, :already_claimed) + + {:error, %Req.TransportError{} = reason} -> + resolve_ambiguous_claim(store, key, body, encoded, reason) + + result -> + result + end + end + end + + # A lost claim response is recoverable only for this boot's exact stored state. + defp resolve_ambiguous_claim(store, key, body, encoded, reason) do + case read_matching_owned_state(store, key, body, encoded) do + {:ok, etag} -> {:ok, {:claimed, etag}} + :error -> {:error, reason} end end diff --git a/lib/durable_server/object_store.ex b/lib/durable_server/object_store.ex index 9f173c0..f53a05f 100644 --- a/lib/durable_server/object_store.ex +++ b/lib/durable_server/object_store.ex @@ -753,7 +753,9 @@ defmodule DurableServer.ObjectStore do method: :put, url: key, body: body, - retry: false + retry: :transient, + max_retries: 2, + receive_timeout: @default_retry_attempt_timeout ) do {:ok, %{status: status, headers: headers}} when status >= 200 and status < 300 -> diff --git a/test/durable_server/object_store_claim_test.exs b/test/durable_server/object_store_claim_test.exs new file mode 100644 index 0000000..c0d9929 --- /dev/null +++ b/test/durable_server/object_store_claim_test.exs @@ -0,0 +1,200 @@ +defmodule DurableServer.ObjectStoreClaimTest do + use ExUnit.Case, async: true + + alias DurableServer.Backends.ObjectStore, as: ObjectStoreBackend + alias DurableServer.{Meta, ObjectStore, StorageBackend, StoredState} + + @moduletag :capture_log + + describe "initial claim retries" do + test "retries transient failures with fresh signatures" do + store = + object_store([ + %Req.TransportError{reason: :closed}, + %Req.Response{status: 503}, + %Req.Response{status: 200, headers: %{"etag" => ["claimed-etag"]}} + ]) + + assert {:ok, {:claimed, "claimed-etag"}} = + ObjectStore.try_claim(store, "server/one", "state") + + for _attempt <- 1..3 do + assert_receive {:request, %Req.Request{method: :put} = request} + assert Req.Request.get_header(request, "if-none-match") == ["*"] + assert [authorization] = Req.Request.get_header(request, "authorization") + refute authorization =~ "SignedHeaders=authorization;" + end + end + + test "stops after two retries" do + store = object_store(List.duplicate(%Req.TransportError{reason: :closed}, 3)) + + assert {:error, %Req.TransportError{reason: :closed}} = + ObjectStore.try_claim(store, "server/one", "state") + + for _attempt <- 1..3 do + assert_receive {:request, %Req.Request{method: :put}} + end + end + + test "does not retry a permanent response" do + store = object_store([%Req.Response{status: 403}]) + + assert {:error, %Req.Response{status: 403}} = + ObjectStore.try_claim(store, "server/one", "state") + end + + test "does not retry a claim conflict" do + store = object_store([%Req.Response{status: 412}]) + + assert {:error, :already_claimed} = ObjectStore.try_claim(store, "server/one", "state") + end + end + + describe "initial claim recovery" do + test "adopts the committed claim when a lost response is followed by a conflict" do + attempted = stored_state(%{value: 1}) + + store = + object_store([ + %Req.TransportError{reason: :closed}, + %Req.Response{status: 412}, + stored_response(attempted) + ]) + + backend = StorageBackend.new(ObjectStoreBackend, store) + + assert {:ok, {:claimed, "committed-etag"}} = + StorageBackend.try_claim(backend, "server/one", attempted) + + assert_receive {:request, %Req.Request{method: :put}} + assert_receive {:request, %Req.Request{method: :put}} + assert_receive {:request, %Req.Request{method: :get} = read} + assert Req.Request.get_header(read, "x-tigris-consistent") == ["true"] + end + + test "adopts the committed claim when all PUT responses are lost" do + attempted = stored_state(%{value: 1}) + + store = + object_store([ + %Req.TransportError{reason: :closed}, + %Req.TransportError{reason: :closed}, + %Req.TransportError{reason: :closed}, + stored_response(attempted) + ]) + + backend = StorageBackend.new(ObjectStoreBackend, store) + + assert {:ok, {:claimed, "committed-etag"}} = + StorageBackend.try_claim(backend, "server/one", attempted) + end + + test "does not adopt a claim owned by another boot" do + attempted = stored_state(%{value: 1}) + persisted = stored_state(%{value: 1}, %{node_ref: 124}) + store = object_store([%Req.Response{status: 412}, stored_response(persisted)]) + backend = StorageBackend.new(ObjectStoreBackend, store) + + assert {:error, :already_claimed} = + StorageBackend.try_claim(backend, "server/one", attempted) + end + + test "does not adopt different state owned by the same boot" do + attempted = stored_state(%{value: 1}) + persisted = stored_state(%{value: 2}) + store = object_store([%Req.Response{status: 412}, stored_response(persisted)]) + backend = StorageBackend.new(ObjectStoreBackend, store) + + assert {:error, :already_claimed} = + StorageBackend.try_claim(backend, "server/one", attempted) + end + + for field <- [:pid, :node_ref, :node_str] do + test "does not adopt a matching claim without #{field}" do + attempted = stored_state(%{value: 1}, %{unquote(field) => nil}) + store = object_store([%Req.Response{status: 412}, stored_response(attempted)]) + backend = StorageBackend.new(ObjectStoreBackend, store) + + assert {:error, :already_claimed} = + StorageBackend.try_claim(backend, "server/one", attempted) + end + end + + test "does not adopt generic values after a conflict" do + store = object_store([%Req.Response{status: 412}]) + backend = StorageBackend.new(ObjectStoreBackend, store) + + assert {:error, :already_claimed} = + StorageBackend.try_claim(backend, "server/one", %{value: 1}) + end + + test "preserves the transport error when the claim was not stored" do + store = + object_store([ + %Req.TransportError{reason: :closed}, + %Req.TransportError{reason: :closed}, + %Req.TransportError{reason: :closed}, + %Req.Response{status: 404} + ]) + + backend = StorageBackend.new(ObjectStoreBackend, store) + + assert {:error, %Req.TransportError{reason: :closed}} = + StorageBackend.try_claim(backend, "server/one", stored_state(%{value: 1})) + end + + test "preserves the conflict when reading the claim fails" do + store = object_store([%Req.Response{status: 412}, %Req.Response{status: 403}]) + backend = StorageBackend.new(ObjectStoreBackend, store) + + assert {:error, :already_claimed} = + StorageBackend.try_claim(backend, "server/one", stored_state(%{value: 1})) + end + end + + defp stored_state(state, owner_overrides \\ %{}) do + meta = + struct!( + Meta, + Map.merge( + %{ + module: __MODULE__, + permanent: false, + pid: self(), + status: :running, + node_ref: 123, + node_str: "test@node" + }, + owner_overrides + ) + ) + + %StoredState{vsn: 1, state: state, meta: meta} + end + + defp stored_response(state) do + {:ok, encoded} = ObjectStoreBackend.encode(%ObjectStore{}, state) + %Req.Response{status: 200, body: encoded, headers: %{"etag" => ["committed-etag"]}} + end + + defp object_store(responses) do + parent = self() + responses = start_supervised!({Agent, fn -> responses end}) + + adapter = fn request -> + send(parent, {:request, request}) + response = Agent.get_and_update(responses, fn [response | rest] -> {response, rest} end) + {request, response} + end + + ObjectStore.new( + bucket: "test-bucket", + access_key_id: "test-access-key", + secret_access_key: "test-secret-key", + s3_endpoint: "http://s3.test", + default_region: "us-east-1", + req_opts: [adapter: adapter, retry_delay: 0, retry_log_level: false] + ) + end +end