Skip to content

Commit 0dcbcbd

Browse files
committed
feat: Add portable actor observability
Match JavaScript telemetry and actor diagnostics using native Active Support notifications and Active Record scopes. Isolate subscriber failures while preserving application exceptions, and keep private error text out of default logs. Refs cardmagic/solid-objects-js#42
1 parent c18b82b commit 0dcbcbd

32 files changed

Lines changed: 582 additions & 24 deletions

‎CHANGELOG.md‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,10 @@
11
# Changelog
22

3+
## Unreleased
4+
5+
- Add portable telemetry, isolated observer hooks, metric definitions, and bounded authorized actor diagnostics matching JavaScript.
6+
7+
38
## 0.16.0 - 2026-09-23
49

510
- Find a message whose reference a caller lost.

‎README.md‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -162,6 +162,7 @@ Exactly once is not hiding in a more advanced configuration. Read the
162162
- [Five-minute Rails guide](https://solidobjects.dev/5min/rails)
163163
- [Choosing Solid Objects](docs/fit.md)
164164
- [Operations and recovery](docs/operations.md)
165+
- [Observability and diagnostics](docs/observability.md)
165166
- [Reminders](docs/reminders.md)
166167
- [Reactive ERB](docs/realtime.md)
167168
- [Detailed architecture](docs/architecture.md)

‎docs/observability.md‎

Lines changed: 123 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,123 @@
1+
# Portable observability
2+
3+
Configure `instrumentation(event)` to receive structured events. No exporter SDK
4+
is required. The same JSON envelope is emitted by Ruby, SQLite, PostgreSQL, MySQL,
5+
and the Durable Objects host. Existing JavaScript `name`, `occurredAt`, and
6+
`attributes` fields remain available. Ruby's Active Support notifications remain
7+
available with their existing snake_case payloads.
8+
9+
## Schema version 1
10+
11+
Every event has `schemaVersion`, `name` (prefixed with `solid_objects.`),
12+
`occurredAt` (UTC ISO 8601), `adapter`, `actorType`, `actorId`, `incarnation`,
13+
`revision`, `messageId`, `attempt`, `attributes`, and `metrics`.
14+
Unavailable identifiers are null; `attempt` is zero outside a message attempt.
15+
An incarnation identifies a persisted actor instance, independently of its lease
16+
generation. Revisions and IDs are strings. Process-wide events have null actor
17+
identity. A message ID or revision correlates actor work where applicable.
18+
19+
Only known scalar metadata fields enter the portable envelope. Arguments, state,
20+
results, credentials, backtraces, exception text, nested provider data, and unknown
21+
attributes are excluded. Actor IDs remain correlation data: applications should
22+
use opaque actor identifiers and apply their own retention policy to event logs.
23+
Events and metric samples are immutable. Throwing observers, rejected observer
24+
promises, and a failing instrumentation error logger cannot change a turn's
25+
result. Delivery is best effort and synchronous callbacks should be short;
26+
JavaScript does not await exporters. Telemetry is not a durable audit trail.
27+
28+
| Event | Meaning |
29+
| --- | --- |
30+
| `activation.started/completed/failed` | Local actor activation hook lifecycle |
31+
| `message.started/completed/failed/rejected` | One attempt's execution outcome |
32+
| `message.retry` | Failed attempt durably queued for another attempt |
33+
| `dead_letter.created` | Message exhausted retries or failed permanently |
34+
| `mailbox.depth` | On-demand diagnostic sample; exact depth only if not truncated |
35+
| `reminder.enqueued` | Due reminder dispatch; lateness is measured from due time |
36+
| `outbox.age` | Delivery observation; age is time since the item's current availability time |
37+
| `recovery.reclaimed` | A previously claimed, interrupted message begins another attempt |
38+
| `recovery.completed/failed` | Durable effect recovery callback commits or enters the dead-letter queue |
39+
| `snapshot.read` | Authorized snapshot constructed without exposing its contents |
40+
| `realtime.connected/disconnected` | Actor subscription added or removed |
41+
42+
Events describe local observations. Concurrent deletion, crashes, and failed
43+
exporters can omit events. Never infer exactly-once delivery from event counts.
44+
Additional existing runtime events retain their names.
45+
46+
## Metrics and tracing
47+
48+
Metrics are sample descriptions. Exporting them is opt-in: the runtime does not
49+
register meters, allocate per-actor metric series, or install a vendor SDK.
50+
51+
| Name | Kind | Unit | Aggregation |
52+
| --- | --- | --- | --- |
53+
| `solid_objects.events` | counter | `1` | Sum one per event |
54+
| `solid_objects.duration` | histogram | `ms` | Distribution of observed attempt duration |
55+
| `solid_objects.reminder.lateness` | histogram | `ms` | Distribution of reminder dispatch delay |
56+
| `solid_objects.outbox.age` | histogram | `ms` | Distribution of delivery delay since availability |
57+
| `solid_objects.mailbox.depth` | gauge | `1` | Last exact sampled actor depth; omit truncated samples |
58+
59+
Labels contain only event name, adapter family, and declared actor type. Keep the
60+
actor type registry finite. Never add actor ID, incarnation, message ID, operation
61+
arguments, request IDs, or error text to metric labels. A gauge without an actor
62+
label represents the most recently observed actor; it is not total fleet backlog.
63+
For tracing, correlate start/outcome events using `adapter`, `incarnation`,
64+
`messageId`, and `attempt`, and close or expire spans when no outcome arrives.
65+
66+
## Actor observers and diagnostics
67+
68+
```ruby
69+
cart = ShoppingCart.ref("demo-cart")
70+
stop = cart.observe(authorization_context: operator) do |event|
71+
logger.info(event.to_json)
72+
end
73+
summary = cart.diagnostics(authorization_context: operator, limit: 50)
74+
stop.call
75+
```
76+
77+
Ruby uses `reference.observe(authorization_context:) { |event| ... }` and
78+
`reference.diagnostics(authorization_context:, limit: 50)`. Stop observing by
79+
calling the returned proc. `on` filters one event name, such as `message.retry`.
80+
Observers receive only this actor's events in the current runtime/process; they
81+
are not subscriptions to workers on other hosts. Dispose them when the caller's
82+
session ends or authorization is revoked. JS limits local observers to 1,000.
83+
For remote Durable Objects, configure `instrumentation` on the actor host and
84+
filter by actor identity there; process-local reference observers raise
85+
`UnsupportedCapability`. Remote `reference.diagnostics` is supported.
86+
87+
Both APIs default to denied. Set `authorizeAdministration` / `authorize_administration`
88+
to allow action `observe` or `inspect`, resource `actor_diagnostics`, and resource ID
89+
`JSON.stringify([actorType, actorId])`. Ruby receives a symbol action. Authorization
90+
runs before reading summaries or registering observers; possessing an actor ID
91+
confers no permission.
92+
93+
Diagnostics read at most `limit + 1` rows per queue source, with a hard limit of
94+
100. Each category returns `sampled`, `truncated`, and `oldestAgeMilliseconds`.
95+
The last value measures nonnegative time since availability (or terminal failure
96+
for recovery callbacks); future reminders have zero age. No payloads or row
97+
identifiers are returned. Samples are observations across several queries,
98+
not an atomic fleet snapshot. Large queues can still require database scanning;
99+
the bound limits materialized rows and response size, not query execution time.
100+
101+
Categories are mailbox (ready and claimed), outbox (pending and processing effects
102+
and broadcasts), reminders (scheduled and paused), retries (failed messages still
103+
eligible to run), and recoveryFailures (dead internal effect recovery
104+
callback messages). Durable Objects does not implement process-heartbeat effect recovery;
105+
its recoveryFailures category is empty. Recovery failures are durable records,
106+
not a history of transient database or exporter exceptions.
107+
108+
## A query that works across adapters
109+
110+
Write one envelope per line to `events.jsonl`, using the same instrumentation hook
111+
with SQLite and PostgreSQL. This query reports completed attempts by adapter:
112+
113+
```sh
114+
jq -s 'map(select(.name == "solid_objects.message.completed"))
115+
| group_by(.adapter)
116+
| map({adapter: .[0].adapter, completed: length,
117+
mean_ms: (map(.attributes.durationMilliseconds) | add / length)})' events.jsonl
118+
```
119+
120+
A dashboard can chart the event counter by adapter and outcome and the duration,
121+
reminder lateness, and outbox age distributions using exactly the same fields.
122+
Compare workloads with the same actor types and sampling policy. Counts measure
123+
observations and cannot replace a database query for authoritative queue state.

‎docs/roadmap.md‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,9 @@
22

33
## Implemented and tested
44

5+
- Portable telemetry with a shared JSON schema, metric samples, isolated observer
6+
callbacks, and bounded authorized actor diagnostics. See [observability](observability.md).
7+
58
- Rails engine, install generator, migration, and CLI
69
- Explicit actor registry, references, JSON state, and state migrations
710
- Fluent direct synchronous RPC, configured `sync`, and durable `async`

‎lib/solid_objects.rb‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,8 @@
1616
require "solid_objects/callable_keywords"
1717
require "solid_objects/configuration"
1818
require "solid_objects/instrumentation"
19+
require "solid_objects/telemetry"
20+
require "solid_objects/diagnostics"
1921
require "solid_objects/log_subscriber"
2022
require "solid_objects/serialization"
2123
require "solid_objects/context"

‎lib/solid_objects/activation.rb‎

Lines changed: 6 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -24,15 +24,19 @@ def initialize(lease:)
2424
@actor = build_actor(instance)
2525
@last_used_at = monotonic_now
2626
@pass_exhausted = false
27+
SolidObjects.instrument(:"activation.started", instance_id: instance.id, actor_type: instance.actor_type, actor_id: instance.actor_id, generation: lease.generation)
2728
actor.activate
2829
SolidObjects.instrument(
29-
:"activation.started",
30+
:"activation.completed",
3031
instance_id: instance.id,
3132
actor_type: instance.actor_type,
3233
actor_id: instance.actor_id,
3334
owner_id: lease.owner_id,
3435
generation: lease.generation
3536
)
37+
rescue => error
38+
SolidObjects.instrument(:"activation.failed", instance_id: lease.instance_id, actor_type: instance&.actor_type, actor_id: instance&.actor_id, error_class: error.class.name)
39+
raise
3640
end
3741

3842
# @rbs () -> Integer
@@ -104,8 +108,7 @@ def deactivate
104108
actor_id: actor.actor_id,
105109
owner_id: lease.owner_id,
106110
generation: lease.generation,
107-
error_class: error.class.name,
108-
error_message: error.message
111+
error_class: error.class.name
109112
)
110113
SolidObjects.configuration.logger.error(
111114
"SolidObjects activation deactivation failed " \

‎lib/solid_objects/actor_channel.rb‎

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -49,6 +49,7 @@ def subscribed
4949
end
5050
refresh_outdated_components(snapshot)
5151
transmit_state_payloads(snapshot)
52+
SolidObjects.instrument(:"realtime.connected", actor_type: reference.actor_type, actor_id: reference.actor_id, instance_id: snapshot.instance_id, revision: snapshot.revision)
5253
rescue KeyError,
5354
JSON::ParserError,
5455
InvalidStreamToken,
@@ -57,6 +58,13 @@ def subscribed
5758
reject_and_report(reject_reason(error), actor_type:, actor_id:, error:)
5859
end
5960

61+
# @rbs () -> void
62+
def unsubscribed
63+
return unless reference
64+
65+
SolidObjects.instrument(:"realtime.disconnected", actor_type: reference.actor_type, actor_id: reference.actor_id)
66+
end
67+
6068
private
6169

6270
attr_reader :reference,

‎lib/solid_objects/configuration.rb‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -40,6 +40,7 @@ class Configuration
4040
# @rbs @broadcast_worker_count: Integer
4141
# @rbs @reminder_scheduler_count: Integer
4242
# @rbs @connects_to: Hash[Symbol, untyped]?
43+
# @rbs @instrumentation: Proc?
4344
# @rbs @logger: untyped
4445
# @rbs @stream_signing_secret: String?
4546
# @rbs @broadcast_adapter: Proc?
@@ -95,6 +96,7 @@ class Configuration
9596
:reminder_scheduler_count,
9697
:connects_to,
9798
:logger,
99+
:instrumentation,
98100
:stream_signing_secret,
99101
:broadcast_adapter,
100102
:wake_up_adapter,
@@ -159,6 +161,7 @@ def initialize
159161
@component_path_resolver = nil
160162
@component_authorization_context = ->(controller:) { controller }
161163
@payload_authorization_context = ->(connection:) { connection }
164+
@instrumentation = nil
162165
@logger = if defined?(Rails) && Rails.respond_to?(:logger) && Rails.logger
163166
Rails.logger
164167
else

‎lib/solid_objects/diagnostics.rb‎

Lines changed: 86 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,86 @@
1+
# rbs_inline: enabled
2+
3+
module SolidObjects
4+
class Diagnostics
5+
# @rbs @reference: Reference
6+
# @rbs (Reference) -> void
7+
def initialize(reference)
8+
@reference = reference
9+
end
10+
11+
# @rbs (?limit: Integer, ?authorization_context: untyped) -> Hash[String, untyped]
12+
def summary(limit: 100, authorization_context: nil)
13+
authorize!(:inspect, authorization_context:)
14+
raise ArgumentError, "diagnostic limit must be an integer between 1 and 100" unless limit.is_a?(Integer) && limit.between?(1, 100)
15+
16+
instance = Instance.find_by(actor_type: reference.actor_type, actor_id: reference.actor_id)
17+
now = SolidObjects.database_adapter.database_now
18+
instance_id = instance&.id
19+
mailbox = ReadyMessage.where(instance_id:).order(:available_at).limit(limit + 1).pluck(:available_at) +
20+
ClaimedMessage.where(instance_id:).order(:claimed_at).limit(limit + 1).pluck(:claimed_at)
21+
outbox = Effect.where(instance_id:, status: %w[pending processing]).order(:available_at).limit(limit + 1).pluck(:available_at) +
22+
Broadcast.where(instance_id:, status: %w[pending processing]).order(:available_at).limit(limit + 1).pluck(:available_at)
23+
reminders = Reminder.where(instance_id:, status: %w[scheduled paused]).order(:next_run_at).limit(limit + 1).pluck(:next_run_at)
24+
retries = ReadyMessage.joins(:message).where(instance_id:).where.not(Message.table_name => { error: nil }).order(:available_at).limit(limit + 1).pluck(:available_at)
25+
recovery_messages = Message.where(instance_id:, delivery_mode: "internal").where("idempotency_key LIKE ?", "effect:%:recovery").select(:id)
26+
recovery_failures = DeadLetter.where(instance_id:, message_id: recovery_messages).order(:last_failed_at).limit(limit + 1).pluck(:last_failed_at)
27+
result = {
28+
"actorType" => reference.actor_type,
29+
"actorId" => reference.actor_id,
30+
"incarnation" => instance_id&.to_s,
31+
"revision" => instance&.state_revision&.to_s,
32+
"adapter" => DatabaseAdapter.family(Record.connection).to_s,
33+
"occurredAt" => now.utc.iso8601(3),
34+
"limit" => limit,
35+
"mailbox" => summarize(mailbox, now:, limit:),
36+
"outbox" => summarize(outbox, now:, limit:),
37+
"reminders" => summarize(reminders, now:, limit:),
38+
"retries" => summarize(retries, now:, limit:),
39+
"recoveryFailures" => summarize(recovery_failures, now:, limit:)
40+
}
41+
SolidObjects.instrument(:"mailbox.depth", actor_type: reference.actor_type, actor_id: reference.actor_id, instance_id:, count: result.fetch("mailbox").fetch("sampled"), truncated: result.fetch("mailbox").fetch("truncated"), depth: result.fetch("mailbox").fetch("truncated") ? nil : result.fetch("mailbox").fetch("sampled"))
42+
Serialization.readonly_copy(result)
43+
end
44+
45+
# @rbs (?authorization_context: untyped) { (Hash[String, untyped]) -> untyped } -> Proc
46+
def observe(authorization_context: nil, &block)
47+
authorize!(:observe, authorization_context:)
48+
subscription = ActiveSupport::Notifications.subscribe(/\Asolid_objects\./) do |notification|
49+
payload = notification.payload
50+
next unless payload[:actor_type] == reference.actor_type && payload[:actor_id] == reference.actor_id
51+
52+
begin
53+
block.call(Telemetry.event(notification.name.delete_prefix("solid_objects.").to_sym, payload))
54+
rescue
55+
nil
56+
end
57+
end
58+
-> { ActiveSupport::Notifications.unsubscribe(subscription) }
59+
end
60+
61+
private
62+
63+
attr_reader :reference
64+
65+
# @rbs (Symbol, ?authorization_context: untyped) -> void
66+
def authorize!(action, authorization_context: nil)
67+
allowed = SolidObjects.configuration.authorize_administration.call(
68+
action:,
69+
resource: "actor_diagnostics",
70+
resource_id: [ reference.actor_type, reference.actor_id ].to_json,
71+
authorization_context:
72+
)
73+
raise Unauthorized, "actor diagnostics are not authorized" unless allowed
74+
end
75+
76+
# @rbs (Array[Time?], now: Time, limit: Integer) -> Hash[String, untyped]
77+
def summarize(timestamps, now:, limit:)
78+
oldest = timestamps.compact.min
79+
{
80+
"sampled" => [ timestamps.length, limit ].min,
81+
"truncated" => timestamps.length > limit,
82+
"oldestAgeMilliseconds" => oldest ? [ ((now - oldest) * 1000).round, 0 ].max : nil
83+
}
84+
end
85+
end
86+
end

‎lib/solid_objects/effect_executor.rb‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -137,6 +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)
140141
return deliver_actor_message(effect) if effect.name == ACTOR_MESSAGE_EFFECT
141142

142143
handler = SolidObjects.effect_registry.fetch(effect.name)
@@ -208,6 +209,9 @@ def complete(effect, result)
208209
Mailbox.new.announce(result_message) if result_message
209210
SolidObjects.instrument(
210211
:"effect.completed",
212+
instance_id: effect.instance_id,
213+
actor_type: effect.instance.actor_type,
214+
actor_id: effect.instance.actor_id,
211215
effect_id: effect.effect_id,
212216
effect_name: effect.name,
213217
message_id: effect.message_id,

0 commit comments

Comments
 (0)