feat: add OpenTelemetry support for aggregate snapshots - #60
Conversation
PR SummaryLow Risk Overview Extends aggregate load telemetry to report whether state rebuild started from a snapshot ( Introduces Written by Cursor Bugbot for commit cc6e0f5. This will update automatically on new commits. Configure here. |
WalkthroughThis pull request introduces telemetry and OpenTelemetry instrumentation for aggregate snapshot operations. It adds snapshot lifecycle events (start, stop, exception) with metadata tracking, extends the aggregate populate flow to record snapshot usage, creates OpenTelemetry span handlers for snapshot monitoring, and includes comprehensive test coverage for the new instrumentation. Changes
Sequence DiagramsequenceDiagram
participant App as Application
participant Agg as Aggregate
participant Telem as Telemetry
participant Builder as AggregateStateBuilder
participant OT as OpenTelemetry Handler
App->>Agg: do_take_snapshot()
Agg->>Telem: emit snapshot:start<br/>(uuid, version, snapshot_every...)
Telem->>OT: handle_telemetry_event(:start)
OT->>OT: create span "commanded.aggregate.snapshot"
Agg->>Agg: execute snapshot operation
alt snapshot success
Agg->>Telem: emit snapshot:stop<br/>(with metadata)
Telem->>OT: handle_telemetry_event(:stop)
OT->>OT: end span with status ok
else snapshot error
Agg->>Telem: emit snapshot:exception<br/>(kind, reason, stacktrace)
Telem->>OT: handle_telemetry_event(:exception)
OT->>OT: record exception, set error status
OT->>OT: end span with error
end
App->>Builder: populate(aggregate)
Builder->>Builder: check snapshot exists?
alt snapshot found
Builder->>Builder: load snapshot data<br/>snapshot_used = true
else no snapshot
Builder->>Builder: initialize empty<br/>snapshot_used = false
end
Builder->>Builder: rebuild_from_events(opts)
Builder->>Telem: emit load:stop<br/>(snapshot_used, snapshot_source_version...)
Telem->>OT: update span with<br/>snapshot attributes
Estimated code review effort🎯 3 (Moderate) | ⏱️ ~25 minutes Possibly related PRs
Poem
🚥 Pre-merge checks | ✅ 2 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (2 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches
🧪 Generate unit tests (beta)
📝 Coding Plan
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
e61b7e4 to
a1c4bc8
Compare
There was a problem hiding this comment.
Cursor Bugbot has reviewed your changes and found 1 potential issue.
Bugbot Autofix is OFF. To automatically fix reported issues with cloud agents, have a team admin enable autofix in the Cursor dashboard.
d32c23c to
ae2e7a4
Compare
Signed-off-by: Yordis Prieto <yordis.prieto@gmail.com>
There was a problem hiding this comment.
🧹 Nitpick comments (3)
lib/commanded/opentelemetry/aggregate_snapshot.ex (1)
106-120: Consider consolidating the helper functions to reduce duplication.Both
maybe_add_snapshot_every/2andmaybe_add_snapshot_module_version/2follow the identical pattern. You could use a single generic helper.♻️ Optional: Consolidate with a generic helper
+ defp maybe_add_attribute(attrs, _key, nil), do: attrs + + defp maybe_add_attribute(attrs, key, value) when is_integer(value) do + [{key, value} | attrs] + end + + defp maybe_add_attribute(attrs, _key, _), do: attrs + 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] + maybe_add_attribute(attrs, CommandedAttributes.commanded_snapshot_every(), snapshot_every) 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] + maybe_add_attribute(attrs, CommandedAttributes.commanded_snapshot_module_version(), version) end - - defp maybe_add_snapshot_module_version(attrs, _), do: attrs🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@lib/commanded/opentelemetry/aggregate_snapshot.ex` around lines 106 - 120, Both maybe_add_snapshot_every/2 and maybe_add_snapshot_module_version/2 duplicate the same pattern; replace them with a single generic helper (e.g., maybe_add_optional_integer_attr(attrs, value, attr_key_fn)) that returns attrs unchanged for nil or non-integer and prepends {attr_key_fn.(), value} when value is an integer; update callers of maybe_add_snapshot_every/2 and maybe_add_snapshot_module_version/2 to call the new helper supplying the appropriate key function (CommandedAttributes.commanded_snapshot_every and CommandedAttributes.commanded_snapshot_module_version) and remove the old duplicated functions.test/opentelemetry/snapshotting_postgres_test.exs (1)
74-83: Renamestart_event_store/0to better reflect its purpose.The function doesn't actually start the event store; it only registers an
on_exitcallback for cleanup. The actual event store startup appears to happen elsewhere (likely via application config or the supervisedApp).♻️ Rename for clarity
- defp start_event_store do + defp setup_event_store_cleanup do alias Commanded.EventStore.Adapters.EventStore.Storage config = Storage.config()And update the caller on line 31:
setup do - start_event_store() + setup_event_store_cleanup() start_supervised!({App, snapshotting: %{ExampleAggregate => [snapshot_every: 10]}}) :ok end🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@test/opentelemetry/snapshotting_postgres_test.exs` around lines 74 - 83, The helper start_event_store/0 is misnamed because it only registers an on_exit cleanup and does not start the event store; rename it (for example to register_event_store_cleanup/0) and update its caller accordingly so intent is clear; ensure the function still aliases Commanded.EventStore.Adapters.EventStore.Storage and continues to call Storage.config(), Storage.connect(config) and Storage.reset!(conn, config) inside the on_exit callback, and rename any references to start_event_store/0 to the new function name.test/opentelemetry/aggregate_snapshot_test.exs (1)
156-166: Potential issue: Detaching handlers while iterating may cause issues with shared handler IDs.If the same handler ID is registered for multiple events (which is the case here with
{__MODULE__, :snapshot}), detaching it in the inner loop for the first event will cause subsequent detach calls for the same ID to fail or be no-ops.However, looking at the implementation in
aggregate_snapshot.exline 16, the handler usesattach_manywith a single ID for all events. This means:telemetry.detach/1only needs to be called once.♻️ Simplify handler detachment
defp detach_snapshot_handlers do - for event <- [ - [:commanded, :aggregate, :snapshot, :start], - [:commanded, :aggregate, :snapshot, :stop], - [:commanded, :aggregate, :snapshot, :exception] - ] do - for handler <- :telemetry.list_handlers(event) do - :telemetry.detach(handler.id) - end - end + # The handler uses attach_many with a single ID, so we only need to check one event + for handler <- :telemetry.list_handlers([:commanded, :aggregate, :snapshot, :start]) do + :telemetry.detach(handler.id) + end end🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@test/opentelemetry/aggregate_snapshot_test.exs` around lines 156 - 166, The detach_snapshot_handlers helper iterates events and detaches each handler repeatedly, but the snapshot handler was attached with a single shared ID (via attach_many) so subsequent detaches are redundant/failing; update detach_snapshot_handlers to collect unique handler ids (or directly detach the shared id used by aggregate_snapshot's attach_many, e.g. {__MODULE__, :snapshot}) and call :telemetry.detach/1 only once per unique id (ensure you reference the existing helper detach_snapshot_handlers and the handler id used by aggregate_snapshot's attach_many).
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Nitpick comments:
In `@lib/commanded/opentelemetry/aggregate_snapshot.ex`:
- Around line 106-120: Both maybe_add_snapshot_every/2 and
maybe_add_snapshot_module_version/2 duplicate the same pattern; replace them
with a single generic helper (e.g., maybe_add_optional_integer_attr(attrs,
value, attr_key_fn)) that returns attrs unchanged for nil or non-integer and
prepends {attr_key_fn.(), value} when value is an integer; update callers of
maybe_add_snapshot_every/2 and maybe_add_snapshot_module_version/2 to call the
new helper supplying the appropriate key function
(CommandedAttributes.commanded_snapshot_every and
CommandedAttributes.commanded_snapshot_module_version) and remove the old
duplicated functions.
In `@test/opentelemetry/aggregate_snapshot_test.exs`:
- Around line 156-166: The detach_snapshot_handlers helper iterates events and
detaches each handler repeatedly, but the snapshot handler was attached with a
single shared ID (via attach_many) so subsequent detaches are redundant/failing;
update detach_snapshot_handlers to collect unique handler ids (or directly
detach the shared id used by aggregate_snapshot's attach_many, e.g. {__MODULE__,
:snapshot}) and call :telemetry.detach/1 only once per unique id (ensure you
reference the existing helper detach_snapshot_handlers and the handler id used
by aggregate_snapshot's attach_many).
In `@test/opentelemetry/snapshotting_postgres_test.exs`:
- Around line 74-83: The helper start_event_store/0 is misnamed because it only
registers an on_exit cleanup and does not start the event store; rename it (for
example to register_event_store_cleanup/0) and update its caller accordingly so
intent is clear; ensure the function still aliases
Commanded.EventStore.Adapters.EventStore.Storage and continues to call
Storage.config(), Storage.connect(config) and Storage.reset!(conn, config)
inside the on_exit callback, and rename any references to start_event_store/0 to
the new function name.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: CHILL
Plan: Pro
Run ID: 884af9ec-d669-425e-b6e9-bcb09fe1a29d
📒 Files selected for processing (12)
lib/commanded/aggregates/aggregate.exlib/commanded/aggregates/aggregate_state_builder.exlib/commanded/opentelemetry.exlib/commanded/opentelemetry/aggregate_populate.exlib/commanded/opentelemetry/aggregate_snapshot.exlib/commanded/opentelemetry/commanded_attributes.extest/aggregates/aggregate_telemetry_test.exstest/opentelemetry/aggregate_populate_test.exstest/opentelemetry/aggregate_snapshot_test.exstest/opentelemetry/snapshotting_postgres_test.exstest/support/factory.extest/support/opentelemetry_case.ex
There was a problem hiding this comment.
🧹 Nitpick comments (3)
test/aggregates/snapshotting_postgres_test.exs (2)
55-57: Redundant assertions after pattern match.Lines 44-53 already pattern-match and verify
snapshot.data.itemsandsnapshot.data.last_indexvalues via the%ExampleAggregate{items: [1, 2, 3, 4, 5, 6, 7, 8, 9, 10], last_index: 10}binding. These subsequent assertions are redundant.🧹 Suggested cleanup
} = snapshot - assert is_struct(snapshot.data, ExampleAggregate) - assert snapshot.data.items == Enum.to_list(1..10) - assert snapshot.data.last_index == 10 + # Pattern match above already validates structure and values end🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@test/aggregates/snapshotting_postgres_test.exs` around lines 55 - 57, The test repeats assertions already enforced by the earlier pattern match against %ExampleAggregate{items: [1,2,3,4,5,6,7,8,9,10], last_index: 10}; remove the redundant lines asserting is_struct(snapshot.data, ExampleAggregate) and the equality checks snapshot.data.items == Enum.to_list(1..10) and snapshot.data.last_index == 10 so the test relies on the existing pattern match for those validations.
74-82: Misleading function name:start_event_storeonly sets up teardown.The function doesn't actually start the event store—it only registers an
on_exitcallback to reset storage. Consider renaming tosetup_event_store_cleanuporregister_event_store_resetto better reflect its behavior.📝 Suggested rename
- defp start_event_store do + defp setup_event_store_cleanup do alias Commanded.EventStore.Adapters.EventStore.Storage config = Storage.config()And update the call in setup:
setup do - start_event_store() + setup_event_store_cleanup() start_supervised!({App, snapshotting: %{ExampleAggregate => [snapshot_every: 10]}})🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@test/aggregates/snapshotting_postgres_test.exs` around lines 74 - 82, The helper start_event_store/0 is misleading because it only registers teardown logic; rename the function (e.g., setup_event_store_cleanup or register_event_store_reset) and update all usages in the test module to the new name; inside the renamed function keep the existing alias Commanded.EventStore.Adapters.EventStore.Storage, Storage.config(), the on_exit(fn -> {:ok, conn} = Storage.connect(config); Storage.reset!(conn, config) end) body unchanged so teardown continues to run.test/opentelemetry/aggregate_snapshot_test.exs (1)
51-56: Consider removing redundant setup block.This nested setup (lines 52-56) duplicates the module-level setup (lines 12-17) since both call
detach_snapshot_handlers()followed byAggregateSnapshot.setup(). The module-level setup already ensures clean handler state before each test.🔧 Suggested simplification
describe "snapshot spans" do - setup do - detach_snapshot_handlers() - AggregateSnapshot.setup() - :ok - end - test "creates span with correct attributes" do🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@test/opentelemetry/aggregate_snapshot_test.exs` around lines 51 - 56, The nested setup in the "snapshot spans" describe block is redundant because the module-level setup already calls detach_snapshot_handlers() and AggregateSnapshot.setup(); remove the inner setup block (the anonymous setup that calls detach_snapshot_handlers() and AggregateSnapshot.setup()) so tests rely on the module-level setup instead, ensuring you only keep the top-level setup that invokes detach_snapshot_handlers() and AggregateSnapshot.setup().
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Nitpick comments:
In `@test/aggregates/snapshotting_postgres_test.exs`:
- Around line 55-57: The test repeats assertions already enforced by the earlier
pattern match against %ExampleAggregate{items: [1,2,3,4,5,6,7,8,9,10],
last_index: 10}; remove the redundant lines asserting is_struct(snapshot.data,
ExampleAggregate) and the equality checks snapshot.data.items ==
Enum.to_list(1..10) and snapshot.data.last_index == 10 so the test relies on the
existing pattern match for those validations.
- Around line 74-82: The helper start_event_store/0 is misleading because it
only registers teardown logic; rename the function (e.g.,
setup_event_store_cleanup or register_event_store_reset) and update all usages
in the test module to the new name; inside the renamed function keep the
existing alias Commanded.EventStore.Adapters.EventStore.Storage,
Storage.config(), the on_exit(fn -> {:ok, conn} = Storage.connect(config);
Storage.reset!(conn, config) end) body unchanged so teardown continues to run.
In `@test/opentelemetry/aggregate_snapshot_test.exs`:
- Around line 51-56: The nested setup in the "snapshot spans" describe block is
redundant because the module-level setup already calls
detach_snapshot_handlers() and AggregateSnapshot.setup(); remove the inner setup
block (the anonymous setup that calls detach_snapshot_handlers() and
AggregateSnapshot.setup()) so tests rely on the module-level setup instead,
ensuring you only keep the top-level setup that invokes
detach_snapshot_handlers() and AggregateSnapshot.setup().
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: CHILL
Plan: Pro
Run ID: aeddb4ec-2a0f-4bbf-a238-a296bc5e712b
📒 Files selected for processing (12)
lib/commanded/aggregates/aggregate.exlib/commanded/aggregates/aggregate_state_builder.exlib/commanded/opentelemetry.exlib/commanded/opentelemetry/aggregate_populate.exlib/commanded/opentelemetry/aggregate_snapshot.exlib/commanded/opentelemetry/commanded_attributes.extest/aggregates/aggregate_telemetry_test.exstest/aggregates/snapshotting_postgres_test.exstest/opentelemetry/aggregate_populate_test.exstest/opentelemetry/aggregate_snapshot_test.exstest/support/factory.extest/support/opentelemetry_case.ex
🚧 Files skipped from review as they are similar to previous changes (1)
- lib/commanded/opentelemetry/aggregate_populate.ex

Adds OpenTelemetry support for aggregate snapshot operations.
Changes
[:commanded, :aggregate, :snapshot]start/stop/exception events when aggregates take snapshots:load :stopmetadata withsnapshot_usedandsnapshot_source_versionso OTel spans can report when a snapshot was used during populateCommanded.OpenTelemetry.AggregateSnapshotattaches to snapshot telemetry and creates spans with application, aggregate UUID/version, snapshot_every, and snapshot_module_version attributesConfiguration
Commanded.OpenTelemetry.setup(aggregate_snapshot: :disabled)to disable snapshot tracingSigned-off-by: Yordis Prieto yordis.prieto@gmail.com