Skip to content

Commit 57d0b6a

Browse files
authored
fix: make take snapshot sync call (#20)
1 parent 437287f commit 57d0b6a

5 files changed

Lines changed: 54 additions & 21 deletions

File tree

‎lib/commanded/aggregates/aggregate.ex‎

Lines changed: 31 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -274,9 +274,9 @@ defmodule Commanded.Aggregates.Aggregate do
274274
end
275275

276276
@doc false
277-
def take_snapshot(application, aggregate_module, aggregate_uuid) do
277+
def take_snapshot(application, aggregate_module, aggregate_uuid, timeout \\ 5_000) do
278278
name = via_name(application, aggregate_module, aggregate_uuid)
279-
GenServer.cast(name, :take_snapshot)
279+
GenServer.call(name, :take_snapshot, timeout)
280280
end
281281

282282
@doc false
@@ -315,11 +315,14 @@ defmodule Commanded.Aggregates.Aggregate do
315315

316316
@doc false
317317
@impl GenServer
318-
def handle_cast(:take_snapshot, %Aggregate{} = state), do: do_take_snapshot(state)
318+
def handle_call(:take_snapshot, _from, %Aggregate{} = state) do
319+
case do_take_snapshot(state) do
320+
{:ok, state} ->
321+
{:reply, :ok, state}
319322

320-
@impl GenServer
321-
def handle_cast({:take_snapshot, lifespan_timeout}, %Aggregate{} = state) do
322-
do_take_snapshot(%Aggregate{state | lifespan_timeout: lifespan_timeout})
323+
{:error, _reason} = error ->
324+
{:reply, error, state}
325+
end
323326
end
324327

325328
@doc false
@@ -353,7 +356,7 @@ defmodule Commanded.Aggregates.Aggregate do
353356

354357
response =
355358
if Snapshotting.snapshot_required?(snapshotting, aggregate_version) do
356-
:ok = GenServer.cast(self(), {:take_snapshot, lifespan_timeout})
359+
send(self(), {:take_snapshot, lifespan_timeout})
357360

358361
# Don't reply with a lifetime because we just asked for a snapshot to
359362
# be taken. When it finishes, it will set the timeout.
@@ -410,6 +413,20 @@ defmodule Commanded.Aggregates.Aggregate do
410413
end
411414
end
412415

416+
@doc false
417+
@impl GenServer
418+
def handle_info({:take_snapshot, lifespan_timeout}, %Aggregate{} = state) do
419+
state = %Aggregate{state | lifespan_timeout: lifespan_timeout}
420+
421+
case do_take_snapshot(state) do
422+
{:ok, state} ->
423+
noreply_with_lifespan(state)
424+
425+
{:error, _error} ->
426+
noreply_with_lifespan(state)
427+
end
428+
end
429+
413430
@doc false
414431
@impl GenServer
415432
def handle_info(:timeout, %Aggregate{} = state) do
@@ -622,18 +639,15 @@ defmodule Commanded.Aggregates.Aggregate do
622639

623640
Logger.debug(describe(state) <> " recording snapshot")
624641

625-
state =
626-
case Snapshotting.take_snapshot(snapshotting, aggregate_version, aggregate_state) do
627-
{:ok, snapshotting} ->
628-
%Aggregate{state | snapshotting: snapshotting}
642+
case Snapshotting.take_snapshot(snapshotting, aggregate_version, aggregate_state) do
643+
{:ok, snapshotting} ->
644+
{:ok, %Aggregate{state | snapshotting: snapshotting}}
629645

630-
{:error, error} ->
631-
Logger.warning(describe(state) <> " snapshot failed due to: " <> inspect(error))
646+
{:error, reason} = error ->
647+
Logger.warning(describe(state) <> " snapshot failed due to: " <> inspect(reason))
632648

633-
state
634-
end
635-
636-
noreply_with_lifespan(state)
649+
error
650+
end
637651
end
638652

639653
defp telemetry_wrong_expected_version(context, from, state) do

‎lib/commanded/commands/dispatcher.ex‎

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -67,6 +67,20 @@ defmodule Commanded.Commands.Dispatcher do
6767
struct(Pipeline, Map.from_struct(payload))
6868
end
6969

70+
# Ignoring this dialyzer warning until dialyxir catches up on Elixir 1.18.x-otp-28
71+
# Please file a bug in https://github.com/jeremyjh/dialyxir/issues with this message.
72+
# Unknown error occurred:
73+
# %FunctionClauseError{
74+
# module: Dialyxir.Warnings.CallWithoutOpaque,
75+
# function: :format_long,
76+
# arity: 1,
77+
# kind: nil,
78+
# args: nil,
79+
# clauses: nil
80+
# }
81+
# Open issue in dialyxir:
82+
# https://github.com/jeremyjh/dialyxir/issues/568
83+
@dialyzer {:nowarn_function, execute: 3}
7084
defp execute(%Pipeline{} = pipeline, %Payload{} = payload, %ExecutionContext{} = context) do
7185
%Pipeline{application: application, assigns: %{aggregate_uuid: aggregate_uuid}} = pipeline
7286
%Payload{aggregate_module: aggregate_module, timeout: timeout} = payload

‎test/application/application_supervisor_test.exs‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,9 @@ defmodule Commanded.Aggregates.ApplicationSupervisorTest do
2828
end
2929

3030
defp assert_hibernated(pid) do
31-
assert Process.info(pid, :current_function) == {:current_function, {:erlang, :hibernate, 3}}
31+
pre_otp_28 = {:current_function, {:erlang, :hibernate, 3}}
32+
post_otp_28 = {:current_function, {:gen_server, :loop_hibernate, 4}}
33+
34+
assert Process.info(pid, :current_function) in [pre_otp_28, post_otp_28]
3235
end
3336
end

‎test/event/event_handler_batch_telemetry_test.exs‎

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -85,8 +85,7 @@ defmodule Commanded.Event.EventHandlerBatchTelemetryTest do
8585
recorded_events = EventFactory.map_to_recorded_events(events, 1, metadata: metadata)
8686
state = setup_state(ErrorHandlingBatchHandler)
8787

88-
{:stop, _, _} =
89-
Handler.handle_info({:events, recorded_events}, state)
88+
{:stop, _, _} = Handler.handle_info({:events, recorded_events}, state)
9089

9190
assert_receive {[:commanded, :event, :batch, :start], _measurements, _metadata}
9291
refute_received {[:commanded, :event, :batch, :stop], _measurements, _metadata}

‎test/event/handler_init_test.exs‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -116,7 +116,10 @@ defmodule Commanded.Event.HandlerInitTest do
116116
end
117117

118118
defp assert_hibernated(pid) do
119-
assert Process.info(pid, :current_function) == {:current_function, {:erlang, :hibernate, 3}}
119+
pre_otp_28 = {:current_function, {:erlang, :hibernate, 3}}
120+
post_otp_28 = {:current_function, {:gen_server, :loop_hibernate, 4}}
121+
122+
assert Process.info(pid, :current_function) in [pre_otp_28, post_otp_28]
120123
end
121124

122125
defp send_subscribed(handler) do

0 commit comments

Comments
 (0)