Skip to content

Commit 7149e9c

Browse files
committed
fix: Isolate outbox telemetry measurements
Keep optional database measurements from failing durable delivery. Cover broadcast age and report message duration only on terminal outcomes so Ruby and JavaScript emit matching samples.
1 parent 0dcbcbd commit 7149e9c

8 files changed

Lines changed: 85 additions & 7 deletions

File tree

‎docs/observability.md‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -92,6 +92,9 @@ confers no permission.
9292

9393
Diagnostics read at most `limit + 1` rows per queue source, with a hard limit of
9494
100. Each category returns `sampled`, `truncated`, and `oldestAgeMilliseconds`.
95+
The limit applies to each combined category: one effect plus one broadcast with
96+
`limit: 1` returns `sampled: 1, truncated: true`, even when both source queries
97+
returned all their rows. The extra row proves that the category exceeds its cap.
9598
The last value measures nonnegative time since availability (or terminal failure
9699
for recovery callbacks); future reminders have zero age. No payloads or row
97100
identifiers are returned. Samples are observations across several queries,

‎lib/solid_objects/broadcast_executor.rb‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -41,6 +41,7 @@ def run_once
4141
broadcast = claim_next
4242
return false unless broadcast
4343

44+
Telemetry.outbox(broadcast)
4445
broadcast_adapter.call(broadcast)
4546
complete(broadcast)
4647
true

‎lib/solid_objects/effect_executor.rb‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -137,7 +137,7 @@ def claim_next
137137

138138
# @rbs (Effect) -> untyped
139139
def deliver(effect)
140-
SolidObjects.instrument(:"outbox.age", instance_id: effect.instance_id, actor_type: effect.instance.actor_type, actor_id: effect.instance.actor_id, message_id: effect.message_id, attempt: effect.attempt_count, age_milliseconds: [ ((database_adapter.database_now - effect.available_at) * 1000).round, 0 ].max)
140+
Telemetry.outbox(effect)
141141
return deliver_actor_message(effect) if effect.name == ACTOR_MESSAGE_EFFECT
142142

143143
handler = SolidObjects.effect_registry.fetch(effect.name)

‎lib/solid_objects/executor.rb‎

Lines changed: 10 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -161,7 +161,7 @@ def complete(result, observable_changes, state_after:, state_changed:)
161161
end
162162
SolidObjects.instrument_after_commit(:"recovery.completed", **instrumentation_payload) if recovery_message?
163163
report_large_state(state_after.byte_size)
164-
SolidObjects.instrument_after_commit(:"message.completed", **instrumentation_payload, revision: message.sequence)
164+
SolidObjects.instrument_after_commit(:"message.completed", **instrumentation_payload, revision: message.sequence, duration_milliseconds: elapsed_milliseconds)
165165
SolidObjects.wake_up.signal
166166
rescue CommittedTransactionError
167167
raise
@@ -424,6 +424,7 @@ def fail_message(error)
424424
:"message.failed",
425425
**instrumentation_payload,
426426
error_class: error.class.name,
427+
duration_milliseconds: elapsed_milliseconds,
427428
dead:
428429
)
429430
SolidObjects.instrument_after_commit(:"recovery.failed", **instrumentation_payload) if dead && recovery_message?
@@ -463,7 +464,8 @@ def reject_message(rejection)
463464
SolidObjects.instrument_after_commit(
464465
:"message.rejected",
465466
**instrumentation_payload,
466-
code: rejection.code
467+
code: rejection.code,
468+
duration_milliseconds: elapsed_milliseconds
467469
)
468470
SolidObjects.wake_up.signal
469471
end
@@ -530,6 +532,11 @@ def recovery_message?
530532
message.delivery_mode == "internal" && message.idempotency_key.to_s.start_with?("effect:") && message.idempotency_key.to_s.end_with?(":recovery")
531533
end
532534

535+
# @rbs () -> (Integer | Float)
536+
def elapsed_milliseconds
537+
((::Process.clock_gettime(::Process::CLOCK_MONOTONIC) - @started_at) * 1000).round(3)
538+
end
539+
533540
# @rbs () -> Hash[Symbol, untyped]
534541
def instrumentation_payload
535542
{
@@ -539,8 +546,7 @@ def instrumentation_payload
539546
actor_id: message.actor_id,
540547
sequence: message.sequence,
541548
attempt: message.attempt_count,
542-
request_id: message.request_id,
543-
duration_milliseconds: ((::Process.clock_gettime(::Process::CLOCK_MONOTONIC) - @started_at) * 1000).round(3)
549+
request_id: message.request_id
544550
}
545551
end
546552
end

‎lib/solid_objects/telemetry.rb‎

Lines changed: 21 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -6,12 +6,31 @@ module Telemetry
66
actorType actorId instanceId incarnation revision messageId requestId attempt sequence
77
operation deliveryMode generation durationMilliseconds latenessMilliseconds ageMilliseconds
88
depth count errorName outcome retryable status effectId effectName reminderId occurrence
9-
truncated broadcastId code role reason processId processKind ownerId componentCount byteCount
9+
outboxKind truncated broadcastId code role reason processId processKind ownerId componentCount byteCount
1010
thresholdBytes previousRunAt nextRunAt name commitAction payload
1111
].freeze
1212
ALIASES = { "errorClass" => "errorName", "stateRevision" => "revision", "durationMs" => "durationMilliseconds" }.freeze
1313

1414
class << self
15+
# @rbs (Effect | Broadcast) -> void
16+
def outbox(record)
17+
return unless SolidObjects.configuration.instrumentation || ActiveSupport::Notifications.notifier.listening?("solid_objects.outbox.age")
18+
19+
instance = record.instance
20+
SolidObjects.instrument(
21+
:"outbox.age",
22+
instance_id: record.instance_id,
23+
actor_type: instance.actor_type,
24+
actor_id: instance.actor_id,
25+
message_id: record.message_id,
26+
attempt: record.attempt_count,
27+
outbox_kind: record.is_a?(Effect) ? "effect" : "broadcast",
28+
age_milliseconds: [ ((SolidObjects.database_adapter.database_now - record.available_at) * 1000).round, 0 ].max
29+
)
30+
rescue
31+
nil
32+
end
33+
1534
# @rbs (Symbol, Hash[Symbol, untyped]) -> void
1635
def emit(name, payload)
1736
observer = SolidObjects.configuration.instrumentation
@@ -27,7 +46,7 @@ def event(name, payload)
2746
attributes = safe_attributes(payload)
2847
adapter = DatabaseAdapter.family(Record.connection).to_s
2948
event_name = "solid_objects.#{name}"
30-
labels = { "event" => event_name, "adapter" => adapter, "actorType" => attributes.fetch("actorType", "") }
49+
labels = { "event" => event_name, "adapter" => adapter, "actorType" => attributes["actorType"].to_s }
3150
metrics = [ { "name" => "solid_objects.events", "kind" => "counter", "unit" => "1", "value" => 1, "labels" => labels } ]
3251
[
3352
[ "durationMilliseconds", "solid_objects.duration", "histogram", "ms" ],

‎sig/generated/lib/solid_objects/executor.rbs‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -108,6 +108,9 @@ module SolidObjects
108108
# @rbs () -> bool
109109
def recovery_message?: () -> bool
110110

111+
# @rbs () -> (Integer | Float)
112+
def elapsed_milliseconds: () -> (Integer | Float)
113+
111114
# @rbs () -> Hash[Symbol, untyped]
112115
def instrumentation_payload: () -> Hash[Symbol, untyped]
113116
end

‎sig/generated/lib/solid_objects/telemetry.rbs‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,9 @@ module SolidObjects
66

77
ALIASES: untyped
88

9+
# @rbs (Effect | Broadcast) -> void
10+
def self.outbox: (Effect | Broadcast) -> void
11+
912
# @rbs (Symbol, Hash[Symbol, untyped]) -> void
1013
def self.emit: (Symbol, Hash[Symbol, untyped]) -> void
1114

‎test/integration/telemetry_test.rb‎

Lines changed: 43 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@ class TelemetryTest < ActiveSupport::TestCase
66
class Counter < SolidObjects::Actor
77
actor_type "telemetry-counter"
88
attribute :count, default: 0
9+
observable :count, broadcast: :value
910

1011
def arrange
1112
emit :telemetry_effect, secret: "private effect"
@@ -40,6 +41,8 @@ def increment
4041
assert_equal 1, event.fetch("attempt")
4142
assert event.fetch("incarnation")
4243
refute_includes event.fetch("metrics").to_json, "actorId"
44+
started = events.find { |entry| entry.fetch("name") == "solid_objects.message.started" }
45+
refute_includes started.fetch("metrics").map { |metric| metric.fetch("name") }, "solid_objects.duration"
4346
end
4447
test "diagnostics and observers require authorization and are bounded" do
4548
reference = Counter.ref("diagnostics")
@@ -118,4 +121,44 @@ def increment
118121
ensure
119122
worker&.stop
120123
end
124+
test "outbox measurement failure cannot fail delivery" do
125+
SolidObjects.configuration.instrumentation = ->(_) { true }
126+
Counter.ref("measurement").arrange
127+
calls = []
128+
SolidObjects.register_effect(:telemetry_effect) {
129+
calls << :delivered
130+
"result"
131+
}
132+
executor = SolidObjects::EffectExecutor.new
133+
executor.define_singleton_method(:claim_next) do
134+
super().tap do |effect|
135+
effect.define_singleton_method(:available_at) { raise "measurement failed" }
136+
end
137+
end
138+
assert executor.run_once
139+
assert_equal [ :delivered ], calls
140+
assert_equal "completed", SolidObjects::Effect.first.status
141+
ensure
142+
executor&.stop
143+
end
144+
145+
test "samples broadcast age and caps the combined outbox category" do
146+
events = []
147+
SolidObjects.configuration.instrumentation = ->(event) { events << event }
148+
SolidObjects.configuration.authorize_administration = ->(**) { true }
149+
SolidObjects.configuration.broadcast_adapter = ->(_) { true }
150+
reference = Counter.ref("broadcast")
151+
reference.increment
152+
reference.arrange
153+
summary = reference.diagnostics(limit: 1)
154+
assert_equal 1, summary.fetch("outbox").fetch("sampled")
155+
assert summary.fetch("outbox").fetch("truncated")
156+
executor = SolidObjects::BroadcastExecutor.new
157+
assert executor.run_once
158+
event = events.find { |entry| entry.fetch("name") == "solid_objects.outbox.age" && entry.fetch("attributes")["outboxKind"] == "broadcast" }
159+
assert event
160+
assert_operator event.fetch("metrics").last.fetch("value"), :>=, 0
161+
ensure
162+
executor&.stop
163+
end
121164
end

0 commit comments

Comments
 (0)