diff --git a/CHANGELOG.md b/CHANGELOG.md index d79a993..97423ed 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,54 @@ ## Unreleased +- `solid_objects.activation.started` now fires before the actor's `activate()` + hook. Before, it fired after a successful hook. The new + `solid_objects.activation.completed` event takes that meaning, and + `solid_objects.activation.failed` reports a failed hook. Ruby changes the same + events. Move a subscriber that reads `activation.started` as a finished + activation to `activation.completed`. +- Reference `observe()` and `on()` reject a missing `onEvent` callback with + `TypeError` before authorization, as Ruby raises `ArgumentError` without a + block. Tests pin the 1,000-observer limit and observer removal on `close()`. + The roadmap no longer lists portable observability as future work. + +- Add the shared telemetry contract `compatibility/telemetry-events.json`. Tests + compare the attribute allowlist and the attribute keys of each core SQL event + with Ruby. Activation events carry `ownerId`, commit action events carry + `activationGeneration`, `reminder.enqueued` carries `messageId` and `attempt`, + and a truncated `mailbox.depth` sample reports `depth: null`. +- Port Ruby's query and observable purity tests for each kind of staged work. + Name `runtime.deadLetters.effects.retry(id)` for a dead transmit, document the + `actor_diagnostics` authorization, and remove the unused `compatibility/ruby.yml`. + +- Pin personalized payload isolation, staged-work rejection, and configured + UTF-8 byte limits with regressions matching Ruby. Share timeout telemetry + fixtures for wait reasons and activation ownership; document Ruby's fixes + and its yielding SQLite wait handler in the parity ledger. + +- Detect observable replacement of staged work even when the intent count is + unchanged. Compare complete intent snapshots around projection evaluation. + +- Preserve reserved JSON keys such as `__proto__` as own data properties without + changing object prototypes. Include them in size limits, state, arguments, and + retained results; share compatibility fixtures with Ruby. +- Record Ruby's matching query/projection purity guards, retained background + results, and the intentional 191-character Ruby / 255-character JS reminder + name limits in the parity ledger. + +- Preserve transmit staging order within a source message through a persisted + effect position, including retries. Schema migration 14 adds `effects.position`; + run `runtime.install()` before starting upgraded workers. Legacy rows retain + their existing ID tie-break because their original order cannot be recovered. +- Reject explicit null transmit arguments, matching Ruby; omission still defaults + to an empty object. Expand the shared wire fixtures and adapter ordering tests. +- Refresh the parity ledger for authorized message reads and retained background + results, bulk redrive/audit records, reminder cancellation, and automatic wake-up. + +- Patch the Cloudflare development tooling's Undici dependency to 7.29.1 to resolve the high-severity WebSocket and TLS advisories reported by CI. + +- Add portable observability envelopes, metric definitions, actor observers, and bounded authorization-aware diagnostics across SQL runtimes and the Durable Objects host. Isolate failing instrumentation and error loggers from actor work. + - Build before `npm publish` reads the manifest. npm validates `bin` against the working tree before `prepack` produces `dist`, so every release logged `No bin file found at dist/executable.js` twice. The published package was diff --git a/README.md b/README.md index e58af1b..51f1724 100644 --- a/README.md +++ b/README.md @@ -257,6 +257,7 @@ There is no exactly-once delivery. Read the - [Choosing Solid Objects](docs/fit.md) - [Public API](docs/api.md) - [Operations and recovery](docs/operations.md) +- [Observability and diagnostics](docs/observability.md) - [Detailed architecture](docs/architecture.md) - [Detailed documentation](docs/) diff --git a/compatibility/json-values.json b/compatibility/json-values.json new file mode 100644 index 0000000..9e2b191 --- /dev/null +++ b/compatibility/json-values.json @@ -0,0 +1,14 @@ +[ + { + "name": "object prototype key", + "value": { "__proto__": { "role": "ordinary data" }, "value": 1 } + }, + { "name": "scalar prototype key", "value": { "__proto__": "ordinary data" } }, + { "name": "null prototype key", "value": { "__proto__": null } }, + { + "name": "nested reserved names", + "value": { + "items": [{ "__proto__": { "constructor": "data" }, "prototype": true, "hasOwnProperty": 1 }] + } + } +] diff --git a/compatibility/ruby.yml b/compatibility/ruby.yml deleted file mode 100644 index 62cad32..0000000 --- a/compatibility/ruby.yml +++ /dev/null @@ -1,8 +0,0 @@ -ruby_repository: https://github.com/cardmagic/solid-objects-ruby -ruby_version: 0.14.0 -ruby_commit: e49c694 -npm_version: 0.14.0 -status: spiritual_parity -releases: - "0.12.0": spiritual_parity - "0.14.0": spiritual_parity diff --git a/compatibility/sync-timeout.json b/compatibility/sync-timeout.json new file mode 100644 index 0000000..a851215 --- /dev/null +++ b/compatibility/sync-timeout.json @@ -0,0 +1,10 @@ +[ + { "rubyReason": "actor_paused", "waitingOn": "actorPaused" }, + { "rubyReason": "activation_held", "waitingOn": "activationHeld" }, + { "rubyReason": "earlier_message", "waitingOn": "earlierMessage" }, + { "rubyReason": "message_claimed", "waitingOn": "messageClaimed" }, + { "rubyReason": "not_yet_available", "waitingOn": "notYetAvailable" }, + { "rubyReason": "ready_unclaimed", "waitingOn": "readyUnclaimed" }, + { "rubyReason": "database_contention", "waitingOn": "databaseContention" }, + { "rubyReason": "unknown", "waitingOn": "unknown" } +] diff --git a/compatibility/telemetry-events.json b/compatibility/telemetry-events.json new file mode 100644 index 0000000..19669db --- /dev/null +++ b/compatibility/telemetry-events.json @@ -0,0 +1,378 @@ +{ + "attributes": [ + "activationGeneration", + "activationOwnerId", + "actorId", + "actorType", + "ageMilliseconds", + "attempt", + "broadcastId", + "broadcastWorkers", + "byteCount", + "code", + "commitAction", + "component", + "componentCount", + "count", + "currentIntervalMilliseconds", + "deliveryMode", + "depth", + "durationMilliseconds", + "effectId", + "effectName", + "effectWorkers", + "errorName", + "failureCount", + "generation", + "idlePollingIntervalMilliseconds", + "incarnation", + "instanceId", + "intervalMilliseconds", + "latenessMilliseconds", + "messageId", + "name", + "nextRunAt", + "occurrence", + "operation", + "outboxKind", + "outcome", + "ownerId", + "payload", + "phase", + "pollingIntervalMilliseconds", + "previousIntervalMilliseconds", + "previousRunAt", + "processId", + "processKind", + "reason", + "reminderId", + "reminderSchedulers", + "requestId", + "retryable", + "revision", + "role", + "sequence", + "status", + "thresholdBytes", + "timeoutMilliseconds", + "truncated", + "waitingOn", + "workers" + ], + "events": [ + { + "name": "activation.completed", + "attributes": ["actorId", "actorType", "generation", "instanceId", "ownerId"], + "javascriptAttributes": [ + "attempt", + "deliveryMode", + "messageId", + "operation", + "requestId", + "sequence" + ] + }, + { + "name": "activation.failed", + "attributes": ["actorId", "actorType", "errorName", "generation", "instanceId", "ownerId"], + "javascriptAttributes": [ + "attempt", + "deliveryMode", + "messageId", + "operation", + "requestId", + "sequence" + ] + }, + { + "name": "activation.started", + "attributes": ["actorId", "actorType", "generation", "instanceId", "ownerId"], + "javascriptAttributes": [ + "attempt", + "deliveryMode", + "messageId", + "operation", + "requestId", + "sequence" + ] + }, + { + "name": "commit_action.completed", + "attributes": [ + "activationGeneration", + "actorId", + "actorType", + "attempt", + "commitAction", + "deliveryMode", + "instanceId", + "messageId", + "operation", + "requestId", + "sequence" + ] + }, + { + "name": "commit_action.failed", + "attributes": [ + "activationGeneration", + "actorId", + "actorType", + "attempt", + "commitAction", + "deliveryMode", + "errorName", + "instanceId", + "messageId", + "operation", + "requestId", + "sequence" + ] + }, + { + "name": "commit_action.started", + "attributes": [ + "activationGeneration", + "actorId", + "actorType", + "attempt", + "commitAction", + "deliveryMode", + "instanceId", + "messageId", + "operation", + "requestId", + "sequence" + ] + }, + { + "name": "dead_letter.created", + "attributes": [ + "actorId", + "actorType", + "attempt", + "deliveryMode", + "instanceId", + "messageId", + "operation", + "requestId", + "sequence" + ] + }, + { + "name": "mailbox.depth", + "attributes": ["actorId", "actorType", "count", "depth", "instanceId", "truncated"] + }, + { + "name": "message.completed", + "attributes": [ + "actorId", + "actorType", + "attempt", + "deliveryMode", + "durationMilliseconds", + "instanceId", + "messageId", + "operation", + "requestId", + "revision", + "sequence" + ] + }, + { + "name": "message.enqueued", + "attributes": [ + "actorId", + "actorType", + "attempt", + "deliveryMode", + "instanceId", + "messageId", + "operation", + "requestId", + "sequence" + ] + }, + { + "name": "message.failed", + "attributes": [ + "actorId", + "actorType", + "attempt", + "deliveryMode", + "durationMilliseconds", + "errorName", + "instanceId", + "messageId", + "operation", + "outcome", + "requestId", + "retryable", + "sequence" + ] + }, + { + "name": "message.rejected", + "attributes": [ + "actorId", + "actorType", + "attempt", + "code", + "deliveryMode", + "durationMilliseconds", + "instanceId", + "messageId", + "operation", + "requestId", + "sequence" + ] + }, + { + "name": "message.retry", + "attributes": [ + "actorId", + "actorType", + "attempt", + "deliveryMode", + "instanceId", + "messageId", + "operation", + "requestId", + "sequence" + ] + }, + { + "name": "message.started", + "attributes": [ + "actorId", + "actorType", + "attempt", + "deliveryMode", + "instanceId", + "messageId", + "operation", + "requestId", + "sequence" + ] + }, + { + "name": "outbox.age", + "match": { + "outboxKind": "broadcast" + }, + "attributes": [ + "actorId", + "actorType", + "ageMilliseconds", + "attempt", + "broadcastId", + "instanceId", + "messageId", + "outboxKind", + "revision" + ] + }, + { + "name": "outbox.age", + "match": { + "outboxKind": "effect" + }, + "attributes": [ + "actorId", + "actorType", + "ageMilliseconds", + "attempt", + "effectId", + "effectName", + "instanceId", + "messageId", + "outboxKind" + ] + }, + { + "name": "payload_broadcast.failed", + "attributes": ["actorId", "actorType", "errorName", "payload"] + }, + { + "name": "polling.interval_changed", + "attributes": [ + "currentIntervalMilliseconds", + "previousIntervalMilliseconds", + "reason", + "role" + ] + }, + { + "name": "realtime.connected", + "attributes": ["actorId", "actorType"] + }, + { + "name": "realtime.disconnected", + "attributes": ["actorId", "actorType"] + }, + { + "name": "recovery.completed", + "attributes": [ + "actorId", + "actorType", + "attempt", + "deliveryMode", + "instanceId", + "messageId", + "operation", + "requestId", + "sequence" + ] + }, + { + "name": "recovery.failed", + "attributes": [ + "actorId", + "actorType", + "attempt", + "deliveryMode", + "instanceId", + "messageId", + "operation", + "requestId", + "sequence" + ] + }, + { + "name": "recovery.reclaimed", + "attributes": [ + "actorId", + "actorType", + "attempt", + "deliveryMode", + "instanceId", + "messageId", + "operation", + "requestId", + "sequence" + ] + }, + { + "name": "reminder.enqueued", + "attributes": [ + "actorId", + "actorType", + "attempt", + "instanceId", + "latenessMilliseconds", + "messageId", + "occurrence", + "operation", + "reminderId" + ] + }, + { + "name": "snapshot.read", + "attributes": ["actorId", "actorType", "instanceId", "revision"] + }, + { + "name": "sync.enqueue_timeout", + "attributes": ["actorId", "actorType", "operation", "timeoutMilliseconds"] + } + ] +} diff --git a/compatibility/transmit-envelopes.json b/compatibility/transmit-envelopes.json index 74aa5d1..c47137a 100644 --- a/compatibility/transmit-envelopes.json +++ b/compatibility/transmit-envelopes.json @@ -84,6 +84,16 @@ "operation": "increment", "arguments": [1] } + }, + { + "name": "null arguments", + "envelope": { + "effectId": "fixture-effect-0008", + "actorType": "transmit-counters", + "actorId": "fixture-counter", + "operation": "increment", + "arguments": null + } } ] } diff --git a/docs/api.md b/docs/api.md index 3b64938..95f6245 100644 --- a/docs/api.md +++ b/docs/api.md @@ -105,7 +105,7 @@ envelope. The envelope carries its name in `invalidations`. A component registry can then refresh a reauthorized endpoint, and the value stays private. `MessageReference` does not retain an invocation's authorization context. -Supply `authorizationContext` to each `status()`, `result()`, and `wait()` call; +Supply `authorizationContext` to each `status()`, `result()`, `outcome()`, and `wait()` call; the runtime reauthorizes the persisted operation every time. Durable results are JSON, so an operation that returns `undefined` or is declared `void` resolves as `null`. @@ -130,6 +130,10 @@ function playerForSession(options: { ### Reminders +Keyed reminder names combine the operation, a colon, and the key. JavaScript +limits the combined name to 255 UTF-16 code units; Ruby limits it to 191 +characters. Use at most 191 ASCII characters for names shared across runtimes. + A reminder is one alarm per actor and name. If you schedule a name that is already armed, the runtime **moves the existing alarm**. It does not add a second one. A reminder is therefore safe to re-arm from a handler that can run @@ -904,6 +908,11 @@ open SAH pool otherwise blocks the next candidate until the worker dies. ## `solid-objects/transmit` +Schema migration 14 persists an effect's staging position within its source message. +Run `runtime.install()` before upgraded workers start. New effects preserve staging +order even when two transmits originate in one turn. Legacy rows use position zero +and their existing ID tie-break; their original order cannot be recovered. + The transactional outbox bridge between a local runtime and a server runtime. An actor stages a transmit intent with `this.transmit()` in the same transaction as its state change. The effect worker drains the @@ -923,10 +932,10 @@ outbox with at-least-once delivery, per-actor order, and retry backoff. `deliver` while offline and the effect retries with backoff. Give a browser runtime a generous `maxAttempts`; an effect that exhausts its attempts during a long offline period lands in dead letters, and - `runtime.deadLetters.retry` re-queues it. + `runtime.deadLetters.effects.retry(id)` re-queues it. - `receiveTransmitEnvelope(options)`: idempotent server ingest. `arguments` is optional in the envelope and defaults to an empty object, matching the - staging side and the Ruby ingest. It enqueues an + staging side and the Ruby ingest. Explicit `null` and arrays are invalid. It enqueues an internal message with `transmit:` as the idempotency key, so a replayed envelope applies once. The host must authenticate the sender before this call; internal delivery skips `authorizeMessage`. The call @@ -1057,3 +1066,16 @@ Behavior: Mounting, authorization actions, CSRF behavior, pages, and extensions are in [Operator dashboard](dashboard.md). + +### Portable telemetry and diagnostics + +`InstrumentationEvent` includes the versioned envelope and immutable `MetricSample` +values. `EventObserver` is a provider-free event callback. An actor reference offers +`observe({ onEvent, authorizationContext })`, `on(name, { onEvent, authorizationContext })`, +and `diagnostics(options)`. Both observer methods resolve to an unsubscribe function. +They reject a missing `onEvent` with `TypeError` before authorization. A runtime +accepts at most 1,000 local observers; the next one rejects with `RangeError`. +`close()` removes every local observer. +`DiagnosticOptions` accepts `authorizationContext` and a `limit` from 1 to 100. +`ActorDiagnostics` holds five bounded `DiagnosticSummary` values, with `sampled`, +`truncated`, and `oldestAgeMilliseconds`. See [observability](observability.md). diff --git a/docs/authorization.md b/docs/authorization.md index c76ef68..1f376b4 100644 --- a/docs/authorization.md +++ b/docs/authorization.md @@ -10,6 +10,12 @@ resource, optional resource ID, and the caller's authorization context. Retry authorization happens before lookup so a denied caller cannot use record IDs as an existence oracle. +Actor diagnostics and actor observers call `authorizeAdministration` with the +resource `actor_diagnostics` and the resource ID +`JSON.stringify([actorType, actorId])`. `diagnostics()` uses the action +`inspect`. `observe()` and `on()` use the action `observe`. The check runs before +any queue read or observer registration. + `runtime.realtime` is transport-neutral. The host application authenticates its WebSocket or stream connection, passes that fresh server-side subject as the session's `authorizationContext`, and forwards incoming protocol messages to diff --git a/docs/observability.md b/docs/observability.md new file mode 100644 index 0000000..a8c93da --- /dev/null +++ b/docs/observability.md @@ -0,0 +1,192 @@ +# Portable observability + +Configure `instrumentation(event)` to receive structured events. No exporter SDK +is required. The same JSON envelope is emitted by Ruby, SQLite, PostgreSQL, MySQL, +and the Durable Objects host. Existing JavaScript `name`, `occurredAt`, and +`attributes` fields remain available. Ruby's Active Support notifications remain +available with their existing snake_case payloads. + +```ts +import { configure } from "solid-objects" + +const runtime = configure({ + ...applicationConfiguration, + instrumentation: (event) => { + console.info(JSON.stringify(event)) + }, +}) +``` + +JavaScript passes `instrumentation` to `configure()`. Ruby sets +`configuration.instrumentation` inside `SolidObjects.configure`. + +## Schema version 1 + +Every event has `schemaVersion`, `name` (prefixed with `solid_objects.`), +`occurredAt` (UTC ISO 8601), `adapter`, `actorType`, `actorId`, `incarnation`, +`revision`, `messageId`, `attempt`, `attributes`, and `metrics`. +Unavailable identifiers are null; `attempt` is zero outside a message attempt. +Ruby publishes the RBS types `SolidObjects::portable_event`, +`SolidObjects::actor_diagnostics`, and `SolidObjects::event_observer`. JavaScript +exports `InstrumentationEvent`, `ActorDiagnostics`, and `EventObserver`. +An incarnation identifies a persisted actor instance, independently of its lease +generation. Revisions and IDs are strings. Process-wide events have null actor +identity. A message ID or revision correlates actor work where applicable. + +Only known scalar metadata fields enter the portable envelope. Arguments, state, +results, credentials, backtraces, exception text, nested provider data, and unknown +attributes are excluded. Actor IDs remain correlation data: applications should +use opaque actor identifiers and apply their own retention policy to event logs. +Events and metric samples are immutable. Throwing observers, rejected observer +promises, and a failing instrumentation error logger cannot change a turn's +result. When an exporter or observer raises, both runtimes log +`solid_objects.instrumentation.failed` with the event name and the error class. +Delivery is best effort and synchronous callbacks should be short; +JavaScript does not await exporters. Telemetry is not a durable audit trail. + +| Event | Meaning | +| ------------------------------------------- | ---------------------------------------------------------------------------- | +| `activation.started/completed/failed` | Local actor activation hook lifecycle | +| `message.started/completed/failed/rejected` | One attempt's execution outcome | +| `message.retry` | Failed attempt durably queued for another attempt | +| `dead_letter.created` | Message exhausted retries or failed permanently | +| `commit_action.started/completed/failed` | One registered commit action inside the commit transaction | +| `mailbox.depth` | On-demand diagnostic sample; `depth` is null when the sample is truncated | +| `reminder.enqueued` | Due reminder dispatch; lateness is measured from due time | +| `outbox.age` | Delivery observation; age is time since the item's current availability time | +| `recovery.reclaimed` | A previously claimed, interrupted message begins another attempt | +| `recovery.completed/failed` | Durable effect recovery callback commits or enters the dead-letter queue | +| `snapshot.read` | Authorized snapshot constructed without exposing its contents | +| `realtime.connected/disconnected` | Actor subscription added or removed | +| `payload_broadcast.failed` | One personalized payload failed; `payload` names it | + +Events describe local observations. Concurrent deletion, crashes, and failed +exporters can omit events. Never infer exactly-once delivery from event counts. +Additional existing runtime events retain their names. + +## Event attributes + +`compatibility/telemetry-events.json` holds the attribute allowlist and the exact +attribute keys of each core event. Both test suites compare the events of the SQL +runtimes with this file. Message events carry `operation` and `deliveryMode`. +`message.failed` also carries a boolean `retryable` and an `outcome` of +`retrying` or `dead`. Commit action events carry `commitAction` and +`activationGeneration`. Activation events carry `generation` and `ownerId`. + +Ruby activates an actor instance before it claims a message, so its activation +events have no message fields. JavaScript activates an actor in the turn that +claims a message, so its activation events also carry the `messageId`, +`requestId`, `sequence`, `attempt`, `operation`, and `deliveryMode` of that +message. The contract file records these keys as JavaScript-only. + +The Durable Objects host sends the same envelope, but some of its events carry +fewer attributes. The contract file does not apply to that host. + +Portable `polling.interval_changed` events carry `previousIntervalMilliseconds`, +`currentIntervalMilliseconds`, and a string `reason`. Ruby converts its native +second-based notification values while preserving the original notification. + +## Synchronous timeout diagnostics + +`solid_objects.sync.timeout` includes `waitingOn`, `activationOwnerId`, and +`activationGeneration` in `attributes`. Generations are decimal strings; +unavailable activation fields are null. Both runtimes use the same wait reasons: +`actorPaused`, `activationHeld`, `earlierMessage`, `messageClaimed`, +`notYetAvailable`, `readyUnclaimed`, `databaseContention`, and `unknown`. +Ruby's exception attributes and Active Support notifications retain their native +snake_case names and reason values. The portable instrumentation envelope uses +the shared camelCase contract. The Durable Objects host does not send +`sync.timeout`. A `call` timeout there reports `waitingOn: "unknown"`. + +## Metrics and tracing + +Metrics are sample descriptions. Exporting them is opt-in: the runtime does not +register meters, allocate per-actor metric series, or install a vendor SDK. + +| Name | Kind | Unit | Aggregation | +| --------------------------------- | --------- | ---- | ------------------------------------------------------ | +| `solid_objects.events` | counter | `1` | Sum one per event | +| `solid_objects.duration` | histogram | `ms` | Distribution of observed attempt duration | +| `solid_objects.reminder.lateness` | histogram | `ms` | Distribution of reminder dispatch delay | +| `solid_objects.outbox.age` | histogram | `ms` | Distribution of delivery delay since availability | +| `solid_objects.mailbox.depth` | gauge | `1` | Last exact sampled actor depth; omit truncated samples | + +Labels contain only event name, adapter family, and declared actor type. Keep the +actor type registry finite. Never add actor ID, incarnation, message ID, operation +arguments, request IDs, or error text to metric labels. A gauge without an actor +label represents the most recently observed actor; it is not total fleet backlog. +For tracing, correlate start/outcome events using `adapter`, `incarnation`, +`messageId`, and `attempt`, and close or expire spans when no outcome arrives. + +## Actor observers and diagnostics + +```ts +const cart = ShoppingCart.ref("demo-cart") +const stop = await cart.observe({ + authorizationContext: operator, + onEvent: (event) => { + console.info(JSON.stringify(event)) + }, +}) +const summary = await cart.diagnostics({ authorizationContext: operator, limit: 50 }) +stop() +``` + +Ruby uses `reference.observe(authorization_context:) { |event| ... }` and +`reference.diagnostics(authorization_context:, limit: 50)`. Stop observing by +calling the returned proc. `on` filters one event name, such as `message.retry`. +Observers receive only this actor's events in the current runtime/process; they +are not subscriptions to workers on other hosts. Dispose them when the caller's +session ends or authorization is revoked. Each runtime or process accepts at most +1,000 local observers. More observers raise `RangeError` in JavaScript and +`ArgumentError` in Ruby. An observer needs a callback or a block. JavaScript +rejects a missing `onEvent` with `TypeError`, and Ruby raises `ArgumentError` +without a block. Both checks run before authorization. + +The JavaScript Durable Objects host does not support process-local reference +observers. On that host, `observe` and `on` raise `UnsupportedCapability`, and +remote `reference.diagnostics` works. To observe a remote actor, configure +`instrumentation` on the actor host and filter the events by actor identity. +Ruby has no Durable Objects host. + +Both APIs default to denied. Set `authorizeAdministration` / `authorize_administration` +to allow action `observe` or `inspect`, resource `actor_diagnostics`, and resource ID +`JSON.stringify([actorType, actorId])`. Ruby receives a symbol action. Authorization +runs before reading summaries or registering observers; possessing an actor ID +confers no permission. + +Diagnostics read at most `limit + 1` rows per queue source, with a hard limit +of 100 rows. Each category returns `sampled`, `truncated`, and +`oldestAgeMilliseconds`. +The limit applies to each combined category: one effect plus one broadcast with +`limit: 1` returns `sampled: 1, truncated: true`, even when both source queries +returned all their rows. The extra row proves that the category exceeds its cap. +The last value measures nonnegative time since availability (or terminal failure +for recovery callbacks); future reminders have zero age. No payloads or row +identifiers are returned. Samples are observations across several queries, +not an atomic fleet snapshot. Large queues can still require database scanning; +the bound limits materialized rows and response size, not query execution time. + +Categories are mailbox (ready and claimed), outbox (pending and processing effects +and broadcasts), reminders (scheduled and paused), retries (failed messages still +eligible to run), and recoveryFailures (dead internal effect recovery +callback messages). Durable Objects does not implement process-heartbeat effect recovery; +its recoveryFailures category is empty. Recovery failures are durable records, +not a history of transient database or exporter exceptions. + +## A query that works across adapters + +Write one envelope per line to `events.jsonl`, using the same instrumentation hook +with SQLite and PostgreSQL. This query reports completed attempts by adapter: + +```sh +jq -s 'map(select(.name == "solid_objects.message.completed")) + | group_by(.adapter) + | map({adapter: .[0].adapter, completed: length, + mean_ms: (map(.attributes.durationMilliseconds) | add / length)})' events.jsonl +``` + +A dashboard can chart the event counter by adapter and outcome and the duration, +reminder lateness, and outbox age distributions using exactly the same fields. +Compare workloads with the same actor types and sampling policy. Counts measure +observations and cannot replace a database query for authoritative queue state. diff --git a/docs/parity.md b/docs/parity.md index 176c1fa..6f33a29 100644 --- a/docs/parity.md +++ b/docs/parity.md @@ -4,18 +4,15 @@ This ledger tracks capability parity with the Ruby `solid_objects` gem. Parity preserves a capability and its correctness or security boundary. It does not copy a Rails API into Node. -Reference: Ruby `solid_objects` 0.14.0. The JavaScript package began at the -Ruby design's `0.12` capability generation; that version number did not imply -earlier JavaScript releases. - -The Node `0.14.0` implementation has capability parity with that reference. Its -relational runtime, correctness boundaries, administration, diagnostics, -operator dashboard, realtime projections, browser behavior, and supported -adapters have native equivalents. Transport- and framework-neutral JavaScript -APIs replace the Rails-specific render surfaces. Three rows below are explicit -scope boundaries that the Ruby reference shares: the partial guard row, the -backpressure row, and the shared planned result-lookup row. They are not missing -Ruby capabilities. +Reference: Ruby `solid_objects` 0.16.0 and the unreleased changes recorded in +both changelogs. The JavaScript package began at the Ruby design's `0.12` +capability generation; that version number did not imply earlier JavaScript releases. + +The SQL runtimes share durable identity, ordered mailboxes, fenced commits, +effects, reminders, authorization, administration, and realtime guarantees. +Transport- and framework-neutral JavaScript APIs replace Rails-specific surfaces. +The partial write-guard and backpressure rows below record their actual limits. +Browser and Cloudflare hosting have separate scope boundaries later in this file. `0.14.0` also adds `runtime.enqueueInternalMessage()`, `runtime.enqueueInternalMessageInTransaction()`, and @@ -32,7 +29,7 @@ such boundary between a gem and its dependents. Operation-reference typing is runtime-specific: TypeScript infers scheduled and transmitted operations from the concrete receiver, and checks literal effect callback names. Ruby offers opt-in RBS generation from declared application types -in [solid-objects-ruby#66](https://github.com/cardmagic/solid-objects-ruby/pull/66). +since Ruby 0.14.7. Both preserve runtime operation validation and global effect/commit-action names; this does not imply automatic TypeScript-style inference in Ruby. @@ -51,60 +48,63 @@ SQL waiters read results and status from one statement. Ruby already checks completion and returns the result from the same loaded message; no Ruby change is needed for the JavaScript stale-result race fix. -| Capability | Status | TypeScript shape or remaining work | -| ------------------------------------------------------------------------------------------------------------- | ------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | -| Actor registry, durable identity, JSON state, and adjacent state migrations | Native | Ordinary classes, static actor types, inferred state, explicit migrations, and isolated runtime context across every actor-instance callback. | -| Fluent committed calls and background delivery | Native | `await reference.operation()` and `reference.send.operation()`. | -| Ordered mailbox, sequence allocation, idempotency, retries, dead letters, leases, renewal, and fenced commits | Native | Relational ready/claimed membership tables, distinct generated request IDs and caller idempotency keys, durable history, adapter-appropriate sequence locking, and PostgreSQL/MySQL row locks held from fence validation through commit. | -| Domain rejection and strict poison ordering | Native | Rejections accept JavaScript identifier-style codes and roll back without retry; invalid codes fail terminally, while retryable failures block later operations until completion or dead-lettering. | -| Bounded activation passes and hot-actor fairness | Native | Configurable turn-count and elapsed-time budgets bound each pass, then move only that actor's already-due memberships behind actors already waiting. | -| Bounded claim candidate scan | Native | A configurable ordered scan continues to another ready actor when a worker loses the first candidate's lease race. | -| Backpressure and payload caps | Partial | Serialization enforces a shared maximum JSON nesting depth, raising `InvalidPayload`, and an optional caller-supplied `maxBytes` limit, raising `PayloadTooLarge`; reminder names are bounded to 255 characters. Distributed per-actor rate limits, global admission control, and cache-capacity eviction are not planned in either runtime. They are hot, request-path, and loss-tolerant, so one durable ordered message per check is the wrong shape. Solid Objects Pro answers them with grouped and ephemeral operations, which [fit](fit.md) describes. | -| Idle activation cache | Native | Long-running workers retain hydrated actors under renewable fenced leases, restore public state after failed turns, and release on timeout, fairness yield, lease loss, or shutdown. | -| Transactional effects and outcome operations | Native | At-least-once handlers receive immutable stable effect, attempt, source-message, and actor identity; success and failure operations also receive the originally staged arguments for correlation. Typed callback envelopes are exported. | -| Actor-to-actor delivery | Native | `sendTo(reference).operation()` stages delivery in the source actor commit. | -| One-shot and recurring reminders | Native | Scheduling, replacement events, catch-up policy, stale-claim recovery, pausing, authorized inspection, and idempotent resume are implemented. | -| Same-database commit actions | Native | Registered actions receive source-message identity, mailbox sequence, activation generation, and the fenced transaction connection. | -| Ambient transaction rejection | Native | Committed calls and message waits fail before blocking when the current async context already owns a transaction on the Solid Objects adapter. | -| Direct application-write isolation during actor code | Partial | `guardApplicationDatabase()` fails closed for operations, projections, migrations, and commit actions; only the supplied fenced commit-action connection may write. Unwrapped clients cannot be intercepted. | -| Committed snapshots | Native | `snapshot()` returns authorized persisted fields and inferred getters from one read-only committed state image; realtime replay reads explicit observables without mailbox history. | -| Actor destruction and incarnation fencing | Native | Authorized cascading deletion creates a fresh instance ID on recreation; an authorized waiter receives `ActorDestroyed` when that incarnation disappears. | -| Result recovery and sync timeout diagnostics | Native | Status, result, and wait reauthorize the stored operation; terminal failure raises structured `MessageFailed`; whole-call adapter deadlines distinguish enqueue, wait, database, activation, and mailbox blockers. | -| Result lookup by request ID and idempotency key | Native | `runtime.findBy({ requestId })` and `reference.findBy({ idempotencyKey })` rebuild a `MessageReference`, authorized with the hook the original call ran. An actor remembers the keys of its own finished turns, so a key lookup separates a pruned message from one that never existed. | +| Capability | Status | TypeScript shape or remaining work | +| ------------------------------------------------------------------------------------------------------------- | ------- | --------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | +| Actor registry, durable identity, JSON state, and adjacent state migrations | Native | Ordinary classes, static actor types, inferred state, explicit migrations, and isolated runtime context across every actor-instance callback. | +| Query and projection purity | Native | Both runtimes reject state mutation and staged effects, recovery checks, commit actions, reminders, or outbound messages in queries, observable projections, and personalized payloads. Projection guards compare the complete staged work, including same-count replacements. Query and observable violations fail terminally with `QueryMutatedState`; payload violations fail only that payload. Snapshot projection reads enforce the same rule. | +| JSON property preservation | Native | Reserved property names such as `__proto__`, `constructor`, and `prototype` remain ordinary data, including through size checks and immutable copies. Both suites consume `compatibility/json-values.json`. | +| Fluent committed calls and background delivery | Native | `await reference.operation()` and `reference.send.operation()`. | +| Ordered mailbox, sequence allocation, idempotency, retries, dead letters, leases, renewal, and fenced commits | Native | Relational ready/claimed membership tables, distinct generated request IDs and caller idempotency keys, durable history, adapter-appropriate sequence locking, and PostgreSQL/MySQL row locks held from fence validation through commit. | +| Domain rejection and strict poison ordering | Native | Rejections accept JavaScript identifier-style codes and roll back without retry; invalid codes fail terminally, while retryable failures block later operations until completion or dead-lettering. | +| Bounded activation passes and hot-actor fairness | Native | Configurable turn-count and elapsed-time budgets bound each pass, then move only that actor's already-due memberships behind actors already waiting. | +| Bounded claim candidate scan | Native | A configurable ordered scan continues to another ready actor when a worker loses the first candidate's lease race. | +| Backpressure and payload caps | Partial | Serialization enforces a shared maximum JSON nesting depth, raising `InvalidPayload`, and an optional caller-supplied `maxBytes` limit, raising `PayloadTooLarge`; keyed reminder names are bounded to 255 characters in JS and 191 in Ruby. These existing database limits remain runtime-specific; use at most 191 ASCII characters for a name shared across runtimes. Distributed per-actor rate limits, global admission control, and cache-capacity eviction are not planned in either runtime. They are hot, request-path, and loss-tolerant, so one durable ordered message per check is the wrong shape. Solid Objects Pro answers them with grouped and ephemeral operations, which [fit](fit.md) describes. | +| Idle activation cache | Native | Long-running workers retain hydrated actors under renewable fenced leases, restore public state after failed turns, and release on timeout, fairness yield, lease loss, or shutdown. | +| Transactional effects and outcome operations | Native | At-least-once handlers receive immutable stable effect, attempt, source-message, and actor identity; success and failure operations also receive the originally staged arguments for correlation. Typed callback envelopes are exported. | +| Actor-to-actor delivery | Native | `sendTo(reference).operation()` stages delivery in the source actor commit. | +| One-shot and recurring reminders | Native | Scheduling, cancellation by name/key/handle, read-your-writes inspection, catch-up policy, stale-claim recovery, pausing, and authorized resume. Cancellation cannot recall an already-enqueued occurrence. | +| Same-database commit actions | Native | Registered actions receive source-message identity, mailbox sequence, activation generation, and the fenced transaction connection. | +| Ambient transaction rejection | Native | Committed calls and message waits fail before blocking when the current async context already owns a transaction on the Solid Objects adapter. | +| Direct application-write isolation during actor code | Partial | `guardApplicationDatabase()` fails closed for operations, projections, migrations, and commit actions; only the supplied fenced commit-action connection may write. Unwrapped clients cannot be intercepted. | +| Committed snapshots | Native | `snapshot()` returns authorized persisted fields and inferred getters from one read-only committed state image; realtime replay reads explicit observables without mailbox history. | +| Actor destruction and incarnation fencing | Native | Authorized cascading deletion creates a fresh instance ID on recreation; an authorized waiter receives `ActorDestroyed` when that incarnation disappears. | +| Result recovery and sync timeout diagnostics | Native | Status, result, outcome, and wait reauthorize the original stored operation and arguments. All delivery modes retain immutable JSON results under the result-size cap. Result reads raise terminal rejection/failure errors; outcome returns them as data. Whole-call deadlines distinguish enqueue, wait, database, activation, and mailbox blockers. | +| Result lookup by request ID and idempotency key | Native | `runtime.findBy({ requestId })` and `reference.findBy({ idempotencyKey })` rebuild a `MessageReference`, authorized with the hook the original call ran. An actor remembers the keys of its own finished turns, so a key lookup separates a pruned message from one that never existed. | Effect callback envelopes are typed with `EffectFailurePayload`, `EffectSuccessPayload`, and `SerializedError` in both SQL and Cloudflare. -Ruby RBS contracts are tracked in cardmagic/solid-objects-ruby#64 and preserve -Ruby field names; this does not change runtime delivery semantics. +Ruby publishes RBS contracts with native Ruby field names; this does not change runtime delivery semantics. ## Operations -| Capability | Status | TypeScript shape or remaining work | -| ----------------------------------------------------------------------------- | ------ | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------ | -| Process registration, heartbeats, stale claim recovery, and graceful shutdown | Native | Runtime roles persist host, PID, runtime versions, draining and stopped transitions, cooperative cancellation, and a bounded shutdown deadline; cleanup recovers stale claims. | -| Failed-role replacement | Native | Built-in and registered roles are rebuilt through their factories with capped backoff; shutdown is the terminal replacement boundary. | -| Additional supervised components | Native | `registerComponent()` builds, validates, runs, and stops application components with the runtime. | -| Dead-letter inspection and retry | Native | `runtime.deadLetters` provides deny-by-default immutable inspection and idempotent durable retry linkage. | -| Reconciliation reads | Native | Authorized cursor pages cover active, quiet, and orphaned instances; bounded state batches are migrated and deeply frozen. | -| Message, process, and opt-in instance retention | Native | Supervised scheduling bounds message and process growth; authorized manual APIs add preview and keep destructive instance expiration explicit. SQL instance pruning uses an index on actor type and update time. | -| Doctor and schema verification | Native | Structured checks cover configuration, schema/version shape, adapter server versions, neutral-context policy probes, live roles, and a targeted round trip. | -| CLI | Native | The packaged executable loads an application runtime and exposes start, diagnostics, processes, dead letters, reminders, and explicit retention pruning as JSON. | -| Operator dashboard | Native | The opt-in `solid-objects/web` export provides Fetch and Node/Connect mounting, authorized runtime views and actions, session-backed CSRF, filtering, paging, charts, and immutable extension hooks. Matches the Ruby dashboard's own documented limits: no audit trail of admin actions, dead-letter retry is one at a time, and pause sets a flag rather than interrupting an in-flight turn. | -| Structured instrumentation | Native | An isolated transport-neutral sink emits immutable lifecycle metadata and structurally excludes application payloads. | -| Large committed state warning | Native | `warnStateBytes` reports one `solid_objects.state.large` event, holding the actor type, actor ID, byte count, and threshold, when a committed image passes a 128 KB soft threshold. The event holds no application state, and it reports after the commit. The Ruby gem carries the same event and the same 5 MB hard default from `0.14.3`, as `warn_state_bytes`. Its threshold defaults to 64 KB rather than 128 KB, because its measured curve falls sooner: it keeps 55% of its empty-state throughput at 13 KB, where this package keeps 98% at 16 KB. | -| Public test helper | Native | `runtime.testing` provides role-selective deterministic draining, explicit-time due-reminder execution, and dependency-ordered reset without relying on cascades. | +| Capability | Status | TypeScript shape or remaining work | +| ----------------------------------------------------------------------------- | ------ | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------ | +| Process registration, heartbeats, stale claim recovery, and graceful shutdown | Native | Runtime roles persist host, PID, runtime versions, draining and stopped transitions, cooperative cancellation, and a bounded shutdown deadline; cleanup recovers stale claims. | +| Failed-role replacement | Native | Built-in and registered roles are rebuilt through their factories with capped backoff; shutdown is the terminal replacement boundary. | +| Additional supervised components | Native | `registerComponent()` builds, validates, runs, and stops application components with the runtime. | +| Dead-letter inspection and retry | Native | `runtime.deadLetters` provides deny-by-default immutable inspection and idempotent durable retry linkage. | +| Reconciliation reads | Native | Authorized cursor pages cover active, quiet, and orphaned instances; bounded state batches are migrated and deeply frozen. | +| Message, process, and opt-in instance retention | Native | Supervised scheduling bounds message and process growth; authorized manual APIs add preview and keep destructive instance expiration explicit. SQL instance pruning uses an index on actor type and update time. | +| Doctor and schema verification | Native | Structured checks cover configuration, schema/version shape, adapter server versions, neutral-context policy probes, live roles, and a targeted round trip. | +| CLI | Native | The packaged executable loads an application runtime and exposes start, diagnostics, processes, dead letters, reminders, and explicit retention pruning as JSON. | +| Operator dashboard | Native | Fetch and Node/Connect mounting, authorized views and actions, CSRF, filtering, paging, charts, and extensions match the Rack dashboard. Both dashboards expose individual retries and pause/resume; bulk redrive is available through the administration APIs. Pause does not interrupt an in-flight turn. | +| Bulk redrive and administration audit | Native | Message, effect, and broadcast scopes support idempotent retry and durable, bounded, cancellable redrive. Retry and task transitions persist administration identity in audit rows. | +| Structured instrumentation | Native | Versioned events, metric samples, isolated observers, and bounded authorized diagnostics match Ruby. Both test suites check SQL event attributes against `compatibility/telemetry-events.json`; only JS activation events add their turn's message fields. Both runtimes log exporter and observer failures, require observer callbacks, and cap local observers at 1,000. Timeout events share wait reasons, activation owner IDs, and string generations; see [observability](observability.md). Durable Objects uses host instrumentation and remote diagnostics, sends fewer attributes, and does not send `sync.timeout`. | +| Large committed state warning | Native | `warnStateBytes` reports one `solid_objects.state.large` event, holding the actor type, actor ID, byte count, and threshold, when a committed image passes a 128 KB soft threshold. The event holds no application state, and it reports after the commit. The Ruby gem carries the same event and the same 5 MB hard default from `0.14.3`, as `warn_state_bytes`. Its threshold defaults to 64 KB rather than 128 KB, because its measured curve falls sooner: it keeps 55% of its empty-state throughput at 13 KB, where this package keeps 98% at 16 KB. | +| Public test helper | Native | `runtime.testing` provides role-selective deterministic draining, explicit-time due-reminder execution, and dependency-ordered reset without relying on cascades. | ## Databases and wake-up -| Capability | Status | TypeScript shape or remaining work | -| ------------------------ | ------ | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | -| SQLite | Native | Uses built-in `node:sqlite`, serialized process-local access, bounded transient writer retries, foreign keys, strict tables, database time, and deadline-bounded access and lock waits. | -| PostgreSQL | Native | Optional `pg` 8.23 peer, bounded pooling, 64-bit schema, row-locked sequences, server checks, and deadline-bounded pool, statement, and lock waits. | -| MySQL | Native | Optional `mysql2` 3.23 peer, bounded pooling, InnoDB schema, row-locked sequences, scoped deadlock retry, and deadline-bounded pool, query, and lock waits. Ruby also tests a second client, `trilogy`; Node has no comparable second MySQL client, so only `mysql2` is tracked here. | -| Durable polling fallback | Native | Every role progresses without a notification service. Effects, reminders, and broadcasts use canonical ordered polling indexes; PostgreSQL and MySQL lock only the selected row with `FOR UPDATE SKIP LOCKED`. Broadcasts and reminders compare separate available and stale-recovery probes, preserve the oldest-first choice, and retry past candidates locked by another claimant. | -| In-process wake-up | Native | A generation-based default adapter prevents claim-to-wait signal loss; commits wake role-specific waiters and polling remains the fallback. | -| PostgreSQL wake-up | Native | `database.wakeUp()` uses one dedicated event-driven client, role-specific `LISTEN/NOTIFY`, generation fencing, reconnectable listeners, and durable polling fallback. | -| Redis wake-up | Native | An optional `redis` peer provides role-specific Pub/Sub over separate lazy publisher/subscriber connections, with bounded failures and durable polling fallback. | +| Capability | Status | TypeScript shape or remaining work | +| --------------------------- | ------ | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | +| SQLite | Native | Uses built-in `node:sqlite`, serialized process-local access, bounded transient writer retries, foreign keys, strict tables, database time, and deadline-bounded access and lock waits. Ruby uses a yielding busy handler for background writes so a waiting thread lets the lock holder commit, including on Rails 7.1/7.2; both preserve synchronous deadlines. | +| PostgreSQL | Native | Optional `pg` 8.23 peer, bounded pooling, 64-bit schema, row-locked sequences, server checks, and deadline-bounded pool, statement, and lock waits. | +| MySQL | Native | Optional `mysql2` 3.23 peer, bounded pooling, InnoDB schema, row-locked sequences, scoped deadlock retry, and deadline-bounded pool, query, and lock waits. Ruby also tests a second client, `trilogy`; Node has no comparable second MySQL client, so only `mysql2` is tracked here. | +| Durable polling fallback | Native | Every role progresses without a notification service. Effects, reminders, and broadcasts use canonical ordered polling indexes; PostgreSQL and MySQL lock only the selected row with `FOR UPDATE SKIP LOCKED`. Broadcasts and reminders compare separate available and stale-recovery probes, preserve the oldest-first choice, and retry past candidates locked by another claimant. | +| In-process wake-up | Native | A generation-based adapter prevents claim-to-wait signal loss. It is the polling fallback or an explicit opt-out from automatic cross-process selection. | +| Automatic wake-up selection | Native | Prefer configured Redis, then proven PostgreSQL notifications, then polling. Probe failures warn and report the resolved capability; MySQL needs Redis for cross-process wake-up. | +| PostgreSQL wake-up | Native | `database.wakeUp()` uses one dedicated event-driven client, role-specific `LISTEN/NOTIFY`, generation fencing, reconnectable listeners, and durable polling fallback. | +| Redis wake-up | Native | An optional `redis` peer provides role-specific Pub/Sub over separate lazy publisher/subscriber connections, with bounded failures and durable polling fallback. | Both runtimes select a wake-up adapter automatically. `wakeUp` takes a name or an adapter and defaults to `"automatic"`, which prefers a configured Redis URL, @@ -152,11 +152,14 @@ with the hook the original call ran, and report absence, an unregistered actor, and a refusal the same way. `outcome` reports the status, result, error, rejection, and attempt count in both. -Three details differ, and all come from the runtimes rather than the feature. -This runtime stores a result for every completed message, so a lookup answers -one for asynchronous work; Ruby stores a result only for `sync` delivery, so a -lookup there answers the status and the error but not the result. This runtime -needed a new unique index on `request_id`, added as schema version 12, because +Both runtimes retain JSON results for synchronous, background, and internal +messages under the configured result-size cap. Every status, result, and outcome +read rechecks authorization against the original operation and arguments. +`result` raises terminal failures; `outcome` exposes them as data. Results discarded +by older Ruby workers cannot be recovered retroactively. + +Two storage details differ. This runtime needed a new unique index on +`request_id`, added as schema version 12, because its table constrained the pair `(actor_type, actor_id, request_id)`; the Ruby schema has carried a global unique index since its first migration. The error record also differs: Ruby's `ErrorRecord` carries `class_name`, `message`, and @@ -189,15 +192,15 @@ because a Durable Object indexes only its own messages, so that form raises ## Realtime and browser behavior -| Capability | Status | TypeScript shape or remaining work | -| ------------------------------------------------------------ | -------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------ | -| Explicit observable projection and durable invalidations | Native | `observables()` is opt-in and invalidation-only by default. `broadcastValue()` sends changed values; `broadcastInvalidation()` explicitly sends only changed names while comparing the real value. Private or subscriber-specific values belong behind invalidation-only component endpoints or in typed payloads. | -| Action Cable channels and signed stream names | Not applicable | `runtime.realtime` provides authenticated transport-neutral sessions; the host owns its HTTP/WebSocket server and authentication. | -| Authorized subscriptions | Native | Each request is denied by default and authorized before actor lookup; sessions replay committed observables and fence ordered durable revisions. Multi-process hosts explicitly bridge their shared transport. | -| Turbo scalar replacement | Not applicable | The browser client exposes invalidations to application rendering code. Framework adapters can be separate packages. | -| Keyed component refresh, morph/replace, and batch coalescing | Native | A typed framework-neutral registry selects explicit dependencies, coalesces batch requests, aborts superseded work, fences each target, and delegates synchronous application strategy to the host. | -| Personalized payload broadcasts | Native | Static typed projections run against committed state under each fresh subscriber context, reauthorize as queries, isolate failures, and carry independent revision fences. | -| Real-browser compatibility suite | Native | Playwright exercises subscription replay over native WebSocket, incarnation/revision fences, payload delivery, component batching, and cancellation in Chromium. | +| Capability | Status | TypeScript shape or remaining work | +| ------------------------------------------------------------ | -------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------ | +| Explicit observable projection and durable invalidations | Native | `observables()` is opt-in and invalidation-only by default. `broadcastValue()` sends changed values; `broadcastInvalidation()` explicitly sends only changed names while comparing the real value. Private or subscriber-specific values belong behind invalidation-only component endpoints or in typed payloads. | +| Action Cable channels and signed stream names | Not applicable | `runtime.realtime` provides authenticated transport-neutral sessions; the host owns its HTTP/WebSocket server and authentication. | +| Authorized subscriptions | Native | Each request is denied by default and authorized before actor lookup; sessions replay committed observables and fence ordered durable revisions. Multi-process hosts explicitly bridge their shared transport. | +| Turbo scalar replacement | Not applicable | The browser client exposes invalidations to application rendering code. Framework adapters can be separate packages. | +| Keyed component refresh, morph/replace, and batch coalescing | Native | A typed framework-neutral registry selects explicit dependencies, coalesces batch requests, aborts superseded work, fences each target, and delegates synchronous application strategy to the host. | +| Personalized payload broadcasts | Native | Both runtimes evaluate each personalized payload on an isolated actor from the same committed snapshot, including generated defaults, reject state mutation and staged work, prevent application database writes, and enforce the configured payload byte limit. They reauthorize as queries, isolate failures, and carry independent revision fences. | +| Real-browser compatibility suite | Native | Playwright exercises subscription replay over native WebSocket, incarnation/revision fences, payload delivery, component batching, and cancellation in Chromium. | ## JavaScript-only: the browser runtime @@ -257,6 +260,8 @@ and [#48](https://github.com/cardmagic/solid-objects-ruby/issues/48)): with `register_transmit` is the staging side. Identifiers differ by runtime idiom, but both sides guarantee the same wire contract: +- transmit siblings preserve staging order within each source message, including retries; +- explicit `arguments: null` is invalid; only omission defaults to an empty object; - envelope keys are camelCase (`effectId`, `actorType`, `actorId`, `operation`, and an optional `arguments` that defaults to an empty object); @@ -264,6 +269,13 @@ runtime idiom, but both sides guarantee the same wire contract: - a replay with changed arguments raises the idempotency conflict on both sides and leaves the first application intact. +JavaScript schema migration 14 adds an effect position for ordering within a +source message. Ruby uses its existing sequential effect primary key. Existing +JS rows receive position zero and retain their earlier ID-based tie-break; +their original staging order cannot be reconstructed. Run `runtime.install()` +before starting upgraded workers. The stronger ordering applies to effects +staged by upgraded workers. + `compatibility/transmit-envelopes.json` is committed to both repositories with a consuming test on each side, so the contract is enforced from both sides of the repository boundary. Manual cross-runtime QA (Node to Rails diff --git a/package.json b/package.json index 8d1487a..96c85c4 100644 --- a/package.json +++ b/package.json @@ -115,8 +115,8 @@ "test": "vitest run", "test:browser": "pnpm run build && playwright test", "test:coverage": "vitest run --coverage", - "test:postgresql": "vitest run test/postgresql.test.ts test/effect-recovery.test.ts test/instance-retention.test.ts", - "test:mysql": "vitest run test/mysql.test.ts test/effect-recovery.test.ts test/instance-retention.test.ts", + "test:postgresql": "vitest run test/postgresql.test.ts test/effect-recovery.test.ts test/instance-retention.test.ts test/transmit-ordering.test.ts", + "test:mysql": "vitest run test/mysql.test.ts test/effect-recovery.test.ts test/instance-retention.test.ts test/transmit-ordering.test.ts", "test:package": "node scripts/release-artifact-smoke.mjs", "test:recovery": "pnpm run build && node examples/failure-recovery/demo.ts", "test:at-least-once": "pnpm run build && node examples/at-least-once/demo.ts", diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 11e0d4c..06d19ca 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -5,6 +5,7 @@ settings: excludeLinksFromLockfile: false overrides: + undici@7.29.0: 7.29.1 sharp@<0.35.4: ^0.35.4 importers: @@ -1082,8 +1083,8 @@ packages: undici-types@7.18.2: resolution: {integrity: sha512-AsuCzffGHJybSaRrmr5eHr81mwJU3kjw6M+uprWvCXiNeN9SOGwQ3Jn8jb8m3Z6izVgknn1R0FTCEAP2QrLY/w==} - undici@7.29.0: - resolution: {integrity: sha512-IDxfleLmmbSskfWSUATiN1nfn2rDuvnMOqb5CWR92iIfojA0Ud+ulOAAEQ57LPr9rWmsreUyf5lwyao+7GNNVw==} + undici@7.29.1: + resolution: {integrity: sha512-RYONW2MeafgYlkVOKYKkA/Ag7BmXqgIWCa8t1m0JcxrQg9pI9lEqRhAOruOBCbAohOa/gkCF+iPi9hrgvTzu6Q==} engines: {node: '>=20.18.1'} unenv@2.0.0-rc.24: @@ -1838,7 +1839,7 @@ snapshots: dependencies: '@cspotcode/source-map-support': 0.8.1 sharp: 0.35.4(@types/node@24.13.3) - undici: 7.29.0 + undici: 7.29.1 workerd: 1.20260903.1 ws: 8.21.0 youch: 4.1.0-beta.10 @@ -2041,7 +2042,7 @@ snapshots: undici-types@7.18.2: {} - undici@7.29.0: {} + undici@7.29.1: {} unenv@2.0.0-rc.24: dependencies: diff --git a/pnpm-workspace.yaml b/pnpm-workspace.yaml index 0f8fd9d..61f7a64 100644 --- a/pnpm-workspace.yaml +++ b/pnpm-workspace.yaml @@ -6,4 +6,5 @@ minimumReleaseAgeExclude: - "@vitest/mocker@4.1.11" - sharp@0.35.4 overrides: + undici@7.29.0: 7.29.1 sharp@<0.35.4: ^0.35.4 diff --git a/roadmap.md b/roadmap.md index 987ff63..bf19d6e 100644 --- a/roadmap.md +++ b/roadmap.md @@ -44,11 +44,3 @@ letters, reminders, outboxes, incarnation, and revision. Fleet-wide operations remain optional where a backend cannot provide them. [Issue #41: portable actor administration](https://github.com/cardmagic/solid-objects-js/issues/41) - -## Portable observability and diagnostics - -Define common structured events, metrics, and diagnostic views for activation -duration, mailbox depth, retries, dead letters, reminder lateness, outbox age, -recovery failures, and realtime subscriptions without leaking provider APIs. - -[Issue #42: portable observability and diagnostics](https://github.com/cardmagic/solid-objects-js/issues/42) diff --git a/src/actor-runtime.ts b/src/actor-runtime.ts index 8f17e18..9ff9ac0 100644 --- a/src/actor-runtime.ts +++ b/src/actor-runtime.ts @@ -1,3 +1,5 @@ +import type { EventObserver } from "./telemetry.js" +import type { ActorDiagnostics, DiagnosticOptions } from "./diagnostics.js" import type { RealtimeManager } from "./realtime.js" import type { Actor, ActorClass } from "./actor.js" import type { @@ -29,6 +31,17 @@ export interface SnapshotWithIncarnation { } export interface ActorRuntime { + diagnostics( + reference: ActorReferenceCore, + options?: DiagnosticOptions, + ): Promise + observe( + reference: ActorReferenceCore, + options: { + onEvent: EventObserver + authorizationContext?: SnapshotOptions["authorizationContext"] + }, + ): Promise<() => void> ref( actorClass: ActorClass, actorId: ActorIdentifier, diff --git a/src/actor.ts b/src/actor.ts index 7a6e718..b1065c0 100644 --- a/src/actor.ts +++ b/src/actor.ts @@ -536,6 +536,10 @@ export abstract class Actor { intentCount(): number { return Object.values(this.#intents).reduce((count, intents) => count + intents.length, 0) } + + intentSnapshot(): string { + return JSON.stringify(this.#intents) + } } function isObservableBroadcast(value: object): value is ObservableBroadcast { diff --git a/src/cloudflare/engine.ts b/src/cloudflare/engine.ts index d582c1a..0613a26 100644 --- a/src/cloudflare/engine.ts +++ b/src/cloudflare/engine.ts @@ -1,3 +1,5 @@ +import { diagnosticLimit, diagnosticSummary } from "../diagnostics.js" +import { telemetryEvent, deliverTelemetry } from "../telemetry.js" import type { Actor, ActorClass, ActorIntents } from "../actor.js" import { withActorContext, withActorProjection, withRuntime } from "../context.js" import { @@ -78,7 +80,9 @@ export class ActorEngine { arguments: jsonObject(input.payload.arguments), }) this.bind(input) - return this.store.atomic(() => normalizeJson(this.enqueue(input))) + const message = await this.store.atomic(() => this.enqueue(input)) + this.emit("message.enqueued", { messageId: message.id, attempt: message.attempt }) + return normalizeJson(message) } if (input.method === "message" || input.method === "lookup") return this.readMessage(input) if (input.method === "administration") return this.administer(input) @@ -100,6 +104,7 @@ export class ActorEngine { String(input.payload.subscriptionId), ) }) + this.emit("realtime.disconnected", { actorType: input.actorType, actorId: input.actorId }) return null } if (input.method === "snapshot") { @@ -135,13 +140,15 @@ export class ActorEngine { }), ) } + if (input.method === "subscribe") + this.emit("realtime.connected", { actorType: input.actorType, actorId: input.actorId }) return this.projection({ input, payloadNames }) } async pump(): Promise { - await this.store.atomic(() => { + const dispatchedReminders = await this.store.atomic(() => { this.store.prune() - this.scheduleReminders() + const dispatched = this.scheduleReminders() if (this.actorRunning) { const head = this.store.head() if (head?.status === "claimed") { @@ -154,7 +161,9 @@ export class ActorEngine { outbox.availableAt = Date.now() + RECOVERY_INTERVAL this.store.saveOutbox(outbox) } + return dispatched }) + for (const reminder of dispatchedReminders) this.emit("reminder.enqueued", reminder) const work: Promise[] = [] if (!this.actorRunning) work.push(this.drainActors()) for (const outbox of this.store.outboxHeads()) { @@ -430,6 +439,7 @@ export class ActorEngine { }) if (stableJson(actorState(actor, definition.stateKeys)) !== before || actor.hasIntents()) throw new QueryMutatedState("snapshot getters must not mutate state or stage work") + this.emit("snapshot.read", { actorType: identity.actorType, actorId: identity.actorId }) return { snapshot, instanceId: instance?.incarnation ?? "0", @@ -513,7 +523,12 @@ export class ActorEngine { this.store.saveMessage(message) } this.cached = undefined - this.emit("actor.destroyed") + this.emit("actor.destroyed", { + actorType: instance.actorType, + actorId: instance.actorId, + incarnation: instance.incarnation, + revision: instance.revision, + }) return true } @@ -551,9 +566,11 @@ export class ActorEngine { readReminders: this.readReminders, }) if (this.cached?.actor !== actor) { + this.emit("activation.started", { messageId: message.id, attempt: message.attempt }) await withActorContext({ actor, runtime: this.runtime }, () => actor.activate()) this.assertCurrent(instance) this.cached = { incarnation: instance.incarnation, actor } + this.emit("activation.completed", { messageId: message.id, attempt: message.attempt }) } } catch (error) { if (error instanceof ActorDestroyed) return @@ -570,19 +587,28 @@ export class ActorEngine { } this.store.saveMessage(message) }) + this.emit("activation.failed", { + messageId: message.id, + attempt: message.attempt, + errorName: errorName(error), + }) this.emit("actor.setup_failed", { errorName: errorName(error) }) return } const stateBefore = deepCopy(actorState(actor, definition.stateKeys)) + const startedAt = performance.now() try { await this.store.atomic(() => { this.assertCurrent(instance) message.status = "claimed" + message.claimedAt = Date.now() message.generation = instance.generation message.attempt += 1 message.availableAt = Date.now() + RECOVERY_INTERVAL this.store.saveMessage(message) }) + if (message.attempt > 1 && message.error === null) + this.emit("recovery.reclaimed", { messageId: message.id, attempt: message.attempt }) this.emit("message.started", { messageId: message.id, attempt: message.attempt }) const evaluated = await evaluateActorTurn({ actor, @@ -634,7 +660,11 @@ export class ActorEngine { this.stage({ instance: current, message, intents, broadcast: evaluated.broadcast }) this.completeReminder(message) }) - this.emit("message.completed", { messageId: message.id }) + this.emit("message.completed", { + messageId: message.id, + attempt: message.attempt, + durationMilliseconds: performance.now() - startedAt, + }) } catch (error) { actor.discardIntents() for (const key of definition.stateKeys) @@ -693,7 +723,13 @@ export class ActorEngine { this.store.saveMessage(message) } }) + if (message.status === "ready") + this.emit("message.retry", { messageId: message.id, attempt: message.attempt }) + if (message.status === "dead") + this.emit("dead_letter.created", { messageId: message.id, attempt: message.attempt }) this.emit("message.failed", { + attempt: message.attempt, + durationMilliseconds: performance.now() - startedAt, messageId: message.id, errorName: errorName(error), status: message.status, @@ -819,9 +855,16 @@ export class ActorEngine { return } try { + this.emit("outbox.age", { + messageId: outbox.messageId, + outboxKind: outbox.kind, + attempt: outbox.attempt + 1, + ageMilliseconds: Math.max(0, Date.now() - outbox.availableAt), + }) await this.store.atomic(() => { if (!this.outboxCurrent(outbox, instance)) throw new ActorDestroyed("outbox was removed") outbox.status = "claimed" + outbox.deliveryAvailableAt = outbox.availableAt outbox.attempt += 1 outbox.availableAt = Date.now() + RECOVERY_INTERVAL this.store.saveOutbox(outbox) @@ -971,9 +1014,10 @@ export class ActorEngine { }) } - private scheduleReminders(): void { + private scheduleReminders(): JsonObject[] { + const dispatched: JsonObject[] = [] const instance = this.store.instance() - if (!instance || instance.paused) return + if (!instance || instance.paused) return dispatched for (const reminder of this.store.rows( "SELECT record FROM reminders WHERE status = 'scheduled' AND due_at <= ? ORDER BY due_at LIMIT ?", [Date.now(), this.settings.maxMessagesPerActivationPass], @@ -991,6 +1035,11 @@ export class ActorEngine { availableAt: reminder.at, }, }) + dispatched.push({ + messageId: message.id, + attempt: message.attempt, + latenessMilliseconds: Math.max(0, Date.now() - reminder.at), + }) message.reminder = { name: reminder.name, generation: reminder.generation } this.store.saveMessage(message) reminder.status = "completed" @@ -1000,6 +1049,7 @@ export class ActorEngine { break } } + return dispatched } private completeReminder(message: Message): void { @@ -1036,13 +1086,15 @@ export class ActorEngine { const action = String(input.payload.action) if ( !(await this.settings.authorizeAdministration({ - action, - resource: `actor:${actorName(input)}`, + action: action === "diagnostics" ? "inspect" : action, + resource: action === "diagnostics" ? "actor_diagnostics" : `actor:${actorName(input)}`, + resourceId: actorName(input), authorizationContext: input.authorizationContext, })) ) throw new Unauthorized("actor administration is not authorized") this.bind(input) + if (action === "diagnostics") return this.diagnostics(input) if (action === "deadLetters") return normalizeJson({ messages: this.store.rows( @@ -1107,21 +1159,75 @@ export class ActorEngine { }) } + private diagnostics(input: HostRequest): JsonObject { + const limit = diagnosticLimit(Number(input.payload.limit ?? 100)) + const instance = this.store.instance() + const now = Date.now() + const queries = { + mailbox: + "SELECT CASE WHEN status = 'claimed' THEN COALESCE(json_extract(record, '$.claimedAt'), json_extract(record, '$.createdAt')) ELSE available_at END AS at FROM messages WHERE status IN ('ready', 'claimed')", + outbox: + "SELECT CASE WHEN status = 'claimed' THEN COALESCE(json_extract(record, '$.deliveryAvailableAt'), available_at) ELSE available_at END AS at FROM outboxes WHERE status IN ('pending', 'claimed')", + reminders: "SELECT due_at AS at FROM reminders WHERE status IN ('scheduled', 'paused')", + retries: + "SELECT available_at AS at FROM messages WHERE status = 'ready' AND json_extract(record, '$.error') IS NOT NULL", + } + const summaries: JsonObject = {} + for (const [name, query] of Object.entries(queries)) { + const rows = this.store.storage.sql + .exec<{ at: number }>(`${query} ORDER BY at LIMIT ?`, limit + 1) + .toArray() + summaries[name] = normalizeJson( + diagnosticSummary({ timestamps: rows.map((row) => row.at), now, limit }), + ) + } + const mailbox = jsonObject(summaries.mailbox) + this.emit("mailbox.depth", { + actorType: input.actorType, + actorId: input.actorId, + count: mailbox.sampled ?? 0, + truncated: mailbox.truncated ?? false, + ...(mailbox.truncated ? {} : { depth: mailbox.sampled ?? 0 }), + }) + return { + actorType: input.actorType, + actorId: input.actorId, + incarnation: instance?.incarnation ?? null, + revision: instance?.revision ?? null, + adapter: "durable-objects", + occurredAt: new Date(now).toISOString(), + limit, + ...summaries, + recoveryFailures: normalizeJson(diagnosticSummary({ timestamps: [], now, limit })), + } + } + private retryDelay(attempt: number): number { const delay = this.settings.retryDelayMilliseconds(attempt) return Number.isFinite(delay) && delay >= 1 ? delay : 1_000 } private emit(name: string, attributes: JsonObject = {}): void { + if (!this.settings.instrumentation) return try { - this.settings.instrumentation?.({ - name: `solid_objects.${name}`, - occurredAt: new Date().toISOString(), - attributes, + const instance = this.store.instance() + const event = telemetryEvent({ + name, + adapter: "durable-objects", + attributes: { + actorType: instance?.actorType ?? null, + actorId: instance?.actorId ?? null, + incarnation: instance?.incarnation ?? null, + revision: instance?.revision ?? null, + ...attributes, + }, }) - } catch { - this.settings.logger.error({ event: "solid_objects.instrumentation.failed", name }) - } + deliverTelemetry({ + observer: this.settings.instrumentation, + event, + logger: this.settings.logger, + }) + } catch {} } } diff --git a/src/cloudflare/records.ts b/src/cloudflare/records.ts index b856cf9..20e80fd 100644 --- a/src/cloudflare/records.ts +++ b/src/cloudflare/records.ts @@ -26,6 +26,7 @@ export interface Message { idempotencyKey: string | null status: MessageStatus attempt: number + claimedAt?: number availableAt: number createdAt: number completedAt: number | null @@ -46,6 +47,7 @@ export interface Outbox { payload: JsonObject status: "pending" | "claimed" | "completed" | "dead" attempt: number + deliveryAvailableAt?: number availableAt: number completedAt: number | null error: JsonObject | null diff --git a/src/cloudflare/runtime.ts b/src/cloudflare/runtime.ts index bebabd3..96c84a8 100644 --- a/src/cloudflare/runtime.ts +++ b/src/cloudflare/runtime.ts @@ -1,3 +1,5 @@ +import type { EventObserver } from "../telemetry.js" +import type { ActorDiagnostics, DiagnosticOptions } from "../diagnostics.js" import "./platform.js" import type { Actor, ActorClass } from "../actor.js" import type { ActorRuntime } from "../actor-runtime.js" @@ -78,6 +80,32 @@ export class CloudflareRuntime implements ActorRuntime { }) } + async observe( + _reference: ActorReferenceCore, + _options: { + onEvent: EventObserver + authorizationContext?: SnapshotOptions["authorizationContext"] + }, + ): Promise<() => void> { + return unsupported( + "process-local observers; configure instrumentation on the actor Durable Object host", + ) + } + + async diagnostics( + reference: ActorReferenceCore, + options: DiagnosticOptions = {}, + ): Promise { + const value = await this.call({ + actorType: reference.actorType, + actorId: reference.actorId, + method: "administration", + authorizationContext: normalizeJson(options.authorizationContext ?? null), + payload: { action: "diagnostics", limit: options.limit ?? 100 }, + }) + return readonlyCopy(value) as ActorDiagnostics & JsonObject + } + async invoke(options: { reference: ActorReferenceCore operation: string diff --git a/src/configuration.ts b/src/configuration.ts index 70b860b..a369d2e 100644 --- a/src/configuration.ts +++ b/src/configuration.ts @@ -1,3 +1,4 @@ +import type { MetricSample } from "./telemetry.js" import { InvalidActor } from "./errors.js" import type { Database } from "./database/types.js" import type { @@ -40,6 +41,15 @@ export interface SubscriptionAuthorizationInput { } export interface InstrumentationEvent { + readonly schemaVersion: 1 + readonly adapter: string + readonly actorType: string | null + readonly actorId: string | null + readonly incarnation: string | null + readonly revision: string | null + readonly messageId: string | null + readonly attempt: number + readonly metrics: readonly MetricSample[] readonly name: string readonly occurredAt: string readonly attributes: DeepReadonly diff --git a/src/diagnostics.ts b/src/diagnostics.ts new file mode 100644 index 0000000..a55eedf --- /dev/null +++ b/src/diagnostics.ts @@ -0,0 +1,114 @@ +import type { SolidObjectsRuntime } from "./runtime.js" +import type { SnapshotOptions } from "./types.js" + +export interface DiagnosticSummary { + readonly sampled: number + readonly truncated: boolean + readonly oldestAgeMilliseconds: number | null +} + +export interface ActorDiagnostics { + readonly actorType: string + readonly actorId: string + readonly incarnation: string | null + readonly revision: string | null + readonly adapter: string + readonly occurredAt: string + readonly limit: number + readonly mailbox: DiagnosticSummary + readonly outbox: DiagnosticSummary + readonly reminders: DiagnosticSummary + readonly retries: DiagnosticSummary + readonly recoveryFailures: DiagnosticSummary +} + +export interface DiagnosticOptions extends SnapshotOptions { + limit?: number +} + +export function diagnosticLimit(limit = 100): number { + if (!Number.isSafeInteger(limit) || limit < 1 || limit > 100) + throw new RangeError("diagnostic limit must be an integer between 1 and 100") + return limit +} + +export async function actorDiagnostics(options: { + runtime: SolidObjectsRuntime + actorType: string + actorId: string + options: DiagnosticOptions +}): Promise { + const { runtime, actorType, actorId } = options + await runtime.authorizeAdministration({ + action: "inspect", + resource: "actor_diagnostics", + resourceId: JSON.stringify([actorType, actorId]), + authorizationContext: options.options.authorizationContext, + }) + const limit = diagnosticLimit(options.options.limit) + const instance = await runtime.repository.findInstanceByIdentity(actorType, actorId) + return runtime.settings.database.connection(async (connection) => { + const now = await connection.nowMilliseconds() + const table = (name: string) => runtime.repository.table(name) + const queries = { + mailbox: `SELECT available_at_ms AS at FROM ${table("ready_messages")} WHERE instance_id = ? UNION ALL SELECT claimed_at_ms AS at FROM ${table("claimed_messages")} WHERE instance_id = ?`, + outbox: `SELECT available_at_ms AS at FROM ${table("effects")} WHERE instance_id = ? AND status IN ('pending', 'processing') UNION ALL SELECT available_at_ms AS at FROM ${table("broadcasts")} WHERE instance_id = ? AND status IN ('pending', 'processing')`, + reminders: `SELECT run_at_ms AS at FROM ${table("reminders")} WHERE instance_id = ? AND status IN ('scheduled', 'paused')`, + retries: `SELECT ready.available_at_ms AS at FROM ${table("ready_messages")} ready JOIN ${table("messages")} message ON message.id = ready.message_id WHERE ready.instance_id = ? AND message.attempt_count > 0 AND message.error IS NOT NULL`, + recoveryFailures: `SELECT dead.created_at_ms AS at FROM ${table("dead_letters")} dead JOIN ${table("messages")} message ON message.id = dead.message_id WHERE dead.instance_id = ? AND message.delivery_mode = 'internal' AND message.idempotency_key LIKE 'effect:%:recovery'`, + } + const summarize = async (name: keyof typeof queries): Promise => { + const parameters = [instance?.id ?? ""] + if (name === "mailbox" || name === "outbox") parameters.push(instance?.id ?? "") + const rows = instance + ? await connection.all<{ at: number | bigint }>( + `${queries[name]} ORDER BY at LIMIT ${limit + 1}`, + parameters, + ) + : [] + return diagnosticSummary({ + timestamps: rows.map((row) => Number(row.at)), + now, + limit, + }) + } + const result: ActorDiagnostics = Object.freeze({ + actorType, + actorId, + incarnation: instance?.id ?? null, + revision: instance ? String(instance.state_revision) : null, + adapter: runtime.settings.database.family, + occurredAt: new Date(now).toISOString(), + limit, + mailbox: await summarize("mailbox"), + outbox: await summarize("outbox"), + reminders: await summarize("reminders"), + retries: await summarize("retries"), + recoveryFailures: await summarize("recoveryFailures"), + }) + runtime.emitInstrumentation("mailbox.depth", { + actorType, + actorId, + instanceId: instance?.id ?? null, + count: result.mailbox.sampled, + truncated: result.mailbox.truncated, + depth: result.mailbox.truncated ? null : result.mailbox.sampled, + }) + return result + }) +} + +export function diagnosticSummary(options: { + timestamps: number[] + now: number + limit: number +}): DiagnosticSummary { + return Object.freeze({ + sampled: Math.min(options.timestamps.length, options.limit), + truncated: options.timestamps.length > options.limit, + oldestAgeMilliseconds: + options.timestamps.length === 0 + ? null + : Math.max(0, options.now - Math.min(...options.timestamps)), + }) +} diff --git a/src/doctor.ts b/src/doctor.ts index 1ed7927..e1159c3 100644 --- a/src/doctor.ts +++ b/src/doctor.ts @@ -78,7 +78,7 @@ const EXPECTED_COLUMNS: Readonly> = { "claimed_at_ms", ], reminders: ["id", "instance_id", "operation", "message_operation", "run_at_ms", "status"], - effects: ["id", "message_id", "instance_id", "name", "status", "available_at_ms"], + effects: ["id", "message_id", "instance_id", "name", "status", "available_at_ms", "position"], effect_recoveries: [ "effect_id", "instance_id", diff --git a/src/index.ts b/src/index.ts index c5029d2..8633ef5 100644 --- a/src/index.ts +++ b/src/index.ts @@ -218,3 +218,6 @@ export { UnsupportedDatabase, } from "./errors.js" export * from "./effect-recovery.js" + +export type { MetricSample, EventObserver } from "./telemetry.js" +export type { ActorDiagnostics, DiagnosticOptions, DiagnosticSummary } from "./diagnostics.js" diff --git a/src/realtime.ts b/src/realtime.ts index 717de3f..f94362c 100644 --- a/src/realtime.ts +++ b/src/realtime.ts @@ -122,6 +122,10 @@ export class RealtimeManager { private removeSession(session: ManagedRealtimeSession): void { for (const [key, sessions] of this.subscriptions) { + if (sessions.has(session)) { + const [actorType, actorId] = JSON.parse(key) as [string, string] + this.runtime.emitInstrumentation("realtime.disconnected", { actorType, actorId }) + } sessions.delete(session) if (sessions.size === 0) this.subscriptions.delete(key) } @@ -130,6 +134,11 @@ export class RealtimeManager { private add(session: ManagedRealtimeSession, subscription: SubscriptionIdentity): void { const key = subscriptionKey(subscription) const sessions = this.subscriptions.get(key) ?? new Set() + if (!sessions.has(session)) + this.runtime.emitInstrumentation("realtime.connected", { + actorType: subscription.actorType, + actorId: subscription.actorId, + }) sessions.add(session) this.subscriptions.set(key, sessions) session.add(subscription) @@ -139,7 +148,11 @@ export class RealtimeManager { const key = subscriptionKey(subscription) const sessions = this.subscriptions.get(key) if (!sessions) return - sessions.delete(session) + if (sessions.delete(session)) + this.runtime.emitInstrumentation("realtime.disconnected", { + actorType: subscription.actorType, + actorId: subscription.actorId, + }) if (sessions.size === 0) this.subscriptions.delete(key) session.remove(subscription) } diff --git a/src/reference.ts b/src/reference.ts index 03ae4de..feb1e44 100644 --- a/src/reference.ts +++ b/src/reference.ts @@ -1,3 +1,5 @@ +import type { EventObserver } from "./telemetry.js" +import type { ActorDiagnostics, DiagnosticOptions } from "./diagnostics.js" import type { Actor, ActorClass } from "./actor.js" import { SyncInsideTransaction, UnknownOperation } from "./errors.js" import type { Outcome } from "./outcome.js" @@ -236,6 +238,34 @@ export class ActorReferenceCore { this.send = createMessageSender(this, {}) } + diagnostics(options: DiagnosticOptions = {}): Promise { + return this.runtime.diagnostics(this, options) + } + + async observe(options: { + onEvent: EventObserver + authorizationContext?: SnapshotOptions["authorizationContext"] + }): Promise<() => void> { + assertObserver(options.onEvent) + return this.runtime.observe(this, options) + } + + async on( + name: string, + options: { + onEvent: EventObserver + authorizationContext?: SnapshotOptions["authorizationContext"] + }, + ): Promise<() => void> { + assertObserver(options.onEvent) + return this.observe({ + ...options, + onEvent: (event) => { + if (event.name === `solid_objects.${name}`) return options.onEvent(event) + }, + }) + } + with(options: InvocationOptions): ActorInvoker { return createInvoker(this, options) } @@ -363,6 +393,12 @@ function createMessageSender( }) as ActorMessageSender } +function assertObserver(onEvent: unknown): void { + if (typeof onEvent !== "function") { + throw new TypeError("an actor observer requires an onEvent callback") + } +} + function assertOperation(operations: ReadonlySet, operation: string): void { if (!operations.has(operation)) { throw new UnknownOperation(`unknown operation ${JSON.stringify(operation)}`) diff --git a/src/repository.ts b/src/repository.ts index 08cf972..153d02a 100644 --- a/src/repository.ts +++ b/src/repository.ts @@ -731,13 +731,13 @@ export class Repository { turn.message.id, ]) - for (const effect of input.intents.effects) { + for (const [position, effect] of input.intents.effects.entries()) { const effectId = effect.id ?? randomUUID() await connection.run( `INSERT INTO ${this.table("effects")} (id, message_id, instance_id, name, arguments, success_operation, failure_operation, - status, max_attempts, available_at_ms) - VALUES (?, ?, ?, ?, ?, ?, ?, 'pending', ?, ?)`, + status, max_attempts, available_at_ms, position) + VALUES (?, ?, ?, ?, ?, ?, ?, 'pending', ?, ?, ?)`, [ effectId, turn.message.id, @@ -748,6 +748,7 @@ export class Repository { effect.failureOperation ?? null, this.settings.maxAttempts, now, + position, ], ) if (effect.recoveryOperation !== undefined || effect.statusOperation !== undefined) { @@ -1842,7 +1843,7 @@ export class Repository { async enqueueReminder( reminder: ReminderRow, options: { nowMilliseconds?: number } = {}, - ): Promise { + ): Promise { return this.settings.database.transaction(async (connection) => { const now = options.nowMilliseconds ?? (await connection.nowMilliseconds()) const claimed = await connection.get( @@ -1858,9 +1859,9 @@ export class Repository { `SELECT id FROM ${this.table("reminders")} WHERE id = ?`, [reminder.id], )) - if (!claimed && !surviving) return false + if (!claimed && !surviving) return undefined if (!claimed) throw new LostActivation("reminder claim no longer matches") - await this.enqueueInTransaction(connection, { + const message = await this.enqueueInTransaction(connection, { actorType: claimed.actor_type, actorId: claimed.actor_id, operation: claimed.message_operation ?? claimed.operation, @@ -1890,7 +1891,7 @@ export class Repository { claimed.claimed_by, ], ) - return true + return message }) } diff --git a/src/runtime.ts b/src/runtime.ts index 935f340..b2107c0 100644 --- a/src/runtime.ts +++ b/src/runtime.ts @@ -1,3 +1,5 @@ +import { actorDiagnostics, type DiagnosticOptions, type ActorDiagnostics } from "./diagnostics.js" +import { telemetryEvent, deliverTelemetry, type EventObserver } from "./telemetry.js" import type { Actor, ActorClass, ObservableProjection } from "./actor.js" import { BroadcastWorker } from "./broadcast-worker.js" import { @@ -209,6 +211,7 @@ export class SolidObjectsRuntime { readonly realtime readonly processes readonly administration + private readonly observers = new Set() private readonly registry = new Map() private readonly effects = new Map() private readonly commitActions = new Map() @@ -833,7 +836,7 @@ export class SolidObjectsRuntime { : (await this.repository.remindersForInstance(instance.id)).map(scheduledReminderOf), }) const stateBefore = stableJson(actorState(actor, registered.definition.stateKeys)) - const intentCount = actor.intentCount() + const intentsBefore = actor.intentSnapshot() const snapshot: Record = { ...state } await withActorProjection({ actor, runtime: this }, async () => { for (const query of registered.definition.queries) { @@ -846,10 +849,16 @@ export class SolidObjectsRuntime { }) if ( stableJson(actorState(actor, registered.definition.stateKeys)) !== stateBefore || - actor.intentCount() !== intentCount + actor.intentSnapshot() !== intentsBefore ) { throw new QueryMutatedState("snapshot getters must not mutate actor state or stage work") } + this.emitInstrumentation("snapshot.read", { + actorType: reference.actorType, + actorId: reference.actorId, + instanceId: instance?.id ?? null, + revision: String(instance?.state_revision ?? 0), + }) return { snapshot: readonlyCopy(snapshot) as ActorSnapshot, instanceId: instance?.id ?? "0", @@ -1276,6 +1285,10 @@ export class SolidObjectsRuntime { (await this.repository.remindersForInstance(turn.instance.id)).map(scheduledReminderOf), }) } catch (error) { + this.emitInstrumentation("activation.failed", { + ...activationInstrumentation(turn), + errorName: error instanceof Error ? error.name : "Error", + }) throw new ActorSetupFailed(error) } const definition = registered.definition @@ -1290,23 +1303,25 @@ export class SolidObjectsRuntime { if (!activated) { try { + this.emitInstrumentation("activation.started", activationInstrumentation(turn)) await withActorContext({ actor, runtime: this }, () => actor.activate()) activated = true - this.emitInstrumentation("activation.started", { - actorType: turn.message.actor_type, - actorId: turn.message.actor_id, - instanceId: turn.instance.id, - generation: String(turn.activationGeneration), - }) + this.emitInstrumentation("activation.completed", activationInstrumentation(turn)) } catch (error) { renewalController.abort() await renewal actor.discardIntents() + this.emitInstrumentation("activation.failed", { + ...activationInstrumentation(turn), + errorName: error instanceof Error ? error.name : "Error", + }) throw new ActorSetupFailed(renewalError ?? error) } } - const startedAt = Date.now() + const startedAt = performance.now() + if (Number(turn.message.attempt_count) > 1 && turn.message.error === null) + this.emitInstrumentation("recovery.reclaimed", messageInstrumentation(turn.message)) this.emitInstrumentation("message.started", messageInstrumentation(turn.message)) try { stateBefore = deepCopy(actorState(actor, definition.stateKeys)) @@ -1343,6 +1358,7 @@ export class SolidObjectsRuntime { throw new UnknownCommitAction(`unknown commit action ${JSON.stringify(intent.name)}`) const attributes = { commitAction: intent.name, + activationGeneration: String(turn.activationGeneration), ...messageInstrumentation(turn.message), } this.emitInstrumentation("commit_action.started", attributes) @@ -1383,10 +1399,13 @@ export class SolidObjectsRuntime { nextRunAt: new Date(replacement.nextRunAtMilliseconds).toISOString(), }) } + if (recoveryMessage(turn.message)) + this.emitInstrumentation("recovery.completed", messageInstrumentation(turn.message)) this.warnAboutLargeState(turn.message, committed) this.emitInstrumentation("message.completed", { ...messageInstrumentation(turn.message), - durationMilliseconds: Date.now() - startedAt, + revision: String(turn.message.sequence), + durationMilliseconds: performance.now() - startedAt, }) return { actor, retainActivation: true, activated } } catch (error) { @@ -1397,7 +1416,7 @@ export class SolidObjectsRuntime { if (error instanceof LostActivation) { this.emitInstrumentation("activation.lost", { ...messageInstrumentation(turn.message), - durationMilliseconds: Date.now() - startedAt, + durationMilliseconds: performance.now() - startedAt, }) return { actor, retainActivation: false, activated } } @@ -1410,7 +1429,7 @@ export class SolidObjectsRuntime { this.emitInstrumentation("message.rejected", { ...messageInstrumentation(turn.message), code: error.code, - durationMilliseconds: Date.now() - startedAt, + durationMilliseconds: performance.now() - startedAt, }) return { actor, retainActivation: activated, activated } } @@ -1423,12 +1442,15 @@ export class SolidObjectsRuntime { ...messageInstrumentation(turn.message), retryable, errorName: error instanceof Error ? error.name : "Error", - durationMilliseconds: Date.now() - startedAt, + durationMilliseconds: performance.now() - startedAt, outcome, }) - if (outcome === "dead") { + if (outcome === "retrying") + this.emitInstrumentation("message.retry", messageInstrumentation(turn.message)) + if (outcome === "dead" && recoveryMessage(turn.message)) + this.emitInstrumentation("recovery.failed", messageInstrumentation(turn.message)) + if (outcome === "dead") this.emitInstrumentation("dead_letter.created", messageInstrumentation(turn.message)) - } return { actor, retainActivation: activated, activated } } } @@ -1474,6 +1496,7 @@ export class SolidObjectsRuntime { await this.callerWorker?.stop() this.realtime.close() await this.closeWakeUp() + this.observers.clear() await this.settings.database.close() clearDefaultRuntime(this) } @@ -1529,28 +1552,65 @@ export class SolidObjectsRuntime { await this.repository.resetForTesting() } + diagnostics( + reference: ActorReferenceCore, + options: DiagnosticOptions = {}, + ): Promise { + return actorDiagnostics({ + runtime: this, + actorType: reference.actorType, + actorId: reference.actorId, + options, + }) + } + + async observe( + reference: ActorReferenceCore, + options: { + onEvent: EventObserver + authorizationContext?: SnapshotOptions["authorizationContext"] + }, + ): Promise<() => void> { + await this.authorizeAdministration({ + action: "observe", + resource: "actor_diagnostics", + resourceId: JSON.stringify([reference.actorType, reference.actorId]), + authorizationContext: options.authorizationContext, + }) + if (this.observers.size >= 1000) + throw new RangeError("at most 1000 local observers may be registered") + const observer: EventObserver = (event) => { + if (event.actorType === reference.actorType && event.actorId === reference.actorId) + return options.onEvent(event) + } + this.observers.add(observer) + return () => { + this.observers.delete(observer) + } + } + emitInstrumentation(name: string, attributes: JsonObject): void { - const instrumentation = this.settings.instrumentation - if (!instrumentation) return + if (!this.settings.instrumentation && this.observers.size === 0) return try { - instrumentation( - Object.freeze({ - name: `solid_objects.${name}`, - occurredAt: new Date().toISOString(), - attributes: readonlyCopy(attributes), - }), - ) - } catch (error) { - this.settings.logger.error({ - event: "solid_objects.instrumentation.failed", - instrumentationEvent: `solid_objects.${name}`, - errorName: error instanceof Error ? error.name : "Error", - }) - } + const event = telemetryEvent({ name, adapter: this.settings.database.family, attributes }) + if (this.settings.instrumentation) + deliverTelemetry({ + observer: this.settings.instrumentation, + event, + logger: this.settings.logger, + }) + for (const observer of this.observers) + deliverTelemetry({ observer, event, logger: this.settings.logger }) + } catch {} } async executeEffect(effect: EffectRow): Promise { try { + this.emitInstrumentation("outbox.age", { + ...effectInstrumentation(effect), + outboxKind: "effect", + ageMilliseconds: Math.max(0, Date.now() - Number(effect.available_at_ms)), + }) const handler = this.effects.get(effect.name) if (!handler) throw new UnknownEffect(`unknown effect ${JSON.stringify(effect.name)}`) const result = normalizeJson( @@ -1592,14 +1652,22 @@ export class SolidObjectsRuntime { if (!actor.operations.has(dispatchOperation)) { throw new UnknownOperation(`unknown reminder operation ${JSON.stringify(dispatchOperation)}`) } - if (!(await this.repository.enqueueReminder(reminder, options))) return + const message = await this.repository.enqueueReminder(reminder, options) + if (!message) return this.wakeUp("actors") this.emitInstrumentation("reminder.enqueued", { reminderId: reminder.id, + messageId: message.id, + attempt: Number(message.attempt_count), actorType: reminder.actor_type, actorId: reminder.actor_id, operation: reminder.operation, + instanceId: reminder.instance_id, + latenessMilliseconds: Math.max( + 0, + (options.nowMilliseconds ?? Date.now()) - Number(reminder.run_at_ms), + ), occurrence: Number(reminder.occurrence), }) } @@ -1607,6 +1675,11 @@ export class SolidObjectsRuntime { async executeBroadcast(broadcast: BroadcastRow): Promise { const deliver = this.settings.broadcast try { + this.emitInstrumentation("outbox.age", { + ...broadcastInstrumentation(broadcast), + outboxKind: "broadcast", + ageMilliseconds: Math.max(0, Date.now() - Number(broadcast.available_at_ms)), + }) const event = readonlyCopy({ actorType: broadcast.actor_type, actorId: broadcast.actor_id, @@ -2152,7 +2225,7 @@ export class SolidObjectsRuntime { state: deepCopy(options.snapshot.state), }) const stateBefore = stableJson(actorState(actor, options.registered.definition.stateKeys)) - const intentCount = actor.intentCount() + const intentsBefore = actor.intentSnapshot() const handler = options.registered.definition.payloads[options.name] if (!handler) { throw new UnknownPayloadBroadcast(`unknown payload broadcast ${options.name}`) @@ -2162,7 +2235,7 @@ export class SolidObjectsRuntime { ) if ( stableJson(actorState(actor, options.registered.definition.stateKeys)) !== stateBefore || - actor.intentCount() !== intentCount + actor.intentSnapshot() !== intentsBefore ) { throw new QueryMutatedState("payload broadcasts must not mutate actor state or stage work") } @@ -2400,8 +2473,17 @@ function actorDestroyedWhileWaiting(): ActorDestroyed { return new ActorDestroyed("actor was destroyed while waiting for its result") } +function recoveryMessage(message: MessageRow): boolean { + return ( + message.delivery_mode === "internal" && + message.idempotency_key?.startsWith("effect:") === true && + message.idempotency_key.endsWith(":recovery") + ) +} + function messageInstrumentation(message: MessageRow): JsonObject { return { + instanceId: message.instance_id, messageId: message.id, requestId: message.request_id, actorType: message.actor_type, @@ -2413,6 +2495,14 @@ function messageInstrumentation(message: MessageRow): JsonObject { } } +function activationInstrumentation(turn: ClaimedTurn): JsonObject { + return { + ...messageInstrumentation(turn.message), + generation: String(turn.activationGeneration), + ownerId: turn.processId, + } +} + function actorMessageContext(message: MessageRow): MessageContext { return Object.freeze({ id: message.id, @@ -2437,6 +2527,7 @@ function restoreActorState(options: { function effectInstrumentation(effect: EffectRow): JsonObject { return { + instanceId: effect.instance_id, effectId: effect.id, effectName: effect.name, messageId: effect.message_id, @@ -2448,6 +2539,7 @@ function effectInstrumentation(effect: EffectRow): JsonObject { function broadcastInstrumentation(broadcast: BroadcastRow): JsonObject { return { + instanceId: broadcast.instance_id, broadcastId: broadcast.id, messageId: broadcast.message_id, actorType: broadcast.actor_type, diff --git a/src/schema.ts b/src/schema.ts index 1ccafae..85f34f8 100644 --- a/src/schema.ts +++ b/src/schema.ts @@ -14,7 +14,8 @@ const INSTANCE_RETENTION_INDEX_VERSION = 10 const DEAD_LETTER_REDRIVE_VERSION = 11 const REQUEST_ID_LOOKUP_VERSION = 12 const REMEMBERED_KEYS_VERSION = 13 -const LATEST_VERSION = REMEMBERED_KEYS_VERSION +const EFFECT_POSITION_VERSION = 14 +const LATEST_VERSION = EFFECT_POSITION_VERSION export const SCHEMA_VERSIONS: readonly number[] = Object.freeze( Array.from({ length: LATEST_VERSION }, (_unused, index) => index + 1), @@ -390,6 +391,16 @@ export async function installSchema(options: { }) } + if (!installedVersions.has(EFFECT_POSITION_VERSION)) { + await addEffectPosition({ connection, family, table: table("effects") }) + await recordMigration({ + connection, + table: table("schema_migrations"), + version: EFFECT_POSITION_VERSION, + schemaIdentity, + }) + } + if (installedVersions.has(POLLING_INDEXES_VERSION)) return const pollingIndexes = [ ["effects", `${prefix}effects_poll`, "status, available_at_ms, id"], @@ -471,12 +482,23 @@ async function installRedrive(options: { }) } +async function addEffectPosition(options: { + connection: DatabaseConnection + family: DatabaseFamily + table: string +}): Promise { + if (await hasColumn({ ...options, column: "position" })) return + await options.connection.run( + `ALTER TABLE ${options.table} ADD COLUMN position INTEGER NOT NULL DEFAULT 0`, + ) +} + async function addFailedAt(options: { connection: DatabaseConnection family: DatabaseFamily table: string }): Promise { - if (await hasFailedAt(options)) return + if (await hasColumn({ ...options, column: "failed_at_ms" })) return const type = options.family === "sqlite" ? "INTEGER" : "BIGINT" await options.connection.run(`ALTER TABLE ${options.table} ADD COLUMN failed_at_ms ${type}`) @@ -485,22 +507,23 @@ async function addFailedAt(options: { ) } -async function hasFailedAt(options: { +async function hasColumn(options: { connection: DatabaseConnection family: DatabaseFamily table: string + column: string }): Promise { if (options.family === "sqlite") { const columns = await options.connection.all<{ name: string }>( `PRAGMA table_info(${options.table})`, ) - return columns.some(({ name }) => name === "failed_at_ms") + return columns.some(({ name }) => name === options.column) } const schema = options.family === "postgresql" ? "current_schema()" : "DATABASE()" const found = await options.connection.get<{ found: number | bigint }>( `SELECT COUNT(*) AS found FROM information_schema.columns - WHERE table_schema = ${schema} AND table_name = ? AND column_name = 'failed_at_ms'`, - [options.table], + WHERE table_schema = ${schema} AND table_name = ? AND column_name = ?`, + [options.table, options.column], ) return Number(found?.found ?? 0) > 0 } diff --git a/src/serialization.ts b/src/serialization.ts index 74b55cc..01301ec 100644 --- a/src/serialization.ts +++ b/src/serialization.ts @@ -54,12 +54,12 @@ function normalize(value: unknown, depth: number): JsonValue { if (Array.isArray(value)) return value.map((item) => normalize(item, depth + 1)) if (isRecord(value)) { - const output: Record = {} - for (const [key, item] of Object.entries(value)) { - if (item === undefined) throw new InvalidPayload(`undefined is not supported at ${key}`) - output[key] = normalize(item, depth + 1) - } - return output + return Object.fromEntries( + Object.entries(value).map(([key, item]) => { + if (item === undefined) throw new InvalidPayload(`undefined is not supported at ${key}`) + return [key, normalize(item, depth + 1)] + }), + ) } throw new InvalidPayload(`${describe(value)} is not JSON-compatible`) diff --git a/src/telemetry.ts b/src/telemetry.ts new file mode 100644 index 0000000..ac37899 --- /dev/null +++ b/src/telemetry.ts @@ -0,0 +1,149 @@ +import type { InstrumentationEvent } from "./configuration.js" +import type { JsonObject, JsonValue, Logger } from "./types.js" +import { readonlyCopy } from "./serialization.js" + +export interface MetricSample { + readonly name: string + readonly kind: "counter" | "gauge" | "histogram" + readonly unit: "1" | "ms" + readonly value: number + readonly labels: Readonly> +} + +export type EventObserver = (event: InstrumentationEvent) => void + +export const portableAttributes: ReadonlySet = new Set([ + "failureCount", + "phase", + "activationOwnerId", + "activationGeneration", + "pollingIntervalMilliseconds", + "idlePollingIntervalMilliseconds", + "currentIntervalMilliseconds", + "component", + "workers", + "effectWorkers", + "broadcastWorkers", + "reminderSchedulers", + "actorType", + "actorId", + "instanceId", + "incarnation", + "revision", + "messageId", + "requestId", + "attempt", + "sequence", + "operation", + "deliveryMode", + "generation", + "durationMilliseconds", + "latenessMilliseconds", + "ageMilliseconds", + "depth", + "count", + "errorName", + "outcome", + "retryable", + "status", + "effectId", + "effectName", + "reminderId", + "occurrence", + "broadcastId", + "code", + "role", + "reason", + "processId", + "processKind", + "ownerId", + "componentCount", + "byteCount", + "thresholdBytes", + "previousRunAt", + "nextRunAt", + "name", + "commitAction", + "outboxKind", + "truncated", + "payload", + "waitingOn", + "timeoutMilliseconds", + "intervalMilliseconds", + "previousIntervalMilliseconds", +]) + +export function telemetryEvent(options: { + name: string + adapter: string + attributes: JsonObject +}): InstrumentationEvent { + const attributes: JsonObject = {} + for (const [key, value] of Object.entries(options.attributes)) { + if ( + portableAttributes.has(key) && + (value === null || ["string", "number", "boolean"].includes(typeof value)) + ) { + attributes[key] = value + } + } + const name = `solid_objects.${options.name}` + const labels = Object.freeze({ + event: name, + adapter: options.adapter, + actorType: String(attributes.actorType ?? ""), + }) + const metrics: MetricSample[] = [ + { name: "solid_objects.events", kind: "counter", unit: "1", value: 1, labels }, + ] + for (const [field, metric, kind, unit] of [ + ["durationMilliseconds", "solid_objects.duration", "histogram", "ms"], + ["latenessMilliseconds", "solid_objects.reminder.lateness", "histogram", "ms"], + ["ageMilliseconds", "solid_objects.outbox.age", "histogram", "ms"], + ["depth", "solid_objects.mailbox.depth", "gauge", "1"], + ] as const) { + const value = attributes[field] + if (typeof value === "number" && Number.isFinite(value)) + metrics.push({ name: metric, kind, unit, value: Math.max(0, value), labels }) + } + return Object.freeze({ + schemaVersion: 1, + name, + occurredAt: new Date().toISOString(), + adapter: options.adapter, + actorType: identifier(attributes.actorType), + actorId: identifier(attributes.actorId), + incarnation: identifier(attributes.incarnation ?? attributes.instanceId), + revision: identifier(attributes.revision), + messageId: identifier(attributes.messageId), + attempt: typeof attributes.attempt === "number" ? attributes.attempt : 0, + attributes: readonlyCopy(attributes), + metrics: Object.freeze(metrics.map((sample) => Object.freeze(sample))), + }) +} + +export function deliverTelemetry(options: { + observer: EventObserver + event: InstrumentationEvent + logger: Logger +}): void { + const failed = (error: ErrorValue) => { + try { + const result = options.logger.error({ + event: "solid_objects.instrumentation.failed", + instrumentationEvent: options.event.name, + errorName: error instanceof Error ? error.name : "Error", + }) + void Promise.resolve(result).catch(() => {}) + } catch {} + } + try { + void Promise.resolve(options.observer(options.event)).catch(failed) + } catch (error) { + failed(error) + } +} + +function identifier(value: JsonValue | undefined): string | null { + return typeof value === "string" || typeof value === "number" ? String(value) : null +} diff --git a/src/transmit.ts b/src/transmit.ts index 981cd26..6897b6b 100644 --- a/src/transmit.ts +++ b/src/transmit.ts @@ -50,7 +50,7 @@ export async function receiveTransmitEnvelope(options: { throw new InvalidPayload(`transmit envelope requires a non-empty ${field}`) } } - const argumentsValue = envelope.arguments ?? {} + const argumentsValue = envelope.arguments === undefined ? {} : envelope.arguments if (!isJsonObject(argumentsValue)) { throw new InvalidPayload("transmit envelope arguments must be a JSON object") } @@ -115,7 +115,7 @@ async function undeliveredEnvelopesThrough(input: { AND messages.actor_type = ? AND messages.actor_id = ? AND messages.sequence <= (SELECT sequence FROM ${messages} WHERE id = ?) - ORDER BY messages.sequence, effects.id`, + ORDER BY messages.sequence, effects.position, effects.id`, [effectName, context.actorType, context.actorId, context.sourceMessageId], ), ) diff --git a/src/turn.ts b/src/turn.ts index 99c5f83..f34d5f7 100644 --- a/src/turn.ts +++ b/src/turn.ts @@ -14,11 +14,11 @@ export function readActorObservables(options: { }): ObservableProjection { const { actor, definition, runtime } = options const before = options.stateJson ?? stableJson(actorState(actor, definition.stateKeys)) - const intentCount = actor.intentCount() + const intentsBefore = actor.intentSnapshot() const projection = withActorProjection({ actor, runtime }, () => actor.observableValues()) if ( stableJson(actorState(actor, definition.stateKeys)) !== before || - actor.intentCount() !== intentCount + actor.intentSnapshot() !== intentsBefore ) { throw new QueryMutatedState("observables must not mutate actor state or stage durable work") } diff --git a/src/wake-up-notification.ts b/src/wake-up-notification.ts index 604ef18..896f79c 100644 --- a/src/wake-up-notification.ts +++ b/src/wake-up-notification.ts @@ -7,7 +7,14 @@ export function notifyWakeUp(options: { role: WakeUpRole }): void { const logFailure = (errorName: string): void => { - options.logger.error({ event: "solid_objects.wake_up.failed", role: options.role, errorName }) + try { + const result = options.logger.error({ + event: "solid_objects.wake_up.failed", + role: options.role, + errorName, + }) + void Promise.resolve(result).catch(() => {}) + } catch {} } try { Promise.resolve(options.adapter.notify(options.role)).catch((error) => diff --git a/test/cloudflare/telemetry.test.ts b/test/cloudflare/telemetry.test.ts new file mode 100644 index 0000000..4438168 --- /dev/null +++ b/test/cloudflare/telemetry.test.ts @@ -0,0 +1,40 @@ +import { env } from "cloudflare:test" +import { describe, expect, it } from "vitest" +import { createRuntime, durableObjects } from "../../src/cloudflare/index.js" +import { Counter, telemetryEvents } from "./worker.js" + +const authorizationContext = "allowed" + +describe("portable Durable Objects telemetry", () => { + it("uses the common envelope and authorizes bounded diagnostics", async () => { + const runtime = createRuntime({ backend: durableObjects({ namespace: env.ACTORS }) }) + const reference = runtime.ref(Counter, "telemetry") + await expect(reference.diagnostics()).rejects.toMatchObject({ name: "Unauthorized" }) + expect(await reference.with({ authorizationContext }).increment()).toBe(1) + await reference.snapshot({ authorizationContext }) + const summary = await reference.diagnostics({ authorizationContext, limit: 1 }) + expect(summary).toMatchObject({ + adapter: "durable-objects", + actorId: "telemetry", + limit: 1, + mailbox: { sampled: 0, truncated: false }, + }) + const event = telemetryEvents.find( + (event) => event.actorId === "telemetry" && event.name === "solid_objects.message.completed", + ) + expect(event).toMatchObject({ + schemaVersion: 1, + adapter: "durable-objects", + actorType: "Counter", + incarnation: expect.any(String), + messageId: expect.any(String), + attempt: 1, + }) + expect(event?.metrics[0]?.labels).toEqual({ + event: "solid_objects.message.completed", + adapter: "durable-objects", + actorType: "Counter", + }) + await expect(reference.diagnostics({ authorizationContext, limit: 101 })).rejects.toThrow() + }) +}) diff --git a/test/cloudflare/worker.ts b/test/cloudflare/worker.ts index 4d4c29f..d6f8cb6 100644 --- a/test/cloudflare/worker.ts +++ b/test/cloudflare/worker.ts @@ -1,3 +1,4 @@ +import type { InstrumentationEvent } from "../../src/configuration.js" import { Actor, broadcastValue, @@ -14,6 +15,7 @@ import { durableObjects, } from "../../src/cloudflare/index.js" +export const telemetryEvents: InstrumentationEvent[] = [] export const gates = new Map void>() export const revokedSessions = new Set() export const deliveries = new Map() @@ -235,6 +237,9 @@ export class Actors extends createDurableObjectsHost({ }, sessions: environment.SESSIONS, }), + instrumentation: (event) => { + telemetryEvents.push(event) + }, authorizeMessage: (input) => input.authorizationContext === "allowed" || (input.authorizationContext === "argument-denied" && input.arguments.amount !== 7), diff --git a/test/instrumentation.test.ts b/test/instrumentation.test.ts index 2306eef..e148914 100644 --- a/test/instrumentation.test.ts +++ b/test/instrumentation.test.ts @@ -1,8 +1,10 @@ import { afterEach, describe, expect, it, vi } from "vitest" import { Actor } from "../src/actor.js" import type { InstrumentationEvent, SolidObjectsConfiguration } from "../src/configuration.js" +import { postgresql } from "../src/database/postgresql.js" import { sqlite } from "../src/database/sqlite.js" import { configure, type SolidObjectsRuntime } from "../src/runtime.js" +import { expectPortableEvents } from "./support/portable-telemetry.js" class InstrumentedActor extends Actor { static override readonly actorType = "InstrumentedActor" @@ -18,14 +20,38 @@ class InstrumentedActor extends Actor { this.reject("not_allowed", { message: "private rejection", details: { secret: this.secret } }) } + arrange(): void { + this.emit("telemetry-effect", { arguments: { secret: this.secret } }) + this.schedule({ at: new Date(0) }).update({ secret: "reminder-private" }) + } + fail(): void { throw new Error(`private failure ${this.secret}`) } + + commit(): void { + this.commitAction("telemetry-action") + } + + commitBadly(): void { + this.commitAction("telemetry-failure") + } +} + +class ActivationFailure extends Actor { + static override readonly actorType = "ActivationFailure" + + protected override async onActivate(): Promise { + throw new Error("private activation failure") + } + + run(): void {} } let runtime: SolidObjectsRuntime | undefined afterEach(async () => { + await runtime?.repository.resetForTesting() await runtime?.close() runtime = undefined }) @@ -75,6 +101,85 @@ describe("structured instrumentation", () => { expect(serialized).not.toContain("result:") }) + it("emits portable SQL lifecycle events that match the shared attribute contract", async () => { + const events: InstrumentationEvent[] = [] + runtime = configuredRuntime({ + maxAttempts: 2, + retryDelayMilliseconds: () => 0, + instrumentation: (event) => { + events.push(event) + }, + }) + await runtime.install() + runtime.registerEffect("telemetry-effect", () => "delivered") + runtime.registerCommitAction("telemetry-action", () => {}) + runtime.registerCommitAction("telemetry-failure", () => { + throw new Error("private commit failure") + }) + const reference = InstrumentedActor.ref("contract") + await reference.update({ secret: "committed" }) + await reference.send.rejectUpdate() + await reference.send.fail() + await reference.send.commit() + await reference.send.commitBadly() + await reference.send.arrange() + await runtime.enqueueInternalMessage({ + actorType: reference.actorType, + actorId: reference.actorId, + operation: "update", + argumentsValue: { secret: "recovered" }, + idempotencyKey: "effect:contract:recovery", + }) + await reference.diagnostics({ limit: 1 }) + await runtime.worker().runUntilIdle() + await ActivationFailure.ref("contract").send.run() + await runtime + .worker() + .runUntilIdle() + .catch(() => {}) + await runtime.effectWorker().runUntilIdle() + await runtime.reminderScheduler().runOnce() + await reference.snapshot() + + expectPortableEvents(events, [ + "activation.started", + "activation.completed", + "activation.failed", + "message.enqueued", + "message.started", + "message.completed", + "message.rejected", + "message.failed", + "message.retry", + "dead_letter.created", + "commit_action.started", + "commit_action.completed", + "commit_action.failed", + "recovery.completed", + "mailbox.depth", + "outbox.age", + "reminder.enqueued", + "snapshot.read", + ]) + expect( + events + .filter((event) => event.name === "solid_objects.message.failed") + .map(({ attributes }) => ({ + retryable: attributes.retryable, + outcome: attributes.outcome, + })), + ).toEqual( + expect.arrayContaining([ + { retryable: true, outcome: "retrying" }, + { retryable: true, outcome: "dead" }, + ]), + ) + expect( + events.find((event) => event.name === "solid_objects.mailbox.depth")?.attributes, + ).toMatchObject({ truncated: true, depth: null }) + expect(JSON.stringify(events)).not.toContain("private") + }) + it("isolates instrumentation failures from durable work", async () => { const logger = { debug: vi.fn(), @@ -99,13 +204,240 @@ describe("structured instrumentation", () => { errorName: "Error", }) }) + it("isolates a failing sink even when its logger also fails", async () => { + const fail = () => { + throw new Error("private sink error") + } + runtime = configuredRuntime({ + instrumentation: fail, + logger: { debug: fail, info: fail, warn: fail, error: fail }, + }) + await runtime.install() + expect(await InstrumentedActor.ref("safe").update({ secret: "committed" })).toBe( + "result:committed", + ) + }) + + it("adds common correlation and bounded metric labels", async () => { + const events: InstrumentationEvent[] = [] + runtime = configuredRuntime({ + instrumentation: (event) => { + events.push(event) + }, + }) + await runtime.install() + await InstrumentedActor.ref("one").update({ secret: "private" }) + const completed = events.find((event) => event.name === "solid_objects.message.completed")! + expect(completed).toMatchObject({ + schemaVersion: 1, + adapter: runtime.settings.database.family, + actorType: InstrumentedActor.actorType, + actorId: "one", + attempt: 1, + incarnation: expect.any(String), + messageId: expect.any(String), + }) + expect(completed.metrics).toEqual( + expect.arrayContaining([ + expect.objectContaining({ + name: "solid_objects.events", + kind: "counter", + unit: "1", + value: 1, + }), + ]), + ) + expect(JSON.stringify(completed.metrics)).not.toContain('"actorId"') + }) + + it("authorizes bounded diagnostics and local observers before accessing an actor", async () => { + runtime = configuredRuntime({ + authorizeAdministration: ({ authorizationContext }) => authorizationContext === "operator", + }) + await runtime.install() + const reference = InstrumentedActor.ref("diagnostics") + await expect(reference.diagnostics()).rejects.toMatchObject({ name: "Unauthorized" }) + await expect(reference.observe({ onEvent: () => {} })).rejects.toMatchObject({ + name: "Unauthorized", + }) + const events: InstrumentationEvent[] = [] + const stop = await reference.on("message.enqueued", { + authorizationContext: "operator", + onEvent: (event) => { + events.push(event) + }, + }) + await reference.send.update({ secret: "first-secret" }) + await reference.send.update({ secret: "second-secret" }) + await InstrumentedActor.ref("other").send.update({ secret: "other-secret" }) + expect(events).toHaveLength(2) + stop() + await reference.send.update({ secret: "third-secret" }) + expect(events).toHaveLength(2) + const summary = await reference.diagnostics({ authorizationContext: "operator", limit: 1 }) + expect(summary.mailbox).toMatchObject({ sampled: 1, truncated: true }) + expect(summary.outbox.sampled).toBe(0) + expect(JSON.stringify(summary)).not.toContain("secret") + expect(Object.isFrozen(summary.mailbox)).toBe(true) + await expect( + reference.diagnostics({ authorizationContext: "operator", limit: 101 }), + ).rejects.toThrow(RangeError) + }) + + it("observes retries, dead letters, snapshots, reminders, outboxes and realtime", async () => { + const events: InstrumentationEvent[] = [] + runtime = configuredRuntime({ + maxAttempts: 2, + retryDelayMilliseconds: () => 1, + authorizeSubscription: () => true, + instrumentation: (event) => { + events.push(event) + }, + }) + await runtime.install() + runtime.registerEffect("telemetry-effect", () => "provider-private") + const reference = InstrumentedActor.ref("lifecycle") + await expect(reference.fail()).rejects.toMatchObject({ name: "MessageFailed" }) + await reference.arrange() + const summary = await reference.diagnostics() + expect(summary.outbox.sampled).toBe(1) + expect(summary.reminders.sampled).toBe(1) + await runtime.effectWorker().runUntilIdle() + await runtime.reminderScheduler().runOnce() + await reference.snapshot() + const session = runtime.realtime.connect({ authorizationContext: "allowed", send: () => {} }) + await session.receive({ + version: 1, + action: "subscribe", + actorType: reference.actorType, + actorId: reference.actorId, + }) + session.close() + expectPortableEvents(events, ["realtime.connected", "realtime.disconnected"]) + expect(events.map((event) => event.name)).toEqual( + expect.arrayContaining([ + "solid_objects.activation.started", + "solid_objects.activation.completed", + "solid_objects.message.retry", + "solid_objects.dead_letter.created", + "solid_objects.reminder.enqueued", + "solid_objects.outbox.age", + "solid_objects.snapshot.read", + "solid_objects.realtime.connected", + "solid_objects.realtime.disconnected", + ]), + ) + expect( + events.find((event) => event.name === "solid_objects.reminder.enqueued")?.metrics, + ).toEqual( + expect.arrayContaining([ + expect.objectContaining({ name: "solid_objects.reminder.lateness", unit: "ms" }), + ]), + ) + expect(JSON.stringify(events)).not.toContain("provider-private") + expect(JSON.stringify(events)).not.toContain("reminder-private") + }) + + it("rejects observers without an onEvent callback", async () => { + runtime = configuredRuntime() + await runtime.install() + const reference = InstrumentedActor.ref("callbackless") + + await expect(reference.observe({} as never)).rejects.toThrow(TypeError) + await expect(reference.on("message.completed", {} as never)).rejects.toThrow(TypeError) + }) + + it("accepts at most 1000 local observers", async () => { + runtime = configuredRuntime() + await runtime.install() + const reference = InstrumentedActor.ref("observer-limit") + const stops = await Promise.all( + Array.from({ length: 1000 }, () => reference.observe({ onEvent: () => {} })), + ) + + const error = await reference.observe({ onEvent: () => {} }).catch((caught: unknown) => caught) + + expect(error).toBeInstanceOf(RangeError) + expect(error).toHaveProperty("message", "at most 1000 local observers may be registered") + stops.pop()?.() + stops.push(await reference.observe({ onEvent: () => {} })) + for (const stop of stops) stop() + }) + + it("removes local observers when the runtime closes", async () => { + runtime = configuredRuntime() + await runtime.install() + const closed = runtime + const identity = { actorType: InstrumentedActor.actorType, actorId: "closed-observer" } + const events: InstrumentationEvent[] = [] + await InstrumentedActor.ref(identity.actorId).observe({ + onEvent: (event) => { + events.push(event) + }, + }) + closed.emitInstrumentation("custom", identity) + + await closed.close() + runtime = undefined + closed.emitInstrumentation("custom", identity) + + expect(events).toHaveLength(1) + }) + + it("isolates asynchronous sinks and rejects unknown metadata", async () => { + const events: InstrumentationEvent[] = [] + runtime = configuredRuntime({ + instrumentation: (event) => { + events.push(event) + }, + }) + await runtime.install() + const reference = InstrumentedActor.ref("async") + const stop = await reference.observe({ + onEvent: async () => { + throw new Error("exporter-private") + }, + }) + expect(await reference.update({ secret: "committed" })).toBe("result:committed") + runtime.emitInstrumentation("custom", { + actorId: "async", + arguments: "private", + state: { secret: "private" }, + response: "private", + password: "private", + }) + expect(events.at(-1)?.attributes).toEqual({ actorId: "async" }) + stop() + }) + + it("reports failed durable recovery callbacks in diagnostics", async () => { + const events: InstrumentationEvent[] = [] + runtime = configuredRuntime({ + instrumentation: (event) => { + events.push(event) + }, + }) + await runtime.install() + const reference = InstrumentedActor.ref("recovery") + await runtime.enqueueInternalMessage({ + actorType: reference.actorType, + actorId: reference.actorId, + operation: "fail", + idempotencyKey: "effect:test:recovery", + }) + await runtime.worker().runUntilIdle() + expect((await reference.diagnostics()).recoveryFailures.sampled).toBe(1) + expectPortableEvents(events, ["recovery.failed"]) + }) }) function configuredRuntime( overrides: Partial = {}, ): SolidObjectsRuntime { return configure({ - database: sqlite({ path: ":memory:" }), + database: process.env.SOLID_OBJECTS_DATABASE_URL?.startsWith("postgresql:") + ? postgresql({ connectionString: process.env.SOLID_OBJECTS_DATABASE_URL }) + : sqlite({ path: ":memory:" }), authorizeMessage: () => true, authorizeQuery: () => true, authorizeDestroy: () => true, diff --git a/test/lifecycle.test.ts b/test/lifecycle.test.ts index 5c75eec..10eae3c 100644 --- a/test/lifecycle.test.ts +++ b/test/lifecycle.test.ts @@ -1,10 +1,11 @@ import { afterEach, describe, expect, it, vi } from "vitest" import { Actor } from "../src/actor.js" -import type { SolidObjectsConfiguration } from "../src/configuration.js" +import type { InstrumentationEvent, SolidObjectsConfiguration } from "../src/configuration.js" import { StateMigrationError } from "../src/errors.js" import type { JsonObject } from "../src/types.js" import { configure, type SolidObjectsRuntime } from "../src/runtime.js" import { sqlite } from "../src/database/sqlite.js" +import { expectPortableEvents } from "./support/portable-telemetry.js" class LifecycleCounter extends Actor { static override readonly actorType = "LifecycleCounter" @@ -677,7 +678,12 @@ describe("runtime lifecycle", () => { }) it("recovers an expired claimed message", async () => { - runtime = configuredRuntime() + const events: InstrumentationEvent[] = [] + runtime = configuredRuntime({ + instrumentation: (event) => { + events.push(event) + }, + }) await runtime.install() const reference = LifecycleCounter.ref("recovered") const message = await reference.send.increment() @@ -697,6 +703,16 @@ describe("runtime lifecycle", () => { expect(await message.result()).toBe(1) const stored = await runtime.repository.findMessage(message.id) expect(Number(stored?.attempt_count)).toBe(2) + expect(events).toEqual( + expect.arrayContaining([ + expect.objectContaining({ + name: "solid_objects.recovery.reclaimed", + attempt: 2, + messageId: message.id, + }), + ]), + ) + expectPortableEvents(events, ["recovery.reclaimed"]) }) it("drains a bounded activation pass before yielding", async () => { diff --git a/test/outboxes.test.ts b/test/outboxes.test.ts index 1a5eb18..c4aa272 100644 --- a/test/outboxes.test.ts +++ b/test/outboxes.test.ts @@ -1,10 +1,15 @@ import { afterEach, describe, expect, it } from "vitest" import { Actor, broadcastInvalidation, broadcastValue } from "../src/actor.js" -import type { BroadcastEvent, SolidObjectsConfiguration } from "../src/configuration.js" +import type { + BroadcastEvent, + InstrumentationEvent, + SolidObjectsConfiguration, +} from "../src/configuration.js" import { NonRetryableError } from "../src/errors.js" import { configure, type SolidObjectsRuntime } from "../src/runtime.js" import { sqlite } from "../src/database/sqlite.js" import type { EffectFailurePayload, EffectSuccessPayload, JsonObject } from "../src/index.js" +import { expectPortableEvents } from "./support/portable-telemetry.js" class Checkout extends Actor { static override readonly actorType = "Checkout" @@ -414,7 +419,11 @@ describe("reminders", () => { describe("observable broadcasts", () => { it("delivers values and invalidation-only observable names", async () => { const events: BroadcastEvent[] = [] + const telemetry: InstrumentationEvent[] = [] runtime = configuredRuntime({ + instrumentation: (event) => { + telemetry.push(event) + }, broadcast: async (event) => { events.push(event) }, @@ -424,6 +433,20 @@ describe("observable broadcasts", () => { await ObservableCounter.ref("counter").increment() expect(await runtime.broadcastWorker().runUntilIdle()).toBe(1) + expect(telemetry).toEqual( + expect.arrayContaining([ + expect.objectContaining({ + name: "solid_objects.outbox.age", + actorId: "counter", + incarnation: expect.any(String), + attributes: expect.objectContaining({ + outboxKind: "broadcast", + ageMilliseconds: expect.any(Number), + }), + }), + ]), + ) + expectPortableEvents(telemetry, ["outbox.age"]) expect(events).toHaveLength(1) expect(events[0]).toMatchObject({ actorType: "ObservableCounter", diff --git a/test/payload-projections.test.ts b/test/payload-projections.test.ts new file mode 100644 index 0000000..445a168 --- /dev/null +++ b/test/payload-projections.test.ts @@ -0,0 +1,193 @@ +import { randomUUID } from "node:crypto" +import { afterEach, expect, it } from "vitest" +import { Actor, type PayloadBroadcasts } from "../src/actor.js" +import type { RealtimeEnvelope } from "../src/browser/index.js" +import type { InstrumentationEvent } from "../src/configuration.js" +import { sqlite } from "../src/database/sqlite.js" +import { configure, type SolidObjectsRuntime } from "../src/runtime.js" +import { expectPortableEvents } from "./support/portable-telemetry.js" + +class PayloadReader extends Actor { + static override readonly actorType = "payload-projection" + static override readonly payloads = { + impure: (actor, action) => { + switch (action) { + case "state": + actor.items.push("unexpected") + break + case "effect": + actor.emit("unexpected") + break + case "recovery": + actor.requestEffectRecovery({ id: "missing-effect" }) + break + case "commit_action": + actor.commitAction("unexpected") + break + case "reminder": + actor.schedule({ at: new Date(Date.now() + 60_000) }).append() + break + case "outbound": + actor.sendTo(PayloadReader.ref("other")).append() + break + case "raise": + actor.items.push("unexpected") + throw new Error("projection failed") + } + return { items: actor.items } + }, + pure: (actor) => ({ items: actor.items }), + text: (_actor, text) => ({ text }), + } satisfies PayloadBroadcasts + + items: string[] = [] + + append(): void { + this.items.push("committed") + } +} + +let runtime: SolidObjectsRuntime | undefined + +class GeneratedDefaults extends Actor { + static override readonly actorType = "payload-generated-defaults" + static override readonly payloads = { + first: (actor) => ({ token: actor.token }), + second: (actor) => ({ token: actor.token }), + } satisfies PayloadBroadcasts + + token = randomUUID() +} + +afterEach(async () => { + await runtime?.close() + runtime = undefined +}) + +it.each(["state", "effect", "recovery", "commit_action", "reminder", "outbound", "raise"])( + "isolates a payload that performs %s from later projections", + async (action) => { + const events: InstrumentationEvent[] = [] + const configured = await startRuntime({ events }) + await configured.ref(PayloadReader, "one").append() + + const delivered = await subscribe({ + authorizationContext: action, + payloads: ["impure", "pure"], + }) + + expect(delivered).toEqual([ + expect.objectContaining({ name: "pure", payload: { items: ["committed"] } }), + ]) + expect(events).toContainEqual( + expect.objectContaining({ + name: "solid_objects.payload_broadcast.failed", + attributes: expect.objectContaining({ + payload: "impure", + errorName: action === "raise" ? "Error" : "QueryMutatedState", + }), + }), + ) + expectPortableEvents(events, ["payload_broadcast.failed"]) + expect(await configured.ref(PayloadReader, "one").snapshot()).toEqual({ items: ["committed"] }) + }, +) + +it.each([0, -1])( + "enforces the configured UTF-8 payload byte boundary with offset %s", + async (offset) => { + const text = "éé" + const events: InstrumentationEvent[] = [] + await startRuntime({ + events, + maxPayloadBytes: Buffer.byteLength(JSON.stringify({ text })) + offset, + }) + + const delivered = await subscribe({ authorizationContext: text, payloads: ["text"] }) + + if (offset === 0) { + expect(delivered).toEqual([expect.objectContaining({ payload: { text } })]) + return + } + expect(delivered).toEqual([]) + expect(events).toContainEqual( + expect.objectContaining({ + name: "solid_objects.payload_broadcast.failed", + attributes: expect.objectContaining({ errorName: "PayloadTooLarge" }), + }), + ) + }, +) + +it("accepts a configured payload limit above one megabyte", async () => { + const text = "x".repeat(1_048_576) + await startRuntime({ events: [], maxPayloadBytes: Buffer.byteLength(JSON.stringify({ text })) }) + + const delivered = await subscribe({ authorizationContext: text, payloads: ["text"] }) + + expect(delivered).toEqual([expect.objectContaining({ payload: { text } })]) +}) + +it("shares generated defaults across payloads from the same snapshot", async () => { + const configured = await startRuntime({ events: [] }) + configured.register(GeneratedDefaults) + const delivered: RealtimeEnvelope[] = [] + const session = configured.realtime.connect({ + authorizationContext: null, + send: (envelope) => { + delivered.push(envelope) + }, + }) + try { + await session.receive({ + version: 1, + action: "subscribe", + actorType: GeneratedDefaults.actorType, + actorId: "new", + payloads: ["first", "second"], + }) + const payloads = delivered.filter((envelope) => envelope.kind === "payload") + expect(payloads).toHaveLength(2) + expect(payloads[0]!.payload).toEqual(payloads[1]!.payload) + } finally { + session.close() + } +}) + +async function startRuntime(options: { events: InstrumentationEvent[]; maxPayloadBytes?: number }) { + runtime = configure({ + database: sqlite({ path: ":memory:" }), + authorizeMessage: () => true, + authorizeQuery: () => true, + authorizeSubscription: () => true, + instrumentation: (event) => { + options.events.push(event) + }, + ...(options.maxPayloadBytes === undefined ? {} : { maxPayloadBytes: options.maxPayloadBytes }), + }) + runtime.register(PayloadReader) + await runtime.install() + return runtime +} + +async function subscribe(options: { authorizationContext: string; payloads: string[] }) { + const delivered: RealtimeEnvelope[] = [] + const session = runtime!.realtime.connect({ + authorizationContext: options.authorizationContext, + send: (envelope) => { + delivered.push(envelope) + }, + }) + try { + await session.receive({ + version: 1, + action: "subscribe", + actorType: PayloadReader.actorType, + actorId: "one", + payloads: options.payloads, + }) + return delivered.filter((envelope) => envelope.kind === "payload") + } finally { + session.close() + } +} diff --git a/test/polling-loop.test.ts b/test/polling-loop.test.ts index 2b3918d..64649da 100644 --- a/test/polling-loop.test.ts +++ b/test/polling-loop.test.ts @@ -2,6 +2,7 @@ import { afterEach, describe, expect, it, vi } from "vitest" import { sqlite } from "../src/database/sqlite.js" import { createRuntime, type SolidObjectsRuntime } from "../src/runtime.js" import type { InstrumentationEvent } from "../src/configuration.js" +import { expectPortableEvents } from "./support/portable-telemetry.js" import { InProcessWakeUpAdapter, type WakeUpAdapter, @@ -86,6 +87,7 @@ describe("idle polling", () => { currentIntervalMilliseconds: 1_000, }, ]) + expectPortableEvents(events, ["polling.interval_changed"]) }) it("backs an idle effect worker off to the configured ceiling", async () => { diff --git a/test/read-only-projections.test.ts b/test/read-only-projections.test.ts new file mode 100644 index 0000000..dcc3c98 --- /dev/null +++ b/test/read-only-projections.test.ts @@ -0,0 +1,174 @@ +import { expect, it } from "vitest" +import { Actor } from "../src/actor.js" +import { sqlite } from "../src/database/sqlite.js" +import { createRuntime, type SolidObjectsRuntime } from "../src/runtime.js" + +class ReplacingProjection extends Actor { + static override readonly actorType = "replacing-projection" + action = "" + + stage({ action }: { action: string }): void { + this.action = action + if (action === "effect") this.emit("original") + if (action === "commit_action") this.commitAction("original") + } + + override observables() { + if (!this.hasIntents()) return { action: this.action } + + this.discardIntents() + if (this.action === "effect") this.emit("replacement") + if (this.action === "commit_action") this.commitAction("replacement") + return { action: this.action } + } +} + +it.each(["effect", "commit_action"])( + "rejects observable replacement of staged %s without retrying or committing", + async (action) => { + const committed: string[] = [] + const runtime = createRuntime({ + database: sqlite({ path: ":memory:" }), + authorizeMessage: () => true, + authorizeQuery: () => true, + maxAttempts: 3, + retryDelayMilliseconds: () => 0, + }) + try { + runtime.register(ReplacingProjection) + runtime.registerCommitAction("original", () => { + committed.push("original") + }) + runtime.registerCommitAction("replacement", () => { + committed.push("replacement") + }) + await runtime.install() + const reference = runtime.ref(ReplacingProjection, "one") + const message = await reference.send.stage({ action }) + await runtime.worker().runUntilIdle() + + expect(await message.outcome()).toMatchObject({ + status: "dead", + error: { name: "QueryMutatedState" }, + attempts: 1, + }) + expect((await reference.snapshot()).action).toBe("") + expect(committed).toEqual([]) + } finally { + await runtime.close() + } + }, +) + +const purityActions = ["effect", "recovery", "commit_action", "reminder", "outbound", "state"] + +class ReadOnlyReader extends Actor { + static override readonly actorType = "read-only-reader" + static queryAction: string | undefined + static projectionAction: string | undefined + items: string[] = [] + + get read(): string[] { + performPurityAction({ actor: this, action: ReadOnlyReader.queryAction }) + return this.items + } + + append(): string[] { + this.items.push("committed") + return this.items + } + + override observables() { + performPurityAction({ actor: this, action: ReadOnlyReader.projectionAction }) + return { projected: this.items } + } +} + +function performPurityAction(options: { actor: ReadOnlyReader; action: string | undefined }): void { + const { actor, action } = options + if (action === "effect") actor.emit("unexpected") + if (action === "recovery") actor.requestEffectRecovery({ id: "missing-effect" }) + if (action === "commit_action") actor.commitAction("unexpected") + if (action === "reminder") actor.schedule({ at: new Date(Date.now() + 60_000) }).append() + if (action === "outbound") actor.sendTo(ReadOnlyReader.ref("other")).append() + if (action === "state") actor.items.push("unexpected") +} + +function purityRuntime(committed: string[]): SolidObjectsRuntime { + const runtime = createRuntime({ + database: sqlite({ path: ":memory:" }), + authorizeMessage: () => true, + authorizeQuery: () => true, + maxAttempts: 3, + retryDelayMilliseconds: () => 0, + }) + runtime.register(ReadOnlyReader) + runtime.registerCommitAction("unexpected", () => { + committed.push("unexpected") + }) + return runtime +} + +async function expectNoCommittedWork(options: { + runtime: SolidObjectsRuntime + committed: readonly string[] +}): Promise { + const { runtime, committed } = options + const rows = (table: string) => + runtime.settings.database.connection((connection) => + connection.all<{ attempt_count?: number | bigint }>( + `SELECT * FROM ${runtime.repository.table(table)}`, + ), + ) + expect((await rows("messages")).map((message) => Number(message.attempt_count))).toEqual([1]) + for (const table of ["effects", "effect_recoveries", "reminders", "broadcasts"]) + expect(await rows(table)).toEqual([]) + expect(committed).toEqual([]) + ReadOnlyReader.queryAction = undefined + ReadOnlyReader.projectionAction = undefined + expect((await runtime.ref(ReadOnlyReader, "one").snapshot()).items).toEqual([]) +} + +it.each(purityActions)( + "rejects a query that stages %s without retrying or committing", + async (action) => { + const committed: string[] = [] + const runtime = purityRuntime(committed) + try { + await runtime.install() + ReadOnlyReader.queryAction = action + + await expect(runtime.ref(ReadOnlyReader, "one").read).rejects.toMatchObject({ + details: { name: "QueryMutatedState" }, + }) + await expectNoCommittedWork({ runtime, committed }) + } finally { + ReadOnlyReader.queryAction = undefined + await runtime.close() + } + }, +) + +it.each(purityActions)( + "rejects an observable that stages %s without retrying or committing", + async (action) => { + const committed: string[] = [] + const runtime = purityRuntime(committed) + try { + await runtime.install() + ReadOnlyReader.projectionAction = action + const message = await runtime.ref(ReadOnlyReader, "one").send.append() + await runtime.worker().runUntilIdle() + + expect(await message.outcome()).toMatchObject({ + status: "dead", + error: { name: "QueryMutatedState" }, + attempts: 1, + }) + await expectNoCommittedWork({ runtime, committed }) + } finally { + ReadOnlyReader.projectionAction = undefined + await runtime.close() + } + }, +) diff --git a/test/reminder-cancellation.test.ts b/test/reminder-cancellation.test.ts index c579481..e09a641 100644 --- a/test/reminder-cancellation.test.ts +++ b/test/reminder-cancellation.test.ts @@ -389,7 +389,7 @@ describe("reminder cancellation", () => { expect(claimed).toBeDefined() await reference.convertByName() - await expect(started.repository.enqueueReminder(claimed!)).resolves.toBe(false) + await expect(started.repository.enqueueReminder(claimed!)).resolves.toBeUndefined() expect(await started.reminderScheduler().runOnce()).toBe(0) }) diff --git a/test/serialization-compatibility.test.ts b/test/serialization-compatibility.test.ts new file mode 100644 index 0000000..3971bb2 --- /dev/null +++ b/test/serialization-compatibility.test.ts @@ -0,0 +1,59 @@ +import { readFileSync } from "node:fs" +import { expect, it } from "vitest" +import { Actor } from "../src/actor.js" +import { sqlite } from "../src/database/sqlite.js" +import { PayloadTooLarge } from "../src/errors.js" +import { createRuntime } from "../src/runtime.js" +import { normalizeJson, readonlyCopy, stableJson } from "../src/serialization.js" +import type { JsonObject } from "../src/types.js" + +const fixtures: { name: string; value: JsonObject }[] = JSON.parse( + readFileSync(new URL("../compatibility/json-values.json", import.meta.url), "utf8"), +) + +for (const { name, value } of fixtures) { + it(`preserves ${name} as JSON data`, () => { + const normalized = normalizeJson(value) + expect(JSON.stringify(normalized)).toBe(JSON.stringify(value)) + expect(Object.getPrototypeOf(normalized)).toBe(Object.prototype) + expect(readonlyCopy(value)).toStrictEqual(value) + expect(JSON.parse(stableJson(value))).toStrictEqual(value) + }) +} + +it("includes reserved keys in the encoded byte limit", () => { + const value = JSON.parse('{"__proto__":"long payload"}') + expect(() => normalizeJson(value, { maxBytes: 2 })).toThrow(PayloadTooLarge) +}) + +class JsonActor extends Actor { + static override readonly actorType = "json-compatibility" + payload: JsonObject = {} + + store({ payload }: { payload: JsonObject }): JsonObject { + this.payload = payload + return payload + } +} + +it("preserves reserved keys through actor arguments, state, and results", async () => { + const runtime = createRuntime({ + database: sqlite({ path: ":memory:" }), + wakeUp: "in_process", + authorizeMessage: () => true, + authorizeQuery: () => true, + }) + try { + await runtime.install() + runtime.register(JsonActor) + const reference = runtime.ref(JsonActor, "one") + for (const { name, value } of fixtures) { + const message = await reference.send.with({ idempotencyKey: name }).store({ payload: value }) + await runtime.worker().runUntilIdle() + expect(await message.result()).toStrictEqual(value) + expect((await reference.snapshot()).payload).toStrictEqual(value) + } + } finally { + await runtime.close() + } +}) diff --git a/test/support/portable-telemetry.ts b/test/support/portable-telemetry.ts new file mode 100644 index 0000000..2e549c4 --- /dev/null +++ b/test/support/portable-telemetry.ts @@ -0,0 +1,43 @@ +import { readFileSync } from "node:fs" +import { expect } from "vitest" +import type { InstrumentationEvent } from "../../src/configuration.js" + +interface PortableEventContract { + readonly name: string + readonly match?: Readonly> + readonly attributes: readonly string[] + readonly javascriptAttributes?: readonly string[] +} + +export const telemetryContract: { + readonly attributes: readonly string[] + readonly events: readonly PortableEventContract[] +} = JSON.parse( + readFileSync(new URL("../../compatibility/telemetry-events.json", import.meta.url), "utf8"), +) + +export function expectPortableAttributes(event: InstrumentationEvent): void { + const name = event.name.replace(/^solid_objects\./, "") + const contract = telemetryContract.events.find( + (entry) => + entry.name === name && + Object.entries(entry.match ?? {}).every(([key, value]) => event.attributes[key] === value), + ) + + expect(contract, `${name} has no portable attribute contract`).toBeDefined() + expect(Object.keys(event.attributes).sort(), `${name} attributes`).toEqual( + [...(contract?.attributes ?? []), ...(contract?.javascriptAttributes ?? [])].sort(), + ) +} + +export function expectPortableEvents( + events: readonly InstrumentationEvent[], + names: readonly string[], +): void { + for (const name of names) { + const matching = events.filter((event) => event.name === `solid_objects.${name}`) + + expect(matching, `expected a solid_objects.${name} event`).not.toHaveLength(0) + for (const event of matching) expectPortableAttributes(event) + } +} diff --git a/test/sync-timeout.test.ts b/test/sync-timeout.test.ts index bd103a4..0b6433b 100644 --- a/test/sync-timeout.test.ts +++ b/test/sync-timeout.test.ts @@ -4,6 +4,7 @@ import type { InstrumentationEvent } from "../src/configuration.js" import { sqlite } from "../src/database/sqlite.js" import { SyncEnqueueTimeout, SyncTimeout } from "../src/errors.js" import { configure, type SolidObjectsRuntime } from "../src/runtime.js" +import { expectPortableEvents } from "./support/portable-telemetry.js" class TimeoutActor extends Actor { static override readonly actorType = "TimeoutActor" @@ -69,6 +70,7 @@ describe("synchronous timeout diagnostics", () => { messageId: error.details.messageId, waitingOn: "activationHeld", activationOwnerId: "blocking-worker", + activationGeneration: String(error.details.activation.generation), }), }), ]), @@ -85,7 +87,8 @@ describe("synchronous timeout diagnostics", () => { }) it("distinguishes a deadline before durable enqueue", async () => { - runtime = configuredRuntime() + const events: InstrumentationEvent[] = [] + runtime = configuredRuntime({ instrumentation: (event) => events.push(event) }) await runtime.install() const enqueue = runtime.repository.enqueue.bind(runtime.repository) vi.spyOn(runtime.repository, "enqueue").mockImplementation(async (input) => { @@ -101,6 +104,11 @@ describe("synchronous timeout diagnostics", () => { } expect(error).toBeInstanceOf(SyncEnqueueTimeout) + expectPortableEvents(events, ["sync.enqueue_timeout"]) + expect( + events.find((event) => event.name === "solid_objects.sync.enqueue_timeout")?.attributes + .timeoutMilliseconds, + ).toBe(1) const messages = await runtime.settings.database.connection((connection) => connection.all(`SELECT id FROM ${runtime?.repository.table("messages")}`), ) diff --git a/test/telemetry-compatibility.test.ts b/test/telemetry-compatibility.test.ts new file mode 100644 index 0000000..0807875 --- /dev/null +++ b/test/telemetry-compatibility.test.ts @@ -0,0 +1,43 @@ +import { readFileSync } from "node:fs" +import { expect, it } from "vitest" +import { portableAttributes, telemetryEvent } from "../src/telemetry.js" +import { telemetryContract } from "./support/portable-telemetry.js" + +const fixtures: { rubyReason: string; waitingOn: string }[] = JSON.parse( + readFileSync(new URL("../compatibility/sync-timeout.json", import.meta.url), "utf8"), +) + +it("matches the shared portable attribute allowlist", () => { + expect([...portableAttributes].sort()).toEqual([...telemetryContract.attributes].sort()) +}) + +it.each(fixtures)("preserves portable $waitingOn timeout diagnostics", ({ waitingOn }) => { + const event = telemetryEvent({ + name: "sync.timeout", + adapter: "sqlite", + attributes: { + waitingOn, + activationOwnerId: "worker-1", + activationGeneration: "7", + arguments: { secret: "private" }, + }, + }) + + expect(event.attributes).toEqual({ + waitingOn, + activationOwnerId: "worker-1", + activationGeneration: "7", + }) +}) + +it("preserves unknown activation fields during database contention", () => { + const attributes = { + waitingOn: "databaseContention", + activationOwnerId: null, + activationGeneration: null, + } + + expect( + telemetryEvent({ name: "sync.timeout", adapter: "sqlite", attributes }).attributes, + ).toEqual(attributes) +}) diff --git a/test/transmit-ordering.test.ts b/test/transmit-ordering.test.ts new file mode 100644 index 0000000..b9b9fa5 --- /dev/null +++ b/test/transmit-ordering.test.ts @@ -0,0 +1,119 @@ +import { afterEach, expect, it, vi } from "vitest" +import { Actor } from "../src/actor.js" +import { sqlite } from "../src/database/sqlite.js" +import { postgresql } from "../src/database/postgresql.js" +import { mysql } from "../src/database/mysql.js" +import { createRuntime, type SolidObjectsRuntime } from "../src/runtime.js" +import { receiveTransmitEnvelope, registerTransmit, TRANSMIT_EFFECT } from "../src/transmit.js" + +class Sender extends Actor { + static override readonly actorType = "OrderedTransmit" + + stage(): void { + const identifiers = vi + .spyOn(globalThis.crypto, "randomUUID") + .mockReturnValueOnce("ffffffff-ffff-4fff-8fff-ffffffffffff") + .mockReturnValueOnce("00000000-0000-4000-8000-000000000000") + try { + this.emit(TRANSMIT_EFFECT, { arguments: { operation: "append", arguments: { value: 1 } } }) + this.emit(TRANSMIT_EFFECT, { arguments: { operation: "append", arguments: { value: 2 } } }) + } finally { + identifiers.mockRestore() + } + } + + stageFluent(): void { + this.transmit().append({ value: 1 }) + this.transmit().append({ value: 2 }) + } + + append(_arguments: { value: number }): void {} +} + +class Receiver extends Actor { + static override readonly actorType = "OrderedTransmit" + values: number[] = [] + + append({ value }: { value: number }): void { + this.values.push(value) + } +} + +const runtimes: SolidObjectsRuntime[] = [] + +afterEach(async () => { + for (const runtime of runtimes) { + await runtime.repository.resetForTesting() + await runtime.close() + } + runtimes.length = 0 +}) + +async function start(prefix: string): Promise { + const connectionString = process.env.SOLID_OBJECTS_DATABASE_URL + const database = connectionString?.startsWith("postgresql:") + ? postgresql({ connectionString }) + : connectionString?.startsWith("mysql:") + ? mysql({ connectionString }) + : sqlite({ path: ":memory:" }) + const runtime = createRuntime({ + database, + tableNamePrefix: prefix, + wakeUp: "in_process", + authorizeMessage: () => true, + authorizeQuery: () => true, + retryDelayMilliseconds: () => 0, + logger: { debug() {}, info() {}, warn() {}, error() {} }, + }) + runtimes.push(runtime) + await runtime.install() + return runtime +} + +it("preserves staging order within a message through delivery retries", async () => { + const sender = await start("transmit_order_sender_") + const receiver = await start("transmit_order_receiver_") + sender.register(Sender) + receiver.register(Receiver) + let failures = 1 + registerTransmit({ + runtime: sender, + deliver: async (envelope) => { + if (envelope.arguments.value === 1 && failures-- > 0) throw new Error("offline") + await receiveTransmitEnvelope({ runtime: receiver, envelope }) + }, + }) + await sender.ref(Sender, "one").stage() + await sender.testing.drain({ roles: ["effects"], maxPasses: 20 }) + await receiver.testing.drain({ roles: ["actors"] }) + expect(await receiver.ref(Receiver, "one").snapshot()).toEqual({ values: [1, 2] }) +}) + +for (const migrationState of ["absent", "interrupted"] as const) { + it(`upgrades ${migrationState} effect positions without changing pending work`, async () => { + const runtime = await start("transmit_order_upgrade_") + runtime.register(Sender) + await runtime.ref(Sender, "legacy").stageFluent() + const effects = runtime.repository.table("effects") + const migrations = runtime.repository.table("schema_migrations") + await runtime.settings.database.connection(async (connection) => { + await connection.run(`DELETE FROM ${migrations} WHERE version = 14`) + if (migrationState === "absent") + await connection.run(`ALTER TABLE ${effects} DROP COLUMN position`) + }) + await runtime.install() + await runtime.install() + await runtime.ref(Sender, "new").stageFluent() + const rows = await runtime.settings.database.connection((connection) => + connection.all<{ actor_id: string; position: number | bigint; status: string }>( + `SELECT instances.actor_id, effects.position, effects.status FROM ${effects} effects JOIN ${runtime.repository.table("instances")} instances ON instances.id = effects.instance_id ORDER BY instances.actor_id, effects.position`, + ), + ) + expect(rows.map((row) => [row.actor_id, Number(row.position), row.status])).toEqual([ + ["legacy", 0, "pending"], + ["legacy", migrationState === "absent" ? 0 : 1, "pending"], + ["new", 0, "pending"], + ["new", 1, "pending"], + ]) + }) +} diff --git a/test/wake-up.test.ts b/test/wake-up.test.ts index ae3eedc..309264f 100644 --- a/test/wake-up.test.ts +++ b/test/wake-up.test.ts @@ -1,3 +1,4 @@ +import { notifyWakeUp } from "../src/wake-up-notification.js" import { afterEach, describe, expect, it, vi } from "vitest" import { Actor } from "../src/actor.js" import { sqlite } from "../src/database/sqlite.js" @@ -237,3 +238,22 @@ async function eventually(condition: () => boolean | Promise): Promise< } throw new Error("condition was not met") } + +it("isolates notification and logger failures", async () => { + const failure = () => { + throw new Error("private failure") + } + const adapter = { + notify: () => Promise.reject(new Error("private notification")), + watch: () => ({ wait: async () => false }), + close: async () => {}, + } + expect(() => + notifyWakeUp({ + adapter, + logger: { debug: failure, info: failure, warn: failure, error: failure }, + role: "actors", + }), + ).not.toThrow() + await Promise.resolve() +})