Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
64 changes: 40 additions & 24 deletions lib/durable_server/backends/object_store.ex
Original file line number Diff line number Diff line change
Expand Up @@ -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}} <-

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

we are pinning on ^encoded but shouldn't we be matching on our own etag instead?

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},
Expand All @@ -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

Expand Down
4 changes: 3 additions & 1 deletion lib/durable_server/object_store.ex
Original file line number Diff line number Diff line change
Expand Up @@ -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 ->
Expand Down
200 changes: 200 additions & 0 deletions test/durable_server/object_store_claim_test.exs
Original file line number Diff line number Diff line change
@@ -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