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
5 changes: 5 additions & 0 deletions lib/commanded/aggregates/aggregate.ex
Original file line number Diff line number Diff line change
Expand Up @@ -71,6 +71,7 @@ defmodule Commanded.Aggregates.Aggregate do
measurements: "%{system_time: integer()}",
metadata: """
%{application: Commanded.Application.t(),
aggregate_module: module(),
aggregate_uuid: String.t(),
aggregate_version: non_neg_integer(),
snapshot_every: non_neg_integer() | nil,
Expand All @@ -84,6 +85,7 @@ defmodule Commanded.Aggregates.Aggregate do
measurements: "%{duration: non_neg_integer()}",
metadata: """
%{application: Commanded.Application.t(),
aggregate_module: module(),
aggregate_uuid: String.t(),
aggregate_version: non_neg_integer(),
snapshot_every: non_neg_integer() | nil,
Expand All @@ -98,6 +100,7 @@ defmodule Commanded.Aggregates.Aggregate do
measurements: "%{duration: non_neg_integer()}",
metadata: """
%{application: Commanded.Application.t(),
aggregate_module: module(),
aggregate_uuid: String.t(),
aggregate_version: non_neg_integer(),
snapshot_every: non_neg_integer() | nil,
Expand Down Expand Up @@ -680,6 +683,7 @@ defmodule Commanded.Aggregates.Aggregate do
defp do_take_snapshot(%Aggregate{} = state) do
%Aggregate{
application: application,
aggregate_module: aggregate_module,
aggregate_uuid: aggregate_uuid,
aggregate_state: aggregate_state,
aggregate_version: aggregate_version,
Expand All @@ -691,6 +695,7 @@ defmodule Commanded.Aggregates.Aggregate do

meta = %{
application: application,
aggregate_module: aggregate_module,
aggregate_uuid: aggregate_uuid,
aggregate_version: aggregate_version,
snapshot_every: snapshot_every,
Expand Down
6 changes: 6 additions & 0 deletions lib/commanded/aggregates/aggregate_state_builder.ex
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ defmodule Commanded.Aggregates.AggregateStateBuilder do
measurements: "%{system_time: integer()}",
metadata: """
%{application: Commanded.Application.t(),
aggregate_module: module(),
aggregate_uuid: String.t(),
aggregate_state: struct(),
aggregate_version: non_neg_integer()}
Expand All @@ -26,6 +27,7 @@ defmodule Commanded.Aggregates.AggregateStateBuilder do
measurements: "%{duration: non_neg_integer(), count: non_neg_integer()}",
metadata: """
%{application: Commanded.Application.t(),
aggregate_module: module(),
aggregate_uuid: String.t(),
aggregate_state: struct(),
aggregate_version: non_neg_integer(),
Expand All @@ -40,6 +42,7 @@ defmodule Commanded.Aggregates.AggregateStateBuilder do
measurements: "%{system_time: integer()}",
metadata: """
%{application: Commanded.Application.t(),
aggregate_module: module(),
aggregate_uuid: String.t(),
aggregate_state: struct(),
aggregate_version: non_neg_integer()}
Expand All @@ -52,6 +55,7 @@ defmodule Commanded.Aggregates.AggregateStateBuilder do
measurements: "%{duration: non_neg_integer(), count: non_neg_integer()}",
metadata: """
%{application: Commanded.Application.t(),
aggregate_module: module(),
aggregate_uuid: String.t(),
aggregate_state: struct(),
aggregate_version: non_neg_integer()}
Expand Down Expand Up @@ -189,13 +193,15 @@ defmodule Commanded.Aggregates.AggregateStateBuilder do
defp telemetry_metadata(%Aggregate{} = state) do
%Aggregate{
application: application,
aggregate_module: aggregate_module,
aggregate_uuid: aggregate_uuid,
aggregate_state: aggregate_state,
aggregate_version: aggregate_version
} = state

%{
application: application,
aggregate_module: aggregate_module,
aggregate_uuid: aggregate_uuid,
aggregate_state: aggregate_state,
aggregate_version: aggregate_version
Expand Down
13 changes: 11 additions & 2 deletions lib/commanded/opentelemetry/aggregate_populate.ex
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ defmodule Commanded.OpenTelemetry.AggregatePopulate do
@moduledoc false

alias Commanded.OpenTelemetry.CommandedAttributes
alias Commanded.OpenTelemetry.Helpers
alias OpenTelemetry.SemConv.Incubating.CodeAttributes
alias OpenTelemetry.SemConv.Incubating.MessagingAttributes
alias OpenTelemetry.Span
Expand Down Expand Up @@ -29,19 +30,23 @@ defmodule Commanded.OpenTelemetry.AggregatePopulate do
meta,
_config
) do
aggregate_module_name = Helpers.module_name(meta.aggregate_module)

attributes = [
{MessagingAttributes.messaging_system(), "commanded"},
{MessagingAttributes.messaging_operation_type(), :receive},
{MessagingAttributes.messaging_operation_name(), "load"},
{MessagingAttributes.messaging_destination_name(), aggregate_module_name},
{CodeAttributes.code_function(), "load"},
{CodeAttributes.code_namespace(), aggregate_module_name},
{CommandedAttributes.commanded_application(), meta.application},
{CommandedAttributes.commanded_aggregate_uuid(), meta.aggregate_uuid},
{CommandedAttributes.commanded_aggregate_version(), meta.aggregate_version}
]

OpentelemetryTelemetry.start_telemetry_span(
@tracer_id,
"commanded.aggregate.load",
"load #{aggregate_module_name}",
meta,
%{
kind: :internal,
Expand Down Expand Up @@ -90,13 +95,17 @@ defmodule Commanded.OpenTelemetry.AggregatePopulate do
meta,
_config
) do
aggregate_module_name = Helpers.module_name(meta.aggregate_module)

attributes = [
# OTel Messaging SemConv
{MessagingAttributes.messaging_system(), "commanded"},
{MessagingAttributes.messaging_operation_type(), :receive},
{MessagingAttributes.messaging_operation_name(), "populate"},
{MessagingAttributes.messaging_destination_name(), aggregate_module_name},
# OTel Code SemConv
{CodeAttributes.code_function(), "populate"},
{CodeAttributes.code_namespace(), aggregate_module_name},
# Commanded-specific
{CommandedAttributes.commanded_application(), meta.application},
{CommandedAttributes.commanded_aggregate_uuid(), meta.aggregate_uuid},
Expand All @@ -105,7 +114,7 @@ defmodule Commanded.OpenTelemetry.AggregatePopulate do

OpentelemetryTelemetry.start_telemetry_span(
@tracer_id,
"commanded.aggregate.populate",
"populate #{aggregate_module_name}",
meta,
%{
kind: :internal,
Expand Down
6 changes: 5 additions & 1 deletion lib/commanded/opentelemetry/aggregate_snapshot.ex
Original file line number Diff line number Diff line change
Expand Up @@ -30,12 +30,16 @@ defmodule Commanded.OpenTelemetry.AggregateSnapshot do
meta,
_config
) do
aggregate_module_name = Helpers.module_name(meta.aggregate_module)

attributes =
[
{MessagingAttributes.messaging_system(), "commanded"},
{MessagingAttributes.messaging_operation_type(), :publish},
{MessagingAttributes.messaging_operation_name(), "snapshot"},
{MessagingAttributes.messaging_destination_name(), aggregate_module_name},
{CodeAttributes.code_function(), "snapshot"},
{CodeAttributes.code_namespace(), aggregate_module_name},
{CommandedAttributes.commanded_application(), meta.application},
{CommandedAttributes.commanded_aggregate_uuid(), meta.aggregate_uuid},
{CommandedAttributes.commanded_aggregate_version(), meta.aggregate_version}
Expand All @@ -45,7 +49,7 @@ defmodule Commanded.OpenTelemetry.AggregateSnapshot do

OpentelemetryTelemetry.start_telemetry_span(
@tracer_id,
"commanded.aggregate.snapshot",
"snapshot #{aggregate_module_name}",
meta,
%{
kind: :internal,
Expand Down
18 changes: 13 additions & 5 deletions test/opentelemetry/aggregate_populate_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -73,7 +73,7 @@ defmodule Commanded.OpenTelemetry.AggregatePopulateTest do

assert_receive {:span,
span(
name: "commanded.aggregate.populate",
name: "populate MockAggregate",
kind: :internal,
attributes: attributes
)},
Expand All @@ -85,7 +85,9 @@ defmodule Commanded.OpenTelemetry.AggregatePopulateTest do
"messaging.system": "commanded",
"messaging.operation.type": :receive,
"messaging.operation.name": "populate",
"messaging.destination.name": "MockAggregate",
"code.function": "populate",
"code.namespace": "MockAggregate",
"commanded.application": MockApp,
"commanded.aggregate.uuid": aggregate_uuid,
"commanded.aggregate.version": 5,
Expand Down Expand Up @@ -123,7 +125,7 @@ defmodule Commanded.OpenTelemetry.AggregatePopulateTest do

assert_receive {:span,
span(
name: "commanded.aggregate.load",
name: "load MockAggregate",
kind: :internal,
attributes: attributes
)},
Expand All @@ -133,7 +135,9 @@ defmodule Commanded.OpenTelemetry.AggregatePopulateTest do
"messaging.system": "commanded",
"messaging.operation.type": :receive,
"messaging.operation.name": "load",
"messaging.destination.name": "MockAggregate",
"code.function": "load",
"code.namespace": "MockAggregate",
"commanded.application": MockApp,
"commanded.aggregate.uuid": aggregate_uuid,
"commanded.aggregate.version": 0,
Expand Down Expand Up @@ -164,7 +168,7 @@ defmodule Commanded.OpenTelemetry.AggregatePopulateTest do

assert_receive {:span,
span(
name: "commanded.aggregate.load",
name: "load MockAggregate",
kind: :internal,
attributes: attributes
)},
Expand Down Expand Up @@ -198,7 +202,7 @@ defmodule Commanded.OpenTelemetry.AggregatePopulateTest do

assert_receive {:span,
span(
name: "commanded.aggregate.populate",
name: "populate MockAggregate",
kind: :internal,
attributes: attributes
)},
Expand All @@ -208,7 +212,9 @@ defmodule Commanded.OpenTelemetry.AggregatePopulateTest do
"messaging.system": "commanded",
"messaging.operation.type": :receive,
"messaging.operation.name": "populate",
"messaging.destination.name": "MockAggregate",
"code.function": "populate",
"code.namespace": "MockAggregate",
"commanded.application": MockApp,
"commanded.aggregate.uuid": aggregate_uuid,
"commanded.aggregate.version": 0,
Expand All @@ -232,7 +238,7 @@ defmodule Commanded.OpenTelemetry.AggregatePopulateTest do

assert_receive {:span,
span(
name: "commanded.aggregate.populate",
name: "populate MockAggregate",
kind: :internal,
attributes: attributes
)},
Expand All @@ -242,7 +248,9 @@ defmodule Commanded.OpenTelemetry.AggregatePopulateTest do
"messaging.system": "commanded",
"messaging.operation.type": :receive,
"messaging.operation.name": "populate",
"messaging.destination.name": "MockAggregate",
"code.function": "populate",
"code.namespace": "MockAggregate",
"commanded.application": MockApp,
"commanded.aggregate.uuid": aggregate_uuid,
"commanded.aggregate.version": 10,
Expand Down
8 changes: 5 additions & 3 deletions test/opentelemetry/aggregate_snapshot_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -72,7 +72,7 @@ defmodule Commanded.OpenTelemetry.AggregateSnapshotTest do

assert_receive {:span,
span(
name: "commanded.aggregate.snapshot",
name: "snapshot MockAggregate",
kind: :internal,
attributes: attributes
)},
Expand All @@ -82,7 +82,9 @@ defmodule Commanded.OpenTelemetry.AggregateSnapshotTest do
"messaging.system": "commanded",
"messaging.operation.type": :publish,
"messaging.operation.name": "snapshot",
"messaging.destination.name": "MockAggregate",
"code.function": "snapshot",
"code.namespace": "MockAggregate",
"commanded.application": MockApp,
"commanded.aggregate.uuid": aggregate_uuid,
"commanded.aggregate.version": 10,
Expand All @@ -108,7 +110,7 @@ defmodule Commanded.OpenTelemetry.AggregateSnapshotTest do

assert_receive {:span,
span(
name: "commanded.aggregate.snapshot",
name: "snapshot MockAggregate",
kind: :internal,
attributes: attributes
)},
Expand Down Expand Up @@ -143,7 +145,7 @@ defmodule Commanded.OpenTelemetry.AggregateSnapshotTest do

assert_receive {:span,
span(
name: "commanded.aggregate.snapshot",
name: "snapshot MockAggregate",
status: {:status, :error, _error_message},
attributes: span_attrs
)},
Expand Down
4 changes: 4 additions & 0 deletions test/support/factory.ex
Original file line number Diff line number Diff line change
Expand Up @@ -601,6 +601,7 @@ defmodule Commanded.TestSupport.Factory do

defaults = [
application: Keyword.get(opts, :application, MockApp),
aggregate_module: Keyword.get(opts, :aggregate_module, MockAggregate),
aggregate_uuid: aggregate_uuid,
aggregate_state: Keyword.get(opts, :aggregate_state, %{}),
aggregate_version: Keyword.get(opts, :aggregate_version, 0)
Expand All @@ -610,6 +611,7 @@ defmodule Commanded.TestSupport.Factory do

%{
application: Keyword.fetch!(opts, :application),
aggregate_module: Keyword.fetch!(opts, :aggregate_module),
aggregate_uuid: Keyword.fetch!(opts, :aggregate_uuid),
aggregate_state: Keyword.fetch!(opts, :aggregate_state),
aggregate_version: Keyword.fetch!(opts, :aggregate_version)
Expand All @@ -630,6 +632,7 @@ defmodule Commanded.TestSupport.Factory do

defaults = [
application: Keyword.get(opts, :application, MockApp),
aggregate_module: Keyword.get(opts, :aggregate_module, MockAggregate),
aggregate_uuid: aggregate_uuid,
aggregate_version: Keyword.get(opts, :aggregate_version, 0),
snapshot_every: Keyword.get(opts, :snapshot_every),
Expand All @@ -640,6 +643,7 @@ defmodule Commanded.TestSupport.Factory do

%{
application: Keyword.fetch!(opts, :application),
aggregate_module: Keyword.fetch!(opts, :aggregate_module),
aggregate_uuid: Keyword.fetch!(opts, :aggregate_uuid),
aggregate_version: Keyword.fetch!(opts, :aggregate_version),
snapshot_every: Keyword.get(opts, :snapshot_every),
Expand Down
Loading