Skip to content

Commit 353e939

Browse files
authored
refactor(dispatcher): align aggregate failure boundary (#98)
Signed-off-by: Yordis Prieto <yordis.prieto@gmail.com>
1 parent 7a07792 commit 353e939

10 files changed

Lines changed: 168 additions & 47 deletions

File tree

‎guides/explanations/commands.md‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -193,7 +193,7 @@ open_account = %OpenAccount{
193193

194194
A command handler has a default timeout of 5 seconds, the same default as a `GenServer.call/3` process call. It must handle the command in this period; otherwise, dispatch fails with an error tuple.
195195

196-
Depending on where the timeout is observed, this may currently be reported as either `{:error, :aggregate_execution_timeout}` or `{:error, :aggregate_execution_failed}`. Both indicate that command execution did not complete within the configured timeout.
196+
If the command does not complete within the configured timeout, dispatch fails with `{:error, :aggregate_execution_timeout}`.
197197

198198
You can configure a different timeout value during command registration by providing a `timeout` option, defined in milliseconds:
199199

‎guides/explanations/fork-differences.md‎

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -52,6 +52,21 @@ This fork maintains an independent release cycle to introduce new features and i
5252

5353
## Features Added ✨
5454

55+
### **Dispatcher Failure Boundary Normalization**
56+
[PR #98](https://github.com/straw-hat-team/commanded/pull/98)
57+
58+
**Changes:**
59+
- Moved aggregate execution exit normalization into `Commanded.Aggregates.Aggregate.execute/5`
60+
- Removed the dispatcher task wrapper used only for `GenServer.call/3` failure isolation
61+
- Made command timeouts deterministic as `{:error, :aggregate_execution_timeout}`
62+
- Preserved abnormal aggregate exit reasons as `{:error, :aggregate_execution_failed, reason}` at the aggregate execution boundary
63+
- Kept public dispatch failures as `{:error, :aggregate_execution_failed}` while preserving `reason` internally for middleware and logging
64+
65+
**Benefits:**
66+
- Keeps failure normalization in one place instead of splitting it across aggregate and dispatcher layers
67+
- Preserves dispatcher safety guarantees without extra runtime infrastructure
68+
- Gives middleware and logs better diagnostic context for abnormal aggregate exits without widening the public dispatch API
69+
5570
### **UUIDv7 Support**
5671
[PR #22](https://github.com/straw-hat-team/commanded/pull/22)
5772

‎lib/commanded/aggregates/aggregate.ex‎

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -258,6 +258,17 @@ defmodule Commanded.Aggregates.Aggregate do
258258
replies, this function returns `{:exit, {:normal, :aggregate_stopped}}`.
259259
`Commanded.Commands.Dispatcher` handles that case by retrying when permitted.
260260
261+
Infrastructure failures are normalized into return values instead of exiting
262+
the caller:
263+
264+
- `{:error, :remote_node_down}` when a remote aggregate node becomes
265+
unavailable during execution.
266+
- `{:error, :aggregate_execution_timeout}` when the command does not
267+
complete within the configured timeout.
268+
- `{:error, :aggregate_execution_failed, reason}` for other
269+
execution-time exits, preserving the underlying exit reason for
270+
middleware and logging.
271+
261272
- `aggregate_version` - the updated version of the aggregate after executing
262273
the command.
263274
- `events` - events produced by the command, can be an empty list.
@@ -282,6 +293,20 @@ defmodule Commanded.Aggregates.Aggregate do
282293

283294
:exit, {:normal, {GenServer, :call, [^name, {:execute_command, ^context}, ^timeout]}} ->
284295
{:exit, {:normal, :aggregate_stopped}}
296+
297+
:exit,
298+
{{:nodedown, _node_name},
299+
{GenServer, :call, [^name, {:execute_command, ^context}, ^timeout]}} ->
300+
{:error, :remote_node_down}
301+
302+
:exit, {:timeout, {GenServer, :call, [^name, {:execute_command, ^context}, ^timeout]}} ->
303+
{:error, :aggregate_execution_timeout}
304+
305+
:exit, {reason, {GenServer, :call, [^name, {:execute_command, ^context}, ^timeout]}} ->
306+
{:error, :aggregate_execution_failed, reason}
307+
308+
:exit, reason ->
309+
{:error, :aggregate_execution_failed, reason}
285310
end
286311
end
287312

‎lib/commanded/application/supervisor.ex‎

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -67,15 +67,13 @@ defmodule Commanded.Application.Supervisor do
6767
end
6868

6969
defp app_child_spec(name, config) do
70-
task_dispatcher_name = Module.concat([name, Commanded.Commands.TaskDispatcher])
7170
aggregates_supervisor_name = Module.concat([name, Commanded.Aggregates.Supervisor])
7271
subscriptions_name = Module.concat([name, Commanded.Subscriptions])
7372
registry_name = Module.concat([name, Commanded.Subscriptions.Registry])
7473
snapshotting = Keyword.get(config, :snapshotting, %{})
7574
hibernate_after = Keyword.get(config, :hibernate_after, :infinity)
7675

7776
[
78-
{Task.Supervisor, name: task_dispatcher_name, hibernate_after: hibernate_after},
7977
{Commanded.Aggregates.Supervisor,
8078
name: aggregates_supervisor_name,
8179
application: name,

‎lib/commanded/commands/dispatcher.ex‎

Lines changed: 1 addition & 28 deletions
Original file line numberDiff line numberDiff line change
@@ -120,34 +120,7 @@ defmodule Commanded.Commands.Dispatcher do
120120
context.metadata
121121
)
122122

123-
task_dispatcher_name = Module.concat([application, Commanded.Commands.TaskDispatcher])
124-
125-
task =
126-
Task.Supervisor.async_nolink(task_dispatcher_name, Aggregate, :execute, [
127-
application,
128-
aggregate_module,
129-
aggregate_uuid,
130-
context,
131-
timeout
132-
])
133-
134-
result =
135-
case Task.yield(task, timeout) || Task.shutdown(task) do
136-
{:ok, result} ->
137-
result
138-
139-
{:exit, {:normal, :aggregate_stopped}} = result ->
140-
result
141-
142-
{:exit, {{:nodedown, _node_name}, {GenServer, :call, _}}} ->
143-
{:error, :remote_node_down}
144-
145-
{:exit, _reason} ->
146-
{:error, :aggregate_execution_failed}
147-
148-
nil ->
149-
{:error, :aggregate_execution_timeout}
150-
end
123+
result = Aggregate.execute(application, aggregate_module, aggregate_uuid, context, timeout)
151124

152125
case result do
153126
{:ok, aggregate_version, events, aggregate_state} ->

‎test/aggregates/execute_command_test.exs‎

Lines changed: 81 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,14 +4,32 @@ defmodule Commanded.Aggregates.ExecuteCommandTest do
44
import Commanded.Helpers.ProcessHelper, only: [shutdown_aggregate: 3]
55

66
alias Commanded.Aggregates.{Aggregate, ExecutionContext}
7+
alias Commanded.Commands.{TimeoutAggregateRoot, TimeoutCommand, TimeoutCommandHandler}
78
alias Commanded.ExampleDomain.{BankAccount, BankApp, OpenAccountHandler}
89
alias Commanded.ExampleDomain.BankAccount.Commands.OpenAccount
910
alias Commanded.ExampleDomain.BankAccount.Events.BankAccountOpened
1011
alias Commanded.Helpers.Wait
11-
alias Commanded.{Registration, UUID}
12+
alias Commanded.TestSupport.RetryStopOnceAggregate
13+
alias Commanded.TestSupport.RetryStopOnceAggregate.Command, as: RetryStopOnceCommand
14+
alias Commanded.{DefaultApp, Registration, UUID}
15+
16+
defmodule CrashCommand do
17+
defstruct [:uuid]
18+
end
19+
20+
defmodule CrashBeforeExecuteAggregate do
21+
alias Commanded.Aggregates.ExecutionContext
22+
23+
defstruct []
24+
25+
def before_execute(_aggregate_state, %ExecutionContext{}), do: Process.exit(self(), :boom)
26+
def execute(%__MODULE__{}, %CrashCommand{}), do: []
27+
end
1228

1329
setup do
1430
start_supervised!(BankApp)
31+
start_supervised!(DefaultApp)
32+
start_supervised!(RetryStopOnceAggregate.Tracker)
1533

1634
:ok
1735
end
@@ -100,6 +118,68 @@ defmodule Commanded.Aggregates.ExecuteCommandTest do
100118
assert state_before == Aggregate.aggregate_state(BankApp, BankAccount, account_number)
101119
end
102120

121+
test "returns aggregate_stopped when aggregate stops after being opened" do
122+
aggregate_uuid = UUID.uuid4()
123+
124+
assert {:ok, ^aggregate_uuid} =
125+
Commanded.Aggregates.Supervisor.open_aggregate(
126+
DefaultApp,
127+
RetryStopOnceAggregate,
128+
aggregate_uuid
129+
)
130+
131+
context = %ExecutionContext{
132+
command: %RetryStopOnceCommand{uuid: aggregate_uuid},
133+
handler: RetryStopOnceAggregate,
134+
function: :execute,
135+
before_execute: :before_execute
136+
}
137+
138+
assert {:exit, {:normal, :aggregate_stopped}} =
139+
Aggregate.execute(DefaultApp, RetryStopOnceAggregate, aggregate_uuid, context)
140+
end
141+
142+
test "returns aggregate_execution_timeout when command execution exceeds the timeout" do
143+
aggregate_uuid = UUID.uuid4()
144+
145+
assert {:ok, ^aggregate_uuid} =
146+
Commanded.Aggregates.Supervisor.open_aggregate(
147+
DefaultApp,
148+
TimeoutAggregateRoot,
149+
aggregate_uuid
150+
)
151+
152+
context = %ExecutionContext{
153+
command: %TimeoutCommand{aggregate_uuid: aggregate_uuid, sleep_in_ms: 200},
154+
handler: TimeoutCommandHandler,
155+
function: :handle
156+
}
157+
158+
assert {:error, :aggregate_execution_timeout} =
159+
Aggregate.execute(DefaultApp, TimeoutAggregateRoot, aggregate_uuid, context, 50)
160+
end
161+
162+
test "returns aggregate_execution_failed with the exit reason when the aggregate exits abnormally" do
163+
aggregate_uuid = UUID.uuid4()
164+
165+
assert {:ok, ^aggregate_uuid} =
166+
Commanded.Aggregates.Supervisor.open_aggregate(
167+
DefaultApp,
168+
CrashBeforeExecuteAggregate,
169+
aggregate_uuid
170+
)
171+
172+
context = %ExecutionContext{
173+
command: %CrashCommand{uuid: aggregate_uuid},
174+
handler: CrashBeforeExecuteAggregate,
175+
function: :execute,
176+
before_execute: :before_execute
177+
}
178+
179+
assert {:error, :aggregate_execution_failed, :boom} =
180+
Aggregate.execute(DefaultApp, CrashBeforeExecuteAggregate, aggregate_uuid, context)
181+
end
182+
103183
describe "command dispatch return" do
104184
alias Commanded.Aggregates.ReturnValue.Command
105185
alias Commanded.Aggregates.ReturnValue.Event

‎test/application/application_supervisor_test.exs‎

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,6 @@ defmodule Commanded.Aggregates.ApplicationSupervisorTest do
1212

1313
test "children are hibernated after inactivity" do
1414
for module <- [
15-
Commanded.Commands.TaskDispatcher,
1615
Commanded.Aggregates.Supervisor,
1716
Commanded.Subscriptions,
1817
Commanded.Subscriptions.Registry

‎test/application/telemetry_test.exs‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -106,16 +106,16 @@ defmodule Commanded.Application.TelemetryTest do
106106
test "emit a single aggregate execute attempt when command execution times out" do
107107
command = %TimeoutCommand{aggregate_uuid: UUID.uuid4(), sleep_in_ms: 2_000}
108108

109-
assert {:error, error} = TimeoutRouter.dispatch(command, application: DefaultApp)
110-
assert error in [:aggregate_execution_failed, :aggregate_execution_timeout]
109+
assert {:error, :aggregate_execution_timeout} =
110+
TimeoutRouter.dispatch(command, application: DefaultApp)
111111

112112
assert_receive {[:commanded, :application, :dispatch, :start], _, _, _}
113113
assert_receive {[:commanded, :aggregate, :execute, :start], _, _, _}
114114

115115
refute_receive {[:commanded, :aggregate, :execute, :start], _, _, _}, 200
116116

117117
assert_receive {[:commanded, :application, :dispatch, :stop], _, _,
118-
%{error: ^error}}
118+
%{error: :aggregate_execution_timeout}}
119119
end
120120

121121
defp attach_telemetry do

‎test/commands/command_timeout_test.exs‎

Lines changed: 2 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -14,11 +14,8 @@ defmodule Commanded.Commands.CommandTimeoutTest do
1414
command = %TimeoutCommand{aggregate_uuid: UUID.uuid4(), sleep_in_ms: 2_000}
1515

1616
# Handler is set to take longer than the configured timeout
17-
case TimeoutRouter.dispatch(command, application: DefaultApp) do
18-
{:error, :aggregate_execution_failed} -> :ok
19-
{:error, :aggregate_execution_timeout} -> :ok
20-
reply -> flunk("received an unexpected response: #{inspect(reply)}")
21-
end
17+
assert {:error, :aggregate_execution_timeout} =
18+
TimeoutRouter.dispatch(command, application: DefaultApp)
2219
end
2320

2421
test "should succeed when handler completes within configured timeout" do

‎test/middleware/middleware_test.exs‎

Lines changed: 40 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@ defmodule Commanded.Middleware.MiddlewareTest do
22
use ExUnit.Case
33

44
import Commanded.Enumerable
5+
import ExUnit.CaptureLog
56

67
alias Commanded.Commands.ExecutionResult
78
alias Commanded.DefaultApp
@@ -21,6 +22,19 @@ defmodule Commanded.Middleware.MiddlewareTest do
2122
alias Commanded.TestSupport.RetryStopOnceAggregate.Command, as: RetryStopOnceCommand
2223
alias Commanded.UUID
2324

25+
defmodule CrashCommand do
26+
defstruct [:uuid]
27+
end
28+
29+
defmodule CrashBeforeExecuteAggregate do
30+
alias Commanded.Aggregates.ExecutionContext
31+
32+
defstruct []
33+
34+
def before_execute(_aggregate_state, %ExecutionContext{}), do: Process.exit(self(), :boom)
35+
def execute(%__MODULE__{}, %CrashCommand{}), do: []
36+
end
37+
2438
defmodule FirstMiddleware do
2539
@behaviour Commanded.Middleware
2640

@@ -82,6 +96,17 @@ defmodule Commanded.Middleware.MiddlewareTest do
8296
before_execute: :before_execute
8397
end
8498

99+
defmodule CrashRouter do
100+
use Commanded.Commands.Router
101+
102+
middleware Commanded.Middleware.Logger
103+
104+
dispatch [CrashCommand],
105+
to: CrashBeforeExecuteAggregate,
106+
identity: :uuid,
107+
before_execute: :before_execute
108+
end
109+
85110
setup do
86111
start_supervised!(CommandAuditMiddleware)
87112
start_supervised!(RetryStopOnceAggregate.Tracker)
@@ -150,12 +175,8 @@ defmodule Commanded.Middleware.MiddlewareTest do
150175
test "should execute middleware failure callback when aggregate process dies" do
151176
command = %Timeout{aggregate_uuid: UUID.uuid4()}
152177

153-
# Force command handling to timeout so the aggregate process is terminated
154-
:ok =
155-
case Router.dispatch(command, application: DefaultApp, timeout: 50) do
156-
{:error, :aggregate_execution_timeout} -> :ok
157-
{:error, :aggregate_execution_failed} -> :ok
158-
end
178+
assert {:error, :aggregate_execution_timeout} =
179+
Router.dispatch(command, application: DefaultApp, timeout: 50)
159180

160181
{dispatched, succeeded, failed} = CommandAuditMiddleware.count_commands()
161182

@@ -164,6 +185,19 @@ defmodule Commanded.Middleware.MiddlewareTest do
164185
assert failed == 1
165186
end
166187

188+
test "should preserve abnormal aggregate exit reason for middleware logging" do
189+
command = %CrashCommand{uuid: UUID.uuid4()}
190+
191+
log =
192+
capture_log(fn ->
193+
assert {:error, :aggregate_execution_failed} =
194+
CrashRouter.dispatch(command, application: DefaultApp)
195+
end)
196+
197+
assert log =~ "failed :aggregate_execution_failed"
198+
assert log =~ "due to: :boom"
199+
end
200+
167201
test "should let a middleware update the metadata" do
168202
command = %IncrementCount{aggregate_uuid: UUID.uuid4(), by: 1}
169203

0 commit comments

Comments
 (0)