Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
51 commits
Select commit Hold shift + click to select a range
b750116
Retry on timeouts
cdegroot Feb 15, 2024
143db2f
Add Aggregate State rebuilding telemetry
cdegroot Feb 15, 2024
a58d242
Put telemetry on dehydration
cdegroot Feb 15, 2024
85c97c3
Mix format
cdegroot Feb 15, 2024
2f67dd6
Setup for CalmWave development (#1)
thomasdziedzic-calmwave Aug 25, 2024
3c084bd
aggregate improvements around retries and telemetry (#2)
thomasdziedzic-calmwave Aug 25, 2024
f9d2280
Pull in initial work by @davydog187.
davydog187 Jan 27, 2022
110f19b
Redefine how acknowledgement works
Aug 17, 2022
450ee05
Drop support for {:error, reason, event}
Sep 18, 2024
fb46ba0
Update docs
Sep 18, 2024
c77eb04
Use delegate_event_to_handler & make confirm_receipt be more generic
Sep 20, 2024
69ce480
Do not retry :skip events
Sep 20, 2024
5c50d3c
Batching support (#3)
fmterrorf Sep 20, 2024
b031bdd
Merge pull request #4 from commanded/master
Jan 15, 2025
134606f
Ran devbox install
Feb 24, 2025
fb7dc0f
Merge branch 'master' into main-calmwave-1.4.7
Feb 24, 2025
a7f486b
Merge branch 'commanded:master' into main-calmwave-1.4.7
cblage Feb 24, 2025
69b4a26
Merge pull request #6 from calmwave-open-source/main-calmwave-1.4.7
Feb 24, 2025
7fd286d
Merge branch 'main-calmwave' into batching-support-updated
Feb 24, 2025
b2ee408
Fix test errors, logging, 1 broken test
Feb 25, 2025
9d199c8
Tests running
Feb 26, 2025
a206f79
Remove TODO
Feb 26, 2025
19bfa9c
Merge pull request #552 from calmwave-open-source/cees/eng-757-invest…
cdegroot Feb 26, 2025
144dd5c
Merge pull request #7 from commanded/master
Feb 26, 2025
d61b6cd
Merge pull request #8 from calmwave-open-source/update-with-master
Feb 26, 2025
2be16d5
Merge branch 'main-calmwave' into batching-support
Feb 26, 2025
58d6dc8
Revert "Retry on timeouts"
drteeth Feb 26, 2025
9ddb686
Update credo and ex_doc
drteeth Feb 26, 2025
b14ef46
Fix test
drteeth Feb 26, 2025
60073ed
Merge pull request #9 from commanded/master
Feb 27, 2025
2e11e27
Merge branch 'main-calmwave' into batching-support
Feb 27, 2025
12cc80b
Add serializer behaviour
Nezteb May 18, 2025
01438fc
Add aggregate behaviour
Nezteb May 18, 2025
59f3450
Cleanup
Nezteb May 18, 2025
166e703
chore: add issue template
yordis Jun 3, 2025
cb99dc9
Make execute/2 callback optional
Nezteb Jun 3, 2025
37d9c67
Merge pull request #630 from yordis/add-issue-template
drteeth Jun 3, 2025
047a2a5
Merge pull request #626 from Nezteb/serializer-behaviour
drteeth Jun 24, 2025
424b359
Merge pull request #627 from Nezteb/aggregate-behaviour
drteeth Jun 24, 2025
e8de0e6
Merge pull request #569 from calmwave-open-source/batching-support
cdegroot Jun 25, 2025
862fca1
Remove devbox stuff that accidentally got merged in
cdegroot Jun 25, 2025
b075bab
Merge pull request #635 from commanded/cdg/cleanup-devbox
cdegroot Jun 25, 2025
75b9937
Fix formatting for older, crankier versions of Elixir
drteeth Jul 28, 2025
c4649c1
Update for 1.18/28 support
drteeth Jul 28, 2025
f5b17ce
Formatting
drteeth Aug 6, 2025
dbb29f1
Suppress CallWithoutOpaque.format_long(nil) until it is sorted out up…
drteeth Aug 6, 2025
bdaddcc
remove FIXME
drteeth Aug 6, 2025
09e6a76
Formatting
drteeth Aug 6, 2025
3c9e3ad
Make Aggregate.take_snapshot/3 a blocking call
satom99 Jul 17, 2025
52171b8
Merge pull request #636 from kantox/KA-3069-make-aggregate-take-snaps…
drteeth Aug 13, 2025
4b828a3
Merge branch 'master' of github.com:commanded/commanded into sync-up
yordis Aug 24, 2025
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
48 changes: 31 additions & 17 deletions lib/commanded/aggregates/aggregate.ex
Original file line number Diff line number Diff line change
Expand Up @@ -274,9 +274,9 @@ defmodule Commanded.Aggregates.Aggregate do
end

@doc false
def take_snapshot(application, aggregate_module, aggregate_uuid) do
def take_snapshot(application, aggregate_module, aggregate_uuid, timeout \\ 5_000) do
name = via_name(application, aggregate_module, aggregate_uuid)
GenServer.cast(name, :take_snapshot)
GenServer.call(name, :take_snapshot, timeout)
end

@doc false
Expand Down Expand Up @@ -315,11 +315,14 @@ defmodule Commanded.Aggregates.Aggregate do

@doc false
@impl GenServer
def handle_cast(:take_snapshot, %Aggregate{} = state), do: do_take_snapshot(state)
def handle_call(:take_snapshot, _from, %Aggregate{} = state) do
case do_take_snapshot(state) do
{:ok, state} ->
{:reply, :ok, state}

@impl GenServer
def handle_cast({:take_snapshot, lifespan_timeout}, %Aggregate{} = state) do
do_take_snapshot(%Aggregate{state | lifespan_timeout: lifespan_timeout})
{:error, _reason} = error ->
{:reply, error, state}
end
end

@doc false
Expand Down Expand Up @@ -353,7 +356,7 @@ defmodule Commanded.Aggregates.Aggregate do

response =
if Snapshotting.snapshot_required?(snapshotting, aggregate_version) do
:ok = GenServer.cast(self(), {:take_snapshot, lifespan_timeout})
send(self(), {:take_snapshot, lifespan_timeout})

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

@doc false
@impl GenServer
def handle_info({:take_snapshot, lifespan_timeout}, %Aggregate{} = state) do
state = %Aggregate{state | lifespan_timeout: lifespan_timeout}

case do_take_snapshot(state) do
{:ok, state} ->
noreply_with_lifespan(state)

{:error, _error} ->
noreply_with_lifespan(state)
end
end

@doc false
@impl GenServer
def handle_info(:timeout, %Aggregate{} = state) do
Expand Down Expand Up @@ -622,18 +639,15 @@ defmodule Commanded.Aggregates.Aggregate do

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

state =
case Snapshotting.take_snapshot(snapshotting, aggregate_version, aggregate_state) do
{:ok, snapshotting} ->
%Aggregate{state | snapshotting: snapshotting}
case Snapshotting.take_snapshot(snapshotting, aggregate_version, aggregate_state) do
{:ok, snapshotting} ->
{:ok, %Aggregate{state | snapshotting: snapshotting}}

{:error, error} ->
Logger.warning(describe(state) <> " snapshot failed due to: " <> inspect(error))
{:error, reason} = error ->
Logger.warning(describe(state) <> " snapshot failed due to: " <> inspect(reason))

state
end

noreply_with_lifespan(state)
error
end
end

defp telemetry_wrong_expected_version(context, from, state) do
Expand Down
14 changes: 14 additions & 0 deletions lib/commanded/commands/dispatcher.ex
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,20 @@ defmodule Commanded.Commands.Dispatcher do
struct(Pipeline, Map.from_struct(payload))
end

# Ignoring this dialyzer warning until dialyxir catches up on Elixir 1.18.x-otp-28
# Please file a bug in https://github.com/jeremyjh/dialyxir/issues with this message.
# Unknown error occurred:
# %FunctionClauseError{
# module: Dialyxir.Warnings.CallWithoutOpaque,
# function: :format_long,
# arity: 1,
# kind: nil,
# args: nil,
# clauses: nil
# }
# Open issue in dialyxir:
# https://github.com/jeremyjh/dialyxir/issues/568
@dialyzer {:nowarn_function, execute: 3}
defp execute(%Pipeline{} = pipeline, %Payload{} = payload, %ExecutionContext{} = context) do
%Pipeline{application: application, assigns: %{aggregate_uuid: aggregate_uuid}} = pipeline
%Payload{aggregate_module: aggregate_module, timeout: timeout} = payload
Expand Down
5 changes: 4 additions & 1 deletion test/application/application_supervisor_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,9 @@ defmodule Commanded.Aggregates.ApplicationSupervisorTest do
end

defp assert_hibernated(pid) do
assert Process.info(pid, :current_function) == {:current_function, {:erlang, :hibernate, 3}}
pre_otp_28 = {:current_function, {:erlang, :hibernate, 3}}
post_otp_28 = {:current_function, {:gen_server, :loop_hibernate, 4}}

assert Process.info(pid, :current_function) in [pre_otp_28, post_otp_28]
end
end
3 changes: 1 addition & 2 deletions test/event/event_handler_batch_telemetry_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -85,8 +85,7 @@ defmodule Commanded.Event.EventHandlerBatchTelemetryTest do
recorded_events = EventFactory.map_to_recorded_events(events, 1, metadata: metadata)
state = setup_state(ErrorHandlingBatchHandler)

{:stop, _, _} =
Handler.handle_info({:events, recorded_events}, state)
{:stop, _, _} = Handler.handle_info({:events, recorded_events}, state)

assert_receive {[:commanded, :event, :batch, :start], _measurements, _metadata}
refute_received {[:commanded, :event, :batch, :stop], _measurements, _metadata}
Expand Down
5 changes: 4 additions & 1 deletion test/event/handler_init_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -116,7 +116,10 @@ defmodule Commanded.Event.HandlerInitTest do
end

defp assert_hibernated(pid) do
assert Process.info(pid, :current_function) == {:current_function, {:erlang, :hibernate, 3}}
pre_otp_28 = {:current_function, {:erlang, :hibernate, 3}}
post_otp_28 = {:current_function, {:gen_server, :loop_hibernate, 4}}

assert Process.info(pid, :current_function) in [pre_otp_28, post_otp_28]
end

defp send_subscribed(handler) do
Expand Down
Loading