diff --git a/CLAUDE.md b/CLAUDE.md index a172a63fad..72c2d3550c 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -261,6 +261,8 @@ When the user asks to track something in a note, store it in `~/.agents/notes/` ## Performance +- Bound UniversalDB ranges before materialization; a stream consumer's item limit does not constrain RocksDB's eager range allocation. See [dead-workflow backfill](docs-internal/engine/gasoline-dead-workflow-backfill.md). + - Use `rivet_perf::{perf_start, perf_finish}` for latency-sensitive async phases, shared I/O wrappers, and suspected backpressure points where slow-tail spans are useful. - Keep `perf_start!` labels bounded for metrics; put high-cardinality values such as IDs, request paths, and byte counts in span fields instead. - Every `PerfMeasure` must end with `perf_finish!` or `perf_abandon!` before leaving scope. diff --git a/Cargo.lock b/Cargo.lock index 5ed497b518..771b681ce1 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1318,7 +1318,7 @@ checksum = "2a2330da5de22e8a3cb63252ce2abb30116bf5265e89c0e01bc17015ce30a476" [[package]] name = "datacenter" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "futures-util", @@ -1377,7 +1377,7 @@ dependencies = [ [[package]] name = "depot" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "async-channel", @@ -1420,7 +1420,7 @@ dependencies = [ [[package]] name = "depot-client-embedded" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "async-trait", @@ -1696,7 +1696,7 @@ checksum = "c34f04666d835ff5d62e058c3995147c06f42fe86ff053337632bca83e42702d" [[package]] name = "epoxy" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "axum 0.8.4", @@ -1741,7 +1741,7 @@ dependencies = [ [[package]] name = "epoxy-protocol" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "rivet-vbare-compiler", @@ -2092,7 +2092,7 @@ dependencies = [ [[package]] name = "gasoline" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "async-stream", @@ -2143,7 +2143,7 @@ dependencies = [ [[package]] name = "gasoline-macros" -version = "2.3.14" +version = "2.3.17" dependencies = [ "proc-macro2", "quote", @@ -2152,7 +2152,7 @@ dependencies = [ [[package]] name = "gasoline-runtime" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "epoxy", @@ -3375,7 +3375,7 @@ dependencies = [ [[package]] name = "namespace" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "epoxy", @@ -3944,7 +3944,7 @@ dependencies = [ [[package]] name = "pegboard" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "base64 0.22.1", @@ -4002,7 +4002,7 @@ dependencies = [ [[package]] name = "pegboard-envoy" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "async-trait", @@ -4051,7 +4051,7 @@ dependencies = [ [[package]] name = "pegboard-gateway" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "async-trait", @@ -4085,7 +4085,7 @@ dependencies = [ [[package]] name = "pegboard-gateway2" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "async-trait", @@ -4120,7 +4120,7 @@ dependencies = [ [[package]] name = "pegboard-gateway3" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "async-trait", @@ -4155,7 +4155,7 @@ dependencies = [ [[package]] name = "pegboard-outbound" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "epoxy", @@ -4182,7 +4182,7 @@ dependencies = [ [[package]] name = "pegboard-runner" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "async-trait", @@ -5054,7 +5054,7 @@ dependencies = [ [[package]] name = "rivet-actor-runtime-socket-protocol" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "rivet-vbare-compiler", @@ -5065,7 +5065,7 @@ dependencies = [ [[package]] name = "rivet-api-builder" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "axum 0.8.4", @@ -5108,7 +5108,7 @@ dependencies = [ [[package]] name = "rivet-api-peer" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "axum 0.8.4", @@ -5144,7 +5144,7 @@ dependencies = [ [[package]] name = "rivet-api-public" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "axum 0.8.4", @@ -5182,7 +5182,7 @@ dependencies = [ [[package]] name = "rivet-api-public-openapi-gen" -version = "2.3.14" +version = "2.3.17" dependencies = [ "rivet-api-public", "serde_json", @@ -5191,7 +5191,7 @@ dependencies = [ [[package]] name = "rivet-api-types" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "gasoline", @@ -5206,7 +5206,7 @@ dependencies = [ [[package]] name = "rivet-api-util" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "axum 0.8.4", @@ -5264,7 +5264,7 @@ dependencies = [ [[package]] name = "rivet-bootstrap" -version = "2.3.14" +version = "2.3.17" dependencies = [ "datacenter", "depot", @@ -5286,7 +5286,7 @@ dependencies = [ [[package]] name = "rivet-build-meta" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "epoxy-protocol", @@ -5300,7 +5300,7 @@ dependencies = [ [[package]] name = "rivet-cache" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "futures-util", @@ -5326,7 +5326,7 @@ dependencies = [ [[package]] name = "rivet-cache-purge" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "futures-util", @@ -5343,14 +5343,14 @@ dependencies = [ [[package]] name = "rivet-cache-result" -version = "2.3.14" +version = "2.3.17" dependencies = [ "rivet-util", ] [[package]] name = "rivet-cli" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anstyle", "anyhow", @@ -5371,7 +5371,7 @@ dependencies = [ [[package]] name = "rivet-config" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "chrono", @@ -5393,7 +5393,7 @@ dependencies = [ [[package]] name = "rivet-config-schema-gen" -version = "2.3.14" +version = "2.3.17" dependencies = [ "rivet-config", "schemars 0.8.22", @@ -5402,7 +5402,7 @@ dependencies = [ [[package]] name = "rivet-container-runner" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "async-trait", @@ -5424,7 +5424,7 @@ dependencies = [ [[package]] name = "rivet-data" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "gasoline", @@ -5437,7 +5437,7 @@ dependencies = [ [[package]] name = "rivet-depot-client" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "async-trait", @@ -5466,11 +5466,11 @@ dependencies = [ [[package]] name = "rivet-depot-client-types" -version = "2.3.14" +version = "2.3.17" [[package]] name = "rivet-depot-protocol" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "rivet-vbare-compiler", @@ -5481,7 +5481,7 @@ dependencies = [ [[package]] name = "rivet-dynamic-config" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "futures-util", @@ -5499,7 +5499,7 @@ dependencies = [ [[package]] name = "rivet-engine" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "async-trait", @@ -5582,7 +5582,7 @@ dependencies = [ [[package]] name = "rivet-env" -version = "2.3.14" +version = "2.3.17" dependencies = [ "lazy_static", "uuid", @@ -5590,7 +5590,7 @@ dependencies = [ [[package]] name = "rivet-envoy-client" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "futures-util", @@ -5622,7 +5622,7 @@ dependencies = [ [[package]] name = "rivet-envoy-protocol" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "hex", @@ -5637,7 +5637,7 @@ dependencies = [ [[package]] name = "rivet-error" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "indoc", @@ -5649,7 +5649,7 @@ dependencies = [ [[package]] name = "rivet-error-macros" -version = "2.3.14" +version = "2.3.17" dependencies = [ "indoc", "proc-macro2", @@ -5660,7 +5660,7 @@ dependencies = [ [[package]] name = "rivet-guard" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "async-trait", @@ -5723,7 +5723,7 @@ dependencies = [ [[package]] name = "rivet-guard-core" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "async-trait", @@ -5773,7 +5773,7 @@ dependencies = [ [[package]] name = "rivet-logs" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "chrono", @@ -5787,7 +5787,7 @@ dependencies = [ [[package]] name = "rivet-metrics" -version = "2.3.14" +version = "2.3.17" dependencies = [ "lazy_static", "prometheus", @@ -5795,7 +5795,7 @@ dependencies = [ [[package]] name = "rivet-metrics-server" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "console-subscriber", @@ -5813,7 +5813,7 @@ dependencies = [ [[package]] name = "rivet-outbound-guard" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "ipnet", @@ -5827,7 +5827,7 @@ dependencies = [ [[package]] name = "rivet-perf" -version = "2.3.14" +version = "2.3.17" dependencies = [ "prometheus", "tokio", @@ -5837,7 +5837,7 @@ dependencies = [ [[package]] name = "rivet-pools" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "clickhouse", @@ -5870,7 +5870,7 @@ dependencies = [ [[package]] name = "rivet-postgres-util" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "rustls", @@ -5881,7 +5881,7 @@ dependencies = [ [[package]] name = "rivet-profiling" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "gasoline", @@ -5896,7 +5896,7 @@ dependencies = [ [[package]] name = "rivet-runner-protocol" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "gasoline", @@ -5913,7 +5913,7 @@ dependencies = [ [[package]] name = "rivet-runtime" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "console-subscriber", @@ -5942,7 +5942,7 @@ dependencies = [ [[package]] name = "rivet-service-manager" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "chrono", @@ -5959,7 +5959,7 @@ dependencies = [ [[package]] name = "rivet-telemetry" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "rivet-config", @@ -5970,7 +5970,7 @@ dependencies = [ [[package]] name = "rivet-term" -version = "2.3.14" +version = "2.3.17" dependencies = [ "console", "derive_builder 0.12.0", @@ -5982,7 +5982,7 @@ dependencies = [ [[package]] name = "rivet-test-deps" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "futures-util", @@ -6000,7 +6000,7 @@ dependencies = [ [[package]] name = "rivet-test-deps-docker" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "portpicker", @@ -6017,7 +6017,7 @@ dependencies = [ [[package]] name = "rivet-test-envoy" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "async-stream", @@ -6033,7 +6033,7 @@ dependencies = [ [[package]] name = "rivet-tracing-utils" -version = "2.3.14" +version = "2.3.17" dependencies = [ "futures-util", "lazy_static", @@ -6043,7 +6043,7 @@ dependencies = [ [[package]] name = "rivet-types" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "gasoline", @@ -6060,7 +6060,7 @@ dependencies = [ [[package]] name = "rivet-universaldb-commit" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "rivet-vbare-compiler", @@ -6071,7 +6071,7 @@ dependencies = [ [[package]] name = "rivet-ups-broadcast" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "futures-util", @@ -6090,7 +6090,7 @@ dependencies = [ [[package]] name = "rivet-ups-protocol" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "base64 0.22.1", @@ -6102,7 +6102,7 @@ dependencies = [ [[package]] name = "rivet-util" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "async-trait", @@ -6130,7 +6130,7 @@ dependencies = [ [[package]] name = "rivet-util-id" -version = "2.3.14" +version = "2.3.17" dependencies = [ "serde", "thiserror 1.0.69", @@ -6141,7 +6141,7 @@ dependencies = [ [[package]] name = "rivet-util-serde" -version = "2.3.14" +version = "2.3.17" dependencies = [ "indexmap 2.14.0", "serde", @@ -6178,7 +6178,7 @@ dependencies = [ [[package]] name = "rivet-version-management" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "futures-util", @@ -6198,7 +6198,7 @@ dependencies = [ [[package]] name = "rivet-workflow-worker" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "datacenter", @@ -6214,7 +6214,7 @@ dependencies = [ [[package]] name = "rivetkit" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "async-trait", @@ -6242,7 +6242,7 @@ dependencies = [ [[package]] name = "rivetkit-actor-persist" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "rivet-vbare-compiler", @@ -6253,7 +6253,7 @@ dependencies = [ [[package]] name = "rivetkit-client" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "axum 0.8.4", @@ -6283,7 +6283,7 @@ dependencies = [ [[package]] name = "rivetkit-client-protocol" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "rivet-vbare-compiler", @@ -6294,7 +6294,7 @@ dependencies = [ [[package]] name = "rivetkit-core" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "async-trait", @@ -6353,7 +6353,7 @@ dependencies = [ [[package]] name = "rivetkit-engine-process" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "libc", @@ -6370,7 +6370,7 @@ dependencies = [ [[package]] name = "rivetkit-inspector-protocol" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "rivet-vbare-compiler", @@ -6381,7 +6381,7 @@ dependencies = [ [[package]] name = "rivetkit-napi" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "async-trait", @@ -6407,7 +6407,7 @@ dependencies = [ [[package]] name = "rivetkit-shared-types" -version = "2.3.14" +version = "2.3.17" dependencies = [ "serde", "serde_json", @@ -6415,7 +6415,7 @@ dependencies = [ [[package]] name = "rivetkit-wasm" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "console_error_panic_hook", @@ -7599,7 +7599,7 @@ dependencies = [ [[package]] name = "test-snapshot-gen" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "async-trait", @@ -8382,7 +8382,7 @@ checksum = "ebc1c04c71510c7f702b52b7c350734c9ff1295c464a03335b00bb84fc54f853" [[package]] name = "universaldb" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "async-trait", @@ -8420,7 +8420,7 @@ dependencies = [ [[package]] name = "universalpubsub" -version = "2.3.14" +version = "2.3.17" dependencies = [ "anyhow", "async-trait", diff --git a/docs-internal/engine/gasoline-dead-workflow-backfill.md b/docs-internal/engine/gasoline-dead-workflow-backfill.md new file mode 100644 index 0000000000..0c6fda71d7 --- /dev/null +++ b/docs-internal/engine/gasoline-dead-workflow-backfill.md @@ -0,0 +1,37 @@ +# Dead-workflow index backfill + +The `gasoline_dead_wf_backfill` workflow rebuilds the dead-workflow index for +existing workflow records. Each activity processes at most 1,000 workflows and +uses the existing three-second early transaction deadline. The activity's +`last_key: Option>` and continuation output remain unchanged. + +RocksDB materializes a range before exposing it through the UniversalDB stream +interface. A consumer-side workflow count does not bound that allocation. The +backfill therefore discovers one workflow with a key-only seek, +reads its name, error, and status markers directly, and advances to the end of +that workflow's data prefix. Output detection uses a separate key-only seek so +any output chunk excludes the workflow without loading its output payload. +Classification point reads are serializable, and an explicit read conflict +covers the output subtree to detect concurrent output insertion. +Input and state chunks are never traversed for classification. + +The returned cursor is the first unprocessed range boundary. Cursors written by +the original implementation, which pointed at an actual first workflow key, +remain valid. A cursor advances only after the workflow's complete classification +and optional index write. The index writes and returned cursor share the same +transaction outcome. An interrupted activity can safely repeat an already +committed chunk because index keys are deterministic. A timeout during reads +leaves that workflow as the next unprocessed range. + +The regression test uses the actual RocksDB-backed `DatabaseDebug` operation. +It seeds a workflow with four MiB of irrelevant input, verifies less than +64 KiB is read to classify it, checks all exclusion markers, and exercises +legacy cursors, one-workflow chunks, and replayed input. Run: + +```sh +RIVET_TEST_DATABASE=filesystem RIVET_TEST_PUBSUB=memory cargo test -p gasoline --lib backfill_skips_large_payloads --locked +``` + +This bounds the backfill's reads, not every engine allocation. It does not +change the RocksDB driver, actor data format, workflow input serialization, +or Envoy protocol. diff --git a/engine/packages/gasoline/src/db/kv/debug.rs b/engine/packages/gasoline/src/db/kv/debug.rs index 2b04271bcd..98060a124e 100644 --- a/engine/packages/gasoline/src/db/kv/debug.rs +++ b/engine/packages/gasoline/src/db/kv/debug.rs @@ -12,7 +12,7 @@ use serde::Serialize; use tracing::Instrument; use universaldb::utils::{FormalChunkedKey, FormalKey, IsolationLevel::*, end_of_key_range}; use universaldb::{ - RangeOption, + KeySelector, RangeOption, options::{ConflictRangeType, StreamingMode}, tuple::{PackResult, TupleDepth, TupleUnpack}, value::Value, @@ -40,6 +40,11 @@ use crate::{ const EARLY_TXN_TIMEOUT: Duration = Duration::from_secs(3); +// These tests seed private workflow keys and exercise the real RocksDB-backed debug operation. +#[cfg(test)] +#[path = "../../../tests/modules/dead_workflow_backfill.rs"] +mod dead_workflow_backfill_tests; + impl DatabaseKv { #[tracing::instrument(level = "debug", skip_all)] async fn get_workflows_inner( @@ -1281,6 +1286,7 @@ impl DatabaseDebug for DatabaseKv { limit: usize, last_key: Option<&[u8]>, ) -> Result<(usize, Option>)> { + ensure!(limit > 0, "backfill workflow limit must be positive"); let last_key = last_key.map(|x| x.to_vec()); self.pools @@ -1289,113 +1295,81 @@ impl DatabaseDebug for DatabaseKv { let last_key = &last_key; async move { let tx = tx.with_subspace(self.subspace.clone()); - + let (start, end) = self + .subspace + .subspace(&keys::workflow::DataSubspaceKey::new()) + .range(); let mut total = 0; - - let entire_subspace_key = keys::workflow::DataSubspaceKey::new(); - let subspace_range = self.subspace.subspace(&entire_subspace_key).range(); - let start = if let Some(last_key) = last_key { - last_key.clone() - } else { - subspace_range.0 - }; - let end = subspace_range.1; - - let mut stream = tx.get_ranges_keyvalues( - RangeOption { - mode: StreamingMode::WantAll, - ..(start.clone(), end).into() - }, - Snapshot, - ); - - // Points at the first key of the workflow currently being scanned. Resuming - // from here rescans that workflow from the start, which is required because a - // workflow is only indexed once all of its keys have been read. - let mut new_last_key = Some(start); - let mut current_workflow_id = None; - let mut name = None; - let mut error = None; - let mut state_matches = true; + // This remains a raw range boundary, accepting previously persisted first-key cursors. + let mut new_last_key = Some(last_key.clone().unwrap_or(start)); let fut = async { - while let Some(entry) = stream.try_next().await? { - let workflow_id = *self.subspace.unpack::(entry.key())?; - - if let Some(curr) = current_workflow_id { - if workflow_id != curr { - // Save if matches query - if let (Some(name), Some(error)) = (name.take(), error.take()) - && state_matches - { - tx.write( - &keys::workflow::DeadIdxKey::new(name, error, curr), - (), - )?; - } - - total += 1; - - // Reset state - new_last_key = Some(entry.key().to_vec()); - state_matches = true; - - // Stop on a workflow boundary so the cursor never points in - // the middle of a workflow's keys - if total >= limit { - return anyhow::Ok(()); - } - } + while total < limit { + let Some(start) = &new_last_key else { break }; + // A RocksDB range stream materializes all values before yielding. Seek only + // the next key so input, state, and output payloads never enter this scan. + let first_key = tx + .get_key( + &KeySelector::first_greater_or_equal(start.as_slice()), + Snapshot, + ) + .await?; + if first_key.is_empty() || first_key.as_slice() >= end.as_slice() { + new_last_key = None; + break; } - - current_workflow_id = Some(workflow_id); - - if let Ok(name_key) = - self.subspace.unpack::(entry.key()) - { - name = Some(name_key.deserialize(entry.value())?); - } else if let Ok(_) = self - .subspace - .unpack::(entry.key()) - { - state_matches = false; - } else if let Ok(_) = self - .subspace - .unpack::(entry.key()) - { - state_matches = false; - } else if let Ok(_) = self - .subspace - .unpack::(entry.key()) - { - state_matches = false; - } else if let Ok(_) = self + let workflow_id = *self.subspace.unpack::(&first_key)?; + let output_range = self .subspace - .unpack::(entry.key()) - { - state_matches = false; - } else if let Ok(error_key) = self - .subspace - .unpack::(entry.key()) - { - error = Some(error_key.deserialize(entry.value())?); - } - } - - // Save the last workflow in the range - if let Some(curr) = current_workflow_id { - if let (Some(name), Some(error)) = (name.take(), error.take()) - && state_matches + .subspace(&keys::workflow::OutputKey::new(workflow_id)) + .range(); + // get_key tracks the returned key, not the empty seek gap. Cover the entire + // output prefix so a concurrent completion cannot leave a stale dead index. + tx.add_conflict_range( + &output_range.0, + &output_range.1, + ConflictRangeType::Read, + )?; + let name_key = keys::workflow::NameKey::new(workflow_id); + let error_key = keys::workflow::ErrorKey::new(workflow_id); + let worker_key = keys::workflow::WorkerIdKey::new(workflow_id); + let wake_key = keys::workflow::HasWakeConditionKey::new(workflow_id); + let silence_key = keys::workflow::SilenceTsKey::new(workflow_id); + let output_selector = + KeySelector::first_greater_or_equal(output_range.0.as_slice()); + let (name, error, worker, wake, silenced, output_key) = tokio::try_join!( + tx.read_opt(&name_key, Serializable), + tx.read_opt(&error_key, Serializable), + tx.exists(&worker_key, Serializable), + tx.exists(&wake_key, Serializable), + tx.exists(&silence_key, Serializable), + tx.get_key(&output_selector, Snapshot), + )?; + let has_output = !output_key.is_empty() + && output_key.as_slice() < output_range.1.as_slice(); + if let (Some(name), Some(error)) = (name, error) + && !worker && !wake && !silenced + && !has_output { - tx.write(&keys::workflow::DeadIdxKey::new(name, error, curr), ())?; + tx.write( + &keys::workflow::DeadIdxKey::new(name, error, workflow_id), + (), + )?; } - + // Do not await between indexing and advancing. A timeout leaves an unfinished + // workflow at the cursor, while completed writes and progress commit together. + new_last_key = Some( + self.subspace + .subspace( + &keys::workflow::DataSubspaceKey::new_with_workflow_id( + workflow_id, + ), + ) + .range() + .1, + ); total += 1; } - - // Reached the end of all workflows - new_last_key = None; - anyhow::Ok(()) }; @@ -1403,7 +1377,6 @@ impl DatabaseDebug for DatabaseKv { Ok(res) => res?, Err(_) => tracing::debug!("timed out reading workflows"), } - Ok((total, new_last_key)) } }) diff --git a/engine/packages/gasoline/tests/modules/dead_workflow_backfill.rs b/engine/packages/gasoline/tests/modules/dead_workflow_backfill.rs new file mode 100644 index 0000000000..bea9e0e411 --- /dev/null +++ b/engine/packages/gasoline/tests/modules/dead_workflow_backfill.rs @@ -0,0 +1,145 @@ +//! Real storage regression for the backfill's memory bound and durable workflow cursor. + +use super::*; +use crate::db::Database; + +fn backfill_read_bytes() -> u64 { + rivet_metrics::REGISTRY + .gather() + .into_iter() + .filter(|family| family.name() == "rivet_udb_operation_bytes") + .flat_map(|family| family.get_metric().to_vec()) + .filter(|metric| { + let has_label = |name, value| { + metric + .get_label() + .iter() + .any(|label| label.name() == name && label.value() == value) + }; + has_label("name", "gas_debug_backfill_dead_workflows") && has_label("direction", "read") + }) + .map(|metric| metric.get_counter().value() as u64) + .sum() +} + +#[tokio::test] +async fn backfill_skips_large_payloads_and_preserves_classification_and_cursors() -> Result<()> { + let deps = rivet_test_deps::TestDeps::new().await?; + let db = ::new(deps.config().clone(), deps.pools().clone()).await?; + assert_eq!(db.backfill_dead_workflows(1, None).await?, (0, None)); + + let mut ids = (0..10) + .map(|_| Id::new_v1(deps.config().dc_label())) + .collect::>(); + ids.sort(); + deps.pools() + .udb()? + .txn("seed_backfill_regression", |tx| { + let ids = ids.clone(); + let subspace = db.subspace.clone(); + async move { + let tx = tx.with_subspace(subspace.clone()); + for (index, id) in ids.iter().copied().enumerate() { + if index != 7 && index != 9 { + tx.write( + &keys::workflow::NameKey::new(id), + "backfill-regression".to_string(), + )?; + } + if index != 8 && index != 9 { + tx.write(&keys::workflow::ErrorKey::new(id), "failed".to_string())?; + } + } + // One workflow contains many pages of irrelevant data. Classification must not load it. + let payload = vec![b'x'; 1024]; + for chunk in 0..4096 { + tx.set( + &subspace.pack(&keys::workflow::InputKey::new(ids[0]).chunk(chunk)), + &payload, + ); + } + tx.write(&keys::workflow::WorkerIdKey::new(ids[2]), ids[2])?; + tx.write(&keys::workflow::HasWakeConditionKey::new(ids[3]), ())?; + tx.write(&keys::workflow::SilenceTsKey::new(ids[4]), 1)?; + // Any output chunk excludes the workflow, even if chunk zero is absent. + let large_value = vec![b'x'; 256 * 1024]; + tx.set( + &subspace.pack(&keys::workflow::OutputKey::new(ids[5]).chunk(3)), + &large_value, + ); + // Discovery must read only the key even when the first and only value is large. + tx.set( + &subspace.pack(&keys::workflow::InputKey::new(ids[9]).chunk(0)), + &large_value, + ); + Ok(()) + } + }) + .await?; + + let before = backfill_read_bytes(); + let (count, mut cursor) = db.backfill_dead_workflows(1, None).await?; + assert_eq!(count, 1); + let bytes = backfill_read_bytes() - before; + assert!(bytes > 0, "backfill byte instrumentation was not observed"); + assert!( + bytes < 64 * 1024, + "classifying one workflow read {bytes} bytes of unrelated payload" + ); + assert!(cursor.is_some()); + // Recreate the database facade as on an activity retry, retaining only the durable cursor. + let db = ::new(deps.config().clone(), deps.pools().clone()).await?; + + // The pre-fix cursor is an actual first key of an unprocessed workflow, not an encoded struct. + let legacy_cursor = db.subspace.pack(&keys::workflow::NameKey::new(ids[1])); + let resumed = db.backfill_dead_workflows(1, Some(&legacy_cursor)).await?; + assert_eq!(resumed.0, 1); + let replayed = db.backfill_dead_workflows(1, Some(&legacy_cursor)).await?; + assert_eq!(resumed, replayed); + + let mut total = count; + for _ in 0..ids.len() + 1 { + let Some(previous) = cursor else { break }; + // Round-trip the persisted bytes through JSON exactly like the existing activity input. + let previous: Vec = serde_json::from_str(&serde_json::to_string(&previous)?)?; + let (count, next) = db.backfill_dead_workflows(1, Some(&previous)).await?; + total += count; + if let Some(next) = &next { + assert!( + next > &previous, + "cursor did not advance at a workflow boundary" + ); + } + cursor = next; + } + assert!(cursor.is_none()); + assert_eq!(total, ids.len()); + let total_bytes = backfill_read_bytes() - before; + assert!( + total_bytes < 64 * 1024, + "backfill discovery or output detection loaded payloads: {total_bytes} bytes" + ); + + let indexed = deps + .pools() + .udb()? + .txn("read_backfill_regression", |tx| { + let subspace = db.subspace.clone(); + async move { + let range = subspace.subspace(&keys::workflow::DeadIdxKey::subspace( + "backfill-regression".to_string(), + )); + tx.get_ranges_keyvalues((&range).into(), Snapshot) + .map(|entry| { + Ok(subspace + .unpack::(entry?.key())? + .workflow_id) + }) + .try_collect::>() + .await + } + }) + .await?; + assert_eq!(indexed, vec![ids[0], ids[1], ids[6]]); + Ok(()) +}