Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 4 additions & 1 deletion docs-internal/engine/rivetkit-telemetry.md
Original file line number Diff line number Diff line change
Expand Up @@ -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:

Expand All @@ -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.

Expand Down Expand Up @@ -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
2 changes: 2 additions & 0 deletions docs/content/docs/general/tracing.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -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).

Expand Down Expand Up @@ -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
49 changes: 49 additions & 0 deletions rivetkit-rust/packages/actor-persist/src/versioned.rs
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
use anyhow::{Result, bail};
use serde::{Deserialize, Serialize};
use vbare::OwnedVersionedData;

use crate::generated::{v1, v2, v3, v4};
Expand Down Expand Up @@ -535,3 +536,51 @@ impl OwnedVersionedData for RunWakeAt {
Vec::<fn(Self) -> Result<Self>>::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<String>,
pub traceparent: Option<String>,
pub tracestate: Option<String>,
}

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<Self::Latest> {
match self {
Self::V1(data) => Ok(data),
}
}

fn deserialize_version(payload: &[u8], version: u16) -> Result<Self> {
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<Vec<u8>> {
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<impl Fn(Self) -> Result<Self>> {
Vec::<fn(Self) -> Result<Self>>::new()
}

fn serialize_converters() -> Vec<impl Fn(Self) -> Result<Self>> {
Vec::<fn(Self) -> Result<Self>>::new()
}
}
56 changes: 55 additions & 1 deletion rivetkit-rust/packages/rivetkit-core/src/actor/context.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};

Expand Down Expand Up @@ -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> {
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<ActorInvocationTraceContext> {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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`.
Expand Down Expand Up @@ -1082,6 +1083,54 @@ pub(crate) async fn persist_run_wake_at(db: &SqliteDb, wake_at: Option<i64>) ->
Ok(())
}

pub(crate) async fn load_workflow_trace(db: &SqliteDb) -> Result<IncomingTraceContext> {
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::<persist_versioned::WorkflowTraceContext>(
&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::WorkflowTraceContext>(
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<Option<String>> {
let result = db
.query(LOAD_INSPECTOR_TOKEN_SQL, None)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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";

Expand Down
4 changes: 4 additions & 0 deletions rivetkit-rust/packages/rivetkit-core/src/actor/keys.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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];

Expand Down
3 changes: 2 additions & 1 deletion rivetkit-rust/packages/rivetkit-core/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
Loading
Loading