fix(event_store): defer stream_forward telemetry until stream consumed - #64
Conversation
PR SummaryMedium Risk Overview Replaces the prior Extends telemetry tests to cover list vs lazy-stream behaviour, early halt, and exception emission, and updates the shared test handler to also subscribe to Written by Cursor Bugbot for commit 57ed763. This will update automatically on new commits. Configure here. |
|
Warning Rate limit exceeded
Your organization is not enrolled in usage-based pricing. Contact your admin to enable usage-based pricing to continue reviews beyond the rate limit, or try again in 2 minutes and 25 seconds. ⌛ How to resolve this issue?After the wait time has elapsed, a review can be triggered using the We recommend that you space out your commits to avoid hitting the rate limit. 🚦 How do rate limits work?CodeRabbit enforces hourly rate limits for each developer per organization. Our paid plans have higher rate limits than the trial, open-source and free plans. In all cases, we re-allow further reviews after a brief timeout. Please see our FAQ for further information. ℹ️ Review info⚙️ Run configurationConfiguration used: Organization UI Review profile: CHILL Plan: Pro Run ID: 📒 Files selected for processing (2)
WalkthroughEventStore.stream_forward/4 now emits explicit telemetry: it emits a start event with span context and times, resolves the adapter and calls adapter.stream_forward/4, emits stop immediately for eager lists or defers stop until enumeration for lazy streams, and emits exception telemetry (with kind/reason/stacktrace) on errors. Changes
Sequence Diagram(s)sequenceDiagram
participant Client as Client
participant ES as EventStore
participant Adapter as Adapter
participant Telemetry as Telemetry
Client->>ES: call stream_forward(...)
ES->>Telemetry: emit [:commanded,:event_store,:stream_forward,:start] (monotonic/system, telemetry_span_context)
ES->>Adapter: resolve & call adapter.stream_forward(...)
alt Adapter returns immediate list
Adapter-->>ES: list
ES->>Telemetry: emit [:...:stream_forward,:stop] (duration)
ES-->>Client: return list
else Adapter returns lazy enumerable
Adapter-->>ES: enumerable
ES-->>Client: return wrapped enumerable (stop deferred)
Client->>ES: enumerate wrapped stream
ES->>Telemetry: emit [:...:stream_forward,:stop] when enumeration completes
end
opt resolution or adapter raises
ES->>Telemetry: emit [:...:stream_forward,:exception] (kind, reason, stacktrace, duration, span context)
ES-->>Client: re-raise exception
end
Estimated code review effort🎯 3 (Moderate) | ⏱️ ~20 minutes Possibly related PRs
Poem
🚥 Pre-merge checks | ✅ 3✅ Passed checks (3 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches🧪 Generate unit tests (beta)
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 |
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.
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@lib/commanded/event_store.ex`:
- Around line 308-321: The telemetry block inside the Stream.transform call can
fail to emit any telemetry if the downstream consumer raises (leaving the
OpenTelemetry span orphaned); update the stream-wrapping logic (the
Stream.transform usage in wrap_stream_forward_telemetry or the function that
creates the stream telemetry) so that exceptions during enumeration emit a
[:commanded, :event_store, :stream_forward, :exception] telemetry event with
duration, monotonic_time, and the original span_context, and always ensure the
span_context is included in the telemetry metadata; implement this by wrapping
the stream enumeration with a construct that catches errors (e.g., convert to a
Stream.resource-based wrapper or add a try/rescue around enumeration) to call
:telemetry.execute for :exception before re-raising, and keep the existing :stop
telemetry on normal completion.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: CHILL
Plan: Pro
Run ID: 2ce22e2e-15e4-4007-bce9-062201c798e3
📒 Files selected for processing (4)
lib/commanded/event_store.extest/application/dynamic_applications_test.exstest/commands/dispatch_command_test.exstest/event_store/stream_forward_deferred_telemetry_test.exs
768eb19 to
44048d2
Compare
There was a problem hiding this comment.
🧹 Nitpick comments (1)
test/event_store/telemetry_test.exs (1)
107-107: Avoid runtime-generated atoms for telemetry handler IDs.The handler ID is created via dynamic atom interpolation with
System.unique_integer/1. Runtime-generated atoms accumulate in the BEAM atom table and are never garbage collected. Since:telemetry.attach_many/4accepts any term, use a non-atom identifier instead.♻️ Suggested change
- handler = :"stream_forward_ex-#{System.unique_integer([:positive])}" + handler = {__MODULE__, :stream_forward_ex, make_ref()}🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@test/event_store/telemetry_test.exs` at line 107, The test currently builds a runtime atom for the telemetry handler ID via handler = :"stream_forward_ex-#{System.unique_integer([:positive])}", which leaks atoms; change the handler to a non-atom term (for example a tuple or string) and pass that to :telemetry.attach_many/4 instead (e.g., use a tuple like {:stream_forward_ex, System.unique_integer([:positive])} or a string interpolation) so the handler id is not an atom while keeping uniqueness for tests.
🤖 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/event_store/telemetry_test.exs`:
- Line 107: The test currently builds a runtime atom for the telemetry handler
ID via handler = :"stream_forward_ex-#{System.unique_integer([:positive])}",
which leaks atoms; change the handler to a non-atom term (for example a tuple or
string) and pass that to :telemetry.attach_many/4 instead (e.g., use a tuple
like {:stream_forward_ex, System.unique_integer([:positive])} or a string
interpolation) so the handler id is not an atom while keeping uniqueness for
tests.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: CHILL
Plan: Pro
Run ID: 0b1c71e8-180d-4792-84fa-1cbfa26c1bb5
📒 Files selected for processing (2)
lib/commanded/event_store.extest/event_store/telemetry_test.exs
🚧 Files skipped from review as they are similar to previous changes (1)
- lib/commanded/event_store.ex
44048d2 to
94de2c0
Compare
There was a problem hiding this comment.
🧹 Nitpick comments (1)
test/event_store/telemetry_test.exs (1)
193-197: Minor inconsistency: missingduration >= 0assertion.The full enumeration test (line 152) asserts both
is_integer(stop_meas.duration)andstop_meas.duration >= 0, but this halted-stream test only asserts the integer check. For consistency and completeness, consider adding the non-negative assertion here as well.Proposed fix
assert_receive {:deferred, [:commanded, :event_store, :stream_forward, :stop], stop_meas, _} assert is_integer(stop_meas.duration) + assert stop_meas.duration >= 0 end🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@test/event_store/telemetry_test.exs` around lines 193 - 197, The deferred stop measurement assertion is missing a non-negative check: after the existing assert_receive that binds stop_meas (matching {:deferred, [:commanded, :event_store, :stream_forward, :stop], stop_meas, _}), add an assertion that stop_meas.duration >= 0 to mirror the full enumeration test; keep the existing is_integer(stop_meas.duration) check and simply append assert stop_meas.duration >= 0.
🤖 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/event_store/telemetry_test.exs`:
- Around line 193-197: The deferred stop measurement assertion is missing a
non-negative check: after the existing assert_receive that binds stop_meas
(matching {:deferred, [:commanded, :event_store, :stream_forward, :stop],
stop_meas, _}), add an assertion that stop_meas.duration >= 0 to mirror the full
enumeration test; keep the existing is_integer(stop_meas.duration) check and
simply append assert stop_meas.duration >= 0.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: CHILL
Plan: Pro
Run ID: 79aab0c4-fd0b-4714-8d6e-99b4e55c2386
📒 Files selected for processing (2)
lib/commanded/event_store.extest/event_store/telemetry_test.exs
✅ Files skipped from review due to trivial changes (1)
- lib/commanded/event_store.ex
94de2c0 to
3703033
Compare
Signed-off-by: Yordis Prieto <yordis.prieto@gmail.com>
3703033 to
57ed763
Compare

Summary
[:commanded, :event_store, :stream_forward, :stop]telemetry until stream enumeration completes for lazy-stream adapters, sodurationreflects actual read-from-store latency:stopimmediately, preserving backwards compatibility[:commanded, :event_store, :stream_forward, :exception]withkind,reason, andstacktracewhen adapter resolution orstream_forwardraises, matching:telemetry.span/3behavior