Skip to content
Merged
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
81 changes: 72 additions & 9 deletions lib/commanded/aggregates/aggregate.ex
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,49 @@ defmodule Commanded.Aggregates.Aggregate do
"""
})

telemetry_event(%{
event: [:commanded, :aggregate, :snapshot, :start],
description: "Emitted when an aggregate begins taking a snapshot",
measurements: "%{system_time: integer()}",
metadata: """
%{application: Commanded.Application.t(),
aggregate_uuid: String.t(),
aggregate_version: non_neg_integer(),
snapshot_every: non_neg_integer() | nil,
snapshot_module_version: non_neg_integer()}
"""
})

telemetry_event(%{
event: [:commanded, :aggregate, :snapshot, :stop],
description: "Emitted when an aggregate completes taking a snapshot",
measurements: "%{duration: non_neg_integer()}",
metadata: """
%{application: Commanded.Application.t(),
aggregate_uuid: String.t(),
aggregate_version: non_neg_integer(),
snapshot_every: non_neg_integer() | nil,
snapshot_module_version: non_neg_integer(),
error: nil | any()}
"""
})

telemetry_event(%{
event: [:commanded, :aggregate, :snapshot, :exception],
description: "Emitted when an aggregate raises during snapshot",
measurements: "%{duration: non_neg_integer()}",
metadata: """
%{application: Commanded.Application.t(),
aggregate_uuid: String.t(),
aggregate_version: non_neg_integer(),
snapshot_every: non_neg_integer() | nil,
snapshot_module_version: non_neg_integer(),
kind: :throw | :error | :exit,
reason: any(),
stacktrace: list()}
"""
})

@moduledoc """
Aggregate is a `GenServer` process used to provide access to an
instance of an event sourced aggregate.
Expand Down Expand Up @@ -628,22 +671,42 @@ defmodule Commanded.Aggregates.Aggregate do

defp do_take_snapshot(%Aggregate{} = state) do
%Aggregate{
application: application,
aggregate_uuid: aggregate_uuid,
aggregate_state: aggregate_state,
aggregate_version: aggregate_version,
snapshotting: snapshotting
snapshotting: %Snapshotting{
snapshot_every: snapshot_every,
snapshot_module_version: snapshot_module_version
}
} = state

Logger.debug(describe(state) <> " recording snapshot")
meta = %{
application: application,
aggregate_uuid: aggregate_uuid,
aggregate_version: aggregate_version,
snapshot_every: snapshot_every,
snapshot_module_version: snapshot_module_version
}

:telemetry.span([:commanded, :aggregate, :snapshot], meta, fn ->
Logger.debug(describe(state) <> " recording snapshot")

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

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

error
end
stop_meta =
Map.put(meta, :error, if(match?({:error, _}, result), do: elem(result, 1), else: nil))

{result, stop_meta}
end)
end

defp telemetry_wrong_expected_version(context, from, state) do
Expand Down
65 changes: 53 additions & 12 deletions lib/commanded/aggregates/aggregate_state_builder.ex
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,9 @@ defmodule Commanded.Aggregates.AggregateStateBuilder do
%{application: Commanded.Application.t(),
aggregate_uuid: String.t(),
aggregate_state: struct(),
aggregate_version: non_neg_integer()}
aggregate_version: non_neg_integer(),
snapshot_used: boolean(),
snapshot_source_version: non_neg_integer() | nil}
"""
})

Expand Down Expand Up @@ -69,29 +71,48 @@ defmodule Commanded.Aggregates.AggregateStateBuilder do
def populate(%Aggregate{} = state) do
%Aggregate{aggregate_module: aggregate_module, snapshotting: snapshotting} = state

aggregate =
{aggregate, snapshot_used, snapshot_source_version} =
case Snapshotting.read_snapshot(snapshotting) do
{:ok, %SnapshotData{source_version: source_version, data: data}} ->
%Aggregate{
agg = %Aggregate{
state
| aggregate_version: source_version,
aggregate_state: data
}

{agg, true, source_version}

{:error, _error} ->
# No snapshot present, or exists but for outdated state, so use initial empty state
%Aggregate{state | aggregate_version: 0, aggregate_state: struct(aggregate_module)}
agg = %Aggregate{
state
| aggregate_version: 0,
aggregate_state: struct(aggregate_module)
}

{agg, false, nil}
end

rebuild_from_events(aggregate)
rebuild_from_events(aggregate,
snapshot_used: snapshot_used,
snapshot_source_version: snapshot_source_version
)
end

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

## Options

* `:snapshot_used` - whether a snapshot was used as initial state (default: `false`)
* `:snapshot_source_version` - version of the snapshot, if used (default: `nil`)
"""
def rebuild_from_events(%Aggregate{} = state) do
def rebuild_from_events(%Aggregate{} = state, opts \\ []) do
snapshot_used = Keyword.get(opts, :snapshot_used, false)
snapshot_source_version = Keyword.get(opts, :snapshot_source_version)

load_prefix = [:commanded, :aggregate, :load]
load_start = Telemetry.start(load_prefix, telemetry_metadata(state))
meta = telemetry_metadata(state)
load_start = Telemetry.start(load_prefix, meta)

%Aggregate{
application: application,
Expand All @@ -107,13 +128,25 @@ defmodule Commanded.Aggregates.AggregateStateBuilder do
@read_event_batch_size
) do
{:error, :stream_not_found} ->
# aggregate does not exist, return initial state
Telemetry.stop(load_prefix, load_start, telemetry_metadata(state), %{count: 0})
Telemetry.stop(
load_prefix,
load_start,
load_stop_metadata(state, snapshot_used, snapshot_source_version),
%{count: 0}
)

{state, 0}

event_stream ->
{state, count} = rebuild_from_event_stream(event_stream, state)
Telemetry.stop(load_prefix, load_start, telemetry_metadata(state), %{count: count})

Telemetry.stop(
load_prefix,
load_start,
load_stop_metadata(state, snapshot_used, snapshot_source_version),
%{count: count}
)

{state, count}
end

Expand Down Expand Up @@ -145,6 +178,14 @@ defmodule Commanded.Aggregates.AggregateStateBuilder do
{state, count}
end

defp load_stop_metadata(aggregate, snapshot_used, snapshot_source_version) do
telemetry_metadata(aggregate)
|> Map.merge(%{
snapshot_used: snapshot_used,
snapshot_source_version: snapshot_source_version
})
end

defp telemetry_metadata(%Aggregate{} = state) do
%Aggregate{
application: application,
Expand Down
14 changes: 14 additions & 0 deletions lib/commanded/opentelemetry.ex
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@ defmodule Commanded.OpenTelemetry do

alias Commanded.OpenTelemetry.Aggregate
alias Commanded.OpenTelemetry.AggregatePopulate
alias Commanded.OpenTelemetry.AggregateSnapshot
alias Commanded.OpenTelemetry.Application, as: OTelApplication
alias Commanded.OpenTelemetry.EventHandler
alias Commanded.OpenTelemetry.EventStore
Expand Down Expand Up @@ -65,6 +66,11 @@ defmodule Commanded.OpenTelemetry do
default: [],
doc: "Aggregate populate tracing configuration. Use `:disabled` to disable."
],
aggregate_snapshot: [
type: {:in, [:disabled, []]},
default: [],
doc: "Aggregate snapshot tracing configuration. Use `:disabled` to disable."
],
application: [
type: {:in, [:disabled, []]},
default: [],
Expand Down Expand Up @@ -120,6 +126,9 @@ defmodule Commanded.OpenTelemetry do
# Disable aggregate populate tracing
Commanded.OpenTelemetry.setup(aggregate_populate: :disabled)

# Disable aggregate snapshot tracing
Commanded.OpenTelemetry.setup(aggregate_snapshot: :disabled)

# Disable event store tracing
Commanded.OpenTelemetry.setup(event_store: :disabled)

Expand All @@ -141,6 +150,11 @@ defmodule Commanded.OpenTelemetry do
_config -> AggregatePopulate.setup()
end

case opts[:aggregate_snapshot] do
:disabled -> :ok
_config -> AggregateSnapshot.setup()
end

case opts[:application] do
:disabled -> :ok
_config -> OTelApplication.setup()
Expand Down
14 changes: 14 additions & 0 deletions lib/commanded/opentelemetry/aggregate_populate.ex
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,20 @@ defmodule Commanded.OpenTelemetry.AggregatePopulate do
meta.aggregate_version
)

Span.set_attribute(
ctx,
CommandedAttributes.commanded_snapshot_used(),
meta[:snapshot_used] || false
)

if snapshot_source_version = meta[:snapshot_source_version] do
Span.set_attribute(
ctx,
CommandedAttributes.commanded_snapshot_source_version(),
snapshot_source_version
)
end

OpentelemetryTelemetry.end_telemetry_span(@tracer_id, meta)
end

Expand Down
121 changes: 121 additions & 0 deletions lib/commanded/opentelemetry/aggregate_snapshot.ex
Original file line number Diff line number Diff line change
@@ -0,0 +1,121 @@
defmodule Commanded.OpenTelemetry.AggregateSnapshot do
@moduledoc false

alias Commanded.OpenTelemetry.CommandedAttributes
alias Commanded.OpenTelemetry.Helpers
alias OpenTelemetry.SemConv.ErrorAttributes
alias OpenTelemetry.SemConv.Incubating.CodeAttributes
alias OpenTelemetry.SemConv.Incubating.MessagingAttributes
alias OpenTelemetry.Span

@tracer_id __MODULE__

def setup do
:ok =
:telemetry.attach_many(
{__MODULE__, :snapshot},
[
[:commanded, :aggregate, :snapshot, :start],
[:commanded, :aggregate, :snapshot, :stop],
[:commanded, :aggregate, :snapshot, :exception]
],
&__MODULE__.handle_telemetry_event/4,
%{}
)
end

def handle_telemetry_event(
[:commanded, :aggregate, :snapshot, :start],
_measurements,
meta,
_config
) do
attributes =
[
{MessagingAttributes.messaging_system(), "commanded"},
{MessagingAttributes.messaging_operation_type(), :publish},
{MessagingAttributes.messaging_operation_name(), "snapshot"},
{CodeAttributes.code_function(), "snapshot"},
{CommandedAttributes.commanded_application(), meta.application},
{CommandedAttributes.commanded_aggregate_uuid(), meta.aggregate_uuid},
{CommandedAttributes.commanded_aggregate_version(), meta.aggregate_version}
]
|> maybe_add_snapshot_every(meta[:snapshot_every])
|> maybe_add_snapshot_module_version(meta[:snapshot_module_version])

OpentelemetryTelemetry.start_telemetry_span(
@tracer_id,
"commanded.aggregate.snapshot",
meta,
%{
kind: :internal,
attributes: attributes
}
)
end

def handle_telemetry_event(
[:commanded, :aggregate, :snapshot, :stop],
_measurements,
meta,
_config
) do
ctx = OpentelemetryTelemetry.set_current_telemetry_span(@tracer_id, meta)

if error = meta[:error] do
Span.set_attribute(
ctx,
ErrorAttributes.error_type(),
Helpers.to_error_type(error, @tracer_id)
)

Span.set_status(ctx, OpenTelemetry.status(:error, Helpers.format_error(error)))
end

OpentelemetryTelemetry.end_telemetry_span(@tracer_id, meta)
end

def handle_telemetry_event(
[:commanded, :aggregate, :snapshot, :exception],
_measurements,
%{kind: kind, reason: reason, stacktrace: stacktrace} = meta,
_config
) do
ctx = OpentelemetryTelemetry.set_current_telemetry_span(@tracer_id, meta)

Span.set_attribute(ctx, :"erlang.exception.kind", kind)

exception = Exception.normalize(kind, reason, stacktrace)

Span.set_attribute(
ctx,
ErrorAttributes.error_type(),
Helpers.to_error_type(exception, @tracer_id)
)

Span.record_exception(ctx, exception, stacktrace)

Span.set_status(
ctx,
OpenTelemetry.status(:error, Exception.format_banner(kind, reason, stacktrace))
)

OpentelemetryTelemetry.end_telemetry_span(@tracer_id, meta)
end

defp maybe_add_snapshot_every(attrs, nil), do: attrs

defp maybe_add_snapshot_every(attrs, snapshot_every) when is_integer(snapshot_every) do
[{CommandedAttributes.commanded_snapshot_every(), snapshot_every} | attrs]
end

defp maybe_add_snapshot_every(attrs, _), do: attrs

defp maybe_add_snapshot_module_version(attrs, nil), do: attrs

defp maybe_add_snapshot_module_version(attrs, version) when is_integer(version) do
[{CommandedAttributes.commanded_snapshot_module_version(), version} | attrs]
end

defp maybe_add_snapshot_module_version(attrs, _), do: attrs
end
Loading
Loading