From 465da0dabaa3f27ebf5c1142b2af552bc883636e Mon Sep 17 00:00:00 2001 From: Sree Narayanan Date: Wed, 2 Sep 2026 02:40:32 +0400 Subject: [PATCH 1/6] docs(rivetkit): add actor tracing docs --- .claude/reference/testing.md | 1 + CLAUDE.md | 1 + docs-internal/engine/rivetkit-telemetry.md | 124 +++++ docs/content/docs/general/tracing.mdx | 188 +++++++ docs/sidebar.json | 4 + .../docs/general-tracing/application-span.ts | 23 + examples/docs/general-tracing/sdk.ts | 4 + examples/docs/package.json | 2 + pnpm-lock.yaml | 517 ++++++++++++++++++ 9 files changed, 864 insertions(+) create mode 100644 docs-internal/engine/rivetkit-telemetry.md create mode 100644 docs/content/docs/general/tracing.mdx create mode 100644 examples/docs/general-tracing/application-span.ts create mode 100644 examples/docs/general-tracing/sdk.ts diff --git a/.claude/reference/testing.md b/.claude/reference/testing.md index 32a03f3380..8f34506212 100644 --- a/.claude/reference/testing.md +++ b/.claude/reference/testing.md @@ -47,6 +47,7 @@ For RivetKit runtime or parity bugs, use `rivetkit-typescript/packages/rivetkit` - Keep RivetKit test fixtures scoped to the engine-only runtime. - Prefer targeted integration tests under `rivetkit-typescript/packages/rivetkit/tests/` over shared multi-driver matrices. +- A span and its parent can arrive in different OTLP export batches, so a trace test that waits for the child and then asserts its `parentSpanId` is racy. Wait on a predicate over the whole exported span list until both are present, then assert the relationship. ## Frontend testing diff --git a/CLAUDE.md b/CLAUDE.md index 93cad1e4c6..c0679ee549 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -394,6 +394,7 @@ Load these only when the task touches the topic. - **[SQLite VFS parity](docs-internal/engine/sqlite-vfs.md)** — native Rust VFS ↔ WASM TypeScript VFS 1:1 parity rule, v2 storage keys, chunk layout, delete/truncate strategy. Read before touching either VFS. - **[SQLite optimizations](docs-internal/engine/SQLITE_OPTIMIZATIONS.md)** — brief tracker for SQLite cold-read, VFS, storage, preload, and benchmark optimization ideas. - **[TLS trust roots](docs-internal/engine/tls-trust-roots.md)** — rustls native+webpki union rationale, which clients use which backend. +- **[RivetKit telemetry](docs-internal/engine/rivetkit-telemetry.md)** — Core-owned invocation and SQLite spans, ray semantics, persisted trace origins, native OTLP export. Read before touching actor tracing or log correlation. - **[Sleep sequence](docs-internal/engine/sleep-sequence.md)** — engine lifecycle authority, `keepAwake` vs `waitUntil` semantics, grace deadline shutdown-token abort, `can_arm_sleep_timer` vs `can_finalize_sleep` predicates. Read before touching sleep/destroy lifecycle. ### Agent procedural (`.claude/reference/`) diff --git a/docs-internal/engine/rivetkit-telemetry.md b/docs-internal/engine/rivetkit-telemetry.md new file mode 100644 index 0000000000..d060b6f775 --- /dev/null +++ b/docs-internal/engine/rivetkit-telemetry.md @@ -0,0 +1,124 @@ +# RivetKit telemetry + +Internal reference for actor traces and log correlation. Core owns telemetry behavior. Runtime adapters activate Core's context in the host language. + +See [NAPI bridge](napi-bridge.md) for binding conventions and [Core internals](rivetkit-core-internals.md) for actor dispatch and lifecycle wiring. + +## Ownership + +Telemetry crosses two runtimes but keeps one trace: + +```text +client + headers: ray ID + W3C trace context + → rivetkit-core + invocation spans + native operation spans + → NAPI + activates the invocation context in TypeScript + → application spans +``` + +- `rivetkit-core::telemetry` owns spans and completion +- NAPI only translates context between Core and TypeScript +- Rust and TypeScript export through their own OpenTelemetry SDKs + +Core spans use the `rivetkit::telemetry` tracing target. Log layers exclude this target, and the export layer excludes unrelated diagnostic spans. + +## Telemetry surfaces + +| Work | Span | Kind | Relationship | +| --- | --- | --- | --- | +| Action | `{actor}/{action}` | `server` | Child of incoming context | +| Schedule | `{actor}/{action}` | `internal` | New trace linked to its origin | +| Raw HTTP | `{actor}/onRequest` | `server` | Child of incoming context | +| Queue send | `{actor}/queue.send` | `producer` | Child of incoming context | +| Queue receive | `{actor}/queue.receive` | `consumer` | Linked to the send origin | +| Actor call | `{callee}/{action}` | `client` | Child of application or invocation span | +| SQLite | `rivet.sqlite.{operation}` | `internal` | Child of application or invocation span | + +Core records these attributes: + +| Scope | Attributes | +| --- | --- | +| Actor | `rivet.actor.id`, `rivet.actor.name`, `rivet.actor.key` | +| Correlation | `rivet.ray.id` | +| Invocation | `rivet.invocation.type`, `otel.status_code`, `error.type` | +| Action | `rivet.action.name` | +| HTTP | `http.request.method`, `http.response.status_code` | +| Queue | `rivet.queue.name` | +| SQLite | `rivet.operation.system`, `rivet.operation.name` | + +Raw HTTP spans use `onRequest`, never the request path. Handler errors use their `group.code` as `error.type`. A 5xx response uses the status code. Abandoned SQLite and actor-call tracking uses `actor.operation_abandoned` to represent an unknown outcome. + +## Logs + +Actor loggers include actor identity and ray ID. A valid span also adds `traceId` and `spanId`. Do not retain an invocation logger for unrelated work. + +## Context propagation + +```text +incoming headers + x-rivet-ray-id ───────────────────────────────┐ + traceparent + tracestate ──→ invocation span ├─→ actor call / HTTP / queue + │ +application span ──────────────────────────────┘ preferred parent + +schedule or queue send ── stores origin ── later execution links to origin +``` + +Core accepts correlation headers on actions, raw HTTP requests, and queue sends. + +- Ray IDs match `[A-Za-z0-9_-]` and contain 1–30 characters +- RivetKit propagates ray IDs but does not create them +- Invalid W3C context starts a root span without rejecting the request +- Explicit HTTP trace headers override generated headers as one pair +- Per-call context overrides static client telemetry headers +- Application spans take precedence over the invocation span as outbound parents + +The Engine gateway supplies its guard ray ID when the caller sends none. External clients read `rivet.ray.id` from OpenTelemetry baggage; a configured `x-rivet-ray-id` header is the fallback in both clients. + +The Engine keeps two ray values per request and never merges them: + +- `ray_id` is the guard's own `Id`. The Engine uses it for everything it does internally: `ctx.with_ray`, `api-public`, gasoline operations, and workflow and signal records. A caller cannot choose it. +- `external_ray_id` is the caller's validated `x-rivet-ray-id`, or the guard's `Id` as a string when the caller sends none. The guard forwards it to the actor, returns it in the response header, and uses it in WebSocket close frames. + +The `guard_request` span and the guard access logs record both fields, so either value leads to the other. The caller's value stays a string because an Engine `Id` embeds a datacenter label that a client cannot produce. + +Rust owns ray ID validation in `rivetkit-client-protocol::telemetry_headers`. Core and the Rust client use `TraceContextPropagator`. TypeScript handles context in `common/otel-context.ts`. + +Caller trace context provides correlation, not identity or authorization. Strip incoming correlation headers at an untrusted boundary when callers must not select these values. + +## Invocation lifetime + +```text +dispatch + → start span and timer + → run handler with context + → send reply + → wait for tracked waitUntil work + → end span +``` + +`ActorInvocation` completes once. A reply dropped unsent finishes the span with `otel.status_code` `ERROR` and `error.type` `actor.dropped_reply`. A dispatch the actor task refuses is answered before any span exists. `c.keepAwake` remains part of the handler, while `waitUntil` may extend the span beyond the reply. + +Schedules and queue messages store their ray ID and W3C trace context in the `ray_id`, `traceparent`, and `tracestate` columns of their own rows, added by internal schema migration v2. Each schedule fire starts a new linked trace. Each queue receipt links to the send origin. + +## Native export + +The `native-runtime` feature enables Core's exporter. Hosts attach `telemetry::export::layer()` and call `shutdown_best_effort()` during shutdown. + +- An OTLP endpoint enables export +- Protocol selection supports `grpc`, `http/protobuf`, and `http/json` +- Export uses a bounded background queue and never fails actor work +- NAPI forwards OpenTelemetry SDK warnings, such as dropped spans, to the JavaScript logger through `setTelemetryLogSink`. `createRegistry()` installs the sink and `shutdownTelemetry()` releases it. While a sink is installed the Rust log layers skip those warnings, and without one they print through the Rust log layers, so each warning is visible once + +## Data policy + +Record actor identity, invocation type, HTTP method and status, correlation IDs, operation names, and error identity. Do not record arguments, results, connection parameters, SQL text or bindings, actor state, arbitrary headers, or raw error messages. + +## Gaps + +- WebSocket handlers, lifecycle hooks, connection callbacks, KV, and actor-state operations have no dedicated spans +- WebSocket action messages and inspector actions do not inherit caller context +- Actor creation ray IDs do not reach the actor runtime +- Wasm does not export host spans or expose Core invocation context to TypeScript diff --git a/docs/content/docs/general/tracing.mdx b/docs/content/docs/general/tracing.mdx new file mode 100644 index 0000000000..77983be2a4 --- /dev/null +++ b/docs/content/docs/general/tracing.mdx @@ -0,0 +1,188 @@ +--- +title: "OpenTelemetry" +description: "Export traces from Rivet Actors, add application spans, and follow work across actors with ray IDs." +skill: true +--- + +Rivet gives you automatic instrumentation of actor actions, HTTP handlers, actor-to-actor calls, SQLite operations, queues, and scheduled actions. + +This helps you investigate anything that looks slow, walk through complex actor-to-actor flows, and get end-to-end context into the life of your actor. + +
+ + + + + + YOUR ACTOR PROCESS + + + RivetKit + automatic actor spans + + + Your code + your application spans + + + + + + + OTEL_EXPORTER_OTLP_TRACES_ENDPOINT + + + OTLP collector + Jaeger, Tempo, Honeycomb + + +
+ +## Quickstart + +Set an OTLP exporter within the process that runs your actors, next to `RIVET_ENDPOINT`, and redeploy: + +```bash +OTEL_SERVICE_NAME=my-rad-app +OTEL_EXPORTER_OTLP_TRACES_ENDPOINT=http://collector:4318/v1/traces +OTEL_EXPORTER_OTLP_TRACES_PROTOCOL=http/protobuf +``` + +`OTEL_EXPORTER_OTLP_TRACES_ENDPOINT` acts as a feature flag; setting it enables tracing. See [Configure the exporter](#configure-the-exporter) for more options. + +You should now see spans named `{actor}/{action}` in your observability platform under `OTEL_SERVICE_NAME`. + +## What gets traced + +Each span shows how long an operation took and whether it failed. RivetKit records: + +- Actions and HTTP handlers, with the operations they perform shown underneath +- Calls through `c.client()`, including time spent routing, waking the other actor, and retrying +- SQLite operations, without recording SQL text, bindings, or results +- Queue sends and receipts, with a link from each receipt to its sender +- Scheduled actions, in a new trace linked to the work that scheduled them + +The full list of span names and attributes lives in [`telemetry.rs`](https://github.com/rivet-dev/rivet/blob/main/rivetkit-rust/packages/rivetkit-core/src/telemetry.rs), and the SQLite operations in [`sqlite/mod.rs`](https://github.com/rivet-dev/rivet/blob/main/rivetkit-rust/packages/rivetkit-core/src/actor/sqlite/mod.rs). + +## Add application spans + +You can trace your own actor code with `@opentelemetry/api`. Your spans will nest with RivetKit spans in the same trace. + +Here's a quick guide on how to set this up: + + + + +```bash +npm install @opentelemetry/api @opentelemetry/sdk-node +``` + + + + +Start the SDK before your RivetKit registry: + + + + + + +Wrap the code you want to measure in `startActiveSpan()`: + + + + + + +Calling `query` produces this trace: + +``` +counter/query +└── counter.query application span + └── rivet.sqlite.execute traced by RivetKit +``` + + +Use a fixed span name, such as `generate_response`. Avoid building the name from changing values, like `generate_response_${requestId}`. Put those values in span attributes instead. + + +## Follow work across actors with ray IDs + +A ray ID follows a request across actor calls and scheduled work. Rivet supplies one automatically, or you can provide your own. + +Each blue box shows an actor running an action. RivetKit records each execution as a **span**, with its duration and outcome. + +
+ + + + + + ONE RAY ID + + + Caller + + + x-rivet-ray-id + optional + + + Rivet + + + + + Action: appMention + Actor: slackThread + + + c.client() + + + Action: summarize + Actor: summarizer + + + c.schedule.after(60s) + + + Action: postDigest + Actor: slackThread + + Same ray ID across this work + ray ID: my-rivet-ray-id + + +
+ +- Search for `rivet.ray.id` in traces or `rayId` in [actor logs](/actors/docs/general/logging) to find related work. +- To use your own ID, set `rivet.ray.id` in [OpenTelemetry baggage](https://opentelemetry.io/docs/concepts/signals/baggage/) for client calls, or send `x-rivet-ray-id` on raw HTTP requests. + + +A ray ID must be 1–30 characters long and may contain only letters, digits, hyphens (`-`), and underscores (`_`). Invalid values are ignored. Use ray IDs to find related traces, not to authenticate requests. + + +## Configure the exporter + +| Variable | Purpose | +| --- | --- | +| `OTEL_SERVICE_NAME` | Service name reported with spans | +| `OTEL_EXPORTER_OTLP_TRACES_ENDPOINT` | Trace collector endpoint | +| `OTEL_EXPORTER_OTLP_TRACES_PROTOCOL` | `http/protobuf`, `http/json`, or `grpc` | +| `OTEL_EXPORTER_OTLP_HEADERS` | Headers sent with every export, such as auth | +| `OTEL_TRACES_SAMPLER`, `OTEL_TRACES_SAMPLER_ARG` | Sampling | +| `OTEL_SDK_DISABLED=true` | Turn telemetry off | + +For more settings, see the [OTLP exporter configuration](https://opentelemetry.io/docs/languages/sdk-configuration/otlp-exporter/) and [SDK configuration](https://opentelemetry.io/docs/languages/sdk-configuration/general/). + +RivetKit sends spans in the background. Actor requests continue if the collector is slow or unavailable. Spans can be lost when exports fail or the buffer fills; export failures and dropped spans are reported in warnings. + +## Known limitations + +- WebSocket handlers, lifecycle hooks, connection callbacks, KV operations, and actor state operations do not have dedicated spans +- Actions called over `.connect()` start a new trace instead of joining the caller's trace +- Effect spans are not traced. Use `@effect/opentelemetry` +- RivetKit does not add `rivet.ray.id` to your application spans +- The Wasm runtime is not traced +- A hard process exit can lose spans still waiting in the export queue diff --git a/docs/sidebar.json b/docs/sidebar.json index 13992a1f4c..41fe7a2f7c 100644 --- a/docs/sidebar.json +++ b/docs/sidebar.json @@ -180,6 +180,10 @@ "title": "Logging", "href": "/actors/docs/general/logging" }, + { + "title": "OpenTelemetry", + "href": "/actors/docs/general/tracing" + }, { "title": "Errors", "href": "/actors/docs/errors" diff --git a/examples/docs/general-tracing/application-span.ts b/examples/docs/general-tracing/application-span.ts new file mode 100644 index 0000000000..acc960320a --- /dev/null +++ b/examples/docs/general-tracing/application-span.ts @@ -0,0 +1,23 @@ +import { SpanStatusCode, trace } from "@opentelemetry/api"; +import { actor } from "rivetkit"; +import { db } from "rivetkit/db"; + +const tracer = trace.getTracer("counter"); + +export const counter = actor({ + state: {}, + db: db(), + actions: { + query: async (c) => + tracer.startActiveSpan("counter.query", async (span) => { + try { + return await c.db.execute("SELECT 1 AS value"); + } catch (error) { + span.setStatus({ code: SpanStatusCode.ERROR }); + throw error; + } finally { + span.end(); + } + }), + }, +}); diff --git a/examples/docs/general-tracing/sdk.ts b/examples/docs/general-tracing/sdk.ts new file mode 100644 index 0000000000..1528d6b5fb --- /dev/null +++ b/examples/docs/general-tracing/sdk.ts @@ -0,0 +1,4 @@ +import { NodeSDK } from "@opentelemetry/sdk-node"; + +const sdk = new NodeSDK(); +sdk.start(); diff --git a/examples/docs/package.json b/examples/docs/package.json index 769324ddf4..d38b92a77b 100644 --- a/examples/docs/package.json +++ b/examples/docs/package.json @@ -27,6 +27,8 @@ "pg": "^8.14.1", "vitest": "^3.0.9", "pino": "^9.7.0", + "@opentelemetry/api": "^1.1.0", + "@opentelemetry/sdk-node": "^0.211.0", "react": "^19.2.0", "react-dom": "^19.2.0" }, diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index abd0c91ef9..0d035b8640 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -849,6 +849,12 @@ importers: '@hono/node-server': specifier: ^1.14.1 version: 1.19.13(hono@4.11.9) + '@opentelemetry/api': + specifier: ^1.1.0 + version: 1.9.0 + '@opentelemetry/sdk-node': + specifier: ^0.211.0 + version: 0.211.0(@opentelemetry/api@1.9.0) '@rivetkit/cloudflare-workers': specifier: workspace:* version: link:../../rivetkit-typescript/packages/cloudflare-workers @@ -6335,6 +6341,15 @@ packages: '@modelcontextprotocol/sdk': optional: true + '@grpc/grpc-js@1.14.4': + resolution: {integrity: sha512-k9Dj3DV/itK9D06Y8f190Qgop7/Ui+D0njFV3LHMPwPT75DpXLQohE9Wmz0QElrJnzsjB7KPWiKJbOl7IPDArQ==} + engines: {node: '>=12.10.0'} + + '@grpc/proto-loader@0.8.1': + resolution: {integrity: sha512-wtF6h+DY6M3YaDBPAmvuuA6jV8Sif9MjtOI5euKFWRgCDl5PeDpPsHR9u2l6St5ceY8AZgoNDww5+HvEsXFsGg==} + engines: {node: '>=6'} + hasBin: true + '@hey-api/client-fetch@0.5.7': resolution: {integrity: sha512-hLpID6NCs8+stbz935UyvyGOXY44oLBSOy7ZEpwXxj977A/0U41iihDQllDoCJrxtbe06DnDgwPOn6/xnRJ71w==} deprecated: Starting with v0.73.0, this package is bundled directly inside @hey-api/openapi-ts. @@ -6935,6 +6950,9 @@ packages: '@jridgewell/trace-mapping@0.3.9': resolution: {integrity: sha512-3Belt6tdc8bPgAtbcmdtNJlirVoTmEb5e2gC94PnkwEW9jI6CAHUeoG85tjWP5WquqfavoMtMwiG4P926ZKKuQ==} + '@js-sdsl/ordered-map@4.4.2': + resolution: {integrity: sha512-iUKgm52T8HOE/makSxjqoWhe95ZJA1/G1sYsGev2JDKUSS14KAgg1LHb+Ba+IPow0xflbnSkOsZcO08C7w1gYw==} + '@keyv/serialize@1.1.1': resolution: {integrity: sha512-dXn3FZhPv0US+7dtJsIi2R+c7qWYiReoEh5zUntWCf4oSpMNib8FDhSoed6m3QyZdx5hK7iLFkYk3rNxwt8vTA==} @@ -7515,40 +7533,200 @@ packages: '@open-draft/until@2.1.0': resolution: {integrity: sha512-U69T3ItWHvLwGg5eJ0n3I62nWuE6ilHlmz7zM0npLBRvPRd7e6NYmg54vvRtP5mZG7kZqZCFVdsTWo7BPtBujg==} + '@opentelemetry/api-logs@0.211.0': + resolution: {integrity: sha512-swFdZq8MCdmdR22jTVGQDhwqDzcI4M10nhjXkLr1EsIzXgZBqm4ZlmmcWsg3TSNf+3mzgOiqveXmBLZuDi2Lgg==} + engines: {node: '>=8.0.0'} + '@opentelemetry/api@1.9.0': resolution: {integrity: sha512-3giAOQvZiH5F9bMlMiv8+GSPMeqg0dbaeo58/0SlA9sxSqZhnUtxzX9/2FzyhS9sWQf5S0GJE0AKBrFqjpeYcg==} engines: {node: '>=8.0.0'} + '@opentelemetry/configuration@0.211.0': + resolution: {integrity: sha512-PNsCkzsYQKyv8wiUIsH+loC4RYyblOaDnVASBtKS22hK55ToWs2UP6IsrcfSWWn54wWTvVe2gnfwz67Pvrxf2Q==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.9.0 + '@opentelemetry/context-async-hooks@2.11.0': resolution: {integrity: sha512-Tr79DyWI8itsBdg+jH+opjfrwLzX+erk1/ExkIwhWoAVjVrJIn2y5+cGjTC0Vy8fyNIA/y8wuJPZwr1T3xCZeQ==} engines: {node: ^18.19.0 || >=20.6.0} peerDependencies: '@opentelemetry/api': '>=1.0.0 <1.10.0' + '@opentelemetry/context-async-hooks@2.5.0': + resolution: {integrity: sha512-uOXpVX0ZjO7heSVjhheW2XEPrhQAWr2BScDPoZ9UDycl5iuHG+Usyc3AIfG6kZeC1GyLpMInpQ6X5+9n69yOFw==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': '>=1.0.0 <1.10.0' + '@opentelemetry/core@2.11.0': resolution: {integrity: sha512-7YP44XH0tV6+Mb54x2YGf84i7yi+31MBZlE8JwvozkxyTvXbSp10X7cI7YE49ChJ3shMJoBmCJF3+1QFBJctGA==} engines: {node: ^18.19.0 || >=20.6.0} peerDependencies: '@opentelemetry/api': '>=1.0.0 <1.10.0' + '@opentelemetry/core@2.5.0': + resolution: {integrity: sha512-ka4H8OM6+DlUhSAZpONu0cPBtPPTQKxbxVzC4CzVx5+K4JnroJVBtDzLAMx4/3CDTJXRvVFhpFjtl4SaiTNoyQ==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': '>=1.0.0 <1.10.0' + + '@opentelemetry/exporter-logs-otlp-grpc@0.211.0': + resolution: {integrity: sha512-UhOoWENNqyaAMP/dL1YXLkXt6ZBtovkDDs1p4rxto9YwJX1+wMjwg+Obfyg2kwpcMoaiIFT3KQIcLNW8nNGNfQ==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/exporter-logs-otlp-http@0.211.0': + resolution: {integrity: sha512-c118Awf1kZirHkqxdcF+rF5qqWwNjJh+BB1CmQvN9AQHC/DUIldy6dIkJn3EKlQnQ3HmuNRKc/nHHt5IusN7mA==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/exporter-logs-otlp-proto@0.211.0': + resolution: {integrity: sha512-kMvfKMtY5vJDXeLnwhrZMEwhZ2PN8sROXmzacFU/Fnl4Z79CMrOaL7OE+5X3SObRYlDUa7zVqaXp9ZetYCxfDQ==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/exporter-metrics-otlp-grpc@0.211.0': + resolution: {integrity: sha512-D/U3G8L4PzZp8ot5hX9wpgbTymgtLZCiwR7heMe4LsbGV4OdctS1nfyvaQHLT6CiGZ6FjKc1Vk9s6kbo9SWLXQ==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/exporter-metrics-otlp-http@0.211.0': + resolution: {integrity: sha512-lfHXElPAoDSPpPO59DJdN5FLUnwi1wxluLTWQDayqrSPfWRnluzxRhD+g7rF8wbj1qCz0sdqABl//ug1IZyWvA==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/exporter-metrics-otlp-proto@0.211.0': + resolution: {integrity: sha512-61iNbffEpyZv/abHaz3BQM3zUtA2kVIDBM+0dS9RK68ML0QFLRGYa50xVMn2PYMToyfszEPEgFC3ypGae2z8FA==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/exporter-prometheus@0.211.0': + resolution: {integrity: sha512-cD0WleEL3TPqJbvxwz5MVdVJ82H8jl8mvMad4bNU24cB5SH2mRW5aMLDTuV4614ll46R//R3RMmci26mc2L99g==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/exporter-trace-otlp-grpc@0.211.0': + resolution: {integrity: sha512-eFwx4Gvu6LaEiE1rOd4ypgAiWEdZu7Qzm2QNN2nJqPW1XDeAVH1eNwVcVQl+QK9HR/JCDZ78PZgD7xD/DBDqbw==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/exporter-trace-otlp-http@0.211.0': + resolution: {integrity: sha512-F1Rv3JeMkgS//xdVjbQMrI3+26e5SXC7vXA6trx8SWEA0OUhw4JHB+qeHtH0fJn46eFItrYbL5m8j4qi9Sfaxw==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/exporter-trace-otlp-proto@0.211.0': + resolution: {integrity: sha512-DkjXwbPiqpcPlycUojzG2RmR0/SIK8Gi9qWO9znNvSqgzrnAIE9x2n6yPfpZ+kWHZGafvsvA1lVXucTyyQa5Kg==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/exporter-zipkin@2.5.0': + resolution: {integrity: sha512-bk9VJgFgUAzkZzU8ZyXBSWiUGLOM3mZEgKJ1+jsZclhRnAoDNf+YBdq+G9R3cP0+TKjjWad+vVrY/bE/vRR9lA==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.0.0 + + '@opentelemetry/instrumentation@0.211.0': + resolution: {integrity: sha512-h0nrZEC/zvI994nhg7EgQ8URIHt0uDTwN90r3qQUdZORS455bbx+YebnGeEuFghUT0HlJSrLF4iHw67f+odY+Q==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/otlp-exporter-base@0.211.0': + resolution: {integrity: sha512-bp1+63V8WPV+bRI9EQG6E9YID1LIHYSZVbp7f+44g9tRzCq+rtw/o4fpL5PC31adcUsFiz/oN0MdLISSrZDdrg==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/otlp-grpc-exporter-base@0.211.0': + resolution: {integrity: sha512-mR5X+N4SuphJeb7/K7y0JNMC8N1mB6gEtjyTLv+TSAhl0ZxNQzpSKP8S5Opk90fhAqVYD4R0SQSAirEBlH1KSA==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/otlp-transformer@0.211.0': + resolution: {integrity: sha512-julhCJ9dXwkOg9svuuYqqjXLhVaUgyUvO2hWbTxwjvLXX2rG3VtAaB0SzxMnGTuoCZizBT7Xqqm2V7+ggrfCXA==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/propagator-b3@2.5.0': + resolution: {integrity: sha512-g10m4KD73RjHrSvUge+sUxUl8m4VlgnGc6OKvo68a4uMfaLjdFU+AULfvMQE/APq38k92oGUxEzBsAZ8RN/YHg==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': '>=1.0.0 <1.10.0' + + '@opentelemetry/propagator-jaeger@2.5.0': + resolution: {integrity: sha512-t70ErZCncAR/zz5AcGkL0TF25mJiK1FfDPEQCgreyAHZ+mRJ/bNUiCnImIBDlP3mSDXy6N09DbUEKq0ktW98Hg==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': '>=1.0.0 <1.10.0' + '@opentelemetry/resources@2.11.0': resolution: {integrity: sha512-Ie7+8q8MDF4FAEQCKVMTx3ReUvxiIAgIiiW3c9JdmP8+HMcDy20puT+AHjexnExgnbvBxjQ9fjkFDWrikJ2jQA==} engines: {node: ^18.19.0 || >=20.6.0} peerDependencies: '@opentelemetry/api': '>=1.3.0 <1.10.0' + '@opentelemetry/resources@2.5.0': + resolution: {integrity: sha512-F8W52ApePshpoSrfsSk1H2yJn9aKjCrbpQF1M9Qii0GHzbfVeFUB+rc3X4aggyZD8x9Gu3Slua+s6krmq6Dt8g==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': '>=1.3.0 <1.10.0' + + '@opentelemetry/sdk-logs@0.211.0': + resolution: {integrity: sha512-O5nPwzgg2JHzo59kpQTPUOTzFi0Nv5LxryG27QoXBciX3zWM3z83g+SNOHhiQVYRWFSxoWn1JM2TGD5iNjOwdA==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': '>=1.4.0 <1.10.0' + + '@opentelemetry/sdk-metrics@2.5.0': + resolution: {integrity: sha512-BeJLtU+f5Gf905cJX9vXFQorAr6TAfK3SPvTFqP+scfIpDQEJfRaGJWta7sJgP+m4dNtBf9y3yvBKVAZZtJQVA==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': '>=1.9.0 <1.10.0' + + '@opentelemetry/sdk-node@0.211.0': + resolution: {integrity: sha512-+s1eGjoqmPCMptNxcJJD4IxbWJKNLOQFNKhpwkzi2gLkEbCj6LzSHJNhPcLeBrBlBLtlSpibM+FuS7fjZ8SSFQ==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': '>=1.3.0 <1.10.0' + '@opentelemetry/sdk-trace-base@2.11.0': resolution: {integrity: sha512-H19x/TX/LZdqiYOjM7fqtSxwlplC5pgelavqbQdHbhdq0q/AI/TGkM2dfGuuynTXmJPeF2HoZVoPDu+TGoW78A==} engines: {node: ^18.19.0 || >=20.6.0} peerDependencies: '@opentelemetry/api': '>=1.3.0 <1.10.0' + '@opentelemetry/sdk-trace-base@2.5.0': + resolution: {integrity: sha512-VzRf8LzotASEyNDUxTdaJ9IRJ1/h692WyArDBInf5puLCjxbICD6XkHgpuudis56EndyS7LYFmtTMny6UABNdQ==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': '>=1.3.0 <1.10.0' + '@opentelemetry/sdk-trace-node@2.11.0': resolution: {integrity: sha512-CuvCMJmZxswhNLlM2LfuLOW3h3fZujA4hsG4B+Sz4dX2zvaXO8Ng74cnDHWD64gLszTlhiG3c0iNUjj4g+0/sA==} engines: {node: ^18.19.0 || >=20.6.0} peerDependencies: '@opentelemetry/api': '>=1.0.0 <1.10.0' + '@opentelemetry/sdk-trace-node@2.5.0': + resolution: {integrity: sha512-O6N/ejzburFm2C84aKNrwJVPpt6HSTSq8T0ZUMq3xT2XmqT4cwxUItcL5UWGThYuq8RTcbH8u1sfj6dmRci0Ow==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': '>=1.0.0 <1.10.0' + '@opentelemetry/sdk-trace@2.11.0': resolution: {integrity: sha512-fFnTqGm8/G73GQVnxYi7LXa1ZVYEUvgL6XI1LpvV0bPC7WQ/ZGgKxCSl8FnlZBKto9JHHEFTO6s6CUpvvtwFrA==} engines: {node: ^18.19.0 || >=20.6.0} @@ -7674,12 +7852,21 @@ packages: '@protobufjs/codegen@2.0.4': resolution: {integrity: sha512-YyFaikqM5sH0ziFZCN3xDC7zeGaB/d0IUb9CATugHWbd1FRFwWwt4ld4OYMPWu5a3Xe01mGAULCdqhMlPl29Jg==} + '@protobufjs/codegen@2.0.5': + resolution: {integrity: sha512-zgXFLzW3Ap33e6d0Wlj4MGIm6Ce8O89n/apUaGNB/jx+hw+ruWEp7EwGUshdLKVRCxZW12fp9r40E1mQrf/34g==} + '@protobufjs/eventemitter@1.1.0': resolution: {integrity: sha512-j9ednRT81vYJ9OfVuXG6ERSTdEL1xVsNgqpkxMsbIabzSo3goCjDIveeGv5d03om39ML71RdmrGNjG5SReBP/Q==} + '@protobufjs/eventemitter@1.1.1': + resolution: {integrity: sha512-vW1GmwMZNnL+gMRaovlh9yZX74kc+TTU3FObkkurpMaRtBfLP3ldjS9KQWlwZgraRE0+dheEEoAxdzcJQ8eXZg==} + '@protobufjs/fetch@1.1.0': resolution: {integrity: sha512-lljVXpqXebpsijW71PZaCYeIcE5on1w5DlQy5WH6GLbFryLUrBD4932W/E2BSpfRJWseIL4v/KPgBFxDOIdKpQ==} + '@protobufjs/fetch@1.1.1': + resolution: {integrity: sha512-GpptLrs57adMSuHi3VNj0mAF8dwh36LMaYF6XyJ6JMWlVsc+t42tm1HSEDmOs3A8fC9yyeisgLhsTVQokOZ0zw==} + '@protobufjs/float@1.0.2': resolution: {integrity: sha512-Ddb+kVXlXst9d+R9PfTIxh1EdNkgoRe5tOX6t01f1lYWOvJnSPDBlG241QLzcyPdoNTsblLUdujGSE4RzrTZGQ==} @@ -7695,6 +7882,9 @@ packages: '@protobufjs/utf8@1.1.0': resolution: {integrity: sha512-Vvn3zZrhQZkkBE8LSuW3em98c0FwgO4nxzv6OdSxPKJIEKY2bGbHn+mhGIPerzI4twdxaP8/0+06HBpwf345Lw==} + '@protobufjs/utf8@1.1.2': + resolution: {integrity: sha512-b1UQwcEZ4yCnMCD8DAL1VlbvBJE9/IX4FTIp7BG1xYpf29SLazLSrqUkj4w7Y5y7cCVP6E5tcqqcI0xemPkHug==} + '@publint/pack@0.1.4': resolution: {integrity: sha512-HDVTWq3H0uTXiU0eeSQntcVUTPP3GamzeXI41+x7uU9J65JgWQh3qWZHblR1i0npXfFtF+mxBiU2nJH8znxWnQ==} engines: {node: '>=18'} @@ -11025,6 +11215,11 @@ packages: resolution: {integrity: sha512-5cvg6CtKwfgdmVqY1WIiXKc3Q1bkRqGLi+2W/6ao+6Y7gu/RCwRuAhGEzh5B4KlszSuTLgZYuqFqo5bImjNKng==} engines: {node: '>= 0.6'} + acorn-import-attributes@1.9.5: + resolution: {integrity: sha512-n02Vykv5uA3eHGM/Z2dQrcD56kL8TyDb2p1+0P83PClMnC/nc+anbQRhIOWnSq4Ke/KvDPrY3C9hDtC/A3eHnQ==} + peerDependencies: + acorn: ^8 + acorn-import-phases@1.0.4: resolution: {integrity: sha512-wKmbr/DDiIXzEOiWrTTUcDm24kQ2vGfZQvM2fwg2vXqR5uW6aapr7ObPtj1th32b9u90/Pf4AItvdTh42fBmVQ==} engines: {node: '>=10.13.0'} @@ -14120,6 +14315,9 @@ packages: resolution: {integrity: sha512-TR3KfrTZTYLPB6jUjfx6MF9WcWrHL9su5TObK4ZkYgBdWKPOFoSoQIdEuTuR82pmtxH2spWG9h6etwfr1pLBqQ==} engines: {node: '>=6'} + import-in-the-middle@2.0.6: + resolution: {integrity: sha512-3vZV3jX0XRFW3EJDTwzWoZa+RH1b8eTTx6YOCjglrLyPuepwoBti1k3L2dKwdCUrnVEfc5CuRuGstaC/uQJJaw==} + import-lazy@4.0.0: resolution: {integrity: sha512-rKtvo6a868b5Hu3heneU+L4yEQ4jYKLtjpnPeUdK7h0yzXGmyBTypknlkCvHFBqfX9YlorEiMM6Dnq/5atfHkw==} engines: {node: '>=8'} @@ -15528,6 +15726,9 @@ packages: resolution: {integrity: sha512-t9VmxaqrmANnEOBhpSDI6HD192Ge48k8vmWqQQL7hSFEqHEYwZbbsu49+aKLWZeRvFs3j1pMhXOqqF4kPlvjkQ==} engines: {node: '>=18.0.0'} + module-details-from-path@1.0.4: + resolution: {integrity: sha512-EGWKgxALGMgzvxYF1UyGTy0HXX/2vHLkw6+NvDKW2jypWbHpjQuj4UMcqQWXHERJhVGKikolT06G3bcKe4fi7w==} + monaco-editor@0.55.1: resolution: {integrity: sha512-jz4x+TJNFHwHtwuV9vA9rMujcZRb0CEilTEwG2rRSpe/A7Jdkuj8xPKttCgOh+v/lkHy7HsZ64oj+q3xoAFl9A==} @@ -16533,6 +16734,14 @@ packages: resolution: {integrity: sha512-CvexbZtbov6jW2eXAvLukXjXUW1TzFaivC46BpWc/3BpcCysb5Vffu+B3XHMm8lVEuy2Mm4XGex8hBSg1yapPg==} engines: {node: '>=12.0.0'} + protobufjs@7.6.6: + resolution: {integrity: sha512-dYDWdjSl5RNb7SgPxGQcRU+GtvP7s2fpkrY0r432PcOIaZ0/rBcxEZnQN67iJhFuQiVw754JDoPruPCNdGsbjg==} + engines: {node: '>=12.0.0'} + + protobufjs@8.0.0: + resolution: {integrity: sha512-jx6+sE9h/UryaCZhsJWbJtTEy47yXoGNYI4z8ZaRncM0zBKeRqjO2JEcOUYwrYGb1WLhXM1FfMzW3annvFv0rw==} + engines: {node: '>=12.0.0'} + proxy-addr@2.0.7: resolution: {integrity: sha512-llQsMLSUDUPT44jdrU/O37qlnifitDP+ZwrmmZcoSKyLKvtZxpyV0n2/bD/N4tBAAZ/gJEdZU7KMraoK1+XYAg==} engines: {node: '>= 0.10'} @@ -16986,6 +17195,10 @@ packages: resolution: {integrity: sha512-Xf0nWe6RseziFMu+Ap9biiUbmplq6S9/p+7w7YXP/JBHhrUDDUhwa+vANyubuqfZWTveU//DYVGsDG7RKL/vEw==} engines: {node: '>=0.10.0'} + require-in-the-middle@8.0.1: + resolution: {integrity: sha512-QT7FVMXfWOYFbeRBF6nu+I6tr2Tf3u0q8RIEjNob/heKY/nh7drD/k7eeMFmSQgnTtCzLDcCu/XEnpW2wk4xCQ==} + engines: {node: '>=9.3.0 || >=8.10.0 <9.0.0'} + requireg@0.2.2: resolution: {integrity: sha512-nYzyjnFcPNGR3lx9lwPPPnuQxv6JWEZd2Ci0u9opN7N5zUEPIhY/GbL3vMGOr2UXwEg9WwSyV9X9Y/kLFgPsOg==} engines: {node: '>= 4.0.0'} @@ -22134,6 +22347,18 @@ snapshots: - supports-color - utf-8-validate + '@grpc/grpc-js@1.14.4': + dependencies: + '@grpc/proto-loader': 0.8.1 + '@js-sdsl/ordered-map': 4.4.2 + + '@grpc/proto-loader@0.8.1': + dependencies: + lodash.camelcase: 4.3.0 + long: 5.3.2 + protobufjs: 7.6.6 + yargs: 17.7.2 + '@hey-api/client-fetch@0.5.7': {} '@hono/node-server@1.14.2(hono@4.12.33)': @@ -22793,6 +23018,8 @@ snapshots: '@jridgewell/resolve-uri': 3.1.2 '@jridgewell/sourcemap-codec': 1.5.5 + '@js-sdsl/ordered-map@4.4.2': {} + '@keyv/serialize@1.1.1': {} '@ladle/react-context@1.0.1(react-dom@19.1.0(react@19.1.0))(react@19.1.0)': @@ -23753,23 +23980,240 @@ snapshots: '@open-draft/until@2.1.0': {} + '@opentelemetry/api-logs@0.211.0': + dependencies: + '@opentelemetry/api': 1.9.0 + '@opentelemetry/api@1.9.0': {} + '@opentelemetry/configuration@0.211.0(@opentelemetry/api@1.9.0)': + dependencies: + '@opentelemetry/api': 1.9.0 + '@opentelemetry/core': 2.5.0(@opentelemetry/api@1.9.0) + yaml: 2.9.0 + '@opentelemetry/context-async-hooks@2.11.0(@opentelemetry/api@1.9.0)': dependencies: '@opentelemetry/api': 1.9.0 + '@opentelemetry/context-async-hooks@2.5.0(@opentelemetry/api@1.9.0)': + dependencies: + '@opentelemetry/api': 1.9.0 + '@opentelemetry/core@2.11.0(@opentelemetry/api@1.9.0)': dependencies: '@opentelemetry/api': 1.9.0 '@opentelemetry/semantic-conventions': 1.40.0 + '@opentelemetry/core@2.5.0(@opentelemetry/api@1.9.0)': + dependencies: + '@opentelemetry/api': 1.9.0 + '@opentelemetry/semantic-conventions': 1.40.0 + + '@opentelemetry/exporter-logs-otlp-grpc@0.211.0(@opentelemetry/api@1.9.0)': + dependencies: + '@grpc/grpc-js': 1.14.4 + '@opentelemetry/api': 1.9.0 + '@opentelemetry/core': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/otlp-exporter-base': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/otlp-grpc-exporter-base': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/otlp-transformer': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/sdk-logs': 0.211.0(@opentelemetry/api@1.9.0) + + '@opentelemetry/exporter-logs-otlp-http@0.211.0(@opentelemetry/api@1.9.0)': + dependencies: + '@opentelemetry/api': 1.9.0 + '@opentelemetry/api-logs': 0.211.0 + '@opentelemetry/core': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/otlp-exporter-base': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/otlp-transformer': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/sdk-logs': 0.211.0(@opentelemetry/api@1.9.0) + + '@opentelemetry/exporter-logs-otlp-proto@0.211.0(@opentelemetry/api@1.9.0)': + dependencies: + '@opentelemetry/api': 1.9.0 + '@opentelemetry/api-logs': 0.211.0 + '@opentelemetry/core': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/otlp-exporter-base': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/otlp-transformer': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/resources': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/sdk-logs': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/sdk-trace-base': 2.5.0(@opentelemetry/api@1.9.0) + + '@opentelemetry/exporter-metrics-otlp-grpc@0.211.0(@opentelemetry/api@1.9.0)': + dependencies: + '@grpc/grpc-js': 1.14.4 + '@opentelemetry/api': 1.9.0 + '@opentelemetry/core': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/exporter-metrics-otlp-http': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/otlp-exporter-base': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/otlp-grpc-exporter-base': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/otlp-transformer': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/resources': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/sdk-metrics': 2.5.0(@opentelemetry/api@1.9.0) + + '@opentelemetry/exporter-metrics-otlp-http@0.211.0(@opentelemetry/api@1.9.0)': + dependencies: + '@opentelemetry/api': 1.9.0 + '@opentelemetry/core': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/otlp-exporter-base': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/otlp-transformer': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/resources': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/sdk-metrics': 2.5.0(@opentelemetry/api@1.9.0) + + '@opentelemetry/exporter-metrics-otlp-proto@0.211.0(@opentelemetry/api@1.9.0)': + dependencies: + '@opentelemetry/api': 1.9.0 + '@opentelemetry/core': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/exporter-metrics-otlp-http': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/otlp-exporter-base': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/otlp-transformer': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/resources': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/sdk-metrics': 2.5.0(@opentelemetry/api@1.9.0) + + '@opentelemetry/exporter-prometheus@0.211.0(@opentelemetry/api@1.9.0)': + dependencies: + '@opentelemetry/api': 1.9.0 + '@opentelemetry/core': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/resources': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/sdk-metrics': 2.5.0(@opentelemetry/api@1.9.0) + + '@opentelemetry/exporter-trace-otlp-grpc@0.211.0(@opentelemetry/api@1.9.0)': + dependencies: + '@grpc/grpc-js': 1.14.4 + '@opentelemetry/api': 1.9.0 + '@opentelemetry/core': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/otlp-exporter-base': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/otlp-grpc-exporter-base': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/otlp-transformer': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/resources': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/sdk-trace-base': 2.5.0(@opentelemetry/api@1.9.0) + + '@opentelemetry/exporter-trace-otlp-http@0.211.0(@opentelemetry/api@1.9.0)': + dependencies: + '@opentelemetry/api': 1.9.0 + '@opentelemetry/core': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/otlp-exporter-base': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/otlp-transformer': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/resources': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/sdk-trace-base': 2.5.0(@opentelemetry/api@1.9.0) + + '@opentelemetry/exporter-trace-otlp-proto@0.211.0(@opentelemetry/api@1.9.0)': + dependencies: + '@opentelemetry/api': 1.9.0 + '@opentelemetry/core': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/otlp-exporter-base': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/otlp-transformer': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/resources': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/sdk-trace-base': 2.5.0(@opentelemetry/api@1.9.0) + + '@opentelemetry/exporter-zipkin@2.5.0(@opentelemetry/api@1.9.0)': + dependencies: + '@opentelemetry/api': 1.9.0 + '@opentelemetry/core': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/resources': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/sdk-trace-base': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/semantic-conventions': 1.40.0 + + '@opentelemetry/instrumentation@0.211.0(@opentelemetry/api@1.9.0)': + dependencies: + '@opentelemetry/api': 1.9.0 + '@opentelemetry/api-logs': 0.211.0 + import-in-the-middle: 2.0.6 + require-in-the-middle: 8.0.1 + transitivePeerDependencies: + - supports-color + + '@opentelemetry/otlp-exporter-base@0.211.0(@opentelemetry/api@1.9.0)': + dependencies: + '@opentelemetry/api': 1.9.0 + '@opentelemetry/core': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/otlp-transformer': 0.211.0(@opentelemetry/api@1.9.0) + + '@opentelemetry/otlp-grpc-exporter-base@0.211.0(@opentelemetry/api@1.9.0)': + dependencies: + '@grpc/grpc-js': 1.14.4 + '@opentelemetry/api': 1.9.0 + '@opentelemetry/core': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/otlp-exporter-base': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/otlp-transformer': 0.211.0(@opentelemetry/api@1.9.0) + + '@opentelemetry/otlp-transformer@0.211.0(@opentelemetry/api@1.9.0)': + dependencies: + '@opentelemetry/api': 1.9.0 + '@opentelemetry/api-logs': 0.211.0 + '@opentelemetry/core': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/resources': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/sdk-logs': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/sdk-metrics': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/sdk-trace-base': 2.5.0(@opentelemetry/api@1.9.0) + protobufjs: 8.0.0 + + '@opentelemetry/propagator-b3@2.5.0(@opentelemetry/api@1.9.0)': + dependencies: + '@opentelemetry/api': 1.9.0 + '@opentelemetry/core': 2.5.0(@opentelemetry/api@1.9.0) + + '@opentelemetry/propagator-jaeger@2.5.0(@opentelemetry/api@1.9.0)': + dependencies: + '@opentelemetry/api': 1.9.0 + '@opentelemetry/core': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/resources@2.11.0(@opentelemetry/api@1.9.0)': dependencies: '@opentelemetry/api': 1.9.0 '@opentelemetry/core': 2.11.0(@opentelemetry/api@1.9.0) '@opentelemetry/semantic-conventions': 1.40.0 + '@opentelemetry/resources@2.5.0(@opentelemetry/api@1.9.0)': + dependencies: + '@opentelemetry/api': 1.9.0 + '@opentelemetry/core': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/semantic-conventions': 1.40.0 + + '@opentelemetry/sdk-logs@0.211.0(@opentelemetry/api@1.9.0)': + dependencies: + '@opentelemetry/api': 1.9.0 + '@opentelemetry/api-logs': 0.211.0 + '@opentelemetry/core': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/resources': 2.5.0(@opentelemetry/api@1.9.0) + + '@opentelemetry/sdk-metrics@2.5.0(@opentelemetry/api@1.9.0)': + dependencies: + '@opentelemetry/api': 1.9.0 + '@opentelemetry/core': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/resources': 2.5.0(@opentelemetry/api@1.9.0) + + '@opentelemetry/sdk-node@0.211.0(@opentelemetry/api@1.9.0)': + dependencies: + '@opentelemetry/api': 1.9.0 + '@opentelemetry/api-logs': 0.211.0 + '@opentelemetry/configuration': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/context-async-hooks': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/core': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/exporter-logs-otlp-grpc': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/exporter-logs-otlp-http': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/exporter-logs-otlp-proto': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/exporter-metrics-otlp-grpc': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/exporter-metrics-otlp-http': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/exporter-metrics-otlp-proto': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/exporter-prometheus': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/exporter-trace-otlp-grpc': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/exporter-trace-otlp-http': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/exporter-trace-otlp-proto': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/exporter-zipkin': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/instrumentation': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/propagator-b3': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/propagator-jaeger': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/resources': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/sdk-logs': 0.211.0(@opentelemetry/api@1.9.0) + '@opentelemetry/sdk-metrics': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/sdk-trace-base': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/sdk-trace-node': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/semantic-conventions': 1.40.0 + transitivePeerDependencies: + - supports-color + '@opentelemetry/sdk-trace-base@2.11.0(@opentelemetry/api@1.9.0)': dependencies: '@opentelemetry/api': 1.9.0 @@ -23778,6 +24222,13 @@ snapshots: '@opentelemetry/sdk-trace': 2.11.0(@opentelemetry/api@1.9.0) '@opentelemetry/semantic-conventions': 1.40.0 + '@opentelemetry/sdk-trace-base@2.5.0(@opentelemetry/api@1.9.0)': + dependencies: + '@opentelemetry/api': 1.9.0 + '@opentelemetry/core': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/resources': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/semantic-conventions': 1.40.0 + '@opentelemetry/sdk-trace-node@2.11.0(@opentelemetry/api@1.9.0)': dependencies: '@opentelemetry/api': 1.9.0 @@ -23785,6 +24236,13 @@ snapshots: '@opentelemetry/core': 2.11.0(@opentelemetry/api@1.9.0) '@opentelemetry/sdk-trace-base': 2.11.0(@opentelemetry/api@1.9.0) + '@opentelemetry/sdk-trace-node@2.5.0(@opentelemetry/api@1.9.0)': + dependencies: + '@opentelemetry/api': 1.9.0 + '@opentelemetry/context-async-hooks': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/core': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/sdk-trace-base': 2.5.0(@opentelemetry/api@1.9.0) + '@opentelemetry/sdk-trace@2.11.0(@opentelemetry/api@1.9.0)': dependencies: '@opentelemetry/api': 1.9.0 @@ -23887,13 +24345,21 @@ snapshots: '@protobufjs/codegen@2.0.4': {} + '@protobufjs/codegen@2.0.5': {} + '@protobufjs/eventemitter@1.1.0': {} + '@protobufjs/eventemitter@1.1.1': {} + '@protobufjs/fetch@1.1.0': dependencies: '@protobufjs/aspromise': 1.1.2 '@protobufjs/inquire': 1.1.0 + '@protobufjs/fetch@1.1.1': + dependencies: + '@protobufjs/aspromise': 1.1.2 + '@protobufjs/float@1.0.2': {} '@protobufjs/inquire@1.1.0': {} @@ -23904,6 +24370,8 @@ snapshots: '@protobufjs/utf8@1.1.0': {} + '@protobufjs/utf8@1.1.2': {} + '@publint/pack@0.1.4': {} '@radix-ui/number@1.1.1': {} @@ -28130,6 +28598,10 @@ snapshots: mime-types: 3.0.2 negotiator: 1.0.0 + acorn-import-attributes@1.9.5(acorn@8.16.0): + dependencies: + acorn: 8.16.0 + acorn-import-phases@1.0.4(acorn@8.16.0): dependencies: acorn: 8.16.0 @@ -31463,6 +31935,13 @@ snapshots: parent-module: 1.0.1 resolve-from: 4.0.0 + import-in-the-middle@2.0.6: + dependencies: + acorn: 8.16.0 + acorn-import-attributes: 1.9.5(acorn@8.16.0) + cjs-module-lexer: 2.2.0 + module-details-from-path: 1.0.4 + import-lazy@4.0.0: {} imurmurhash@0.1.4: {} @@ -33304,6 +33783,8 @@ snapshots: modern-tar@0.7.7: {} + module-details-from-path@1.0.4: {} + monaco-editor@0.55.1: dependencies: dompurify: 3.2.7 @@ -34497,6 +34978,35 @@ snapshots: '@types/node': 22.19.15 long: 5.3.2 + protobufjs@7.6.6: + dependencies: + '@protobufjs/aspromise': 1.1.2 + '@protobufjs/base64': 1.1.2 + '@protobufjs/codegen': 2.0.5 + '@protobufjs/eventemitter': 1.1.1 + '@protobufjs/fetch': 1.1.1 + '@protobufjs/float': 1.0.2 + '@protobufjs/path': 1.1.2 + '@protobufjs/pool': 1.1.0 + '@protobufjs/utf8': 1.1.2 + '@types/node': 22.19.15 + long: 5.3.2 + + protobufjs@8.0.0: + dependencies: + '@protobufjs/aspromise': 1.1.2 + '@protobufjs/base64': 1.1.2 + '@protobufjs/codegen': 2.0.4 + '@protobufjs/eventemitter': 1.1.0 + '@protobufjs/fetch': 1.1.0 + '@protobufjs/float': 1.0.2 + '@protobufjs/inquire': 1.1.0 + '@protobufjs/path': 1.1.2 + '@protobufjs/pool': 1.1.0 + '@protobufjs/utf8': 1.1.0 + '@types/node': 22.19.15 + long: 5.3.2 + proxy-addr@2.0.7: dependencies: forwarded: 0.2.0 @@ -35082,6 +35592,13 @@ snapshots: require-from-string@2.0.2: {} + require-in-the-middle@8.0.1: + dependencies: + debug: 4.4.3(supports-color@8.1.1) + module-details-from-path: 1.0.4 + transitivePeerDependencies: + - supports-color + requireg@0.2.2: dependencies: nested-error-stacks: 2.0.1 From 982eb970b6039cb51e704f64d848fa2a326fcb7b Mon Sep 17 00:00:00 2001 From: Sree Narayanan Date: Thu, 17 Sep 2026 15:30:35 +0400 Subject: [PATCH 2/6] feat(rivetkit): add tracing for workflow runs and steps --- .../packages/actor-persist/src/versioned.rs | 48 ++++ .../rivetkit-core/src/actor/context.rs | 56 +++- .../src/actor/internal_storage/mod.rs | 50 +++- .../src/actor/internal_storage/queries.rs | 1 + .../packages/rivetkit-core/src/actor/keys.rs | 4 + .../packages/rivetkit-core/src/lib.rs | 3 +- .../packages/rivetkit-core/src/telemetry.rs | 261 +++++++++++++++++- .../rivetkit-core/tests/sql_efficiency.rs | 15 + .../packages/rivetkit-napi/index.d.ts | 19 ++ .../packages/rivetkit-napi/index.js | 4 +- .../rivetkit-napi/src/actor_context.rs | 91 ++++++ .../driver-test-suite/registry-static.ts | 7 +- .../fixtures/driver-test-suite/telemetry.ts | 51 ++++ .../packages/rivetkit/src/actor/config.ts | 38 +++ .../rivetkit/src/registry/napi-runtime.ts | 35 ++- .../packages/rivetkit/src/registry/native.ts | 9 +- .../packages/rivetkit/src/registry/runtime.ts | 5 + .../rivetkit/src/registry/wasm-runtime.ts | 13 + .../packages/rivetkit/src/workflow/driver.ts | 32 ++- .../tests/driver/actor-telemetry.test.ts | 177 ++++++++++++ .../packages/workflow-engine/src/context.ts | 17 +- .../packages/workflow-engine/src/driver.ts | 39 +++ .../packages/workflow-engine/src/index.ts | 110 ++++++-- .../packages/workflow-engine/src/keys.ts | 1 + 24 files changed, 1033 insertions(+), 53 deletions(-) diff --git a/rivetkit-rust/packages/actor-persist/src/versioned.rs b/rivetkit-rust/packages/actor-persist/src/versioned.rs index d64bd2d9e0..12341020b1 100644 --- a/rivetkit-rust/packages/actor-persist/src/versioned.rs +++ b/rivetkit-rust/packages/actor-persist/src/versioned.rs @@ -1,4 +1,5 @@ use anyhow::{Result, bail}; +use serde::{Deserialize, Serialize}; use vbare::OwnedVersionedData; use crate::generated::{v1, v2, v3, v4}; @@ -535,3 +536,50 @@ impl OwnedVersionedData for RunWakeAt { Vec:: Result>::new() } } + +/// The span of the last workflow run, which the next run links to. +#[derive(Debug, Default, Serialize, Deserialize)] +pub struct WorkflowTraceContextV1 { + pub traceparent: Option, + pub tracestate: Option, +} + +pub enum WorkflowTraceContext { + V1(WorkflowTraceContextV1), +} + +impl OwnedVersionedData for WorkflowTraceContext { + type Latest = WorkflowTraceContextV1; + + fn wrap_latest(latest: Self::Latest) -> Self { + Self::V1(latest) + } + + fn unwrap_latest(self) -> Result { + match self { + Self::V1(data) => Ok(data), + } + } + + fn deserialize_version(payload: &[u8], version: u16) -> Result { + match version { + 1 => Ok(Self::V1(serde_bare::from_slice(payload)?)), + _ => bail!("invalid workflow trace context version: {version}"), + } + } + + fn serialize_version(self, version: u16) -> Result> { + match (self, version) { + (Self::V1(data), 1) => serde_bare::to_vec(&data).map_err(Into::into), + (_, version) => bail!("unexpected workflow trace context version: {version}"), + } + } + + fn deserialize_converters() -> Vec Result> { + Vec:: Result>::new() + } + + fn serialize_converters() -> Vec Result> { + Vec:: Result>::new() + } +} diff --git a/rivetkit-rust/packages/rivetkit-core/src/actor/context.rs b/rivetkit-rust/packages/rivetkit-core/src/actor/context.rs index 849a4b831e..79c6cd5b16 100644 --- a/rivetkit-rust/packages/rivetkit-core/src/actor/context.rs +++ b/rivetkit-rust/packages/rivetkit-core/src/actor/context.rs @@ -54,7 +54,8 @@ use crate::inspector::{Inspector, InspectorSnapshot}; use crate::sqlite::SqliteDb; use crate::telemetry::{ ActorInvocationTelemetry, ActorInvocationTraceContext, ActorTelemetryIdentity, - OutboundCallInvocation, + IncomingTraceContext, OutboundCallInvocation, WorkflowRunInvocation, WorkflowRunOutcome, + WorkflowStepSpan, }; use crate::types::{ActorKey, ConnId, ListOpts, format_actor_key}; @@ -305,6 +306,59 @@ impl ActorContext { .start_outbound_call(actor_name, action_name) } + /// Opens one workflow run as its own invocation, linked to the span the + /// previous run persisted. + #[doc(hidden)] + pub async fn start_workflow_span(&self) -> WorkflowRunInvocation { + let previous = if tracing::enabled!(target: "rivetkit::telemetry", tracing::Level::INFO) { + internal_storage::load_workflow_trace(&self.0.sql) + .await + .unwrap_or_else(|error| { + tracing::warn!( + actor_id = %self.actor_id(), + ?error, + "failed to load the previous workflow run span, so this run will not link to it" + ); + IncomingTraceContext::default() + }) + } else { + IncomingTraceContext::default() + }; + WorkflowRunInvocation::start(self, previous) + } + + /// Closes a workflow run and persists its span for the next run to link to. + #[doc(hidden)] + pub async fn finish_workflow_span( + &self, + run: WorkflowRunInvocation, + outcome: WorkflowRunOutcome, + ) { + let Some(trace_context) = run.finish(outcome) else { + return; + }; + if let Err(error) = + internal_storage::persist_workflow_trace(&self.0.sql, trace_context).await + { + tracing::warn!( + actor_id = %self.actor_id(), + ?error, + "failed to persist the workflow run span, so the next run will not link to it" + ); + } + } + + /// Opens the span for one step attempt, or nothing when this handle serves + /// no invocation or tracing is off. + #[doc(hidden)] + pub fn start_workflow_step_span( + &self, + step_name: &str, + attempt: u32, + ) -> Option { + WorkflowStepSpan::start(self, step_name, attempt) + } + /// Returns correlation for the invocation this handle serves, absent when /// the handle is not bound to one or tracing is disabled. pub fn invocation_trace_context(&self) -> Option { diff --git a/rivetkit-rust/packages/rivetkit-core/src/actor/internal_storage/mod.rs b/rivetkit-rust/packages/rivetkit-core/src/actor/internal_storage/mod.rs index aefcaf171f..e7bf252af8 100644 --- a/rivetkit-rust/packages/rivetkit-core/src/actor/internal_storage/mod.rs +++ b/rivetkit-rust/packages/rivetkit-core/src/actor/internal_storage/mod.rs @@ -7,7 +7,7 @@ use rivetkit_actor_persist::versioned as persist_versioned; use crate::actor::connection::{ PersistedConnection, PersistedSubscription, encode_persisted_connection, }; -use crate::actor::keys::make_workflow_key; +use crate::actor::keys::{WORKFLOW_TRACE_CONTEXT_KEY, make_workflow_key}; use crate::actor::messages::WorkflowKvWrite; use crate::actor::persist::{ decode_latest_with_embedded_version, encode_latest_with_embedded_version, @@ -35,6 +35,7 @@ const QUEUE_MESSAGE_IDS_PER_QUERY: usize = 128; const WORKFLOW_KV_VALUE_LIMIT: usize = 256 * 1024; pub(crate) const RUN_WAKE_AT_META_KEY: &str = "run_wake_at"; const RUN_WAKE_AT_VERSION: u16 = 1; +const WORKFLOW_TRACE_CONTEXT_VERSION: u16 = 1; /// Depot rejects SQLite commits that dirty more than `MAX_COMMIT_RAW_DIRTY_BYTES` /// (320 pages * 4 KiB = 1.3 MiB) in `engine/packages/depot/src/conveyer/constants.rs`. @@ -1082,6 +1083,53 @@ pub(crate) async fn persist_run_wake_at(db: &SqliteDb, wake_at: Option) -> Ok(()) } +pub(crate) async fn load_workflow_trace(db: &SqliteDb) -> Result { + let result = db + .query( + LOAD_WORKFLOW_KV_SQL, + Some(vec![BindParam::Blob(WORKFLOW_TRACE_CONTEXT_KEY.to_vec())]), + ) + .await + .context("load workflow trace context")?; + let Some(row) = result.rows.first() else { + return Ok(IncomingTraceContext::default()); + }; + let payload = read_blob(row, 0, "workflow trace context")?; + let stored = decode_latest_with_embedded_version::( + &payload, + "workflow trace context", + )?; + Ok(IncomingTraceContext { + traceparent: stored.traceparent, + tracestate: stored.tracestate, + ..Default::default() + }) +} + +pub(crate) async fn persist_workflow_trace( + db: &SqliteDb, + trace_context: IncomingTraceContext, +) -> Result<()> { + let payload = encode_latest_with_embedded_version::( + persist_versioned::WorkflowTraceContextV1 { + traceparent: trace_context.traceparent, + tracestate: trace_context.tracestate, + }, + WORKFLOW_TRACE_CONTEXT_VERSION, + "workflow trace context", + )?; + db.execute( + UPSERT_WORKFLOW_KV_SQL, + Some(vec![ + BindParam::Blob(WORKFLOW_TRACE_CONTEXT_KEY.to_vec()), + BindParam::Blob(payload), + ]), + ) + .await + .context("persist workflow trace context")?; + Ok(()) +} + pub(crate) async fn load_inspector_token(db: &SqliteDb) -> Result> { let result = db .query(LOAD_INSPECTOR_TOKEN_SQL, None) diff --git a/rivetkit-rust/packages/rivetkit-core/src/actor/internal_storage/queries.rs b/rivetkit-rust/packages/rivetkit-core/src/actor/internal_storage/queries.rs index 14c164810b..c55742b2b5 100644 --- a/rivetkit-rust/packages/rivetkit-core/src/actor/internal_storage/queries.rs +++ b/rivetkit-rust/packages/rivetkit-core/src/actor/internal_storage/queries.rs @@ -91,6 +91,7 @@ pub(crate) const UPSERT_QUEUE_NEXT_ID_SQL: &str = "INSERT INTO _rivet_runtime (i pub(crate) const UPSERT_LAST_PUSHED_ALARM_SQL: &str = "INSERT INTO _rivet_runtime (id, last_pushed_alarm, inspector_token, queue_next_id) VALUES (1, ?, ?, ?) ON CONFLICT(id) DO UPDATE SET last_pushed_alarm = excluded.last_pushed_alarm"; pub(crate) const UPSERT_RUN_WAKE_AT_SQL: &str = "INSERT INTO _rivet_meta (key, value) VALUES (?, ?) ON CONFLICT(key) DO UPDATE SET value = excluded.value"; pub(crate) const UPSERT_INSPECTOR_TOKEN_SQL: &str = "INSERT INTO _rivet_runtime (id, last_pushed_alarm, inspector_token, queue_next_id) VALUES (1, ?, ?, ?) ON CONFLICT(id) DO UPDATE SET inspector_token = excluded.inspector_token"; +pub(crate) const LOAD_WORKFLOW_KV_SQL: &str = "SELECT value FROM _rivet_wf_kv WHERE key = ?"; pub(crate) const LOAD_META_TEXT_SQL: &str = "SELECT value FROM _rivet_meta WHERE key = ?"; pub(crate) const UPSERT_META_TEXT_SQL: &str = "INSERT INTO _rivet_meta (key, value) VALUES (?, ?) ON CONFLICT(key) DO UPDATE SET value = excluded.value"; diff --git a/rivetkit-rust/packages/rivetkit-core/src/actor/keys.rs b/rivetkit-rust/packages/rivetkit-core/src/actor/keys.rs index 41b57c6ce3..d81937dc2a 100644 --- a/rivetkit-rust/packages/rivetkit-core/src/actor/keys.rs +++ b/rivetkit-rust/packages/rivetkit-core/src/actor/keys.rs @@ -58,6 +58,10 @@ pub const QUEUE_MESSAGES_PREFIX: [u8; 3] = [ ]; // Prefix for workflow v1 storage: [6, 1, ...workflow_key]. pub const WORKFLOW_STORAGE_PREFIX: [u8; 2] = [WORKFLOW_PREFIX[0], WORKFLOW_STORAGE_VERSION]; +/// Workflow storage key for the last run's trace context. The last two bytes +/// are the tuple encoding of 5, which the workflow engine reserves for core. +pub const WORKFLOW_TRACE_CONTEXT_KEY: [u8; 4] = + [WORKFLOW_PREFIX[0], WORKFLOW_STORAGE_VERSION, 0x15, 5]; // Prefix for trace v1 storage: [7, 1, ...trace_key]. pub const TRACES_STORAGE_PREFIX: [u8; 2] = [TRACES_PREFIX[0], TRACES_STORAGE_VERSION]; diff --git a/rivetkit-rust/packages/rivetkit-core/src/lib.rs b/rivetkit-rust/packages/rivetkit-core/src/lib.rs index 2fe5edf08b..c8a2b81b97 100644 --- a/rivetkit-rust/packages/rivetkit-core/src/lib.rs +++ b/rivetkit-rust/packages/rivetkit-core/src/lib.rs @@ -24,7 +24,8 @@ pub mod telemetry; #[doc(hidden)] pub use telemetry::{ ActorInvocationSpanContext, ActorInvocationTelemetry, ActorInvocationTraceContext, - IncomingTraceContext, OutboundCallInvocation, + IncomingTraceContext, OutboundCallInvocation, WorkflowRunInvocation, WorkflowRunOutcome, + WorkflowStepOutcome, WorkflowStepSpan, }; #[cfg(any(test, feature = "test-support"))] pub mod testing; diff --git a/rivetkit-rust/packages/rivetkit-core/src/telemetry.rs b/rivetkit-rust/packages/rivetkit-core/src/telemetry.rs index b18683d50a..1437984248 100644 --- a/rivetkit-rust/packages/rivetkit-core/src/telemetry.rs +++ b/rivetkit-rust/packages/rivetkit-core/src/telemetry.rs @@ -66,6 +66,9 @@ const REQUEST_INVOCATION_NAME: &str = "onRequest"; /// Name a queue send invocation is reported under. The queue itself is an attribute. const QUEUE_SEND_INVOCATION_NAME: &str = "queue.send"; +/// Name a workflow run invocation is reported under. +const WORKFLOW_RUN_INVOCATION_NAME: &str = "workflow"; + /// What an invocation ran, which decides its name and the attributes that /// identify it on the span. #[derive(Clone, Copy, Debug)] @@ -73,6 +76,7 @@ enum InvocationSubject<'a> { Action(&'a str), Request { method: &'a str }, QueueSend { queue: &'a str }, + WorkflowRun, } impl<'a> InvocationSubject<'a> { @@ -81,6 +85,7 @@ impl<'a> InvocationSubject<'a> { Self::Action(name) => name, Self::Request { .. } => REQUEST_INVOCATION_NAME, Self::QueueSend { .. } => QUEUE_SEND_INVOCATION_NAME, + Self::WorkflowRun => WORKFLOW_RUN_INVOCATION_NAME, } } @@ -89,6 +94,7 @@ impl<'a> InvocationSubject<'a> { Self::Action(name) => span.record("rivet.action.name", name), Self::Request { method } => span.record("http.request.method", method), Self::QueueSend { queue } => span.record("rivet.queue.name", queue), + Self::WorkflowRun => span, }; } } @@ -105,6 +111,7 @@ enum InvocationType { Scheduled, Request, QueueSend, + Workflow, } impl InvocationType { @@ -114,6 +121,7 @@ impl InvocationType { Self::Scheduled => "scheduled", Self::Request => "request", Self::QueueSend => "queue_send", + Self::Workflow => "workflow", } } } @@ -123,12 +131,14 @@ impl InvocationType { /// Every clone of a handle shares one invocation. `application_span` is the /// span the host runtime had active when it resolved this handle; Core cannot /// see the host's span stack, so spans opened through the handle parent there -/// when it is set and to the invocation span otherwise. +/// when it is set and to the invocation span otherwise. `step_span` is the +/// workflow step the handle runs inside. #[doc(hidden)] #[derive(Clone, Debug)] pub struct ActorInvocationTelemetry { inner: Arc, application_span: Option, + step_span: Option, } /// Identity fields that do not change while an actor is alive. Built once per @@ -248,6 +258,92 @@ pub struct OutboundCallInvocation { context: Option, } +/// How one run of the workflow function ended. `Sleeping` covers a timed +/// sleep, a wait for a message, and a retry backoff. +#[doc(hidden)] +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum WorkflowRunOutcome { + Completed, + Sleeping, + Evicted, + Failed, + Cancelled, +} + +impl WorkflowRunOutcome { + pub fn parse(value: &str) -> anyhow::Result { + match value { + "completed" => Ok(Self::Completed), + "sleeping" => Ok(Self::Sleeping), + "evicted" => Ok(Self::Evicted), + "failed" => Ok(Self::Failed), + "cancelled" => Ok(Self::Cancelled), + other => anyhow::bail!("unknown workflow run outcome `{other}`"), + } + } + + fn as_label(self) -> &'static str { + match self { + Self::Completed => "completed", + Self::Sleeping => "sleeping", + Self::Evicted => "evicted", + Self::Failed => "failed", + Self::Cancelled => "cancelled", + } + } + + fn is_error(self) -> bool { + match self { + Self::Completed | Self::Sleeping | Self::Evicted => false, + Self::Failed | Self::Cancelled => true, + } + } +} + +/// How one step attempt ended. `Retry` will be tried again, `Failed` will not. +#[doc(hidden)] +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum WorkflowStepOutcome { + Ok, + Retry, + Failed, +} + +impl WorkflowStepOutcome { + pub fn parse(value: &str) -> anyhow::Result { + match value { + "ok" => Ok(Self::Ok), + "retry" => Ok(Self::Retry), + "failed" => Ok(Self::Failed), + other => anyhow::bail!("unknown workflow step outcome `{other}`"), + } + } + + fn as_label(self) -> &'static str { + match self { + Self::Ok => "ok", + Self::Retry => "retry", + Self::Failed => "failed", + } + } +} + +/// One run of the workflow function. Dropping this without finishing records +/// the run as abandoned. +#[doc(hidden)] +pub struct WorkflowRunInvocation { + ctx: ActorContext, + invocation: ActorInvocation, +} + +/// One attempt at one workflow step. Dropping this without finishing records +/// the attempt as abandoned. +#[doc(hidden)] +pub struct WorkflowStepSpan { + ctx: ActorContext, + span: Option, +} + impl ActorInvocation { pub(crate) fn start_action( ctx: &ActorContext, @@ -283,6 +379,23 @@ impl ActorInvocation { ) } + /// Starts a new trace for one workflow run and links it to the `previous` + /// run. + fn start_workflow_run(ctx: &ActorContext, previous: IncomingTraceContext) -> Self { + let previous_run = parse_remote_parent( + previous.traceparent.as_deref(), + previous.tracestate.as_deref(), + ); + Self::start( + ctx, + InvocationSubject::WorkflowRun, + InvocationType::Workflow, + None, + None, + previous_run, + ) + } + /// Starts the invocation for one message sent into `queue_name` from /// outside the actor. It ends when the send is acknowledged, or when the /// sender's wait for a completion ends. @@ -347,6 +460,7 @@ impl ActorInvocation { http.request.method = tracing::field::Empty, http.response.status_code = tracing::field::Empty, rivet.queue.name = tracing::field::Empty, + rivet.workflow.run.outcome = tracing::field::Empty, otel.status_code = tracing::field::Empty, error.type = tracing::field::Empty, ); @@ -400,6 +514,53 @@ impl ActorInvocation { } } +impl WorkflowRunInvocation { + pub(crate) fn start(ctx: &ActorContext, previous: IncomingTraceContext) -> Self { + let invocation = ActorInvocation::start_workflow_run(ctx, previous); + Self { + ctx: ctx + .clone() + .with_invocation_telemetry(Some(invocation.telemetry())), + invocation, + } + } + + /// The context workflow code runs under. + pub fn ctx(&self) -> ActorContext { + self.ctx.clone() + } + + /// Returns the span the next run links to, or nothing when tracing is off. + pub(crate) fn finish(self, outcome: WorkflowRunOutcome) -> Option { + let mut trace_context = None; + self.invocation.telemetry.finish_with(|span| { + span.record("rivet.workflow.run.outcome", outcome.as_label()); + span.record( + "otel.status_code", + if outcome.is_error() { "ERROR" } else { "OK" }, + ); + let headers = otel_span_context_of(span) + .map(|span_context| w3c_trace_headers(&span_context)) + .unwrap_or_default(); + trace_context = Some(IncomingTraceContext { + traceparent: headers.traceparent, + tracestate: headers.tracestate, + ..Default::default() + }); + }); + trace_context + } +} + +impl Drop for WorkflowRunInvocation { + fn drop(&mut self) { + self.invocation.telemetry.finish_with(|span| { + span.record("otel.status_code", "ERROR"); + span.record("error.type", OPERATION_ABANDONED_ERROR_TYPE); + }); + } +} + impl Drop for ActorInvocation { fn drop(&mut self) { self.telemetry.finish_dropped(); @@ -423,6 +584,7 @@ impl ActorInvocationTelemetry { identity, }), application_span: None, + step_span: None, } } @@ -448,19 +610,21 @@ impl ActorInvocationTelemetry { Self { inner: self.inner.clone(), application_span, + step_span: self.step_span.clone(), } } /// The context a span opened through this handle parents to: the - /// application span when set, else the invocation span while it is open. + /// application span when set, else the workflow step it runs inside, else + /// the invocation span while it is open. fn parent_context(&self) -> Option { let state = self.inner.state.lock(); if state.finished && state.pending_work == 0 { return None; } - match &self.application_span { - Some(application_span) => { - Some(Context::new().with_remote_span_context(application_span.clone())) + match self.application_span.as_ref().or(self.step_span.as_ref()) { + Some(span_context) => { + Some(Context::new().with_remote_span_context(span_context.clone())) } None => state.span.as_ref().map(tracing::Span::context), } @@ -487,7 +651,10 @@ impl ActorInvocationTelemetry { } state.span.clone() }; - let span = span.and_then(|span| span_context_of(&span)); + let span = match &self.step_span { + Some(step_span) => w3c_span_context(step_span), + None => span.and_then(|span| span_context_of(&span)), + }; Some(ActorInvocationTraceContext { ray_id: self.inner.ray_id.clone(), @@ -549,6 +716,39 @@ impl ActorInvocationTelemetry { }) } + /// Opens the span for one step attempt. Spans opened through the returned + /// handle parent to it. + pub(crate) fn start_workflow_step( + &self, + step_name: &str, + attempt: u32, + ) -> Option<(tracing::Span, Self)> { + let parent = self.parent_context()?; + let span = tracing::info_span!( + target: "rivetkit::telemetry", + parent: None, + "rivet.workflow.step", + otel.name = %format!("{}/{}", self.inner.identity.actor_name, step_name), + otel.kind = "internal", + rivet.workflow.step.name = %step_name, + rivet.workflow.step.attempt = attempt, + rivet.workflow.step.outcome = tracing::field::Empty, + rivet.ray.id = self.inner.ray_id.as_deref(), + rivet.actor.id = %self.inner.identity.actor_id, + rivet.actor.name = %self.inner.identity.actor_name, + rivet.actor.key = %self.inner.identity.actor_key, + otel.status_code = tracing::field::Empty, + error.type = tracing::field::Empty, + ); + span.set_parent(parent); + let step_telemetry = Self { + inner: self.inner.clone(), + application_span: None, + step_span: otel_span_context_of(&span), + }; + Some((span, step_telemetry)) + } + pub(crate) fn start_sqlite(&self, operation: SqliteOperation) -> Option { let parent = self.parent_context()?; let (span_name, operation_name) = operation.names(); @@ -637,6 +837,42 @@ impl Drop for OutboundCallInvocation { } } +impl WorkflowStepSpan { + pub(crate) fn start(ctx: &ActorContext, step_name: &str, attempt: u32) -> Option { + let (span, step_telemetry) = ctx + .invocation_telemetry()? + .start_workflow_step(step_name, attempt)?; + Some(Self { + ctx: ctx.clone().with_invocation_telemetry(Some(step_telemetry)), + span: Some(span), + }) + } + + /// The context the step's callback runs under. + pub fn ctx(&self) -> ActorContext { + self.ctx.clone() + } + + /// `error` is what the step threw. Its group and code become `error.type`. + pub fn finish(mut self, outcome: WorkflowStepOutcome, error: Option<&anyhow::Error>) { + let Some(span) = self.span.take() else { + return; + }; + span.record("rivet.workflow.step.outcome", outcome.as_label()); + record_outcome(&span, error); + } +} + +impl Drop for WorkflowStepSpan { + fn drop(&mut self) { + let Some(span) = self.span.take() else { + return; + }; + span.record("otel.status_code", "ERROR"); + span.record("error.type", OPERATION_ABANDONED_ERROR_TYPE); + } +} + impl SqliteOperationSpan { pub(crate) fn span(&self) -> tracing::Span { self.span.as_ref().expect("sqlite span is present").clone() @@ -701,12 +937,12 @@ fn w3c_trace_headers(span_context: &SpanContext) -> OwnedTraceHeaders { } /// An action or a raw HTTP request is entered from outside the actor, a -/// scheduled fire originates inside it, and a queue send produces a message -/// the actor consumes later. +/// scheduled fire and a workflow run originate inside it, and a queue send +/// produces a message the actor consumes later. fn otel_kind(invocation_type: InvocationType) -> &'static str { match invocation_type { InvocationType::Action | InvocationType::Request => "server", - InvocationType::Scheduled => "internal", + InvocationType::Scheduled | InvocationType::Workflow => "internal", InvocationType::QueueSend => "producer", } } @@ -747,7 +983,12 @@ pub(crate) fn start_queue_receive(ctx: &ActorContext, message: &QueueMessage) -> .invocation_telemetry() .and_then(|telemetry| telemetry.parent_context().map(|parent| (telemetry, parent))); if let Some((telemetry, parent)) = invocation_parent { - if let Some(ray_id) = telemetry.inner.ray_id.as_deref() { + let ray_id = telemetry + .inner + .ray_id + .as_deref() + .or(message.trace_context.ray_id.as_deref()); + if let Some(ray_id) = ray_id { span.record("rivet.ray.id", ray_id); } span.set_parent(parent); diff --git a/rivetkit-rust/packages/rivetkit-core/tests/sql_efficiency.rs b/rivetkit-rust/packages/rivetkit-core/tests/sql_efficiency.rs index eedba335de..faa559e36f 100644 --- a/rivetkit-rust/packages/rivetkit-core/tests/sql_efficiency.rs +++ b/rivetkit-rust/packages/rivetkit-core/tests/sql_efficiency.rs @@ -116,6 +116,9 @@ fn fixture(row_count: usize) -> Connection { let mut user_kv_insert = tx .prepare("INSERT INTO _rivet_user_kv (key, value) VALUES (?, x'01')") .expect("prepare user kv seed"); + let mut workflow_kv_insert = tx + .prepare("INSERT INTO _rivet_wf_kv (key, value) VALUES (?, x'01')") + .expect("prepare workflow kv seed"); let mut schedule_insert = tx .prepare("INSERT INTO _rivet_schedule_events (event_id, trigger_at, action, args, kind, cron_expression, timezone, interval_ms, last_started_at, max_history) VALUES (?, ?, 'run', x'01', ?, NULL, NULL, NULL, NULL, 100)") .expect("prepare schedule seed"); @@ -134,6 +137,9 @@ fn fixture(row_count: usize) -> Connection { user_kv_insert .execute([key.as_bytes()]) .expect("seed user kv"); + workflow_kv_insert + .execute([key.as_bytes()]) + .expect("seed workflow kv"); let event_id = if index % 3 == 0 { format!("at:{key}") } else { @@ -157,6 +163,7 @@ fn fixture(row_count: usize) -> Connection { drop(conn_state_insert); drop(queue_insert); drop(user_kv_insert); + drop(workflow_kv_insert); drop(schedule_insert); drop(history_insert); tx.commit().expect("commit fixture seed"); @@ -464,6 +471,14 @@ fn query_catalog() -> Vec { params: vec![], expectation: indexed(None, &["_rivet_runtime"]), }, + QueryCase { + id: "workflow.kv_get", + sql: internal_storage::LOAD_WORKFLOW_KV_SQL.into(), + params: vec![Value::Blob( + crate::actor::keys::WORKFLOW_TRACE_CONTEXT_KEY.to_vec(), + )], + expectation: indexed(None, &["_rivet_wf_kv"]), + }, QueryCase { id: "runtime.run_wake", sql: internal_storage::LOAD_RUN_WAKE_AT_SQL.into(), diff --git a/rivetkit-typescript/packages/rivetkit-napi/index.d.ts b/rivetkit-typescript/packages/rivetkit-napi/index.d.ts index 5d3db31820..4885ebda7f 100644 --- a/rivetkit-typescript/packages/rivetkit-napi/index.d.ts +++ b/rivetkit-typescript/packages/rivetkit-napi/index.d.ts @@ -338,6 +338,9 @@ export declare class ActorContext { * case the caller sends its own context as before. */ startCallSpan(actorName: string, actionName: string): OutboundCall | null + startWorkflowSpan(): Promise + /** Returns nothing outside a workflow run or when tracing is off. */ + startWorkflowStepSpan(stepName: string, attempt: number): WorkflowStepSpan | null provisionActorRuntimeSocket(): Promise schedule(): Schedule queue(): Queue @@ -407,6 +410,22 @@ export declare class OutboundCall { */ finish(error?: string | undefined | null): void } +/** + * One open workflow run. Collecting it without `finish` records the run as + * abandoned. + */ +export declare class WorkflowSpan { + /** The context workflow code runs under. */ + ctx(): ActorContext + finish(outcome: string): Promise +} +/** One open attempt at one workflow step. */ +export declare class WorkflowStepSpan { + /** The context the step's callback runs under. */ + ctx(): ActorContext + /** `error` is what the step threw, as the bridge encodes it. */ + finish(outcome: string, error?: string | undefined | null): void +} export declare class NapiActorFactory { constructor(callbacks: object, config?: JsActorConfig | undefined | null) } diff --git a/rivetkit-typescript/packages/rivetkit-napi/index.js b/rivetkit-typescript/packages/rivetkit-napi/index.js index 36aa76ee75..0d39d47667 100644 --- a/rivetkit-typescript/packages/rivetkit-napi/index.js +++ b/rivetkit-typescript/packages/rivetkit-napi/index.js @@ -310,12 +310,14 @@ if (!nativeBinding) { throw new Error(`Failed to load native binding`) } -const { ActorContext, decodeInspectorRequest, encodeInspectorResponse, OutboundCall, NapiActorFactory, CancellationToken, ConnHandle, JsNativeDatabase, JsSqliteTransaction, JsActorStateTransaction, HttpResponseBodyStream, HttpRequestBodyStream, Kv, Queue, QueueMessage, CoreRegistry, Schedule, WebSocket, setTelemetryLogSink, shutdownTelemetry } = nativeBinding +const { ActorContext, decodeInspectorRequest, encodeInspectorResponse, OutboundCall, WorkflowSpan, WorkflowStepSpan, NapiActorFactory, CancellationToken, ConnHandle, JsNativeDatabase, JsSqliteTransaction, JsActorStateTransaction, HttpResponseBodyStream, HttpRequestBodyStream, Kv, Queue, QueueMessage, CoreRegistry, Schedule, WebSocket, setTelemetryLogSink, shutdownTelemetry } = nativeBinding module.exports.ActorContext = ActorContext module.exports.decodeInspectorRequest = decodeInspectorRequest module.exports.encodeInspectorResponse = encodeInspectorResponse module.exports.OutboundCall = OutboundCall +module.exports.WorkflowSpan = WorkflowSpan +module.exports.WorkflowStepSpan = WorkflowStepSpan module.exports.NapiActorFactory = NapiActorFactory module.exports.CancellationToken = CancellationToken module.exports.ConnHandle = ConnHandle diff --git a/rivetkit-typescript/packages/rivetkit-napi/src/actor_context.rs b/rivetkit-typescript/packages/rivetkit-napi/src/actor_context.rs index 4a35b25c51..765a7593c0 100644 --- a/rivetkit-typescript/packages/rivetkit-napi/src/actor_context.rs +++ b/rivetkit-typescript/packages/rivetkit-napi/src/actor_context.rs @@ -22,6 +22,8 @@ use rivetkit_core::{ ActorWorkKind, ConnHandle as CoreConnHandle, KeepAwakeRegion, OutboundCallInvocation as CoreOutboundCallInvocation, Request as CoreRequest, RequestSaveOpts, StateDelta, WebSocketCallbackRegion, WorkflowKvWrite, + WorkflowRunInvocation as CoreWorkflowRunInvocation, WorkflowRunOutcome, WorkflowStepOutcome, + WorkflowStepSpan as CoreWorkflowStepSpan, }; use scc::HashMap as SccHashMap; use tokio::sync::mpsc::UnboundedSender; @@ -354,6 +356,32 @@ impl ActorContext { }) } + #[napi] + pub async fn start_workflow_span(&self) -> WorkflowSpan { + let span = self.inner.start_workflow_span().await; + WorkflowSpan { + ctx: span.ctx(), + shared: self.shared.clone(), + span: Mutex::new(Some(span)), + } + } + + /// Returns nothing outside a workflow run or when tracing is off. + #[napi] + pub fn start_workflow_step_span( + &self, + step_name: String, + attempt: u32, + ) -> Option { + self.inner + .start_workflow_step_span(&step_name, attempt) + .map(|step| WorkflowStepSpan { + ctx: step.ctx(), + shared: self.shared.clone(), + step: Some(step), + }) + } + #[napi] pub async fn provision_actor_runtime_socket( &self, @@ -1167,3 +1195,66 @@ impl OutboundCall { invocation.finish(error.as_ref()); } } + +/// One open workflow run. Collecting it without `finish` records the run as +/// abandoned. +#[napi] +pub struct WorkflowSpan { + ctx: CoreActorContext, + shared: Arc, + span: Mutex>, +} + +#[napi] +impl WorkflowSpan { + /// The context workflow code runs under. + #[napi] + pub fn ctx(&self) -> ActorContext { + ActorContext { + inner: self.ctx.clone(), + shared: self.shared.clone(), + } + } + + #[napi] + pub async fn finish(&self, outcome: String) -> napi::Result<()> { + let outcome = WorkflowRunOutcome::parse(&outcome).map_err(napi_anyhow_error)?; + let Some(span) = self.span.lock().take() else { + return Ok(()); + }; + self.ctx.finish_workflow_span(span, outcome).await; + Ok(()) + } +} + +/// One open attempt at one workflow step. +#[napi] +pub struct WorkflowStepSpan { + ctx: CoreActorContext, + shared: Arc, + step: Option, +} + +#[napi] +impl WorkflowStepSpan { + /// The context the step's callback runs under. + #[napi] + pub fn ctx(&self) -> ActorContext { + ActorContext { + inner: self.ctx.clone(), + shared: self.shared.clone(), + } + } + + /// `error` is what the step threw, as the bridge encodes it. + #[napi] + pub fn finish(&mut self, outcome: String, error: Option) -> napi::Result<()> { + let outcome = WorkflowStepOutcome::parse(&outcome).map_err(napi_anyhow_error)?; + let Some(step) = self.step.take() else { + return Ok(()); + }; + let error = error.map(anyhow_error_from_js_reason); + step.finish(outcome, error.as_ref()); + Ok(()) + } +} diff --git a/rivetkit-typescript/packages/rivetkit/fixtures/driver-test-suite/registry-static.ts b/rivetkit-typescript/packages/rivetkit/fixtures/driver-test-suite/registry-static.ts index f928fbd6f5..d14589a4de 100644 --- a/rivetkit-typescript/packages/rivetkit/fixtures/driver-test-suite/registry-static.ts +++ b/rivetkit-typescript/packages/rivetkit/fixtures/driver-test-suite/registry-static.ts @@ -152,7 +152,11 @@ import { import { lifecycleObserver, startStopRaceActor } from "./start-stop-race"; import { stateZodCoercionActor } from "./state-zod-coercion"; import { statelessActor } from "./stateless"; -import { telemetryActor, telemetryRunConsumerActor } from "./telemetry"; +import { + telemetryActor, + telemetryRunConsumerActor, + workflowTracedActor, +} from "./telemetry"; import { driverCtxActor, dynamicVarActor, @@ -201,6 +205,7 @@ export const registry = setup({ use: { telemetryActor, telemetryRunConsumerActor, + workflowTracedActor, // From counter.ts counter, // From counter-conn.ts diff --git a/rivetkit-typescript/packages/rivetkit/fixtures/driver-test-suite/telemetry.ts b/rivetkit-typescript/packages/rivetkit/fixtures/driver-test-suite/telemetry.ts index aab1c54151..0ba4cf6f31 100644 --- a/rivetkit-typescript/packages/rivetkit/fixtures/driver-test-suite/telemetry.ts +++ b/rivetkit-typescript/packages/rivetkit/fixtures/driver-test-suite/telemetry.ts @@ -2,6 +2,7 @@ import { trace } from "@opentelemetry/api"; import { NodeTracerProvider } from "@opentelemetry/sdk-trace-node"; import { actor, queue, UserError } from "rivetkit"; import { db } from "@/common/database/mod"; +import { workflow } from "@/workflow/mod"; // Only a traced runtime gets a JavaScript tracer, so the other driver // fixtures keep running without an OpenTelemetry context manager. @@ -121,3 +122,53 @@ export const telemetryRunConsumerActor = actor({ }, actions: {}, }); + +/** Waits for two messages and logs when it sleeps, so a test can send the second one after the sleep. */ +export const workflowTracedActor = actor({ + state: { chargeAttempts: 0, wakes: 0 }, + db: db(), + onWake: (c) => { + c.state.wakes += 1; + }, + queues: { + approve: jobSchema, + resume: jobSchema, + }, + onSleep: (c) => { + c.log.warn( + { slept_actor_key: c.key[0] }, + "workflow traced actor slept", + ); + }, + run: workflow(async (ctx) => { + await ctx.queue.next("wait-approve", { names: ["approve"] }); + await ctx.step("reserve-stock", async (c) => { + await c.db.execute("SELECT 'reserve-stock' AS step"); + }); + await ctx.queue.next("wait-resume", { names: ["resume"] }); + await ctx.step({ + name: "charge-card", + maxRetries: 3, + retryBackoffBase: 10, + retryBackoffMax: 10, + run: async (c) => { + c.state.chargeAttempts += 1; + if (c.state.chargeAttempts <= 2) { + throw new UserError("card declined", { + code: "card_declined", + }); + } + }, + }); + await ctx.step("notify", async (c) => { + const client = c.client(); + await client.telemetryRunConsumerActor + .getOrCreate(c.key) + .send("runJobs", { id: "workflow-notify" }); + }); + }), + actions: { getWakes: (c) => c.state.wakes }, + options: { + sleepTimeout: 50, + }, +}); diff --git a/rivetkit-typescript/packages/rivetkit/src/actor/config.ts b/rivetkit-typescript/packages/rivetkit/src/actor/config.ts index 66bb857800..a83473a35d 100644 --- a/rivetkit-typescript/packages/rivetkit/src/actor/config.ts +++ b/rivetkit-typescript/packages/rivetkit/src/actor/config.ts @@ -447,6 +447,37 @@ export interface ActorContext< [key: string]: any; } +/** @experimental How one run of a workflow function ended. */ +export type WorkflowOutcome = + | "completed" + | "sleeping" + | "evicted" + | "failed" + | "cancelled"; + +/** @experimental How one workflow step attempt ended. `retry` will be tried again, `failed` will not. */ +export type WorkflowStepOutcome = "ok" | "retry" | "failed"; + +/** @experimental One open run of a workflow function, traced as its own invocation. */ +export interface WorkflowSpan { + /** Runs `body` attributed to this run. */ + run(body: () => Promise): Promise; + /** + * Called by a workflow engine before each attempt at a step. Returns + * nothing when tracing is off. + */ + startStep(name: string, attempt: number): WorkflowStepSpan | undefined; + finish(outcome: WorkflowOutcome): Promise; +} + +/** @experimental One open attempt at one workflow step. */ +export interface WorkflowStepSpan { + /** Runs `body` attributed to this attempt. */ + run(body: () => Promise): Promise; + /** `error` is what the step callback threw. */ + finish(outcome: WorkflowStepOutcome, error?: unknown): void; +} + /** @experimental */ export interface ActorRun { /** @@ -454,6 +485,13 @@ export interface ActorRun { * ensures this actor's run handler is active, starting it if inactive. */ setWakeAt(timestamp: number | null): Promise; + /** @experimental Called by a workflow engine each time it runs the workflow function. */ + startWorkflowSpan(): Promise; + /** + * @experimental Runs `body` outside the current workflow run, so a workflow + * engine's own history storage stays out of the trace. + */ + runOutsideWorkflowSpan(body: () => T): T; } export type ActionContext< diff --git a/rivetkit-typescript/packages/rivetkit/src/registry/napi-runtime.ts b/rivetkit-typescript/packages/rivetkit/src/registry/napi-runtime.ts index d3f9d2044f..c83eaead5a 100644 --- a/rivetkit-typescript/packages/rivetkit/src/registry/napi-runtime.ts +++ b/rivetkit-typescript/packages/rivetkit/src/registry/napi-runtime.ts @@ -8,6 +8,8 @@ import type { HttpResponseBodyStream as NativeHttpResponseBodyStream, WebSocket as NativeWebSocket, } from "@rivetkit/rivetkit-napi"; +import type { WorkflowSpan } from "@/actor/config"; +import { encodeErrorForBridge } from "@/actor/errors"; import type { ActorInvocationTraceContext } from "@/common/actor-telemetry-context"; import { readActiveTraceHeaders, @@ -599,7 +601,10 @@ export class NapiCoreRuntime implements CoreRuntime { } runWithActorInvocationContext(ctx: ActorContextHandle, run: () => T): T { - const nativeCtx = asNativeActorContext(ctx); + return this.#runAs(asNativeActorContext(ctx), run); + } + + #runAs(nativeCtx: NativeActorContext, run: () => T): T { const span = nativeCtx.invocationTraceContext()?.span; return this.#invocationContext.run(nativeCtx, () => runWithActorInvocationSpan(span, run), @@ -631,6 +636,34 @@ export class NapiCoreRuntime implements CoreRuntime { }; } + async startWorkflowSpan(ctx: ActorContextHandle): Promise { + const span = await asNativeActorContext(ctx).startWorkflowSpan(); + const spanCtx = span.ctx(); + return { + run: (body) => this.#runAs(spanCtx, body), + startStep: (name, attempt) => { + const step = spanCtx.startWorkflowStepSpan(name, attempt); + if (!step) return undefined; + const stepCtx = step.ctx(); + return { + run: (body) => this.#runAs(stepCtx, body), + finish: (outcome, error) => + step.finish( + outcome, + error === undefined + ? undefined + : encodeErrorForBridge(error), + ), + }; + }, + finish: (outcome) => span.finish(outcome), + }; + } + + runOutsideActorInvocationContext(run: () => T): T { + return this.#invocationContext.exit(run); + } + actorName(ctx: ActorContextHandle): string { return asNativeActorContext(ctx).name(); } diff --git a/rivetkit-typescript/packages/rivetkit/src/registry/native.ts b/rivetkit-typescript/packages/rivetkit/src/registry/native.ts index 776bb8edab..6eb0dc71b6 100644 --- a/rivetkit-typescript/packages/rivetkit/src/registry/native.ts +++ b/rivetkit-typescript/packages/rivetkit/src/registry/native.ts @@ -9,6 +9,7 @@ import { type ActorCronEveryOptions, type ActorCronSetOptions, type ActorLogger, + type ActorRun, type ActorSchedule, CONN_STATE_MANAGER_SYMBOL, type CronFire, @@ -2726,7 +2727,7 @@ export class ActorContextHandleAdapter { #queue?: NativeQueueAdapter; #request?: Request; #schedule?: NativeScheduleAdapter; - #run?: { setWakeAt(timestamp: number | null): Promise }; + #run?: ActorRun; #runHandlerConfigured: boolean; #onStateChange?: NativeOnStateChangeHandler; #stateEnabled: boolean; @@ -2907,6 +2908,12 @@ export class ActorContextHandleAdapter { this.#runtime.actorSetRunWakeAt(this.#ctx, timestamp), ); }, + startWorkflowSpan: () => + callNative(() => + this.#runtime.startWorkflowSpan(this.#ctx), + ), + runOutsideWorkflowSpan: (body) => + this.#runtime.runOutsideActorInvocationContext(body), }; } return this.#run; diff --git a/rivetkit-typescript/packages/rivetkit/src/registry/runtime.ts b/rivetkit-typescript/packages/rivetkit/src/registry/runtime.ts index 2a26b034d2..fd5b4e4476 100644 --- a/rivetkit-typescript/packages/rivetkit/src/registry/runtime.ts +++ b/rivetkit-typescript/packages/rivetkit/src/registry/runtime.ts @@ -1,3 +1,4 @@ +import type { WorkflowSpan } from "@/actor/config"; import type { ActorInvocationSpanContext, ActorInvocationTraceContext, @@ -599,6 +600,10 @@ export interface CoreRuntime { actorName: string, actionName: string, ): RuntimeOutboundCall | undefined; + /** Backs `ActorRun.startWorkflowSpan`. */ + startWorkflowSpan(ctx: ActorContextHandle): Promise; + /** Backs `ActorRun.runOutsideWorkflowSpan`. */ + runOutsideActorInvocationContext(run: () => T): T; actorName(ctx: ActorContextHandle): string; actorKey(ctx: ActorContextHandle): RuntimeActorKeySegment[]; actorRegion(ctx: ActorContextHandle): string; diff --git a/rivetkit-typescript/packages/rivetkit/src/registry/wasm-runtime.ts b/rivetkit-typescript/packages/rivetkit/src/registry/wasm-runtime.ts index 2e6d2a00bb..622c7730ee 100644 --- a/rivetkit-typescript/packages/rivetkit/src/registry/wasm-runtime.ts +++ b/rivetkit-typescript/packages/rivetkit/src/registry/wasm-runtime.ts @@ -1,3 +1,4 @@ +import type { WorkflowSpan } from "@/actor/config"; import { decodeBridgeRivetError, RivetError } from "@/actor/errors"; import type { ActorInvocationTraceContext } from "@/common/actor-telemetry-context"; import type { @@ -558,6 +559,18 @@ export class WasmCoreRuntime implements CoreRuntime { return undefined; } + async startWorkflowSpan(_ctx: ActorContextHandle): Promise { + return { + run: (body) => body(), + startStep: () => undefined, + finish: async () => {}, + }; + } + + runOutsideActorInvocationContext(run: () => T): T { + return run(); + } + actorName(ctx: ActorContextHandle): string { return callHandle(asWasmActorContext(ctx), "name"); } diff --git a/rivetkit-typescript/packages/rivetkit/src/workflow/driver.ts b/rivetkit-typescript/packages/rivetkit/src/workflow/driver.ts index feb3b2cc9d..4d7ebd4422 100644 --- a/rivetkit-typescript/packages/rivetkit/src/workflow/driver.ts +++ b/rivetkit-typescript/packages/rivetkit/src/workflow/driver.ts @@ -4,6 +4,7 @@ import type { KVWrite, Message, WorkflowMessageDriver, + WorkflowTelemetryDriver, } from "@rivetkit/workflow-engine"; import type { RunContext } from "@/actor/config"; import type { AnyStaticActorInstance } from "@/actor/definition"; @@ -323,6 +324,7 @@ export class ActorWorkflowDriver implements EngineDriver { readonly atomicBatch = true; readonly workerPollInterval = 100; readonly messageDriver: WorkflowMessageDriver; + readonly telemetry: WorkflowTelemetryDriver; #actor: AnyStaticActorInstance; #runCtx: RunContext; #storage: WorkflowStorage; @@ -334,45 +336,51 @@ export class ActorWorkflowDriver implements EngineDriver { this.#actor = actor; this.#runCtx = runCtx; this.messageDriver = new ActorWorkflowMessageDriver(actor, runCtx); + this.telemetry = { + startSpan: () => runCtx.run.startWorkflowSpan(), + }; this.#storage = new WorkflowStorage(runtimeDbFromContext(runCtx)); } + /** Keeps the engine's history storage out of the workflow run's trace. */ + #untraced(run: () => Promise): Promise { + return this.#runCtx.internalKeepAwake( + this.#runCtx.run.runOutsideWorkflowSpan(run), + ); + } + async get(key: Uint8Array): Promise { - return await this.#runCtx.internalKeepAwake(this.#storage.get(key)); + return await this.#untraced(() => this.#storage.get(key)); } async set(key: Uint8Array, value: Uint8Array): Promise { - await this.#runCtx.internalKeepAwake(this.#storage.set(key, value)); + await this.#untraced(() => this.#storage.set(key, value)); } async delete(key: Uint8Array): Promise { - await this.#runCtx.internalKeepAwake(this.#storage.delete(key)); + await this.#untraced(() => this.#storage.delete(key)); } async batchDelete(keys: Uint8Array[]): Promise { - await this.#runCtx.internalKeepAwake(this.#storage.batchDelete(keys)); + await this.#untraced(() => this.#storage.batchDelete(keys)); } async deletePrefix(prefix: Uint8Array): Promise { - await this.#runCtx.internalKeepAwake( - this.#storage.deletePrefix(prefix), - ); + await this.#untraced(() => this.#storage.deletePrefix(prefix)); } async deleteRange(start: Uint8Array, end: Uint8Array): Promise { - await this.#runCtx.internalKeepAwake( - this.#storage.deleteRange(start, end), - ); + await this.#untraced(() => this.#storage.deleteRange(start, end)); } async list(prefix: Uint8Array): Promise { - return await this.#runCtx.internalKeepAwake(this.#storage.list(prefix)); + return await this.#untraced(() => this.#storage.list(prefix)); } async batch(writes: KVWrite[]): Promise { if (writes.length === 0) return; - await this.#runCtx.internalKeepAwake( + await this.#untraced(() => this.#actor.stateManager.saveStateAndWorkflowBatch(writes), ); } diff --git a/rivetkit-typescript/packages/rivetkit/tests/driver/actor-telemetry.test.ts b/rivetkit-typescript/packages/rivetkit/tests/driver/actor-telemetry.test.ts index ef96f2308c..ea91115492 100644 --- a/rivetkit-typescript/packages/rivetkit/tests/driver/actor-telemetry.test.ts +++ b/rivetkit-typescript/packages/rivetkit/tests/driver/actor-telemetry.test.ts @@ -869,6 +869,183 @@ describeDriverMatrix( } } }, 60_000); + describe("workflow", () => { + const actorName = "workflowTracedActor"; + let actorKey: string; + let approveRayId: string; + let resumeRayId: string; + let spans: ExportedSpan[]; + let runs: ExportedSpan[]; + + const named = (name: string) => + spans + .filter( + (span) => + span.name === name && + span.attributes["rivet.actor.key"] === actorKey, + ) + .sort((a, b) => + Number(a.endTimeUnixNano - b.endTimeUnixNano), + ); + const attribute = (found: ExportedSpan[], key: string) => + found.map((span) => span.attributes[key]); + + beforeAll(async () => { + actorKey = `workflow-traced-${crypto.randomUUID()}`; + approveRayId = `approve-${crypto.randomUUID().slice(0, 8)}`; + resumeRayId = `resume-${crypto.randomUUID().slice(0, 8)}`; + const workflowActor = + traced.client.workflowTracedActor.getOrCreate([ + actorKey, + ]); + await withRayBaggage(approveRayId, () => + workflowActor.send("approve", { id: "approve" }), + ); + // The actor process logs from its sleep hook, and calling an action + // to ask would count as activity and keep the actor awake. + await vi.waitFor( + () => { + expect( + traced.runtime.getRuntimeOutput?.(), + ).toContain(`slept_actor_key=${actorKey}`); + }, + { timeout: 30_000, interval: 100 }, + ); + await withRayBaggage(resumeRayId, () => + workflowActor.send("resume", { id: "resume" }), + ); + spans = await waitForSpans( + traceExports, + "the completed workflow run and the queue send it caused", + (exported) => { + spans = exported; + return ( + attribute( + named(`${actorName}/workflow`), + "rivet.workflow.run.outcome", + ).includes("completed") && + named("telemetryRunConsumerActor/queue.send") + .length > 0 + ); + }, + 30_000, + ); + runs = named(`${actorName}/workflow`); + }, 90_000); + + test("chains every run to the one before it across the actor sleeping", async () => { + const wakes = await traced.client.workflowTracedActor + .getOrCreate([actorKey]) + .getWakes(); + expect(wakes).toBeGreaterThanOrEqual(2); + expect(runs.length).toBeGreaterThanOrEqual(4); + expect(runs[0].links).toEqual([]); + for (const [index, run] of runs.entries()) { + expect(run.parentSpanId).toBeUndefined(); + expect(run.attributes["rivet.invocation.type"]).toBe( + "workflow", + ); + if (index === 0) continue; + const previous = runs[index - 1]; + expect(run.traceId).not.toBe(previous.traceId); + expect(run.links).toEqual([ + { + traceId: previous.traceId, + spanId: previous.spanId, + }, + ]); + } + const outcomes = attribute( + runs, + "rivet.workflow.run.outcome", + ); + const completedAt = outcomes.indexOf("completed"); + expect(new Set(outcomes.slice(0, completedAt))).toEqual( + new Set(["sleeping"]), + ); + expect(runs[completedAt].statusCode).toBe(OTLP_STATUS_OK); + }); + + test("reports a completed step once, with only its own work under it", () => { + const reserve = named(`${actorName}/reserve-stock`); + expect(reserve).toHaveLength(1); + expect(reserve[0].attributes).toMatchObject({ + "rivet.workflow.step.name": "reserve-stock", + "rivet.workflow.step.attempt": "1", + "rivet.workflow.step.outcome": "ok", + }); + const sqliteParents = spans + .filter((span) => span.name.startsWith("rivet.sqlite.")) + .map((span) => span.parentSpanId); + expect( + sqliteParents.filter((id) => id === reserve[0].spanId), + ).toHaveLength(1); + for (const run of runs) { + expect(sqliteParents).not.toContain(run.spanId); + } + }); + + test("reports each attempt of a retried step under its own run", () => { + const attempts = named(`${actorName}/charge-card`); + expect( + attempts.map((attempt) => [ + attempt.attributes["rivet.workflow.step.attempt"], + attempt.attributes["rivet.workflow.step.outcome"], + attempt.attributes["error.type"], + attempt.statusCode, + ]), + ).toEqual([ + ["1", "retry", "user.card_declined", OTLP_STATUS_ERROR], + ["2", "retry", "user.card_declined", OTLP_STATUS_ERROR], + ["3", "ok", undefined, OTLP_STATUS_OK], + ]); + const runIds = runs.map((run) => run.spanId); + const parents = attempts.map( + (attempt) => attempt.parentSpanId, + ); + expect(new Set(parents).size).toBe(3); + for (const parent of parents) { + expect(runIds).toContain(parent); + } + }); + + test("keeps a run's ray unchanged when it receives a message sent under another ray", () => { + const ray = (name: string) => + attribute(named(name), "rivet.ray.id"); + const sends = named(`${actorName}/queue.send`); + for (const receive of named(`${actorName}/queue.receive`)) { + const send = sends.find( + (candidate) => + candidate.attributes["rivet.ray.id"] === + receive.attributes["rivet.ray.id"], + ); + expect(receive.links).toEqual([ + { traceId: send?.traceId, spanId: send?.spanId }, + ]); + } + expect(ray(`${actorName}/queue.receive`).sort()).toEqual( + [approveRayId, resumeRayId].sort(), + ); + expect(new Set(ray(`${actorName}/workflow`))).toEqual( + new Set([undefined]), + ); + expect(ray(`${actorName}/reserve-stock`)).toEqual([undefined]); + expect(ray(`${actorName}/charge-card`)).toEqual([ + undefined, + undefined, + undefined, + ]); + }); + + test("puts a queue send made inside a step in that step's trace", () => { + const [notify] = named(`${actorName}/notify`); + const [send] = named( + "telemetryRunConsumerActor/queue.send", + ); + expect(send.traceId).toBe(notify.traceId); + expect(send.parentSpanId).toBe(notify.spanId); + }); + }); }); }, { diff --git a/rivetkit-typescript/packages/workflow-engine/src/context.ts b/rivetkit-typescript/packages/workflow-engine/src/context.ts index c4ee278bf7..8ba8eb63ce 100644 --- a/rivetkit-typescript/packages/workflow-engine/src/context.ts +++ b/rivetkit-typescript/packages/workflow-engine/src/context.ts @@ -1,5 +1,5 @@ import type { Logger } from "pino"; -import type { EngineDriver } from "./driver.js"; +import type { EngineDriver, WorkflowSpan } from "./driver.js"; import { extractErrorInfo, getErrorEventTag, @@ -345,6 +345,7 @@ export class WorkflowContextImpl implements WorkflowContextInterface { private abortController: AbortController; private currentLocation: Location; private visitedKeys: Set; + private runSpan?: WorkflowSpan; private mode: "forward" | "rollback"; private rollbackActions?: RollbackAction[]; private rollbackCheckpointSet: boolean; @@ -369,6 +370,7 @@ export class WorkflowContextImpl implements WorkflowContextInterface { onError?: WorkflowErrorHandler, logger?: Logger, visitedKeys?: Set, + runSpan?: WorkflowSpan, ) { this.currentLocation = location; this.abortController = abortController ?? new AbortController(); @@ -379,6 +381,7 @@ export class WorkflowContextImpl implements WorkflowContextInterface { this.onError = onError; this.logger = logger; this.visitedKeys = visitedKeys ?? new Set(); + this.runSpan = runSpan; } get abortSignal(): AbortSignal { @@ -435,6 +438,7 @@ export class WorkflowContextImpl implements WorkflowContextInterface { this.onError, this.logger, this.visitedKeys, + this.runSpan, ); } @@ -906,13 +910,19 @@ export class WorkflowContextImpl implements WorkflowContextInterface { // Get timeout configuration const timeout = config.timeout ?? DEFAULT_STEP_TIMEOUT; + const stepSpan = this.runSpan?.startStep( + config.name, + metadata.attempts, + ); + try { // Execute with timeout const output = await this.executeWithTimeout( - config.run(), + stepSpan ? stepSpan.run(() => config.run()) : config.run(), timeout, config.name, ); + stepSpan?.finish("ok"); if (entry.kind.type === "step") { entry.kind.data.output = output; @@ -951,6 +961,7 @@ export class WorkflowContextImpl implements WorkflowContextInterface { // Timeout errors are treated as critical by default. Steps opt // into retrying on timeout with retryOnTimeout: true. if (error instanceof StepTimeoutError && !config.retryOnTimeout) { + stepSpan?.finish("failed", error); metadata.status = "exhausted"; metadata.error = String(error); await this.notifyStepError(config, metadata.attempts, error, { @@ -970,6 +981,7 @@ export class WorkflowContextImpl implements WorkflowContextInterface { error instanceof CriticalError || error instanceof RollbackError ) { + stepSpan?.finish("failed", error); metadata.status = "exhausted"; metadata.error = String(error); await this.notifyStepError(config, metadata.attempts, error, { @@ -989,6 +1001,7 @@ export class WorkflowContextImpl implements WorkflowContextInterface { } const willRetry = metadata.attempts <= maxRetries; + stepSpan?.finish(willRetry ? "retry" : "failed", error); metadata.status = willRetry ? "failed" : "exhausted"; metadata.error = String(error); diff --git a/rivetkit-typescript/packages/workflow-engine/src/driver.ts b/rivetkit-typescript/packages/workflow-engine/src/driver.ts index c9905a87ba..48dfffc59c 100644 --- a/rivetkit-typescript/packages/workflow-engine/src/driver.ts +++ b/rivetkit-typescript/packages/workflow-engine/src/driver.ts @@ -16,6 +16,42 @@ export interface KVWrite { value: Uint8Array; } +/** How one run of the workflow function ended. */ +export type WorkflowOutcome = + | "completed" + | "sleeping" + | "evicted" + | "failed" + | "cancelled"; + +/** How one step attempt ended. `retry` will be tried again, `failed` will not. */ +export type WorkflowStepOutcome = "ok" | "retry" | "failed"; + +/** One open run of the workflow function. */ +export interface WorkflowSpan { + /** Runs `body` attributed to this run. */ + run(body: () => Promise): Promise; + /** Returns nothing when the attempt is not being traced. */ + startStep(name: string, attempt: number): WorkflowStepSpan | undefined; + finish(outcome: WorkflowOutcome): Promise; +} + +/** One open attempt at one step. */ +export interface WorkflowStepSpan { + /** Runs `body` attributed to this attempt. */ + run(body: () => Promise): Promise; + /** `error` is what the step callback threw. */ + finish(outcome: WorkflowStepOutcome, error?: unknown): void; +} + +/** + * Receives the start and end of each workflow run and step attempt. A + * completed step replayed from history is not an attempt and is not reported. + */ +export interface WorkflowTelemetryDriver { + startSpan(): Promise; +} + /** * The engine driver provides the KV and scheduling interface. * Implementations must provide these methods to integrate with different backends. @@ -117,4 +153,7 @@ export interface EngineDriver { messageNames: string[], abortSignal: AbortSignal, ): Promise; + + /** Present when the host traces workflows. */ + readonly telemetry?: WorkflowTelemetryDriver; } diff --git a/rivetkit-typescript/packages/workflow-engine/src/index.ts b/rivetkit-typescript/packages/workflow-engine/src/index.ts index 1aea778736..88a19e2b8a 100644 --- a/rivetkit-typescript/packages/workflow-engine/src/index.ts +++ b/rivetkit-typescript/packages/workflow-engine/src/index.ts @@ -12,7 +12,16 @@ export { WorkflowContextImpl, } from "./context.js"; // Driver -export type { EngineDriver, KVEntry, KVWrite } from "./driver.js"; +export type { + EngineDriver, + KVEntry, + KVWrite, + WorkflowOutcome, + WorkflowSpan, + WorkflowStepOutcome, + WorkflowStepSpan, + WorkflowTelemetryDriver, +} from "./driver.js"; export { extractErrorInfo } from "./error-utils.js"; // Errors export { @@ -143,7 +152,7 @@ import { } from "../schemas/serde.js"; import { type RollbackAction, WorkflowContextImpl } from "./context.js"; // Main workflow runner -import type { EngineDriver } from "./driver.js"; +import type { EngineDriver, WorkflowOutcome, WorkflowSpan } from "./driver.js"; import { extractErrorInfo, getErrorEventTag, @@ -288,6 +297,7 @@ async function executeRollback( historyNotifier?: HistoryNotifier, onError?: RunWorkflowOptions["onError"], logger?: Logger, + runSpan?: WorkflowSpan, ): Promise { const rollbackActions: RollbackAction[] = []; const ctx = new WorkflowContextImpl( @@ -303,6 +313,8 @@ async function executeRollback( historyNotifier, onError, logger, + undefined, + runSpan, ); try { @@ -511,16 +523,19 @@ async function executeLiveWorkflow( let lastResult: WorkflowResult | undefined; while (true) { - const result = await executeWorkflow( - workflowId, - workflowFn, - input, - driver, - messageDriver, - abortController, - onHistoryUpdated, - onError, - logger, + const result = await traceRun(driver, (runSpan) => + executeWorkflow( + workflowId, + workflowFn, + input, + driver, + messageDriver, + abortController, + onHistoryUpdated, + onError, + logger, + runSpan, + ), ); lastResult = result; @@ -642,16 +657,19 @@ export function runWorkflow( options.onError, logger, ) - : executeWorkflow( - workflowId, - workflowFn, - input, - driver, - messageDriver, - abortController, - options.onHistoryUpdated, - options.onError, - logger, + : traceRun(driver, (runSpan) => + executeWorkflow( + workflowId, + workflowFn, + input, + driver, + messageDriver, + abortController, + options.onHistoryUpdated, + options.onError, + logger, + runSpan, + ), ); return { @@ -886,6 +904,49 @@ function findReplayBoundaryEntry( return boundary; } +function runOutcomeFromState(state: WorkflowState): WorkflowOutcome { + switch (state) { + case "completed": + return "completed"; + case "sleeping": + return "sleeping"; + case "failed": + return "failed"; + case "cancelled": + return "cancelled"; + case "pending": + case "running": + case "rolling_back": + return "evicted"; + } +} + +/** Internal: Report one workflow run to the driver's telemetry, when it has any. */ +async function traceRun( + driver: EngineDriver, + run: (span?: WorkflowSpan) => Promise>, +): Promise> { + if (!driver.telemetry) { + return await run(); + } + + const span = await driver.telemetry.startSpan(); + let outcome: WorkflowOutcome = "failed"; + try { + const result = await span.run(() => run(span)); + outcome = runOutcomeFromState(result.state); + return result; + } catch (error) { + // A run throws EvictedError only when the workflow was already cancelled. + if (error instanceof EvictedError) { + outcome = "cancelled"; + } + throw error; + } finally { + await span.finish(outcome); + } +} + /** * Internal: Execute the workflow and return the result. */ @@ -899,6 +960,7 @@ async function executeWorkflow( onHistoryUpdated?: (history: WorkflowHistorySnapshot) => void, onError?: RunWorkflowOptions["onError"], logger?: Logger, + runSpan?: WorkflowSpan, ): Promise> { const storage = await loadStorage(driver); const historyNotifier: HistoryNotifier = onHistoryUpdated @@ -953,6 +1015,7 @@ async function executeWorkflow( historyNotifier, onError, logger, + runSpan, ); } catch (error) { if (error instanceof EvictedError) { @@ -986,6 +1049,8 @@ async function executeWorkflow( historyNotifier, onError, logger, + undefined, + runSpan, ); storage.state = "running"; @@ -1064,6 +1129,7 @@ async function executeWorkflow( historyNotifier, onError, logger, + runSpan, ); } catch (rollbackError) { if (rollbackError instanceof EvictedError) { diff --git a/rivetkit-typescript/packages/workflow-engine/src/keys.ts b/rivetkit-typescript/packages/workflow-engine/src/keys.ts index d5e5b83c1e..a83e2ac01d 100644 --- a/rivetkit-typescript/packages/workflow-engine/src/keys.ts +++ b/rivetkit-typescript/packages/workflow-engine/src/keys.ts @@ -14,6 +14,7 @@ export const KEY_PREFIX = { HISTORY: 2, // History entries: [2, ...locationSegments] WORKFLOW: 3, // Workflow metadata: [3, field] ENTRY_METADATA: 4, // Entry metadata: [4, entryId] + // 5 is reserved for RivetKit's workflow trace context. } as const; // Workflow metadata field identifiers From 4d34551ca77ce1cb10e5fe6910375177230ebe71 Mon Sep 17 00:00:00 2001 From: Sree Narayanan Date: Mon, 21 Sep 2026 21:26:07 +0400 Subject: [PATCH 3/6] fix(rivetkit): add trace ids to logs written inside a workflow --- .../fixtures/driver-test-suite/telemetry.ts | 1 + .../rivetkit/src/registry/napi-runtime.ts | 4 ++ .../packages/rivetkit/src/registry/native.ts | 38 +++++++++++-------- .../packages/rivetkit/src/registry/runtime.ts | 2 + .../rivetkit/src/registry/wasm-runtime.ts | 4 ++ .../tests/driver/actor-telemetry.test.ts | 12 ++++++ 6 files changed, 46 insertions(+), 15 deletions(-) diff --git a/rivetkit-typescript/packages/rivetkit/fixtures/driver-test-suite/telemetry.ts b/rivetkit-typescript/packages/rivetkit/fixtures/driver-test-suite/telemetry.ts index 0ba4cf6f31..62390cf1cb 100644 --- a/rivetkit-typescript/packages/rivetkit/fixtures/driver-test-suite/telemetry.ts +++ b/rivetkit-typescript/packages/rivetkit/fixtures/driver-test-suite/telemetry.ts @@ -143,6 +143,7 @@ export const workflowTracedActor = actor({ run: workflow(async (ctx) => { await ctx.queue.next("wait-approve", { names: ["approve"] }); await ctx.step("reserve-stock", async (c) => { + c.log.warn({ workflow_log_key: c.key[0] }, "reserving stock"); await c.db.execute("SELECT 'reserve-stock' AS step"); }); await ctx.queue.next("wait-resume", { names: ["resume"] }); diff --git a/rivetkit-typescript/packages/rivetkit/src/registry/napi-runtime.ts b/rivetkit-typescript/packages/rivetkit/src/registry/napi-runtime.ts index c83eaead5a..f822e70a0e 100644 --- a/rivetkit-typescript/packages/rivetkit/src/registry/napi-runtime.ts +++ b/rivetkit-typescript/packages/rivetkit/src/registry/napi-runtime.ts @@ -664,6 +664,10 @@ export class NapiCoreRuntime implements CoreRuntime { return this.#invocationContext.exit(run); } + currentInvocationScope(): object | undefined { + return this.#invocationContext.getStore(); + } + actorName(ctx: ActorContextHandle): string { return asNativeActorContext(ctx).name(); } diff --git a/rivetkit-typescript/packages/rivetkit/src/registry/native.ts b/rivetkit-typescript/packages/rivetkit/src/registry/native.ts index 6eb0dc71b6..c9beb2d91a 100644 --- a/rivetkit-typescript/packages/rivetkit/src/registry/native.ts +++ b/rivetkit-typescript/packages/rivetkit/src/registry/native.ts @@ -2709,6 +2709,9 @@ class NativeConnectionMap implements ReadonlyMap { readonly [Symbol.toStringTag] = "NativeConnectionMap"; } +/** Logger cache key for code that runs outside any invocation. */ +const NO_INVOCATION_SCOPE = {}; + export class ActorContextHandleAdapter { #runtime: CoreRuntime; #ctx: ActorContextHandle; @@ -2723,7 +2726,7 @@ export class ActorContextHandleAdapter { #db?: unknown; #dispatchCancelToken?: CancellationTokenHandle; #kv?: NativeKvAdapter; - #log?: ActorLogger; + #logByInvocationScope = new WeakMap(); #queue?: NativeQueueAdapter; #request?: Request; #schedule?: NativeScheduleAdapter; @@ -2972,21 +2975,26 @@ export class ActorContextHandleAdapter { ); } + /** Cached per invocation, because the `run` context outlives workflow runs and steps. */ get log() { - if (!this.#log) { - const invocation = this.#invocationTraceContext(); - this.#log = logger().child({ - actorId: this.actorId, - actorName: this.name, - actorKey: this.key, - ...(invocation?.rayId && { rayId: invocation.rayId }), - ...(invocation?.span && { - traceId: invocation.span.traceId, - spanId: invocation.span.spanId, - }), - }); - } - return this.#log; + const scope = + this.#runtime.currentInvocationScope() ?? NO_INVOCATION_SCOPE; + const cached = this.#logByInvocationScope.get(scope); + if (cached) return cached; + + const invocation = this.#invocationTraceContext(); + const log = logger().child({ + actorId: this.actorId, + actorName: this.name, + actorKey: this.key, + ...(invocation?.rayId && { rayId: invocation.rayId }), + ...(invocation?.span && { + traceId: invocation.span.traceId, + spanId: invocation.span.spanId, + }), + }); + this.#logByInvocationScope.set(scope, log); + return log; } get abortSignal(): AbortSignal { diff --git a/rivetkit-typescript/packages/rivetkit/src/registry/runtime.ts b/rivetkit-typescript/packages/rivetkit/src/registry/runtime.ts index fd5b4e4476..24825dfe0e 100644 --- a/rivetkit-typescript/packages/rivetkit/src/registry/runtime.ts +++ b/rivetkit-typescript/packages/rivetkit/src/registry/runtime.ts @@ -604,6 +604,8 @@ export interface CoreRuntime { startWorkflowSpan(ctx: ActorContextHandle): Promise; /** Backs `ActorRun.runOutsideWorkflowSpan`. */ runOutsideActorInvocationContext(run: () => T): T; + /** Identity of the invocation the caller runs in, usable as a cache key. */ + currentInvocationScope(): object | undefined; actorName(ctx: ActorContextHandle): string; actorKey(ctx: ActorContextHandle): RuntimeActorKeySegment[]; actorRegion(ctx: ActorContextHandle): string; diff --git a/rivetkit-typescript/packages/rivetkit/src/registry/wasm-runtime.ts b/rivetkit-typescript/packages/rivetkit/src/registry/wasm-runtime.ts index 622c7730ee..052c0fb591 100644 --- a/rivetkit-typescript/packages/rivetkit/src/registry/wasm-runtime.ts +++ b/rivetkit-typescript/packages/rivetkit/src/registry/wasm-runtime.ts @@ -571,6 +571,10 @@ export class WasmCoreRuntime implements CoreRuntime { return run(); } + currentInvocationScope(): object | undefined { + return undefined; + } + actorName(ctx: ActorContextHandle): string { return callHandle(asWasmActorContext(ctx), "name"); } diff --git a/rivetkit-typescript/packages/rivetkit/tests/driver/actor-telemetry.test.ts b/rivetkit-typescript/packages/rivetkit/tests/driver/actor-telemetry.test.ts index ea91115492..0ca0ffa3c6 100644 --- a/rivetkit-typescript/packages/rivetkit/tests/driver/actor-telemetry.test.ts +++ b/rivetkit-typescript/packages/rivetkit/tests/driver/actor-telemetry.test.ts @@ -985,6 +985,18 @@ describeDriverMatrix( } }); + test("logs written inside a step carry the step's trace ids", () => { + const [reserve] = named(`${actorName}/reserve-stock`); + const line = traced.runtime + .getRuntimeOutput?.() + .split("\n") + .find((candidate) => + candidate.includes(`workflow_log_key=${actorKey}`), + ); + expect(line).toContain(`traceId=${reserve.traceId}`); + expect(line).toContain(`spanId=${reserve.spanId}`); + }); + test("reports each attempt of a retried step under its own run", () => { const attempts = named(`${actorName}/charge-card`); expect( From 1b235f5e6c63e7d59276178145bd988d02d1a177 Mon Sep 17 00:00:00 2001 From: Sree Narayanan Date: Tue, 22 Sep 2026 16:24:09 +0400 Subject: [PATCH 4/6] feat(rivetkit): use queue message ray in workflow runs --- .../packages/actor-persist/src/versioned.rs | 3 +- .../rivetkit-core/src/actor/context.rs | 2 +- .../src/actor/internal_storage/mod.rs | 3 +- .../packages/rivetkit-core/src/telemetry.rs | 70 +++++++++++++------ .../tests/driver/actor-telemetry.test.ts | 27 ++++--- 5 files changed, 71 insertions(+), 34 deletions(-) diff --git a/rivetkit-rust/packages/actor-persist/src/versioned.rs b/rivetkit-rust/packages/actor-persist/src/versioned.rs index 12341020b1..3f677d7de5 100644 --- a/rivetkit-rust/packages/actor-persist/src/versioned.rs +++ b/rivetkit-rust/packages/actor-persist/src/versioned.rs @@ -537,9 +537,10 @@ impl OwnedVersionedData for RunWakeAt { } } -/// The span of the last workflow run, which the next run links to. +/// The span and ray of the last workflow run, which the next run continues. #[derive(Debug, Default, Serialize, Deserialize)] pub struct WorkflowTraceContextV1 { + pub ray_id: Option, pub traceparent: Option, pub tracestate: Option, } diff --git a/rivetkit-rust/packages/rivetkit-core/src/actor/context.rs b/rivetkit-rust/packages/rivetkit-core/src/actor/context.rs index 79c6cd5b16..14a33ee4a6 100644 --- a/rivetkit-rust/packages/rivetkit-core/src/actor/context.rs +++ b/rivetkit-rust/packages/rivetkit-core/src/actor/context.rs @@ -327,7 +327,7 @@ impl ActorContext { WorkflowRunInvocation::start(self, previous) } - /// Closes a workflow run and persists its span for the next run to link to. + /// Closes a workflow run and persists its span and ray for the next run. #[doc(hidden)] pub async fn finish_workflow_span( &self, diff --git a/rivetkit-rust/packages/rivetkit-core/src/actor/internal_storage/mod.rs b/rivetkit-rust/packages/rivetkit-core/src/actor/internal_storage/mod.rs index e7bf252af8..50fcb50b4a 100644 --- a/rivetkit-rust/packages/rivetkit-core/src/actor/internal_storage/mod.rs +++ b/rivetkit-rust/packages/rivetkit-core/src/actor/internal_storage/mod.rs @@ -1100,9 +1100,9 @@ pub(crate) async fn load_workflow_trace(db: &SqliteDb) -> Result Result<()> { let payload = encode_latest_with_embedded_version::( persist_versioned::WorkflowTraceContextV1 { + ray_id: trace_context.ray_id, traceparent: trace_context.traceparent, tracestate: trace_context.tracestate, }, diff --git a/rivetkit-rust/packages/rivetkit-core/src/telemetry.rs b/rivetkit-rust/packages/rivetkit-core/src/telemetry.rs index 1437984248..a589d6b20e 100644 --- a/rivetkit-rust/packages/rivetkit-core/src/telemetry.rs +++ b/rivetkit-rust/packages/rivetkit-core/src/telemetry.rs @@ -153,7 +153,8 @@ pub(crate) struct ActorTelemetryIdentity { #[derive(Debug)] struct InvocationInner { - ray_id: Option, + /// A workflow run takes the ray of each queue message it receives. + takes_message_ray: bool, // This lock is used from Drop paths, and its guard never crosses an await. state: Mutex, identity: Arc, @@ -161,6 +162,7 @@ struct InvocationInner { #[derive(Debug)] struct InvocationState { + ray_id: Option, span: Option, finished: bool, pending_work: usize, @@ -390,7 +392,7 @@ impl ActorInvocation { ctx, InvocationSubject::WorkflowRun, InvocationType::Workflow, - None, + previous.ray_id, None, previous_run, ) @@ -480,7 +482,7 @@ impl ActorInvocation { }; Self { - telemetry: ActorInvocationTelemetry::new(ray_id, span, identity), + telemetry: ActorInvocationTelemetry::new(invocation_type, ray_id, span, identity), } } @@ -530,10 +532,15 @@ impl WorkflowRunInvocation { self.ctx.clone() } - /// Returns the span the next run links to, or nothing when tracing is off. + /// Returns the span and ray the next run continues, or nothing when + /// tracing is off. pub(crate) fn finish(self, outcome: WorkflowRunOutcome) -> Option { + let ray_id = self.invocation.telemetry.ray_id(); let mut trace_context = None; self.invocation.telemetry.finish_with(|span| { + if let Some(ray_id) = ray_id.as_deref() { + span.record("rivet.ray.id", ray_id); + } span.record("rivet.workflow.run.outcome", outcome.as_label()); span.record( "otel.status_code", @@ -543,9 +550,9 @@ impl WorkflowRunInvocation { .map(|span_context| w3c_trace_headers(&span_context)) .unwrap_or_default(); trace_context = Some(IncomingTraceContext { + ray_id, traceparent: headers.traceparent, tracestate: headers.tracestate, - ..Default::default() }); }); trace_context @@ -554,7 +561,11 @@ impl WorkflowRunInvocation { impl Drop for WorkflowRunInvocation { fn drop(&mut self) { + let ray_id = self.invocation.telemetry.ray_id(); self.invocation.telemetry.finish_with(|span| { + if let Some(ray_id) = ray_id.as_deref() { + span.record("rivet.ray.id", ray_id); + } span.record("otel.status_code", "ERROR"); span.record("error.type", OPERATION_ABANDONED_ERROR_TYPE); }); @@ -569,14 +580,16 @@ impl Drop for ActorInvocation { impl ActorInvocationTelemetry { fn new( + invocation_type: InvocationType, ray_id: Option, span: Option, identity: Arc, ) -> Self { Self { inner: Arc::new(InvocationInner { - ray_id, + takes_message_ray: matches!(invocation_type, InvocationType::Workflow), state: Mutex::new(InvocationState { + ray_id, span, finished: false, pending_work: 0, @@ -588,6 +601,18 @@ impl ActorInvocationTelemetry { } } + fn ray_id(&self) -> Option { + self.inner.state.lock().ray_id.clone() + } + + /// Records the ray on a span opened inside this invocation. The ray is + /// borrowed under the lock, so this allocates nothing. + fn record_ray(&self, span: &tracing::Span) { + if let Some(ray_id) = self.inner.state.lock().ray_id.as_deref() { + span.record("rivet.ray.id", ray_id); + } + } + /// Returns a handle for the same invocation whose spans parent to the /// application span identified by `traceparent` and `tracestate`. Absent /// or invalid context, or the invocation span itself, which the host sees @@ -644,22 +669,19 @@ impl ActorInvocationTelemetry { /// Returns correlation fields only while this actor invocation is active. #[doc(hidden)] pub fn trace_context(&self) -> Option { - let span = { + let (ray_id, span) = { let state = self.inner.state.lock(); if state.finished && state.pending_work == 0 { return None; } - state.span.clone() + (state.ray_id.clone(), state.span.clone()) }; let span = match &self.step_span { Some(step_span) => w3c_span_context(step_span), None => span.and_then(|span| span_context_of(&span)), }; - Some(ActorInvocationTraceContext { - ray_id: self.inner.ray_id.clone(), - span, - }) + Some(ActorInvocationTraceContext { ray_id, span }) } /// Trace context that work caused by this invocation records: the invocation's @@ -674,7 +696,7 @@ impl ActorInvocationTelemetry { let parent_span = parent.span(); let headers = w3c_trace_headers(parent_span.span_context()); IncomingTraceContext { - ray_id: self.inner.ray_id.clone(), + ray_id: self.ray_id(), traceparent: headers.traceparent, tracestate: headers.tracestate, } @@ -704,10 +726,11 @@ impl ActorInvocationTelemetry { otel.kind = "client", rivet.actor.name = %actor_name, rivet.action.name = %action_name, - rivet.ray.id = self.inner.ray_id.as_deref(), + rivet.ray.id = tracing::field::Empty, otel.status_code = tracing::field::Empty, error.type = tracing::field::Empty, ); + self.record_ray(&span); span.set_parent(parent); let context = span_context_of(&span); Some(OutboundCallInvocation { @@ -733,13 +756,14 @@ impl ActorInvocationTelemetry { rivet.workflow.step.name = %step_name, rivet.workflow.step.attempt = attempt, rivet.workflow.step.outcome = tracing::field::Empty, - rivet.ray.id = self.inner.ray_id.as_deref(), + rivet.ray.id = tracing::field::Empty, rivet.actor.id = %self.inner.identity.actor_id, rivet.actor.name = %self.inner.identity.actor_name, rivet.actor.key = %self.inner.identity.actor_key, otel.status_code = tracing::field::Empty, error.type = tracing::field::Empty, ); + self.record_ray(&span); span.set_parent(parent); let step_telemetry = Self { inner: self.inner.clone(), @@ -760,13 +784,14 @@ impl ActorInvocationTelemetry { otel.kind = "internal", rivet.operation.system = "sqlite", rivet.operation.name = operation_name, - rivet.ray.id = self.inner.ray_id.as_deref(), + rivet.ray.id = tracing::field::Empty, rivet.actor.id = %self.inner.identity.actor_id, rivet.actor.name = %self.inner.identity.actor_name, rivet.actor.key = %self.inner.identity.actor_key, otel.status_code = tracing::field::Empty, error.type = tracing::field::Empty, ); + self.record_ray(&span); span.set_parent(parent); Some(SqliteOperationSpan { span: Some(span) }) } @@ -983,14 +1008,15 @@ pub(crate) fn start_queue_receive(ctx: &ActorContext, message: &QueueMessage) -> .invocation_telemetry() .and_then(|telemetry| telemetry.parent_context().map(|parent| (telemetry, parent))); if let Some((telemetry, parent)) = invocation_parent { - let ray_id = telemetry - .inner - .ray_id - .as_deref() - .or(message.trace_context.ray_id.as_deref()); - if let Some(ray_id) = ray_id { + let message_ray_id = message.trace_context.ray_id.as_deref(); + let mut state = telemetry.inner.state.lock(); + if let (true, Some(message_ray_id)) = (telemetry.inner.takes_message_ray, message_ray_id) { + state.ray_id = Some(message_ray_id.to_owned()); + } + if let Some(ray_id) = state.ray_id.as_deref().or(message_ray_id) { span.record("rivet.ray.id", ray_id); } + drop(state); span.set_parent(parent); } else if let Some(ray_id) = &message.trace_context.ray_id { span.record("rivet.ray.id", ray_id); diff --git a/rivetkit-typescript/packages/rivetkit/tests/driver/actor-telemetry.test.ts b/rivetkit-typescript/packages/rivetkit/tests/driver/actor-telemetry.test.ts index 0ca0ffa3c6..71199dc336 100644 --- a/rivetkit-typescript/packages/rivetkit/tests/driver/actor-telemetry.test.ts +++ b/rivetkit-typescript/packages/rivetkit/tests/driver/actor-telemetry.test.ts @@ -985,7 +985,7 @@ describeDriverMatrix( } }); - test("logs written inside a step carry the step's trace ids", () => { + test("logs written inside a step carry the step's trace ids and ray", () => { const [reserve] = named(`${actorName}/reserve-stock`); const line = traced.runtime .getRuntimeOutput?.() @@ -995,6 +995,7 @@ describeDriverMatrix( ); expect(line).toContain(`traceId=${reserve.traceId}`); expect(line).toContain(`spanId=${reserve.spanId}`); + expect(line).toContain(`rayId=${approveRayId}`); }); test("reports each attempt of a retried step under its own run", () => { @@ -1021,7 +1022,7 @@ describeDriverMatrix( } }); - test("keeps a run's ray unchanged when it receives a message sent under another ray", () => { + test("gives a run the ray of the message that woke it and keeps it for later runs", () => { const ray = (name: string) => attribute(named(name), "rivet.ray.id"); const sends = named(`${actorName}/queue.send`); @@ -1038,15 +1039,23 @@ describeDriverMatrix( expect(ray(`${actorName}/queue.receive`).sort()).toEqual( [approveRayId, resumeRayId].sort(), ); - expect(new Set(ray(`${actorName}/workflow`))).toEqual( - new Set([undefined]), - ); - expect(ray(`${actorName}/reserve-stock`)).toEqual([undefined]); + expect(ray(`${actorName}/reserve-stock`)).toEqual([ + approveRayId, + ]); expect(ray(`${actorName}/charge-card`)).toEqual([ - undefined, - undefined, - undefined, + resumeRayId, + resumeRayId, + resumeRayId, ]); + const runRays = ray(`${actorName}/workflow`); + const firstWithRay = runRays.findIndex( + (id) => id !== undefined, + ); + expect(runRays[firstWithRay]).toBe(approveRayId); + expect(runRays.slice(firstWithRay)).not.toContain( + undefined, + ); + expect(runRays.at(-1)).toBe(resumeRayId); }); test("puts a queue send made inside a step in that step's trace", () => { From da881f981fe97786c4db133f1851a2ef6ce5d885 Mon Sep 17 00:00:00 2001 From: Sree Narayanan Date: Tue, 22 Sep 2026 14:49:13 +0400 Subject: [PATCH 5/6] docs(rivetkit): add workflow tracing to the tracing docs --- docs-internal/engine/rivetkit-telemetry.md | 5 ++++- docs/content/docs/general/tracing.mdx | 2 ++ 2 files changed, 6 insertions(+), 1 deletion(-) diff --git a/docs-internal/engine/rivetkit-telemetry.md b/docs-internal/engine/rivetkit-telemetry.md index d060b6f775..0eaa5eb0ff 100644 --- a/docs-internal/engine/rivetkit-telemetry.md +++ b/docs-internal/engine/rivetkit-telemetry.md @@ -35,6 +35,8 @@ Core spans use the `rivetkit::telemetry` tracing target. Log layers exclude this | Queue receive | `{actor}/queue.receive` | `consumer` | Linked to the send origin | | Actor call | `{callee}/{action}` | `client` | Child of application or invocation span | | SQLite | `rivet.sqlite.{operation}` | `internal` | Child of application or invocation span | +| Workflow run | `{actor}/workflow` | `internal` | New trace linked to the previous run | +| Workflow step | `{actor}/{step}` | `internal` | Child of the run span, one per attempt | Core records these attributes: @@ -47,6 +49,7 @@ Core records these attributes: | HTTP | `http.request.method`, `http.response.status_code` | | Queue | `rivet.queue.name` | | SQLite | `rivet.operation.system`, `rivet.operation.name` | +| Workflow | `rivet.workflow.run.outcome`, `rivet.workflow.step.name`, `rivet.workflow.step.attempt`, `rivet.workflow.step.outcome` | Raw HTTP spans use `onRequest`, never the request path. Handler errors use their `group.code` as `error.type`. A 5xx response uses the status code. Abandoned SQLite and actor-call tracking uses `actor.operation_abandoned` to represent an unknown outcome. @@ -120,5 +123,5 @@ Record actor identity, invocation type, HTTP method and status, correlation IDs, - WebSocket handlers, lifecycle hooks, connection callbacks, KV, and actor-state operations have no dedicated spans - WebSocket action messages and inspector actions do not inherit caller context -- Actor creation ray IDs do not reach the actor runtime +- Actor creation ray IDs do not reach the actor runtime. A workflow run takes the ray of each queue message it receives and keeps it for later runs - Wasm does not export host spans or expose Core invocation context to TypeScript diff --git a/docs/content/docs/general/tracing.mdx b/docs/content/docs/general/tracing.mdx index 77983be2a4..fe76618623 100644 --- a/docs/content/docs/general/tracing.mdx +++ b/docs/content/docs/general/tracing.mdx @@ -61,6 +61,7 @@ Each span shows how long an operation took and whether it failed. RivetKit recor - SQLite operations, without recording SQL text, bindings, or results - Queue sends and receipts, with a link from each receipt to its sender - Scheduled actions, in a new trace linked to the work that scheduled them +- Workflow runs and steps. Each run is a new trace linked to the run before it The full list of span names and attributes lives in [`telemetry.rs`](https://github.com/rivet-dev/rivet/blob/main/rivetkit-rust/packages/rivetkit-core/src/telemetry.rs), and the SQLite operations in [`sqlite/mod.rs`](https://github.com/rivet-dev/rivet/blob/main/rivetkit-rust/packages/rivetkit-core/src/actor/sqlite/mod.rs). @@ -184,5 +185,6 @@ RivetKit sends spans in the background. Actor requests continue if the collector - Actions called over `.connect()` start a new trace instead of joining the caller's trace - Effect spans are not traced. Use `@effect/opentelemetry` - RivetKit does not add `rivet.ray.id` to your application spans +- Workflow spans carry a `rivet.ray.id` only after the workflow receives a queue message - The Wasm runtime is not traced - A hard process exit can lose spans still waiting in the export queue From f71c35a238b0cdd116f11c2ea93efa4a585852a0 Mon Sep 17 00:00:00 2001 From: Sree Narayanan Date: Tue, 22 Sep 2026 18:25:30 +0400 Subject: [PATCH 6/6] fix(rivetkit): finish a step span after its result is written --- rivetkit-typescript/packages/workflow-engine/src/context.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/rivetkit-typescript/packages/workflow-engine/src/context.ts b/rivetkit-typescript/packages/workflow-engine/src/context.ts index 8ba8eb63ce..b272a8db88 100644 --- a/rivetkit-typescript/packages/workflow-engine/src/context.ts +++ b/rivetkit-typescript/packages/workflow-engine/src/context.ts @@ -922,7 +922,6 @@ export class WorkflowContextImpl implements WorkflowContextInterface { timeout, config.name, ); - stepSpan?.finish("ok"); if (entry.kind.type === "step") { entry.kind.data.output = output; @@ -946,6 +945,7 @@ export class WorkflowContextImpl implements WorkflowContextInterface { await this.flushStorage(); } + stepSpan?.finish("ok"); this.log("debug", { msg: "step completed", step: config.name,