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
4 changes: 2 additions & 2 deletions config/test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,7 @@ config :commanded, TestEventStore,
username: "postgres",
password: "postgres",
database: "eventstore_test",
hostname: "localhost",
hostname: System.get_env("PG_HOST", "localhost"),
pool_size: 5,
pool_overflow: 0

Expand All @@ -49,5 +49,5 @@ config :commanded, Commanded.Projections.Repo,
database: "commanded_projections_test",
username: "postgres",
password: "postgres",
hostname: "localhost",
hostname: System.get_env("PG_HOST", "localhost"),
pool: Ecto.Adapters.SQL.Sandbox
78 changes: 78 additions & 0 deletions guides/explanations/fork-differences.md
Original file line number Diff line number Diff line change
Expand Up @@ -284,3 +284,81 @@ end
**Benefits:**
- All event handlers (projectors, sagas, notification handlers) get processing latency visibility for free
- Enables SLA dashboards and alerting without any application-level code

### **Custom Initial State for Aggregates**
[PR #49](https://github.com/straw-hat-team/commanded/pull/49)

**Changes:**
- Added `initial_state` option to command router's `dispatch` macro
- Allows specifying a module that implements `initial_state/0` callback
- Useful when using protobuf-generated messages as aggregate state

**Usage:**

```elixir
# State module with initial_state/0 callback
defmodule BankAccountState do
defstruct [:account_number, :balance, status: :uninitialized]

def initial_state, do: %__MODULE__{}
end

# Aggregate module (behavior only, no struct)
defmodule BankAccount do
def execute(%BankAccountState{status: :uninitialized}, %OpenAccount{} = cmd) do
%AccountOpened{account_number: cmd.account_number}
end

def apply(%BankAccountState{} = state, %AccountOpened{} = event) do
%BankAccountState{state | account_number: event.account_number, status: :open}
end
end

defmodule MyRouter do
use Commanded.Commands.Router

# Calls BankAccountState.initial_state/0 to create initial state
dispatch [OpenAccount, DepositMoney],
to: BankAccount,
initial_state: BankAccountState,
identity: :account_number
end
```

**Benefits:**
- Decouples aggregate behavior from state representation
- State module controls its own initialization (like `to:` pattern)
- Enables use of protobuf-generated messages as aggregate state instead of being forced to use Elixir structs
- Backwards compatible - if `initial_state` is omitted, `struct(AggregateModule)` is used

**Rationale:**

In the upstream Commanded, the aggregate module serves dual purposes: it defines both the state struct and the behavior (`execute/2` and `apply/2` functions). This coupling becomes problematic when you want to use protobuf-generated modules for state.

Protobuf modules are code-generated and shouldn't be manually modified—any changes would be overwritten on regeneration. To add aggregate behavior to a protobuf struct, you'd need to write custom protobuf extensions or use workarounds, adding complexity to your build pipeline.

By separating the aggregate (behavior) from the state (data structure), you can:

```elixir
# Generated by protobuf - don't modify
defmodule MyApp.Proto.BankAccountState do
use Protobuf, syntax: :proto3
# ... generated fields ...
end

# Your code - aggregate behavior with initial_state returning the protobuf struct
defmodule MyApp.BankAccount do
alias MyApp.Proto.BankAccountState

def initial_state, do: %BankAccountState{status: :STATUS_UNINITIALIZED}

def execute(%BankAccountState{} = state, %OpenAccount{} = cmd), do: ...
def apply(%BankAccountState{} = state, %AccountOpened{} = event), do: ...
end
```

**Why a callback instead of `struct/1`?**

The main reason is to use protobuf-generated messages as aggregate state instead of being forced to use Elixir structs. Protobuf modules are code-generated and have their own initialization semantics. The `initial_state/0` callback lets you return whatever your state module needs—including protobuf messages—rather than relying on `struct(AggregateModule)`, which only works for Elixir structs.

This separation follows the principle that data representation and business logic are distinct concerns that benefit from being in separate modules. It also aligns more closely with the [Functional Decider Pattern](https://thinkbeforecoding.com/post/2021/12/17/functional-event-sourcing-decider), where the aggregate is a set of pure functions (`execute`, `apply`, `initial_state`) operating on state, rather than a stateful object that owns its data structure.
13 changes: 10 additions & 3 deletions lib/application.ex
Original file line number Diff line number Diff line change
Expand Up @@ -214,18 +214,25 @@ defmodule Commanded.Application do
Retrieving aggregate state is done by calling to the opened aggregate,
or querying the event store for an optional state snapshot
and then replaying the aggregate's event stream.

## Options

- `:timeout` - timeout in milliseconds (default: 5000)
- `:initial_state` - module that implements `initial_state/0` callback.
Used when rebuilding state from events. Defaults to `aggregate_module`.

"""
@spec aggregate_state(
aggregate_module :: module(),
aggregate_uuid :: Aggregate.uuid(),
timeout :: integer
timeout_or_opts :: timeout() | Aggregate.aggregate_state_opts()
) :: Aggregate.state()
def aggregate_state(aggregate_module, aggregate_uuid, timeout \\ 5000) do
def aggregate_state(aggregate_module, aggregate_uuid, timeout_or_opts \\ 5000) do
Aggregate.aggregate_state(
__MODULE__,
aggregate_module,
aggregate_uuid,
timeout
timeout_or_opts
)
end

Expand Down
13 changes: 10 additions & 3 deletions lib/commanded.ex
Original file line number Diff line number Diff line change
Expand Up @@ -38,19 +38,26 @@ defmodule Commanded do
Retrieving aggregate state is done by calling to the opened aggregate,
or querying the event store for an optional state snapshot
and then replaying the aggregate's event stream.

## Options

- `:timeout` - timeout in milliseconds (default: 5000)
- `:initial_state` - module that implements `initial_state/0` callback.
Used when rebuilding state from events. Defaults to `aggregate_module`.

"""
@spec aggregate_state(
application :: Commanded.Application.t(),
aggregate_module :: module(),
aggregate_uuid :: Aggregate.uuid(),
timeout :: integer
timeout_or_opts :: timeout() | Aggregate.aggregate_state_opts()
) :: Aggregate.state()
def aggregate_state(application, aggregate_module, aggregate_uuid, timeout \\ 5_000) do
def aggregate_state(application, aggregate_module, aggregate_uuid, timeout_or_opts \\ 5_000) do
Aggregate.aggregate_state(
application,
aggregate_module,
aggregate_uuid,
timeout
timeout_or_opts
)
end
end
60 changes: 59 additions & 1 deletion lib/commanded/aggregates/aggregate.ex
Original file line number Diff line number Diff line change
Expand Up @@ -177,6 +177,9 @@ defmodule Commanded.Aggregates.Aggregate do
@type return_event :: struct() | list(struct()) | {:ok, struct()} | {:ok, list(struct())}
@type no_return_event :: :ok | {:ok, []} | nil | []

@type aggregate_state_opt :: {:timeout, timeout()} | {:initial_state, module()}
@type aggregate_state_opts :: [aggregate_state_opt()]

@doc """
Optionally execute a command against the aggregate. Returns either no event, one event,
a list of events, or an error tuple.
Expand All @@ -194,6 +197,7 @@ defmodule Commanded.Aggregates.Aggregate do
defstruct [
:application,
:aggregate_module,
:initial_state,
:aggregate_uuid,
:aggregate_state,
:snapshotting,
Expand All @@ -208,6 +212,8 @@ defmodule Commanded.Aggregates.Aggregate do

aggregate_module = Keyword.fetch!(aggregate_opts, :aggregate_module)
aggregate_uuid = Keyword.fetch!(aggregate_opts, :aggregate_uuid)
initial_state = Keyword.get(aggregate_opts, :initial_state)
validate_initial_state_module!(initial_state)

unless is_atom(aggregate_module),
do: raise(ArgumentError, message: "aggregate module must be an atom")
Expand All @@ -222,6 +228,7 @@ defmodule Commanded.Aggregates.Aggregate do
state = %Aggregate{
application: application,
aggregate_module: aggregate_module,
initial_state: initial_state,
aggregate_uuid: aggregate_uuid,
snapshotting: Snapshotting.new(application, aggregate_uuid, snapshot_options)
}
Expand Down Expand Up @@ -278,7 +285,17 @@ defmodule Commanded.Aggregates.Aggregate do
end

@doc false
def aggregate_state(application, aggregate_module, aggregate_uuid, timeout \\ 5_000) do
def aggregate_state(application, aggregate_module, aggregate_uuid, timeout_or_opts \\ 5_000)

def aggregate_state(application, aggregate_module, aggregate_uuid, timeout)
when is_integer(timeout) or timeout == :infinity do
aggregate_state(application, aggregate_module, aggregate_uuid, timeout: timeout)
end

def aggregate_state(application, aggregate_module, aggregate_uuid, opts) when is_list(opts) do
timeout = Keyword.get(opts, :timeout, 5_000)
initial_state = Keyword.get(opts, :initial_state)
validate_initial_state_module!(initial_state)
name = via_name(application, aggregate_module, aggregate_uuid)

try do
Expand All @@ -297,6 +314,7 @@ defmodule Commanded.Aggregates.Aggregate do
%Aggregate{
application: application,
aggregate_module: aggregate_module,
initial_state: initial_state,
Comment thread
cursor[bot] marked this conversation as resolved.
aggregate_uuid: aggregate_uuid,
snapshotting: Snapshotting.new(application, aggregate_uuid, snapshot_options)
}
Expand All @@ -308,6 +326,9 @@ defmodule Commanded.Aggregates.Aggregate do
{:ok, result} ->
result

{:exit, reason} ->
exit(reason)

nil ->
exit({:timeout, {GenServer, :call, [name, :aggregate_state, timeout]}})
end
Expand All @@ -320,6 +341,12 @@ defmodule Commanded.Aggregates.Aggregate do
GenServer.call(name, :aggregate_version, timeout)
end

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

@doc false
def take_snapshot(application, aggregate_module, aggregate_uuid, timeout \\ 5_000) do
name = via_name(application, aggregate_module, aggregate_uuid)
Expand Down Expand Up @@ -437,6 +464,14 @@ defmodule Commanded.Aggregates.Aggregate do
reply_with_lifespan(aggregate_version, state)
end

@doc false
@impl GenServer
def handle_call(:initial_state_module, _from, %Aggregate{} = state) do
%Aggregate{initial_state: initial_state} = state

reply_with_lifespan(initial_state, state)
end

@doc false
@impl GenServer
def handle_info({:events, events}, %Aggregate{} = state) do
Expand Down Expand Up @@ -813,4 +848,27 @@ defmodule Commanded.Aggregates.Aggregate do
lifespan_timeout -> {:noreply, state, lifespan_timeout}
end
end

defp validate_initial_state_module!(nil), do: :ok

defp validate_initial_state_module!(initial_state) when is_atom(initial_state) do
case Code.ensure_compiled(initial_state) do
{:module, _} ->
if function_exported?(initial_state, :initial_state, 0) do
:ok
else
raise ArgumentError,
"initial_state module #{inspect(initial_state)} must export initial_state/0 function."
end

{:error, reason} ->
raise ArgumentError,
"initial_state module #{inspect(initial_state)} could not be loaded: #{inspect(reason)}"
end
end

defp validate_initial_state_module!(initial_state) do
raise ArgumentError,
"initial_state must be a module but got: #{inspect(initial_state)}"
end
end
17 changes: 15 additions & 2 deletions lib/commanded/aggregates/aggregate_state_builder.ex
Original file line number Diff line number Diff line change
Expand Up @@ -71,9 +71,13 @@ defmodule Commanded.Aggregates.AggregateStateBuilder do
If the snapshot exists, fetch any subsequent events to rebuild its state.
Otherwise start with the aggregate struct and stream all existing events for
the aggregate from the event store to rebuild its state from those events.

The initial state is determined by the `initial_state` field:
- If a module is provided, `initial_state.initial_state()` is called
- Otherwise, falls back to `struct(aggregate_module)`
"""
def populate(%Aggregate{} = state) do
%Aggregate{aggregate_module: aggregate_module, snapshotting: snapshotting} = state
%Aggregate{snapshotting: snapshotting} = state

{aggregate, snapshot_used, snapshot_source_version} =
case Snapshotting.read_snapshot(snapshotting) do
Expand All @@ -90,7 +94,7 @@ defmodule Commanded.Aggregates.AggregateStateBuilder do
agg = %Aggregate{
state
| aggregate_version: 0,
aggregate_state: struct(aggregate_module)
aggregate_state: create_initial_state(state)
}

{agg, false, nil}
Expand All @@ -102,6 +106,15 @@ defmodule Commanded.Aggregates.AggregateStateBuilder do
)
end

defp create_initial_state(%Aggregate{initial_state: nil, aggregate_module: aggregate_module}) do
struct(aggregate_module)
Comment thread
cursor[bot] marked this conversation as resolved.
end
Comment thread
yordis marked this conversation as resolved.

defp create_initial_state(%Aggregate{initial_state: initial_state})
when is_atom(initial_state) do
initial_state.initial_state()
end

@doc """
Load events from the event store, in batches, to rebuild the aggregate state.

Expand Down
Loading
Loading