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 diff --git a/rivetkit-rust/packages/actor-persist/src/versioned.rs b/rivetkit-rust/packages/actor-persist/src/versioned.rs index d64bd2d9e0..3f677d7de5 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,51 @@ impl OwnedVersionedData for RunWakeAt { Vec:: Result>::new() } } + +/// 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, +} + +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..14a33ee4a6 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 and ray for the next run. + #[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..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 @@ -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,54 @@ 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 { + ray_id: stored.ray_id, + traceparent: stored.traceparent, + tracestate: stored.tracestate, + }) +} + +pub(crate) async fn persist_workflow_trace( + db: &SqliteDb, + trace_context: IncomingTraceContext, +) -> 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, + }, + 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..a589d6b20e 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 @@ -143,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, @@ -151,6 +162,7 @@ struct InvocationInner { #[derive(Debug)] struct InvocationState { + ray_id: Option, span: Option, finished: bool, pending_work: usize, @@ -248,6 +260,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 +381,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, + previous.ray_id, + 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 +462,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, ); @@ -366,7 +482,7 @@ impl ActorInvocation { }; Self { - telemetry: ActorInvocationTelemetry::new(ray_id, span, identity), + telemetry: ActorInvocationTelemetry::new(invocation_type, ray_id, span, identity), } } @@ -400,6 +516,62 @@ 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 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", + 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 { + ray_id, + traceparent: headers.traceparent, + tracestate: headers.tracestate, + }); + }); + trace_context + } +} + +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); + }); + } +} + impl Drop for ActorInvocation { fn drop(&mut self) { self.telemetry.finish_dropped(); @@ -408,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, @@ -423,6 +597,19 @@ impl ActorInvocationTelemetry { identity, }), application_span: None, + step_span: None, + } + } + + 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); } } @@ -448,19 +635,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), } @@ -480,19 +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)), }; - let span = 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 @@ -507,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, } @@ -537,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 { @@ -549,6 +739,40 @@ 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 = 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(), + 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(); @@ -560,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) }) } @@ -637,6 +862,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 +962,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,9 +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 { - if let Some(ray_id) = telemetry.inner.ray_id.as_deref() { + 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-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..62390cf1cb 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,54 @@ 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) => { + 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"] }); + 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..f822e70a0e 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,38 @@ 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); + } + + 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 776bb8edab..c9beb2d91a 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, @@ -2708,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; @@ -2722,11 +2726,11 @@ export class ActorContextHandleAdapter { #db?: unknown; #dispatchCancelToken?: CancellationTokenHandle; #kv?: NativeKvAdapter; - #log?: ActorLogger; + #logByInvocationScope = new WeakMap(); #queue?: NativeQueueAdapter; #request?: Request; #schedule?: NativeScheduleAdapter; - #run?: { setWakeAt(timestamp: number | null): Promise }; + #run?: ActorRun; #runHandlerConfigured: boolean; #onStateChange?: NativeOnStateChangeHandler; #stateEnabled: boolean; @@ -2907,6 +2911,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; @@ -2965,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 2a26b034d2..24825dfe0e 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,12 @@ export interface CoreRuntime { actorName: string, actionName: string, ): RuntimeOutboundCall | undefined; + /** Backs `ActorRun.startWorkflowSpan`. */ + 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 2e6d2a00bb..052c0fb591 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,22 @@ 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(); + } + + currentInvocationScope(): object | undefined { + return undefined; + } + 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..71199dc336 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,204 @@ 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("logs written inside a step carry the step's trace ids and ray", () => { + 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}`); + expect(line).toContain(`rayId=${approveRayId}`); + }); + + 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("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`); + 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(ray(`${actorName}/reserve-stock`)).toEqual([ + approveRayId, + ]); + expect(ray(`${actorName}/charge-card`)).toEqual([ + 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", () => { + 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..b272a8db88 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,10 +910,15 @@ 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, ); @@ -936,6 +945,7 @@ export class WorkflowContextImpl implements WorkflowContextInterface { await this.flushStorage(); } + stepSpan?.finish("ok"); this.log("debug", { msg: "step completed", step: config.name, @@ -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