Skip to content
Draft
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
202 changes: 202 additions & 0 deletions lib/commanded/aggregates/stateless_aggregate.ex
Original file line number Diff line number Diff line change
@@ -0,0 +1,202 @@
defmodule Commanded.Aggregates.StatelessAggregate do
@moduledoc false

alias Commanded.Aggregates.ExecutionContext
alias Commanded.Aggregate.Multi
alias Commanded.Snapshotting
alias Commanded.Aggregates.Aggregate
alias Commanded.Aggregates.AggregateStateBuilder
alias Commanded.Event.Mapper
alias Commanded.EventStore
alias Commanded.Application.Config

require Logger

def execute(application, aggregate_module, aggregate_uuid, context, _timeout) do
aggregate = make_aggregate(application, aggregate_module, aggregate_uuid)
aggregate = AggregateStateBuilder.populate(aggregate)

{result, aggregate} = execute_command(context, aggregate)

aggregate =
if Snapshotting.snapshot_required?(aggregate.snapshotting, aggregate.aggregate_version) do
do_take_snapshot(aggregate)
else
aggregate
end

ExecutionContext.format_reply(result, context, aggregate)
end

defp make_aggregate(application, aggregate_module, aggregate_uuid) do
snapshot_options =
application
|> Config.get(:snapshotting)
|> Kernel.||(%{})
|> Map.get(aggregate_module, [])

snapshotting = Snapshotting.new(application, aggregate_uuid, snapshot_options)

%Aggregate{
aggregate_module: aggregate_module,
snapshotting: snapshotting,
application: application,
aggregate_uuid: aggregate_uuid
}
end

defp before_execute_command(_aggregate_state, %ExecutionContext{before_execute: nil}), do: :ok

defp before_execute_command(aggregate_state, %ExecutionContext{} = context) do
%ExecutionContext{handler: handler, before_execute: before_execute} = context

Kernel.apply(handler, before_execute, [aggregate_state, context])
end

defp execute_command(%ExecutionContext{} = context, %Aggregate{} = aggregate) do
%ExecutionContext{command: command, handler: handler, function: function} = context
%Aggregate{aggregate_state: aggregate_state} = aggregate

Logger.debug(describe(aggregate) <> " executing command: " <> inspect(command))

with :ok <- before_execute_command(aggregate_state, context) do
case Kernel.apply(handler, function, [aggregate_state, command]) do
{:error, _error} = reply ->
{reply, aggregate}

none when none in [:ok, nil, []] ->
{{:ok, []}, aggregate}

%Multi{} = multi ->
case Multi.run(multi) do
{:error, _error} = reply ->
{reply, aggregate}

{aggregate_state, pending_events} ->
persist_events(pending_events, aggregate_state, context, aggregate)
end

{:ok, pending_events} ->
apply_and_persist_events(pending_events, context, aggregate)

pending_events ->
apply_and_persist_events(pending_events, context, aggregate)
end
else
{:error, _error} = reply ->
{reply, aggregate}
end
rescue
error ->
stacktrace = __STACKTRACE__
Logger.error(Exception.format(:error, error, stacktrace))

{{:error, error, stacktrace}, aggregate}
end

defp persist_events(pending_events, aggregate_state, context, %Aggregate{} = aggregate) do
%Aggregate{aggregate_version: expected_version} = aggregate

with :ok <- append_to_stream(pending_events, context, aggregate) do
aggregate_version = expected_version + length(pending_events)

aggregate = %Aggregate{
aggregate
| aggregate_state: aggregate_state,
aggregate_version: aggregate_version
}

{{:ok, pending_events}, aggregate}
else
{:error, :wrong_expected_version} ->
# Fetch missing events from event store
aggregate = AggregateStateBuilder.rebuild_from_events(aggregate)

# Retry command if there are any attempts left
case ExecutionContext.retry(context) do
{:ok, context} ->
Logger.debug(describe(aggregate) <> " wrong expected version, retrying command")

execute_command(context, aggregate)

reply ->
Logger.debug(
describe(aggregate) <> " wrong expected version, but not retrying command"
)

{reply, aggregate}
end

{:error, _error} = reply ->
{reply, aggregate}
end
end

defp apply_and_persist_events(pending_events, context, %Aggregate{} = aggregate) do
%Aggregate{aggregate_module: aggregate_module, aggregate_state: aggregate_state} = aggregate

pending_events = List.wrap(pending_events)
aggregate_state = apply_events(aggregate_module, aggregate_state, pending_events)

persist_events(pending_events, aggregate_state, context, aggregate)
end

defp apply_events(aggregate_module, aggregate_state, events) do
Enum.reduce(events, aggregate_state, &aggregate_module.apply(&2, &1))
end

defp append_to_stream([], _context, _state), do: :ok

defp append_to_stream(pending_events, %ExecutionContext{} = context, %Aggregate{} = state) do
%Aggregate{
application: application,
aggregate_uuid: aggregate_uuid,
aggregate_version: expected_version
} = state

%ExecutionContext{
causation_id: causation_id,
correlation_id: correlation_id,
metadata: metadata
} = context

event_data =
Mapper.map_to_event_data(pending_events,
causation_id: causation_id,
correlation_id: correlation_id,
metadata: metadata
)

EventStore.append_to_stream(application, aggregate_uuid, expected_version, event_data)
end

defp do_take_snapshot(%Aggregate{} = state) do
%Aggregate{
aggregate_state: aggregate_state,
aggregate_version: aggregate_version,
snapshotting: snapshotting
} = state

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

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

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

state
end
end

defp describe(%Aggregate{} = aggregate) do
%Aggregate{
aggregate_module: aggregate_module,
aggregate_uuid: aggregate_uuid,
aggregate_version: aggregate_version
} = aggregate

"#{inspect(aggregate_module)}<#{aggregate_uuid}@#{aggregate_version}>"
end
end
116 changes: 75 additions & 41 deletions lib/commanded/commands/dispatcher.ex
Original file line number Diff line number Diff line change
Expand Up @@ -2,9 +2,10 @@ defmodule Commanded.Commands.Dispatcher do
@moduledoc false
alias Commanded.Aggregates.Aggregate
alias Commanded.Aggregates.ExecutionContext
alias Commanded.Commands.Router
alias Commanded.AggregatesCommands.Router
alias Commanded.Middleware.Pipeline
alias Commanded.Telemetry
alias Commanded.Aggregates.StatelessAggregate

require Logger

Expand All @@ -31,7 +32,8 @@ defmodule Commanded.Commands.Dispatcher do
:metadata,
:retry_attempts,
:returning,
middleware: []
middleware: [],
stateless: false
]
end

Expand Down Expand Up @@ -68,46 +70,10 @@ defmodule Commanded.Commands.Dispatcher do
end

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

{:ok, ^aggregate_uuid} =
Commanded.Aggregates.Supervisor.open_aggregate(
application,
aggregate_module,
aggregate_uuid
)

task_dispatcher_name = Module.concat([application, Commanded.Commands.TaskDispatcher])

task =
Task.Supervisor.async_nolink(task_dispatcher_name, Aggregate, :execute, [
application,
aggregate_module,
aggregate_uuid,
context,
timeout
])

result =
case Task.yield(task, timeout) || Task.shutdown(task) do
{:ok, result} ->
result

{:exit, {:normal, :aggregate_stopped}} = result ->
result

{:exit, {{:nodedown, _node_name}, {GenServer, :call, _}}} ->
{:error, :remote_node_down}
task_dispatcher_name =
Module.concat([pipeline.application, Commanded.Commands.TaskDispatcher])

{:exit, _reason} ->
{:error, :aggregate_execution_failed}

nil ->
{:error, :aggregate_execution_timeout}
end

case result do
case execute_aggregate(task_dispatcher_name, pipeline, payload, context) do
{:ok, aggregate_version, events, aggregate_state} ->
pipeline
|> Pipeline.assign(:aggregate_version, aggregate_version)
Expand Down Expand Up @@ -145,6 +111,74 @@ defmodule Commanded.Commands.Dispatcher do
end
end

defp execute_aggregate(
task_dispatcher_name,
%Pipeline{} = pipeline,
%Payload{stateless: true} = payload,
%ExecutionContext{} = context
) do
%Pipeline{application: application, assigns: %{aggregate_uuid: aggregate_uuid}} = pipeline
%Payload{aggregate_module: aggregate_module, timeout: timeout} = payload

task =
Task.Supervisor.async_nolink(task_dispatcher_name, StatelessAggregate, :execute, [
application,
aggregate_module,
aggregate_uuid,
context,
timeout
])

execute_task(task, timeout)
end

defp execute_aggregate(
task_dispatcher_name,
%Pipeline{} = pipeline,
%Payload{stateless: false} = payload,
%ExecutionContext{} = context
) do
%Pipeline{application: application, assigns: %{aggregate_uuid: aggregate_uuid}} = pipeline
%Payload{aggregate_module: aggregate_module, timeout: timeout} = payload

{:ok, ^aggregate_uuid} =
Commanded.Aggregates.Supervisor.open_aggregate(
application,
aggregate_module,
aggregate_uuid
)

task =
Task.Supervisor.async_nolink(task_dispatcher_name, Aggregate, :execute, [
application,
aggregate_module,
aggregate_uuid,
context,
timeout
])

execute_task(task, timeout)
end

defp execute_task(task, timeout) do
case Task.yield(task, timeout) || Task.shutdown(task) do
{:ok, result} ->
result

{:exit, {:normal, :aggregate_stopped}} = result ->
result

{:exit, {{:nodedown, _node_name}, {GenServer, :call, _}}} ->
{:error, :remote_node_down}

{:exit, _reason} ->
{:error, :aggregate_execution_failed}

nil ->
{:error, :aggregate_execution_timeout}
end
end

defp to_execution_context(%Pipeline{} = pipeline, %Payload{} = payload) do
%Pipeline{command: command, command_uuid: command_uuid, metadata: metadata} = pipeline

Expand Down
3 changes: 2 additions & 1 deletion lib/commanded/commands/router.ex
Original file line number Diff line number Diff line change
Expand Up @@ -583,7 +583,8 @@ defmodule Commanded.Commands.Router do
lifespan: @lifespan,
metadata: metadata,
middleware: @middleware,
retry_attempts: retry_attempts
retry_attempts: retry_attempts,
stateless: Keyword.get(opts, :experimental_stateless, false)
}

Dispatcher.dispatch(payload)
Expand Down