diff --git a/.vale/styles/config/vocabularies/technical/accept.txt b/.vale/styles/config/vocabularies/technical/accept.txt index 4156f480e8a..8cf9b9be539 100644 --- a/.vale/styles/config/vocabularies/technical/accept.txt +++ b/.vale/styles/config/vocabularies/technical/accept.txt @@ -18,6 +18,7 @@ newtype boolean(s?) ddsketch datadog +foldspace otel serde stdin @@ -95,6 +96,7 @@ inlining cloneable callsite enqueue(s|d|ing)? +dequeue(s|d|ing)? hostname backpressure misconfiguration diff --git a/Cargo.lock b/Cargo.lock index f42f82f9742..c5fa4852b9b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -907,7 +907,7 @@ version = "3.1.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "faf9468729b8cbcea668e36183cb69d317348c2e08e994829fb56ebfdfbaac34" dependencies = [ - "windows-sys 0.61.2", + "windows-sys 0.59.0", ] [[package]] @@ -1543,7 +1543,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb" dependencies = [ "libc", - "windows-sys 0.61.2", + "windows-sys 0.59.0", ] [[package]] @@ -1690,6 +1690,56 @@ version = "0.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "77ce24cb58228fbb8aa041425bb1050850ac19177686ea6e0f41a70416f56fdb" +[[package]] +name = "foldspace-core" +version = "0.1.0" +source = "git+https://github.com/DataDog/foldspace?rev=ed9290d5d3e8166c188ba2b94c8ea4f86253eb99#ed9290d5d3e8166c188ba2b94c8ea4f86253eb99" +dependencies = [ + "foldspace-eviction", + "prost", + "serde_json", + "tonic", + "tonic-prost", + "tonic-prost-build", + "zstd", +] + +[[package]] +name = "foldspace-eviction" +version = "0.1.0" +source = "git+https://github.com/DataDog/foldspace?rev=ed9290d5d3e8166c188ba2b94c8ea4f86253eb99#ed9290d5d3e8166c188ba2b94c8ea4f86253eb99" + +[[package]] +name = "foldspace-patterns" +version = "0.1.0" +source = "git+https://github.com/DataDog/foldspace?rev=ed9290d5d3e8166c188ba2b94c8ea4f86253eb99#ed9290d5d3e8166c188ba2b94c8ea4f86253eb99" +dependencies = [ + "foldspace-core", + "foldspace-eviction", + "foldspace-patterns-tokenizer", +] + +[[package]] +name = "foldspace-patterns-tokenizer" +version = "0.1.0" +source = "git+https://github.com/DataDog/foldspace?rev=ed9290d5d3e8166c188ba2b94c8ea4f86253eb99#ed9290d5d3e8166c188ba2b94c8ea4f86253eb99" +dependencies = [ + "regex-automata", + "serde_json", + "thiserror 1.0.69", +] + +[[package]] +name = "foldspace-server" +version = "0.1.0" +source = "git+https://github.com/DataDog/foldspace?rev=ed9290d5d3e8166c188ba2b94c8ea4f86253eb99#ed9290d5d3e8166c188ba2b94c8ea4f86253eb99" +dependencies = [ + "flate2", + "foldspace-core", + "prost", + "zstd", +] + [[package]] name = "form_urlencoded" version = "1.2.2" @@ -3008,7 +3058,7 @@ version = "0.50.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7957b9740744892f114936ab4a57b3f487491bbeafaf8083688b16841a4240e5" dependencies = [ - "windows-sys 0.61.2", + "windows-sys 0.59.0", ] [[package]] @@ -3619,6 +3669,8 @@ dependencies = [ "prettyplease", "prost", "prost-types", + "pulldown-cmark", + "pulldown-cmark-to-cmark", "regex", "syn", "tempfile", @@ -3698,6 +3750,26 @@ dependencies = [ "thiserror 1.0.69", ] +[[package]] +name = "pulldown-cmark" +version = "0.13.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e9f068eba8e7071c5f9511831b44f32c740d5adf574e990f946ddb53db2f314e" +dependencies = [ + "bitflags", + "memchr", + "unicase", +] + +[[package]] +name = "pulldown-cmark-to-cmark" +version = "22.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "50793def1b900256624a709439404384204a5dc3a6ec580281bfaac35e882e90" +dependencies = [ + "pulldown-cmark", +] + [[package]] name = "quanta" version = "0.12.6" @@ -4106,7 +4178,7 @@ dependencies = [ "errno", "libc", "linux-raw-sys 0.12.1", - "windows-sys 0.61.2", + "windows-sys 0.59.0", ] [[package]] @@ -4177,7 +4249,7 @@ dependencies = [ "security-framework", "security-framework-sys", "webpki-root-certs", - "windows-sys 0.61.2", + "windows-sys 0.59.0", ] [[package]] @@ -4326,9 +4398,13 @@ dependencies = [ "facet", "faster-hex", "figment", + "flate2", "float-cmp", "fnv", "foldhash 0.2.0", + "foldspace-core", + "foldspace-patterns", + "foldspace-server", "futures", "hashbrown 0.17.1", "headers", @@ -4377,6 +4453,7 @@ dependencies = [ "test-strategy", "tokio", "tokio-rustls", + "tokio-stream", "tonic", "tower", "tracing", @@ -4711,7 +4788,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5b55fb86dfd3a2f5f76ea78310a88f96c4ea21a3031f8d212443d56123fd0521" dependencies = [ "libc", - "windows-sys 0.61.2", + "windows-sys 0.59.0", ] [[package]] @@ -5241,7 +5318,7 @@ dependencies = [ "getrandom 0.4.2", "once_cell", "rustix 1.1.4", - "windows-sys 0.61.2", + "windows-sys 0.59.0", ] [[package]] @@ -5567,9 +5644,11 @@ dependencies = [ "hyper-util", "percent-encoding", "pin-project", + "rustls-native-certs", "socket2", "sync_wrapper", "tokio", + "tokio-rustls", "tokio-stream", "tower", "tower-layer", @@ -5909,6 +5988,12 @@ dependencies = [ "version_check", ] +[[package]] +name = "unicase" +version = "2.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "dbc4bc3a9f746d862c45cb89d705aa10f187bb96c76001afab07a0d35ce60142" + [[package]] name = "unicode-ident" version = "1.0.24" @@ -6204,7 +6289,7 @@ version = "0.1.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22" dependencies = [ - "windows-sys 0.61.2", + "windows-sys 0.59.0", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index 59bd190dffe..e385a5f75a9 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -114,6 +114,9 @@ bitmask-enum = { version = "2.2", default-features = false } facet = { version = "0.46.0", default-features = false, features = ["std"] } figment = { version = "0.10", default-features = false } foldhash = { version = "0.2", default-features = false, features = ["std"] } +foldspace-core = { git = "https://github.com/DataDog/foldspace", rev = "ed9290d5d3e8166c188ba2b94c8ea4f86253eb99" } +foldspace-patterns = { git = "https://github.com/DataDog/foldspace", rev = "ed9290d5d3e8166c188ba2b94c8ea4f86253eb99", features = ["core-adapter"] } +foldspace-server = { git = "https://github.com/DataDog/foldspace", rev = "ed9290d5d3e8166c188ba2b94c8ea4f86253eb99" } headers = { version = "0.4", default-features = false } http = { version = "1", default-features = false } http-body = { version = "1", default-features = false } @@ -164,6 +167,7 @@ similar-asserts = { version = "2.0", default-features = false } slab = { version = "0.4.12", default-features = false } syn = { version = "2", default-features = false, features = ["full", "parsing", "visit-mut"] } tokio-util = { version = "0.7.18", default-features = false } +tokio-stream = { version = "0.1", default-features = false } tower = { version = "0.5", default-features = false } tracing-subscriber = { version = "0.3", default-features = false } typify = { version = "0.7", default-features = false } @@ -212,6 +216,7 @@ crossbeam-queue = { version = "0.3", default-features = false, features = [ "alloc", ] } float-cmp = { version = "0.10", default-features = false } +flate2 = { version = "1", default-features = false, features = ["rust_backend"] } tower-http = { version = "0.7", default-features = false } bollard = { version = "0.21", default-features = false } home = { version = "0.5", default-features = false } diff --git a/LICENSE-3rdparty.csv b/LICENSE-3rdparty.csv index da407a92257..b1dde9becda 100644 --- a/LICENSE-3rdparty.csv +++ b/LICENSE-3rdparty.csv @@ -127,6 +127,11 @@ flate2,https://github.com/rust-lang/flate2-rs,MIT OR Apache-2.0,"Alex Crichton < float-cmp,https://github.com/mikedilger/float-cmp,MIT,Mike Dilger fnv,https://github.com/servo/rust-fnv,Apache-2.0 OR MIT,Alex Crichton foldhash,https://github.com/orlp/foldhash,Zlib,Orson Peters +foldspace-core,https://github.com/DataDog/foldspace,Apache-2.0,The foldspace-core Authors +foldspace-eviction,https://github.com/DataDog/foldspace,Apache-2.0,The foldspace-eviction Authors +foldspace-patterns,https://github.com/DataDog/foldspace,Apache-2.0,The foldspace-patterns Authors +foldspace-patterns-tokenizer,https://github.com/DataDog/foldspace,Apache-2.0,The foldspace-patterns-tokenizer Authors +foldspace-server,https://github.com/DataDog/foldspace,Apache-2.0,The foldspace-server Authors form_urlencoded,https://github.com/servo/rust-url,MIT OR Apache-2.0,The rust-url developers fs4,https://github.com/al8n/fs4-rs,MIT OR Apache-2.0,"Dan Burkert , Al Liu " futures,https://github.com/rust-lang/futures-rs,MIT OR Apache-2.0,The futures Authors @@ -301,6 +306,8 @@ prost-types,https://github.com/tokio-rs/prost,Apache-2.0,"Dan Burkert protobuf-parse,https://github.com/stepancheg/rust-protobuf/tree/master/protobuf-parse,MIT,Stepan Koltsov protobuf-support,https://github.com/stepancheg/rust-protobuf,MIT,Stepan Koltsov +pulldown-cmark,https://github.com/raphlinus/pulldown-cmark,MIT,"Raph Levien , Marcus Klaas de Vries " +pulldown-cmark-to-cmark,https://github.com/Byron/pulldown-cmark-to-cmark,Apache-2.0,"Sebastian Thiel , Dylan Owen , Alessandro Ogier , Zixian Cai <2891235+caizixian@users.noreply.github.com>, Andrew Lyjak " quanta,https://github.com/metrics-rs/quanta,MIT,Toby Lawrence quick_cache,https://github.com/arthurprs/quick-cache,MIT,Arthur Silva quote,https://github.com/dtolnay/quote,MIT OR Apache-2.0,David Tolnay @@ -455,6 +462,7 @@ typify-impl,https://github.com/oxidecomputer/typify,Apache-2.0,The typify-impl A ucd-trie,https://github.com/BurntSushi/ucd-generate,MIT OR Apache-2.0,Andrew Gallant unarray,https://github.com/cameron1024/unarray,MIT OR Apache-2.0,The unarray Authors uncased,https://github.com/SergioBenitez/uncased,MIT OR Apache-2.0,Sergio Benitez +unicase,https://github.com/seanmonstar/unicase,MIT OR Apache-2.0,Sean McArthur unicode-ident,https://github.com/dtolnay/unicode-ident,(MIT OR Apache-2.0) AND Unicode-3.0,David Tolnay unicode-segmentation,https://github.com/unicode-rs/unicode-segmentation,MIT OR Apache-2.0,"kwantam , Manish Goregaokar " unicode-width,https://github.com/unicode-rs/unicode-width,MIT OR Apache-2.0,"kwantam , Manish Goregaokar " diff --git a/bin/agent-data-plane/src/cli/run.rs b/bin/agent-data-plane/src/cli/run.rs index f84b6d4e6c5..cc80154f646 100644 --- a/bin/agent-data-plane/src/cli/run.rs +++ b/bin/agent-data-plane/src/cli/run.rs @@ -410,7 +410,8 @@ async fn create_topology( &metrics_encoding, &endpoints, ) - .error_context("Failed to configure Datadog forwarder.")?; + .error_context("Failed to configure Datadog forwarder.")? + .with_stateful_logs(config_system.config().domains.logs.stateful.clone()); blueprint.add_forwarder("dd_out", dd_forwarder_config)?; } diff --git a/docs/agent-data-plane/configuration/configuration.md b/docs/agent-data-plane/configuration/configuration.md index 83594e1458c..bdf6a62653d 100644 --- a/docs/agent-data-plane/configuration/configuration.md +++ b/docs/agent-data-plane/configuration/configuration.md @@ -458,6 +458,7 @@ The following settings are specific to ADP and have no equivalent in the core ag | `dogstatsd_string_interner_size_bytes` | Explicit byte budget for context interner | | | `dogstatsd_tcp_port` | TCP listen port for DSD | | | `flush_timeout_secs` | Encoder flush timeout (secs) | | +| `logs_config.use_grpc` | Use stateful Foldspace transport for logs | false | | `memory_limit` | Process memory limit | | | `memory_slop_factor` | Memory accounting slop fraction | 0.25 | | `otlp_allow_context_heap_allocs` | Allow heap allocations for OTLP contexts | | @@ -470,6 +471,10 @@ The following settings are specific to ADP and have no equivalent in the core ag | `otlp_string_interner_size` | OTLP context interner capacity | | | `serializer_max_metrics_per_payload` | Max metrics per payload | | +### `logs_config.use_grpc` + +Enables the initial stateful logs transport. Failed or queued stateful payloads are reconstructed and stored as complete stateless HTTP transactions before entering the normal retry or persisted queues. This does not provide full crash durability: a process crash before conversion can lose unacknowledged stateful payloads, and the persisted queue retains its existing dequeue crash window. + ### `data_plane.otlp.receiver_grpc_endpoint_temporary` Temporary development key for setting ADP's OTLP listen endpoints independently from the Agent's. diff --git a/lib/agent-data-plane-config-system/src/saluki_only.rs b/lib/agent-data-plane-config-system/src/saluki_only.rs index 0a4e00187b7..2af7c9ffd9f 100644 --- a/lib/agent-data-plane-config-system/src/saluki_only.rs +++ b/lib/agent-data-plane-config-system/src/saluki_only.rs @@ -160,6 +160,8 @@ pub struct SalukiOnly { pub apm_config: ApmConfig, /// OTLP receiver and trace knobs (`otlp_config.*`). pub otlp_config: OtlpConfig, + /// Logs transport knobs (`logs_config.*`). + pub logs_config: LogsConfig, /// OTTL span-drop filter (`ottl_filter_config`). pub ottl_filter_config: Option, /// OTTL span-transform processor (`ottl_transform_config`). @@ -213,6 +215,14 @@ pub struct DataPlaneChecks { pub enabled: Option, } +/// `logs_config.*` transport settings absent from the vendored Datadog schema. +#[derive(Clone, Debug, Default, Deserialize)] +#[serde(default)] +pub struct LogsConfig { + /// Whether stateful Foldspace transport is enabled (`logs_config.use_grpc`). + pub use_grpc: Option, +} + /// `apm_config.*` Saluki-only knobs. (The Datadog Agent publishes many other `apm_config.*` keys; /// those are witnessed and ignored here.) #[derive(Clone, Debug, Default, Deserialize)] @@ -418,6 +428,11 @@ impl SalukiOnly { config.shared.metrics_encoding.max_metrics_per_payload = v; } + // domains.logs + if let Some(v) = self.logs_config.use_grpc { + config.domains.logs.stateful.enabled = v; + } + // domains.dogstatsd let dsd = &mut config.domains.dogstatsd; if let Some(v) = self.dogstatsd_tcp_port { @@ -641,6 +656,7 @@ mod tests { "ignore_missing_datadog_fields": true } }, + "logs_config": { "use_grpc": true }, // top-level objects "ottl_filter_config": { "error_mode": "ignore", "traces": { "span": ["attributes[\"a\"] == \"b\""] } }, "ottl_transform_config": { "error_mode": "silent", "trace_statements": ["set(name, \"x\")"] }, @@ -663,6 +679,9 @@ mod tests { assert_eq!(config.shared.metrics_encoding.flush_timeout, Duration::from_secs(7)); assert_eq!(config.shared.metrics_encoding.max_metrics_per_payload, 999); + // domains.logs + assert!(config.domains.logs.stateful.enabled); + // domains.dogstatsd let dsd = &config.domains.dogstatsd; assert_eq!(dsd.listeners.tcp_port, 8126); @@ -786,5 +805,6 @@ mod tests { assert_eq!(otlp.receiver.grpc.endpoint, "localhost:6317"); assert_eq!(otlp.receiver.http.endpoint, "localhost:6318"); assert_eq!(otlp.traces.string_interner_size, DEFAULT_STRING_INTERNER_SIZE_BYTES); + assert!(!config.domains.logs.stateful.enabled); } } diff --git a/lib/agent-data-plane-config/src/domains/logs.rs b/lib/agent-data-plane-config/src/domains/logs.rs new file mode 100644 index 00000000000..dfc4e25f6e2 --- /dev/null +++ b/lib/agent-data-plane-config/src/domains/logs.rs @@ -0,0 +1,19 @@ +//! Logs domain configuration. + +use serde::Serialize; + +/// Resolved logs configuration. +#[derive(Clone, Debug, Default, PartialEq, Serialize)] +pub struct Domain { + /// Stateful Foldspace transport configuration. + pub stateful: StatefulEncoding, +} + +/// Stateful Foldspace transport configuration. +#[derive(Clone, Debug, Default, PartialEq, Serialize)] +pub struct StatefulEncoding { + /// Whether logs use stateful gRPC encoding. + /// + /// Defaults to `false`. Enable this only when the configured logs intake supports Foldspace. + pub enabled: bool, +} diff --git a/lib/agent-data-plane-config/src/domains/mod.rs b/lib/agent-data-plane-config/src/domains/mod.rs index 6db705f8b66..fca531ae751 100644 --- a/lib/agent-data-plane-config/src/domains/mod.rs +++ b/lib/agent-data-plane-config/src/domains/mod.rs @@ -6,6 +6,7 @@ use serde::Serialize; pub mod checks; pub mod dogstatsd; +pub mod logs; pub mod multi_region_failover; pub mod otlp; pub mod traces; @@ -14,6 +15,7 @@ pub mod traces; #[derive(Clone, Debug, Default, PartialEq, Serialize)] pub struct DomainConfiguration { pub dogstatsd: dogstatsd::Domain, + pub logs: logs::Domain, pub otlp: otlp::Domain, pub traces: traces::Domain, pub checks: checks::Domain, diff --git a/lib/datadog-agent/config-overlay-model/src/saluki_keys.rs b/lib/datadog-agent/config-overlay-model/src/saluki_keys.rs index 01510a9f127..dc2c77bae9e 100644 --- a/lib/datadog-agent/config-overlay-model/src/saluki_keys.rs +++ b/lib/datadog-agent/config-overlay-model/src/saluki_keys.rs @@ -22,6 +22,27 @@ pub struct SalukiKey { } pub static SALUKI_KEYS: &[SalukiKey] = &[ + // ── logs.rs ───────────────────────────────────────────────────────────── + SalukiKey { + yaml_path: "logs_config.use_grpc", + description: "Use stateful Foldspace transport for logs", + default: "false", + documentation: Some( + "Enables the initial stateful logs transport. Failed or queued stateful payloads are reconstructed and \ + stored as complete stateless HTTP transactions before entering the normal retry or persisted queues. \ + This does not provide full crash durability: a process crash before conversion can lose unacknowledged \ + stateful payloads, and the persisted queue retains its existing dequeue crash window.", + ), + value_type: "ValueType::Bool", + schema_default: Some("false"), + env_vars: &[], + env_var_override: None, + additional_yaml_paths: &[], + used_by: &["TYPED_CONFIG_SYSTEM"], + test_json: None, + pipeline_affinity: "PipelineAffinity::Pipelines(&[Pipeline::Checks, Pipeline::Otlp])", + filename: "logs.rs", + }, // ── data_plane.rs ──────────────────────────────────────────────────────── SalukiKey { yaml_path: "data_plane.otlp.receiver_grpc_endpoint_temporary", diff --git a/lib/datadog-agent/config-testing/src/config_registry/annotations_index.rs b/lib/datadog-agent/config-testing/src/config_registry/annotations_index.rs index dccb11d4590..86912d828a8 100644 --- a/lib/datadog-agent/config-testing/src/config_registry/annotations_index.rs +++ b/lib/datadog-agent/config-testing/src/config_registry/annotations_index.rs @@ -18,6 +18,7 @@ mod dogstatsd_prefix_filter; mod encoders; mod forwarder; mod get_typed; +mod logs; mod mrf; mod otlp; mod proxy; @@ -40,6 +41,7 @@ pub static SUPPORTED_ANNOTATIONS: LazyLock> = Laz v.extend_from_slice(encoders::ALL); v.extend_from_slice(forwarder::ALL); v.extend_from_slice(get_typed::ALL); + v.extend_from_slice(logs::ALL); v.extend_from_slice(mrf::ALL); v.extend_from_slice(otlp::ALL); v.extend_from_slice(proxy::ALL); diff --git a/lib/datadog-agent/config-testing/src/config_registry/logs.rs b/lib/datadog-agent/config-testing/src/config_registry/logs.rs new file mode 100644 index 00000000000..4cccfec45f4 --- /dev/null +++ b/lib/datadog-agent/config-testing/src/config_registry/logs.rs @@ -0,0 +1,27 @@ +// @generated by build.rs from schema_overlay.yaml — DO NOT EDIT +#[allow(unused_imports)] +use super::schema; +#[allow(unused_imports)] +use super::*; + +static LOGS_CONFIG_USE_GRPC_SCHEMA: SchemaEntry = SchemaEntry { + schema: Schema::Saluki, + yaml_path: "logs_config.use_grpc", + env_vars: &[], + value_type: ValueType::Bool, + default: Some("false"), +}; + +crate::declare_annotations! { + /// `logs_config.use_grpc` + LOGS_CONFIG_USE_GRPC = SalukiAnnotation { + schema: &LOGS_CONFIG_USE_GRPC_SCHEMA, + support_level: SupportLevel::Full, + additional_yaml_paths: &[], + env_var_override: None, + used_by: &[structs::TYPED_CONFIG_SYSTEM], + value_type_override: None, + test_json: None, + pipeline_affinity: PipelineAffinity::Pipelines(&[Pipeline::Checks, Pipeline::Otlp]), + }; +} diff --git a/lib/saluki-components/Cargo.toml b/lib/saluki-components/Cargo.toml index 435166fb74e..749a11df247 100644 --- a/lib/saluki-components/Cargo.toml +++ b/lib/saluki-components/Cargo.toml @@ -31,9 +31,13 @@ ddsketch = { workspace = true } facet = { workspace = true } faster-hex = { workspace = true } figment = { workspace = true } +flate2 = { workspace = true } float-cmp = { workspace = true, features = ["ratio"] } fnv = { workspace = true } foldhash = { workspace = true } +foldspace-core = { workspace = true } +foldspace-patterns = { workspace = true } +foldspace-server = { workspace = true } futures = { workspace = true } hashbrown = { workspace = true } headers = { workspace = true } @@ -84,7 +88,8 @@ tokio = { workspace = true, features = [ "signal", "sync", ] } -tonic = { workspace = true, features = ["router", "server"] } +tokio-stream = { workspace = true } +tonic = { workspace = true, features = ["channel", "router", "server", "tls-native-roots"] } tower = { workspace = true, features = ["limit", "retry", "timeout", "util"] } tracing = { workspace = true } tracing-appender = { workspace = true } diff --git a/lib/saluki-components/src/common/datadog/io.rs b/lib/saluki-components/src/common/datadog/io.rs index b9c8833d9fc..26b06fd0ec0 100644 --- a/lib/saluki-components/src/common/datadog/io.rs +++ b/lib/saluki-components/src/common/datadog/io.rs @@ -14,6 +14,7 @@ use std::{ time::{Duration, Instant}, }; +use agent_data_plane_config::domains::logs::StatefulEncoding; use bytes::Buf; use futures::FutureExt as _; use http::{Request, StatusCode, Uri}; @@ -49,6 +50,7 @@ use super::{ endpoints::{EndpointRoute, EndpointV3Settings, ResolvedEndpoint, RoutableEndpoint, V3EndpointConfig}, middleware::{for_resolved_endpoint, with_allow_arbitrary_tags, with_version_info}, retry_capacity::{TrafficRateWindow, RETRY_QUEUE_CAPACITY_BUCKET_DURATION_SECS}, + stateful_logs::{StatefulEvent, StatefulLogsSender}, telemetry::{ ComponentTelemetry, SharedTransactionQueueTelemetry, TransactionInputTelemetry, TransactionQueueTelemetry, TransactionRetryCounters, TransactionRetryTelemetry, @@ -66,10 +68,36 @@ struct InFlightTransaction { result: R, } +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +enum RetryTransport { + Foldspace, + Http, +} + +struct PendingRetry +where + B: Buf + Clone, +{ + transaction: Transaction, + counters: TransactionRetryCounters, + transport: RetryTransport, +} + type InFlightTransactionResult = InFlightTransaction, RetryCircuitBreakerError>>>>; type InFlightTaskResult = Result, JoinError>; +fn retry_counters_for_transaction( + transaction: &Transaction, retry_telemetry: &mut TransactionRetryTelemetry, endpoint_name: &EndpointNameFn, +) -> TransactionRetryCounters +where + B: Buf + Clone, +{ + let resolved = endpoint_name(transaction.request_uri()); + let logical = resolved.as_deref().unwrap_or_else(|| transaction.request_uri().path()); + retry_telemetry.counters_for(logical) +} + fn prepare_transaction_for_dispatch( pending_txn: PendingTransaction>, retry_telemetry: &mut TransactionRetryTelemetry, endpoint_name: &EndpointNameFn, @@ -81,9 +109,7 @@ where PendingTransaction::HighPriority(txn) => (txn, None), PendingTransaction::LowPriority(txn) => { // Low-priority provenance drives Core Agent-compatible retry accounting, including input overflow. - let resolved = endpoint_name(txn.request_uri()); - let logical = resolved.as_deref().unwrap_or_else(|| txn.request_uri().path()); - let counters = retry_telemetry.counters_for(logical); + let counters = retry_counters_for_transaction(&txn, retry_telemetry, endpoint_name); (txn, Some(counters)) } }; @@ -238,6 +264,7 @@ pub struct TransactionForwarder { endpoints: Vec, endpoint_request_mapper_factory: EndpointRequestMapperFactory, emitter: DiagnosticsEmitter, + stateful_logs: StatefulEncoding, _marker: PhantomData, } @@ -366,10 +393,17 @@ where endpoints, endpoint_request_mapper_factory, emitter, + stateful_logs: StatefulEncoding::default(), _marker: PhantomData, }) } + /// Configures stateful logs transport for endpoint workers. + pub fn with_stateful_logs(mut self, stateful_logs: StatefulEncoding) -> Self { + self.stateful_logs = stateful_logs; + self + } + /// Spawns the I/O task for the forwarder, and any associated endpoint I/O tasks. /// /// Returns a `Handle` that can be used to send transactions to the forwarder, as well as eventually shut it down in @@ -390,6 +424,7 @@ where endpoints, endpoint_request_mapper_factory, emitter, + stateful_logs, _marker, } = self; @@ -408,6 +443,7 @@ where endpoints, endpoint_request_mapper_factory, emitter, + stateful_logs, ), ); @@ -440,6 +476,7 @@ async fn run_io_loop( service: HttpClient, telemetry: ComponentTelemetry, metrics_builder: MetricsBuilder, endpoint_name: Arc, resolved_endpoints: Vec, endpoint_request_mapper_factory: EndpointRequestMapperFactory, emitter: DiagnosticsEmitter, + stateful_logs: StatefulEncoding, ) where B: Body + Buf + Clone + Send + Sync + 'static, B::Data: Send, @@ -476,6 +513,7 @@ async fn run_io_loop( live_config.clone(), service.clone(), telemetry.clone(), + metrics_builder.clone(), txnq_telemetry, retry_telemetry, Arc::clone(&endpoint_name), @@ -483,6 +521,7 @@ async fn run_io_loop( resolved_endpoint, endpoint_request_mapper_factory.clone(), emitter.clone(), + stateful_logs.clone(), ), ); @@ -576,10 +615,10 @@ fn track_transaction_input_for_endpoint( async fn run_endpoint_io_loop( mut txns_rx: mpsc::Receiver>, task_barrier: Arc, context: ComponentContext, config: ForwarderConfiguration, live_config: Option, service: HttpClient, - telemetry: ComponentTelemetry, txnq_telemetry: TransactionQueueTelemetry, + telemetry: ComponentTelemetry, metrics_builder: MetricsBuilder, txnq_telemetry: TransactionQueueTelemetry, mut retry_telemetry: TransactionRetryTelemetry, endpoint_name: Arc, route: EndpointRoute, endpoint: ResolvedEndpoint, endpoint_request_mapper_factory: EndpointRequestMapperFactory, - emitter: DiagnosticsEmitter, + emitter: DiagnosticsEmitter, stateful_logs: StatefulEncoding, ) where B: Body + Buf + Clone + Send + Sync + 'static, B::Data: Send, @@ -620,6 +659,23 @@ async fn run_endpoint_io_loop( ?endpoint_v3_settings, "Starting endpoint I/O task." ); + + let mut stateful_sender = if stateful_logs.enabled && route != EndpointRoute::MetricsPrimary { + match StatefulLogsSender::new( + endpoint.clone(), + &metrics_builder, + &endpoint_domain, + config.request_timeout(), + ) { + Ok(sender) => Some(sender), + Err(error) => { + error!(endpoint_url, %error, "Failed to initialize stateful logs; using stateless transport."); + None + } + } + } else { + None + }; let endpoint_request_mapper = Arc::new(Mutex::new((endpoint_request_mapper_factory)(endpoint))); // Build our endpoint service. @@ -687,12 +743,65 @@ async fn run_endpoint_io_loop( MetaString::from(endpoint_domain.as_str()), config.retry().capacity_time_interval_secs(), ); + if let Some(sender) = stateful_sender.as_mut() { + sender.start().await; + } let mut in_flight = JoinSet::new(); let mut transaction_input_telemetry_by_endpoint = FastHashMap::::default(); + let mut pending_retry = None; let mut done = false; loop { + if !done && pending_retry.is_none() && pending_txns.high_priority_is_empty() { + if let Some(transaction) = pending_txns.pop_low_priority().await { + let counters = + retry_counters_for_transaction(&transaction, &mut retry_telemetry, endpoint_name.as_ref()); + let transport = if stateful_sender + .as_ref() + .is_some_and(|sender| sender.owns_retry_transaction(&transaction)) + { + RetryTransport::Foldspace + } else { + RetryTransport::Http + }; + pending_retry = Some(PendingRetry { + transaction, + counters, + transport, + }); + } + } + + let can_attempt_stateful_retry = pending_retry + .as_ref() + .is_some_and(|retry| retry.transport == RetryTransport::Foldspace) + && stateful_sender + .as_ref() + .is_some_and(StatefulLogsSender::can_attempt_retry); + if !done && can_attempt_stateful_retry { + let retry = pending_retry.take().expect("a Foldspace retry was checked above"); + let sender = stateful_sender + .as_mut() + .expect("Foldspace retry ownership requires a stateful sender"); + match sender + .try_send_retry_transaction(retry.transaction, retry.counters.clone(), &telemetry, &endpoint_domain) + .await + { + Ok(()) => { + debug!(endpoint_url, "Retry sent through an existing Foldspace stream."); + continue; + } + Err(transaction) => { + pending_retry = Some(PendingRetry { + transaction, + counters: retry.counters, + transport: RetryTransport::Foldspace, + }); + } + } + } + select! { // Try and drain the next transaction from our channel, and push it into the pending transactions queue. maybe_txn = txns_rx.recv(), if !done => match maybe_txn { @@ -723,6 +832,17 @@ async fn run_endpoint_io_loop( transaction_size, ); + let txn = match stateful_sender.as_mut() { + Some(sender) => match sender + .try_send_transaction(txn, &mut pending_txns, &telemetry, &endpoint_domain) + .await + { + Ok(()) => continue, + Err(txn) => txn, + }, + None => txn, + }; + match pending_txns.push_high_priority(txn).await { Ok(push_result) => track_queue_drops(&telemetry, &endpoint_domain, push_result), Err(e) => error!(endpoint_url, error = %e, "Failed to enqueue transaction. Events may be permanently lost."), @@ -738,20 +858,31 @@ async fn run_endpoint_io_loop( // While we're not done and there are pending transactions, wait for the service to become ready and then // next the next available pending transaction. - svc = service.ready(), if !done && !pending_txns.is_empty() => match svc { - Ok(svc) => if let Some(pending_txn) = pending_txns.pop().await { - let (metadata, request, retry_counters) = prepare_transaction_for_dispatch( - pending_txn, - &mut retry_telemetry, - endpoint_name.as_ref(), - ); - in_flight.spawn(svc.call(request).map(move |result| InFlightTransaction { - metadata, - retry_counters, - result, - })); + svc = service.ready(), if !done && (!pending_txns.high_priority_is_empty() + || pending_retry.as_ref().is_some_and(|retry| retry.transport == RetryTransport::Http)) => match svc { + Ok(svc) => { + let pending_txn = pending_txns + .pop_high_priority() + .map(PendingTransaction::HighPriority) + .or_else(|| { + pending_retry + .take() + .map(|retry| PendingTransaction::LowPriority(retry.transaction)) + }); + if let Some(pending_txn) = pending_txn { + let (metadata, request, retry_counters) = prepare_transaction_for_dispatch( + pending_txn, + &mut retry_telemetry, + endpoint_name.as_ref(), + ); + in_flight.spawn(svc.call(request).map(move |result| InFlightTransaction { + metadata, + retry_counters, + result, + })); - debug!(endpoint_url, "Request sent."); + debug!(endpoint_url, "Request sent."); + } }, Err(e) => match e { RetryCircuitBreakerError::Service(e) => { @@ -777,10 +908,41 @@ async fn run_endpoint_io_loop( ).await; }, + maybe_event = async { + match stateful_sender.as_mut() { + Some(sender) => sender.next_event().await, + None => std::future::pending::>().await, + } + }, if !done => { + if let Some(event) = maybe_event { + if let Some(sender) = stateful_sender.as_mut() { + sender + .handle_event(event, &mut pending_txns, &telemetry, &endpoint_domain) + .await; + } + } + }, + else => break, } } + if let Some(retry) = pending_retry.take() { + match pending_txns.push_low_priority(retry.transaction).await { + Ok(push_result) => { + retry.counters.increment_requeued(); + track_queue_drops(&telemetry, &endpoint_domain, push_result); + } + Err(e) => { + error!(endpoint_url, error = %e, "Failed to preserve pending retry during shutdown. Events may be permanently lost.") + } + } + } + + if let Some(sender) = stateful_sender.as_mut() { + sender.shutdown(&mut pending_txns, &telemetry, &endpoint_domain).await; + } + // Flush any outstanding transactions in the pending transactions queue, which will potentially enqueue them to disk // if we have disk persistent enabled for the retry queue. match pending_txns.flush().await { @@ -840,7 +1002,7 @@ fn generate_retry_queue_id(context: ComponentContext, endpoint: &ResolvedEndpoin format!("{}/{}/{:x}", context.component_id(), endpoint_host, hash) } -fn track_queue_drops(telemetry: &ComponentTelemetry, domain: &str, push_result: PushResult) { +pub(super) fn track_queue_drops(telemetry: &ComponentTelemetry, domain: &str, push_result: PushResult) { if push_result.had_drops() { saluki_antithesis::sometimes!( true, @@ -917,7 +1079,7 @@ fn build_diagnostics_layer(emitter: DiagnosticsEmitter, endpoint_url: String) -> HttpInspectionLayer::new().with_inspector(StatusCode::FORBIDDEN, forbidden_inspector) } -enum PendingTransaction { +pub(super) enum PendingTransaction { HighPriority(T), LowPriority(T), } @@ -932,8 +1094,9 @@ enum PendingTransaction { /// Ultimately, we use this construction to provide a fast path for new transactions, while limiting the overall number /// of outstanding transactions that are waiting to be processed, with a bias towards preserving the most recent /// transactions so that fresh data can be sent as soon as any temporary networking issues are resolved. -struct PendingTransactions { +pub(super) struct PendingTransactions { high_priority: VecDeque, + max_high_priority: usize, low_priority: RetryQueue, telemetry: TransactionQueueTelemetry, domain: MetaString, @@ -951,6 +1114,7 @@ impl PendingTransactions { ) -> Self { Self { high_priority: VecDeque::with_capacity(max_enqueued), + max_high_priority: max_enqueued, low_priority: retry_queue, telemetry, domain, @@ -961,10 +1125,16 @@ impl PendingTransactions { } } - /// Returns `true` if there are no pending transactions. - /// - /// This includes both the high-priority and low-priority queues. - pub fn is_empty(&self) -> bool { + pub(super) fn high_priority_is_full_with(&self, externally_pending: usize) -> bool { + self.high_priority.len().saturating_add(externally_pending) > self.max_high_priority + } + + pub(super) fn high_priority_is_empty(&self) -> bool { + self.high_priority.is_empty() + } + + #[cfg(test)] + pub(super) fn is_empty(&self) -> bool { self.high_priority.is_empty() && self.low_priority.is_empty() } @@ -974,7 +1144,7 @@ impl PendingTransactions { pub async fn push_high_priority(&mut self, transaction: T) -> Result { self.record_incoming_transaction_size(transaction.size_bytes()).await; - if self.high_priority.len() < self.high_priority.capacity() { + if self.high_priority.len() < self.max_high_priority { self.high_priority.push_back(transaction); self.telemetry.high_prio_queue_insertions().increment(1); @@ -1012,24 +1182,28 @@ impl PendingTransactions { Ok(push_result) } - /// Pops the next transaction from the queue. - /// - /// The high-priority queue is drained first before attempting to pop from the low-priority queue. The returned - /// variant identifies the queue that contained the transaction. - pub async fn pop(&mut self) -> Option> { - // We bias towards handling enqueued transactions first, since those are our "high priority" transactions, and we - // want to keep them flowing as fast as possible. - loop { - if let Some(transaction) = self.high_priority.pop_front() { - self.telemetry.high_prio_queue_removals().increment(1); + pub(super) fn pop_high_priority(&mut self) -> Option { + let transaction = self.high_priority.pop_front()?; + self.telemetry.high_prio_queue_removals().increment(1); - debug!( - high_prio_queue_len = self.high_priority.len(), - "Dequeued pending transaction from high-priority queue." - ); - return Some(PendingTransaction::HighPriority(transaction)); - } + debug!( + high_prio_queue_len = self.high_priority.len(), + "Dequeued pending transaction from high-priority queue." + ); + Some(transaction) + } + #[cfg(test)] + pub(super) async fn pop(&mut self) -> Option> { + if let Some(transaction) = self.pop_high_priority() { + return Some(PendingTransaction::HighPriority(transaction)); + } + + self.pop_low_priority().await.map(PendingTransaction::LowPriority) + } + + pub(super) async fn pop_low_priority(&mut self) -> Option { + loop { let pop_result = self.low_priority.pop().await; let entries_dropped = self.low_priority.take_persisted_entries_dropped(); @@ -1048,7 +1222,7 @@ impl PendingTransactions { low_prio_queue_len = self.low_priority.len(), "Dequeued pending transaction from low-priority queue." ); - return Some(PendingTransaction::LowPriority(transaction)); + return Some(transaction); } Ok(None) => { self.record_retry_queue_size(); @@ -1056,7 +1230,6 @@ impl PendingTransactions { } Err(e) => { error!(error = %e, "Failed to pop transaction from low-priority queue."); - continue; } } } @@ -1449,10 +1622,11 @@ app.datadoghq.com: [key-a, key-b] .expect("overflow push should succeed"); assert!(!push_result.had_drops()); - let first = pending_txns - .pop() - .await - .expect("high-priority transaction should be queued"); + let first = PendingTransaction::HighPriority( + pending_txns + .pop_high_priority() + .expect("high-priority transaction should be queued"), + ); let (metadata, _request, retry_counters) = prepare_transaction_for_dispatch(first, &mut retry_telemetry, &test_logical_endpoint); assert!(retry_counters.is_none()); @@ -1475,7 +1649,12 @@ app.datadoghq.com: [key-a, key-b] None ); - let overflow = pending_txns.pop().await.expect("overflow transaction should be queued"); + let overflow = PendingTransaction::LowPriority( + pending_txns + .pop_low_priority() + .await + .expect("overflow transaction should be queued"), + ); let (metadata, _request, retry_counters) = prepare_transaction_for_dispatch(overflow, &mut retry_telemetry, &test_logical_endpoint); assert!(retry_counters.is_some()); @@ -1557,7 +1736,8 @@ app.datadoghq.com: [key-a, key-b] .await .expect("retry queue push should succeed"); assert!(!push_result.had_drops()); - let pending_txn = pending_txns.pop().await.expect("retry should be queued"); + let pending_txn = + PendingTransaction::LowPriority(pending_txns.pop_low_priority().await.expect("retry should be queued")); let (metadata, request, retry_counters) = prepare_transaction_for_dispatch(pending_txn, &mut retry_telemetry, &test_logical_endpoint); let result = service.call(request).await; @@ -1603,7 +1783,8 @@ app.datadoghq.com: [key-a, key-b] .await .expect("retry queue push should succeed"); assert!(!push_result.had_drops()); - let pending_txn = pending_txns.pop().await.expect("retry should be queued"); + let pending_txn = + PendingTransaction::LowPriority(pending_txns.pop_low_priority().await.expect("retry should be queued")); let (metadata, _request, retry_counters) = prepare_transaction_for_dispatch(pending_txn, &mut retry_telemetry, &test_logical_endpoint); @@ -1621,7 +1802,8 @@ app.datadoghq.com: [key-a, key-b] domain, ) .await; - assert!(pending_txns.is_empty()); + assert!(pending_txns.high_priority_is_empty()); + assert!(pending_txns.low_priority.is_empty()); assert_eq!( recorder.counter(("network_http_requests_retries_total", &retry_metric_tags(domain))), Some(1) @@ -1682,18 +1864,37 @@ app.datadoghq.com: [key-a, key-b] assert!(!push_result.had_drops()); assert_eq!(recorder.gauge("network_http_retry_queue_size"), Some(1.0)); assert_eq!(recorder.gauge("network_http_retry_queue_bytes_per_sec"), Some(20.0)); - assert!(matches!( - pending_txns.pop().await, - Some(PendingTransaction::LowPriority(transaction)) if transaction == "retry" - )); + assert_eq!(pending_txns.pop_low_priority().await.as_deref(), Some("retry")); assert_eq!(recorder.gauge("network_http_retry_queue_size"), Some(0.0)); assert_eq!(recorder.gauge("network_http_retry_queue_bytes_per_sec"), Some(20.0)); - assert!(pending_txns.pop().await.is_none()); + assert!(pending_txns.pop_low_priority().await.is_none()); assert_eq!(recorder.gauge("network_http_retry_queue_size"), Some(0.0)); assert_eq!(recorder.gauge("network_http_retry_queue_bytes_per_sec"), Some(20.0)); } + #[tokio::test] + async fn dedicated_retry_pop_leaves_fresh_transactions_queued() { + let (telemetry, domain) = transaction_queue_telemetry(); + let retry_queue = RetryQueue::new("test".to_string(), 1024); + let mut pending_txns = PendingTransactions::new(4, retry_queue, telemetry, domain, 900); + assert!(!pending_txns + .push_high_priority("fresh".to_string()) + .await + .unwrap() + .had_drops()); + assert!(!pending_txns + .push_low_priority("retry".to_string()) + .await + .unwrap() + .had_drops()); + + assert!(!pending_txns.high_priority_is_empty()); + assert_eq!(pending_txns.pop_low_priority().await.as_deref(), Some("retry")); + assert!(!pending_txns.high_priority_is_empty()); + assert_eq!(pending_txns.pop_high_priority().as_deref(), Some("fresh")); + } + #[tokio::test] async fn retry_queue_capacity_uses_remaining_in_memory_capacity() { let recorder = TestRecorder::default(); diff --git a/lib/saluki-components/src/common/datadog/mod.rs b/lib/saluki-components/src/common/datadog/mod.rs index 2b7db83433c..4ff7189b68d 100644 --- a/lib/saluki-components/src/common/datadog/mod.rs +++ b/lib/saluki-components/src/common/datadog/mod.rs @@ -9,6 +9,7 @@ mod proxy; pub mod request_builder; mod retry; mod retry_capacity; +mod stateful_logs; pub mod telemetry; pub mod transaction; pub mod validation; diff --git a/lib/saluki-components/src/common/datadog/stateful_logs.rs b/lib/saluki-components/src/common/datadog/stateful_logs.rs new file mode 100644 index 00000000000..e9870806730 --- /dev/null +++ b/lib/saluki-components/src/common/datadog/stateful_logs.rs @@ -0,0 +1,2150 @@ +//! Stateful Foldspace logs transport and stateless retry recovery. +//! +//! # Missing +//! +//! - Preserve stateful payloads across a process crash before stateless conversion. +//! - Remove the persisted retry queue's existing dequeue crash window. + +use std::{ + collections::VecDeque, + io::{Read as _, Write as _}, + time::Duration, +}; + +use bytes::Buf; +use chrono::{DateTime, Utc}; +use flate2::{ + read::{GzDecoder, ZlibDecoder}, + write::{GzEncoder, ZlibEncoder}, + Compression, +}; +use foldspace_core::{ + proto::stateful::{ + batch_status, stateful_intake_client::StatefulIntakeClient, BatchStatus, StatefulBatch as ProtoStatefulBatch, + }, + CoreError, DefaultBatchEncoder, DispatchPolicy, LogRecord, MplexEffect, MplexTimerKind, MultiStreamConfig, + MultiStreamCore, PayloadDelivery, PayloadId, PayloadRecovery, ProtoBatchEncoder, SenderId, StatelessLogRecord, + StreamError, StreamId, ZstdBatchCompressor, +}; +use foldspace_patterns::ClusteringPatternExtractor; +#[cfg(test)] +use foldspace_server::ContentEncoding as FoldspaceContentEncoding; +use foldspace_server::StatefulLogsDecoder; +use http::{ + header::{CONTENT_ENCODING, CONTENT_LENGTH}, + HeaderValue, Request, +}; +use saluki_common::collections::FastHashMap; +use saluki_error::{generic_error, ErrorContext as _, GenericError}; +use saluki_metrics::MetricsBuilder; +use serde_json::{Map as JsonMap, Value as JsonValue}; +use stringtheory::MetaString; +use tokio::sync::mpsc; +use tokio_stream::wrappers::ReceiverStream; +use tonic::{metadata::MetadataValue, transport::Channel, Request as TonicRequest}; +use tracing::{debug, error, warn}; +use uuid::Uuid; + +use super::{ + endpoints::ResolvedEndpoint, + io::{track_queue_drops, PendingTransactions}, + telemetry::{ComponentTelemetry, TransactionRetryCounters}, + transaction::{Metadata, Transaction, TransactionBody}, +}; + +const LOGS_INTAKE_PATH: &str = "/api/v2/logs"; +const STATEFUL_SENDERS: usize = 3; +const STATEFUL_BATCH_CAPACITY: usize = usize::MAX; +const STATEFUL_CHANNEL_CAPACITY: usize = 1; +const STATELESS_BATCH_MAX_BYTES: usize = 5 * 1024 * 1024; +const STATE_REQUEST_BYTES: u64 = 5 * 1024 * 1024; +const DUAL_SEND_UUID_FIELD: &str = "dual-send-uuid"; + +type StatefulCore = MultiStreamCore; + +#[derive(Debug)] +struct LogTemplate { + object: JsonMap, +} + +struct RetainedRequest +where + B: Buf + Clone, +{ + metadata: Metadata, + request: Request>, + templates: Vec, +} + +struct RetainedStatelessRequest +where + B: Buf + Clone, +{ + metadata: Metadata, + request: Request>, + retry_counters: TransactionRetryCounters, +} + +enum RetainedPayload +where + B: Buf + Clone, +{ + Stateful(RetainedRequest), + StatelessRetry(RetainedStatelessRequest), +} + +impl RetainedPayload +where + B: Buf + Clone, +{ + fn metadata(&self) -> &Metadata { + match self { + Self::Stateful(retained) => &retained.metadata, + Self::StatelessRetry(retained) => &retained.metadata, + } + } + + fn retry_counters(&self) -> Option<&TransactionRetryCounters> { + match self { + Self::Stateful(_) => None, + Self::StatelessRetry(retained) => Some(&retained.retry_counters), + } + } +} + +impl RetainedRequest +where + B: Buf + Clone, +{ + fn into_transaction(self, logs: Vec) -> Result, GenericError> { + if logs.len() != self.templates.len() { + return Err(generic_error!( + "Foldspace recovery produced {} logs for {} retained templates.", + logs.len(), + self.templates.len() + )); + } + + let mut values = Vec::with_capacity(logs.len()); + for (mut template, log) in self.templates.into_iter().zip(logs) { + template + .object + .insert("message".to_owned(), JsonValue::String(log.message)); + if let Some(uuid) = log.uuid { + template + .object + .insert(DUAL_SEND_UUID_FIELD.to_owned(), JsonValue::String(uuid)); + } + values.push(JsonValue::Object(template.object)); + } + + transaction_from_values(self.metadata, self.request, values) + } +} + +impl RetainedStatelessRequest +where + B: Buf + Clone, +{ + fn into_transaction(self, logs: Vec) -> Result, GenericError> { + let mut values = Vec::with_capacity(logs.len()); + for log in logs { + let original_json = log + .original_json + .ok_or_else(|| generic_error!("Self-contained Foldspace log is missing its original JSON object."))?; + let value: JsonValue = serde_json::from_slice(&original_json) + .error_context("Failed to parse retained self-contained log JSON.")?; + if !value.is_object() { + return Err(generic_error!("Self-contained Foldspace log JSON is not an object.")); + } + values.push(value); + } + + let mut metadata = self.metadata; + metadata.event_count = values.len(); + transaction_from_values(metadata, self.request, values) + } +} + +fn transaction_from_values( + metadata: Metadata, request: Request>, values: Vec, +) -> Result, GenericError> +where + B: Buf + Clone, +{ + let uncompressed = serde_json::to_vec(&values).error_context("Failed to encode recovered stateless logs.")?; + let encoding = content_encoding(request.headers())?; + let body = compress_body(&uncompressed, encoding)?; + let mut request = request.map(|_| TransactionBody::from(body)); + if request.headers().contains_key(CONTENT_LENGTH) { + let content_length = HeaderValue::from_str(&request.body().remaining().to_string()) + .error_context("Failed to update recovered logs content length.")?; + request.headers_mut().insert(CONTENT_LENGTH, content_length); + } + Ok(Transaction::reassemble(metadata, request)) +} + +struct StatelessRecovery +where + B: Buf + Clone, +{ + recovery: PayloadRecovery, + transaction: Transaction, + retry_counters: Option, +} + +enum TransactionEncoding +where + B: Buf + Clone, +{ + Encoded { + payload_id: PayloadId, + effects: Vec, + }, + Passthrough(Box>), +} + +enum PreparedRecovery +where + B: Buf + Clone, +{ + Reconstructed(Box>), + Dropped(PayloadRecovery), + Resolved, + Failed, +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +enum BodyEncoding { + Identity, + Gzip, + Deflate, + Zstd, +} + +fn content_encoding(headers: &http::HeaderMap) -> Result { + let Some(value) = headers.get(CONTENT_ENCODING) else { + return Ok(BodyEncoding::Identity); + }; + match value.to_str().unwrap_or_default().trim().to_ascii_lowercase().as_str() { + "" | "identity" => Ok(BodyEncoding::Identity), + "gzip" => Ok(BodyEncoding::Gzip), + "deflate" => Ok(BodyEncoding::Deflate), + "zstd" => Ok(BodyEncoding::Zstd), + other => Err(generic_error!( + "Stateful logs do not support HTTP content encoding '{}'.", + other + )), + } +} + +fn copy_body(body: &TransactionBody) -> Vec +where + B: Buf + Clone, +{ + let mut body = body.clone(); + let mut bytes = Vec::with_capacity(body.remaining()); + while body.has_remaining() { + let chunk = body.chunk(); + bytes.extend_from_slice(chunk); + let chunk_len = chunk.len(); + body.advance(chunk_len); + } + bytes +} + +fn decompress_body(bytes: &[u8], encoding: BodyEncoding) -> Result, GenericError> { + match encoding { + BodyEncoding::Identity => Ok(bytes.to_vec()), + BodyEncoding::Gzip => { + let mut decoder = GzDecoder::new(bytes); + let mut decoded = Vec::new(); + decoder + .read_to_end(&mut decoded) + .error_context("Failed to decompress gzip logs transaction.")?; + Ok(decoded) + } + BodyEncoding::Deflate => { + let mut decoder = ZlibDecoder::new(bytes); + let mut decoded = Vec::new(); + decoder + .read_to_end(&mut decoded) + .error_context("Failed to decompress deflate logs transaction.")?; + Ok(decoded) + } + BodyEncoding::Zstd => { + zstd::stream::decode_all(bytes).error_context("Failed to decompress zstd logs transaction.") + } + } +} + +fn compress_body(bytes: &[u8], encoding: BodyEncoding) -> Result, GenericError> { + match encoding { + BodyEncoding::Identity => Ok(bytes.to_vec()), + BodyEncoding::Gzip => { + let mut encoder = GzEncoder::new(Vec::new(), Compression::default()); + encoder + .write_all(bytes) + .error_context("Failed to gzip recovered logs transaction.")?; + encoder + .finish() + .error_context("Failed to finish recovered gzip logs transaction.") + } + BodyEncoding::Deflate => { + let mut encoder = ZlibEncoder::new(Vec::new(), Compression::default()); + encoder + .write_all(bytes) + .error_context("Failed to deflate recovered logs transaction.")?; + encoder + .finish() + .error_context("Failed to finish recovered deflate logs transaction.") + } + BodyEncoding::Zstd => { + zstd::stream::encode_all(bytes, 0).error_context("Failed to zstd-compress recovered logs transaction.") + } + } +} + +fn parse_request_logs( + request: &Request>, +) -> Result<(Vec, Vec), GenericError> +where + B: Buf + Clone, +{ + let objects = parse_request_log_objects(request)?; + let mut templates = Vec::with_capacity(objects.len()); + let mut records = Vec::with_capacity(objects.len()); + for mut object in objects { + let record = log_record_from_object(&object, true)?; + object.remove("message"); + object.remove(DUAL_SEND_UUID_FIELD); + records.push(record); + templates.push(LogTemplate { object }); + } + Ok((templates, records)) +} + +fn parse_request_log_records(request: &Request>) -> Result, GenericError> +where + B: Buf + Clone, +{ + parse_request_log_objects(request)? + .iter() + .map(|object| { + let record = log_record_from_object(object, false)?; + let original_json = serde_json::to_vec(object).error_context("Failed to preserve complete retry log.")?; + Ok(StatelessLogRecord::new(record, original_json)) + }) + .collect() +} + +fn parse_request_log_objects( + request: &Request>, +) -> Result>, GenericError> +where + B: Buf + Clone, +{ + let encoding = content_encoding(request.headers())?; + let encoded = copy_body(request.body()); + let decoded = decompress_body(&encoded, encoding)?; + let values: Vec = + serde_json::from_slice(&decoded).error_context("Failed to parse stateless logs transaction.")?; + values + .into_iter() + .map(|value| match value { + JsonValue::Object(object) => Ok(object), + _ => Err(generic_error!("Logs transaction contained a non-object entry.")), + }) + .collect() +} + +fn log_record_from_object(object: &JsonMap, generate_uuid: bool) -> Result { + let timestamp_millis = log_timestamp_millis(object)?; + let message = required_string(object, "message")?; + let status = optional_string(object, "status")?; + let hostname = optional_string(object, "hostname")?; + let service = optional_string(object, "service")?; + let source = optional_string(object, "ddsource")?; + let tags = optional_string(object, "ddtags")?; + let mut uuid = optional_string(object, DUAL_SEND_UUID_FIELD)?; + if uuid.is_none() && generate_uuid { + uuid = Some(Uuid::now_v7().to_string()); + } + + let mut record = LogRecord::new(message.into_bytes(), timestamp_millis); + record.status = status; + record.hostname = hostname; + record.service = service; + record.source = source; + record.tags = tags + .as_deref() + .map(|tags| tags.split(',').map(ToOwned::to_owned).collect()) + .unwrap_or_default(); + record.uuid = uuid; + Ok(record) +} + +fn required_string(object: &JsonMap, field: &str) -> Result { + optional_string(object, field)?.ok_or_else(|| generic_error!("Log entry is missing string field '{}'.", field)) +} + +fn optional_string(object: &JsonMap, field: &str) -> Result, GenericError> { + match object.get(field) { + Some(JsonValue::String(value)) => Ok(Some(value.clone())), + Some(_) => Err(generic_error!("Log field '{}' is not a string.", field)), + None => Ok(None), + } +} + +fn log_timestamp_millis(object: &JsonMap) -> Result { + if let Some(value) = object.get("timestamp") { + return value + .as_i64() + .ok_or_else(|| generic_error!("Log timestamp is not an integer.")); + } + if let Some(value) = object.get("@timestamp") { + let value = value + .as_str() + .ok_or_else(|| generic_error!("Log @timestamp is not a string."))?; + return DateTime::parse_from_rfc3339(value) + .map(|timestamp| timestamp.timestamp_millis()) + .error_context("Failed to parse log @timestamp."); + } + Ok(Utc::now().timestamp_millis()) +} + +#[derive(Debug)] +struct OutboundBatch { + proto: ProtoStatefulBatch, + payload_id: Option, +} + +#[derive(Debug)] +struct InflightBatch { + batch_id: u64, + payload_id: Option, +} + +#[derive(Debug)] +struct SenderStream { + stream_id: StreamId, + outbound: mpsc::Sender, + inflight: Option, + pending: VecDeque, +} + +impl SenderStream { + fn payload_delivery(&self) -> Option<(PayloadId, PayloadDelivery)> { + if let Some(payload_id) = self.inflight.as_ref().and_then(|batch| batch.payload_id) { + return Some((payload_id, PayloadDelivery::Unknown)); + } + self.pending + .iter() + .find_map(|batch| batch.payload_id) + .map(|payload_id| (payload_id, PayloadDelivery::NotSent)) + } + + fn remove_payload(&mut self, payload_id: PayloadId) { + if self + .inflight + .as_ref() + .is_some_and(|batch| batch.payload_id == Some(payload_id)) + { + self.inflight = None; + } + self.pending.retain(|batch| batch.payload_id != Some(payload_id)); + } +} + +#[derive(Debug)] +pub(super) enum StatefulEvent { + Opened { + sender_id: SenderId, + stream_id: StreamId, + outbound: mpsc::Sender, + }, + OpenFailed { + sender_id: SenderId, + stream_id: StreamId, + error: MetaString, + }, + Ack { + sender_id: SenderId, + stream_id: StreamId, + batch_id: u64, + }, + Failed { + sender_id: SenderId, + stream_id: StreamId, + error: MetaString, + failed_payload: Option<(PayloadId, PayloadDelivery)>, + }, + Timer { + sender_id: SenderId, + kind: MplexTimerKind, + }, +} + +#[derive(Clone)] +struct StatefulTelemetry { + converted: metrics::Counter, + conversion_errors: metrics::Counter, + oversized: metrics::Counter, +} + +impl StatefulTelemetry { + fn new(builder: &MetricsBuilder, domain: &str) -> Self { + let domain_tag = ("domain", domain.to_owned()); + Self { + converted: builder + .register_counter_with_tags("stateful_logs_payloads_converted_total", [domain_tag.clone()]), + conversion_errors: builder + .register_counter_with_tags("stateful_logs_conversion_errors_total", [domain_tag.clone()]), + oversized: builder.register_counter_with_tags("stateful_logs_stateless_oversized_total", [domain_tag]), + } + } +} + +pub(super) struct StatefulLogsSender +where + B: Buf + Clone, +{ + core: StatefulCore, + client: StatefulIntakeClient, + encoder: ProtoBatchEncoder, + endpoint: ResolvedEndpoint, + open_timeout: Duration, + max_stateless_batch_bytes: usize, + streams: FastHashMap, + retained: FastHashMap>, + events_tx: mpsc::UnboundedSender, + events_rx: mpsc::UnboundedReceiver, + telemetry: StatefulTelemetry, + disabled: bool, +} + +impl StatefulLogsSender +where + B: Buf + Clone, +{ + pub(super) fn new( + endpoint: ResolvedEndpoint, metrics_builder: &MetricsBuilder, endpoint_domain: &str, request_timeout: Duration, + ) -> Result { + let authority = endpoint + .logs_authority() + .map(ToString::to_string) + .unwrap_or_else(|| endpoint.endpoint().authority().to_owned()); + let grpc_endpoint = format!("{}://{}", endpoint.endpoint().scheme(), authority); + let channel = Channel::from_shared(grpc_endpoint.clone()) + .error_context("Failed to build Foldspace endpoint.")? + .connect_timeout(request_timeout) + .connect_lazy(); + let (events_tx, events_rx) = mpsc::unbounded_channel(); + let core = MultiStreamCore::with_encoder_and_extractor( + MultiStreamConfig { + senders: STATEFUL_SENDERS, + dispatch: DispatchPolicy::Reliable, + batch_capacity: STATEFUL_BATCH_CAPACITY, + max_outstanding_payloads: 1, + fold_catch_up: true, + ..MultiStreamConfig::default() + }, + DefaultBatchEncoder, + ClusteringPatternExtractor::new(), + ); + + debug!(endpoint = %grpc_endpoint, "Configured endpoint-scoped Foldspace sender."); + Ok(Self { + core, + client: StatefulIntakeClient::new(channel), + encoder: ProtoBatchEncoder::new(ZstdBatchCompressor::default()), + endpoint, + open_timeout: request_timeout, + max_stateless_batch_bytes: STATELESS_BATCH_MAX_BYTES, + streams: FastHashMap::default(), + retained: FastHashMap::default(), + events_tx, + events_rx, + telemetry: StatefulTelemetry::new(metrics_builder, endpoint_domain), + disabled: false, + }) + } + + pub(super) async fn start(&mut self) { + let effects = self.core.start(); + self.execute(effects).await; + } + + pub(super) async fn next_event(&mut self) -> Option { + self.events_rx.recv().await + } + + pub(super) async fn try_send_transaction( + &mut self, transaction: Transaction, pending: &mut PendingTransactions>, + component_telemetry: &ComponentTelemetry, endpoint_domain: &str, + ) -> Result<(), Transaction> { + if self.disabled || transaction.request_uri().path() != LOGS_INTAKE_PATH { + return Err(transaction); + } + let (payload_id, effects) = match self.encode_transaction(transaction) { + TransactionEncoding::Encoded { payload_id, effects } => (payload_id, effects), + TransactionEncoding::Passthrough(transaction) => return Err(*transaction), + }; + + let high_priority_full = pending.high_priority_is_full_with(self.retained.len()); + if high_priority_full { + let _ = self + .recover_to_retry( + payload_id, + PayloadDelivery::NotSent, + pending, + component_telemetry, + endpoint_domain, + ) + .await; + } else { + self.execute(effects).await; + } + Ok(()) + } + + pub(super) fn can_attempt_retry(&self) -> bool { + !self.disabled && self.core.can_dispatch_stateless_logs() + } + + pub(super) fn owns_retry_transaction(&self, transaction: &Transaction) -> bool { + transaction.request_uri().path() == LOGS_INTAKE_PATH + } + + pub(super) async fn try_send_retry_transaction( + &mut self, transaction: Transaction, retry_counters: TransactionRetryCounters, + component_telemetry: &ComponentTelemetry, endpoint_domain: &str, + ) -> Result<(), Transaction> { + if !self.can_attempt_retry() || transaction.request_uri().path() != LOGS_INTAKE_PATH { + return Err(transaction); + } + let records = match parse_request_log_records(transaction.request()) { + Ok(records) if !records.is_empty() => records, + Ok(_) | Err(_) => { + component_telemetry.track_permanently_failed_transaction(transaction.metadata(), None, endpoint_domain); + return Ok(()); + } + }; + let log_count = records.len(); + let dispatch = match self + .core + .try_dispatch_stateless_logs(records, self.max_stateless_batch_bytes) + { + Ok(dispatched) => dispatched, + Err(rejected) => { + drop(rejected.into_logs()); + return Err(transaction); + } + }; + let (payload_id, effects, dropped_logs) = dispatch.into_parts(); + if dropped_logs > 0 { + self.telemetry.oversized.increment(dropped_logs as u64); + let mut dropped_metadata = transaction.metadata().clone(); + dropped_metadata.event_count = dropped_logs; + dropped_metadata.data_point_count = 0; + component_telemetry.track_permanently_failed_transaction(&dropped_metadata, None, endpoint_domain); + } + let Some(payload_id) = payload_id else { + return Ok(()); + }; + let sends_on_existing_stream = !effects.is_empty() + && effects.iter().all(|effect| match effect { + MplexEffect::SendBatch { sender_id, batch, .. } => self + .streams + .get(&sender_id.get()) + .is_some_and(|stream| stream.stream_id == batch.stream), + _ => false, + }); + if !sends_on_existing_stream { + error!( + payload_id = payload_id.get(), + "Foldspace accepted a stateless retry without an existing adapter stream." + ); + self.disabled = true; + if let Some(sender_id) = self.core.payload_sender(payload_id) { + if let Ok(recovery) = self + .core + .begin_recovery(sender_id, payload_id, PayloadDelivery::NotSent) + { + let effects = self.core.complete_recovery(recovery); + self.execute(effects).await; + } + } + return Err(transaction); + } + let (mut metadata, request) = transaction.into_parts(); + metadata.event_count = log_count.saturating_sub(dropped_logs); + self.retained.insert( + payload_id.get(), + RetainedPayload::StatelessRetry(RetainedStatelessRequest { + metadata, + request: request.map(|_| TransactionBody::from(Vec::new())), + retry_counters, + }), + ); + self.execute(effects).await; + Ok(()) + } + + fn encode_transaction(&mut self, transaction: Transaction) -> TransactionEncoding { + let (metadata, request) = transaction.into_parts(); + let (templates, records) = match parse_request_logs(&request) { + Ok((templates, records)) if !records.is_empty() => (templates, records), + Ok(_) | Err(_) => { + return TransactionEncoding::Passthrough(Box::new(Transaction::reassemble(metadata, request))) + } + }; + let retained = RetainedRequest { + metadata, + request: request.map(|_| TransactionBody::from(Vec::new())), + templates, + }; + for record in &records { + let effects = self.core.push_log(record); + debug_assert!( + effects.is_empty(), + "one HTTP transaction must map to one Foldspace payload" + ); + } + let (payload_id, effects) = self.core.flush_with_payload_id(); + let payload_id = payload_id.expect("a non-empty logs transaction must flush one payload"); + self.retained + .insert(payload_id.get(), RetainedPayload::Stateful(retained)); + TransactionEncoding::Encoded { payload_id, effects } + } + + pub(super) async fn handle_event( + &mut self, event: StatefulEvent, pending: &mut PendingTransactions>, + component_telemetry: &ComponentTelemetry, endpoint_domain: &str, + ) { + match event { + StatefulEvent::Opened { + sender_id, + stream_id, + outbound, + } => { + self.streams.insert( + sender_id.get(), + SenderStream { + stream_id, + outbound, + inflight: None, + pending: VecDeque::new(), + }, + ); + let effects = self.core.handle_stream_opened(sender_id, stream_id); + self.execute(effects).await; + } + StatefulEvent::OpenFailed { + sender_id, + stream_id, + error, + } => { + let effects = self + .core + .handle_stream_error(sender_id, stream_id, StreamError::new(error.as_ref())); + self.execute(effects).await; + } + StatefulEvent::Ack { + sender_id, + stream_id, + batch_id, + } => { + let Some(stream) = self.streams.get_mut(&sender_id.get()) else { + return; + }; + if stream.stream_id != stream_id { + return; + } + let payload_id = if stream.inflight.as_ref().is_some_and(|batch| batch.batch_id == batch_id) { + stream.inflight.take().and_then(|batch| batch.payload_id) + } else { + None + }; + let effects = self.core.handle_ack(sender_id, stream_id, batch_id); + if let Some(payload_id) = + payload_id.filter(|payload_id| self.core.payload_sender(*payload_id).is_none()) + { + if let Some(retained) = self.retained.remove(&payload_id.get()) { + component_telemetry.track_successful_transaction(retained.metadata(), endpoint_domain); + } + } + self.execute(effects).await; + self.try_send(sender_id).await; + } + StatefulEvent::Failed { + sender_id, + stream_id, + error, + failed_payload, + } => { + let Some(stream) = self.streams.get(&sender_id.get()) else { + return; + }; + if stream.stream_id != stream_id { + return; + } + let stream = self + .streams + .remove(&sender_id.get()) + .expect("Foldspace stream was checked above"); + let recovery = failed_payload.or_else(|| stream.payload_delivery()); + let effects = self + .core + .handle_stream_error(sender_id, stream_id, StreamError::new(error.as_ref())); + if let Some((payload_id, delivery)) = recovery { + if !self + .recover_to_retry(payload_id, delivery, pending, component_telemetry, endpoint_domain) + .await + { + return; + } + } + self.execute(effects).await; + } + StatefulEvent::Timer { sender_id, kind } => { + let effects = self.core.handle_timer(sender_id, kind); + self.execute(effects).await; + } + } + } + + pub(super) async fn shutdown( + &mut self, pending: &mut PendingTransactions>, component_telemetry: &ComponentTelemetry, + endpoint_domain: &str, + ) { + let payloads = self.retained.keys().copied().map(PayloadId).collect::>(); + for payload_id in payloads { + let Some(sender_id) = self.core.payload_sender(payload_id) else { + continue; + }; + let delivery = self + .streams + .get(&sender_id.get()) + .and_then(SenderStream::payload_delivery) + .filter(|(candidate, _)| *candidate == payload_id) + .map_or(PayloadDelivery::NotSent, |(_, delivery)| delivery); + if delivery == PayloadDelivery::Unknown { + if let Some(stream) = self.streams.remove(&sender_id.get()) { + let _ = self.core.handle_stream_error( + sender_id, + stream.stream_id, + StreamError::new("graceful shutdown"), + ); + } + } + let _ = self + .recover_to_retry(payload_id, delivery, pending, component_telemetry, endpoint_domain) + .await; + } + self.streams.clear(); + } + + async fn recover_to_retry( + &mut self, payload_id: PayloadId, delivery: PayloadDelivery, pending: &mut PendingTransactions>, + component_telemetry: &ComponentTelemetry, endpoint_domain: &str, + ) -> bool { + let recovered = + match self.prepare_stateless_recovery(payload_id, delivery, component_telemetry, endpoint_domain) { + PreparedRecovery::Reconstructed(recovered) => *recovered, + PreparedRecovery::Dropped(recovery) => { + let effects = self.core.complete_recovery(recovery); + self.execute(effects).await; + return true; + } + PreparedRecovery::Resolved => return true, + PreparedRecovery::Failed => return false, + }; + let metadata = recovered.transaction.metadata().clone(); + match pending.push_low_priority(recovered.transaction).await { + Ok(push_result) => { + if delivery == PayloadDelivery::NotSent { + if let Some(counters) = recovered.retry_counters.as_ref() { + counters.increment_requeued(); + } + } + self.telemetry.converted.increment(1); + track_queue_drops(component_telemetry, endpoint_domain, push_result); + } + Err(error) => { + self.telemetry.oversized.increment(1); + component_telemetry.track_permanently_failed_transaction(&metadata, None, endpoint_domain); + error!( + payload_id = payload_id.get(), + %error, + "Recovered stateless payload exceeds retry queue limits; dropping it without stateful fallback." + ); + } + } + let effects = self.core.complete_recovery(recovered.recovery); + self.execute(effects).await; + true + } + + fn prepare_stateless_recovery( + &mut self, payload_id: PayloadId, delivery: PayloadDelivery, component_telemetry: &ComponentTelemetry, + endpoint_domain: &str, + ) -> PreparedRecovery { + let Some(sender_id) = self.core.payload_sender(payload_id) else { + return PreparedRecovery::Resolved; + }; + let recovery = match self.core.begin_recovery(sender_id, payload_id, delivery) { + Ok(recovery) => recovery, + Err(error) => { + self.disabled = true; + self.telemetry.conversion_errors.increment(1); + error!(payload_id = payload_id.get(), %error, "Failed to begin Foldspace stateless recovery."); + return PreparedRecovery::Failed; + } + }; + for stream in self.streams.values_mut() { + stream.remove_payload(payload_id); + } + let Some(retained) = self.retained.remove(&payload_id.get()) else { + self.telemetry.conversion_errors.increment(1); + error!( + payload_id = payload_id.get(), + "Foldspace recovery lost its request metadata." + ); + return PreparedRecovery::Dropped(recovery); + }; + let metadata = retained.metadata().clone(); + let retry_counters = retained.retry_counters().cloned(); + let transaction = match retained { + RetainedPayload::Stateful(retained) => StatefulLogsDecoder::new() + .decode_recovery(&recovery) + .error_context("Failed to decode retained Foldspace payload.") + .and_then(|logs| retained.into_transaction(logs)), + RetainedPayload::StatelessRetry(retained) => StatefulLogsDecoder::new() + .decode_recovery(&recovery) + .error_context("Failed to decode remaining self-contained Foldspace retry batches.") + .and_then(|logs| retained.into_transaction(logs)), + }; + let transaction = match transaction { + Ok(transaction) => transaction, + Err(error) => { + self.telemetry.conversion_errors.increment(1); + component_telemetry.track_permanently_failed_transaction(&metadata, None, endpoint_domain); + error!(payload_id = payload_id.get(), %error, "Failed to reconstruct Foldspace payload; dropping it."); + return PreparedRecovery::Dropped(recovery); + } + }; + PreparedRecovery::Reconstructed(Box::new(StatelessRecovery { + recovery, + transaction, + retry_counters, + })) + } + + async fn execute(&mut self, effects: Vec) { + let mut effects = VecDeque::from(effects); + while let Some(effect) = effects.pop_front() { + match effect { + MplexEffect::OpenStream { sender_id, stream_id } => { + if let Err(error) = self.open_stream(sender_id, stream_id) { + effects.extend(self.core.handle_stream_error( + sender_id, + stream_id, + StreamError::new(error.to_string()), + )); + } + } + MplexEffect::SendBatch { + sender_id, + payload_id, + batch, + } => match self.encoder.encode(&batch) { + Ok(proto) => { + if let Some(stream) = self.streams.get_mut(&sender_id.get()) { + stream.pending.push_back(OutboundBatch { proto, payload_id }); + self.try_send(sender_id).await; + } + } + Err(error) => { + let _ = self.events_tx.send(StatefulEvent::Failed { + sender_id, + stream_id: batch.stream, + error: MetaString::from(format!("failed to encode Foldspace batch: {error:?}")), + failed_payload: payload_id.map(|payload_id| (payload_id, PayloadDelivery::NotSent)), + }); + } + }, + MplexEffect::CloseStream { sender_id, .. } => { + self.streams.remove(&sender_id.get()); + } + MplexEffect::ScheduleTimer { sender_id, timer } => { + let events_tx = self.events_tx.clone(); + tokio::spawn(async move { + tokio::time::sleep(timer.after).await; + let _ = events_tx.send(StatefulEvent::Timer { + sender_id, + kind: timer.kind, + }); + }); + } + MplexEffect::ReportError { sender_id, error } => match error { + CoreError::StreamFailed(error) => { + warn!(sender_id = sender_id.get(), error = ?error, "Foldspace stream failed."); + } + error => { + error!( + sender_id = sender_id.get(), + ?error, + "Foldspace sender reported a protocol error." + ); + if let Some(stream) = self.streams.get(&sender_id.get()) { + let _ = self.events_tx.send(StatefulEvent::Failed { + sender_id, + stream_id: stream.stream_id, + error: MetaString::from_static("unrecoverable Foldspace protocol error"), + failed_payload: None, + }); + } + } + }, + } + } + } + + fn open_stream(&mut self, sender_id: SenderId, stream_id: StreamId) -> Result<(), GenericError> { + let (outbound, receiver) = mpsc::channel(STATEFUL_CHANNEL_CAPACITY); + let mut request = TonicRequest::new(ReceiverStream::new(receiver)); + let api_key = MetadataValue::try_from(self.endpoint.api_key()) + .error_context("Foldspace API key is not valid gRPC metadata.")?; + request.metadata_mut().insert("dd-api-key", api_key); + request + .metadata_mut() + .insert("dd-content-encoding", MetadataValue::from_static("zstd")); + let state_request_bytes = MetadataValue::try_from(STATE_REQUEST_BYTES.to_string()) + .error_context("Foldspace state request limit is not valid gRPC metadata.")?; + request + .metadata_mut() + .insert("dd-state-request-bytes", state_request_bytes); + let mut client = self.client.clone(); + let events_tx = self.events_tx.clone(); + let open_timeout = self.open_timeout; + tokio::spawn(async move { + match tokio::time::timeout(open_timeout, client.stateful_stream(request)).await { + Ok(Ok(response)) => { + if events_tx + .send(StatefulEvent::Opened { + sender_id, + stream_id, + outbound, + }) + .is_ok() + { + spawn_response_reader(sender_id, stream_id, response.into_inner(), events_tx); + } + } + Ok(Err(error)) => { + let _ = events_tx.send(StatefulEvent::OpenFailed { + sender_id, + stream_id, + error: MetaString::from(format!("failed to open Foldspace stream: {error}")), + }); + } + Err(_) => { + let _ = events_tx.send(StatefulEvent::OpenFailed { + sender_id, + stream_id, + error: MetaString::from_static("timed out opening Foldspace stream"), + }); + } + } + }); + Ok(()) + } + + async fn try_send(&mut self, sender_id: SenderId) { + let Some(stream) = self.streams.get_mut(&sender_id.get()) else { + return; + }; + if stream.inflight.is_some() { + return; + } + let Some(batch) = stream.pending.pop_front() else { + return; + }; + let batch_id = u64::from(batch.proto.batch_id); + let payload_id = batch.payload_id; + match stream.outbound.send(batch.proto).await { + Ok(()) => { + if let Some(counters) = payload_id + .and_then(|payload_id| self.retained.get(&payload_id.get())) + .and_then(RetainedPayload::retry_counters) + { + counters.increment_retries(); + } + stream.inflight = Some(InflightBatch { batch_id, payload_id }); + } + Err(error) => { + let error_message = error.to_string(); + stream.pending.push_front(OutboundBatch { + proto: error.0, + payload_id, + }); + let _ = self.events_tx.send(StatefulEvent::Failed { + sender_id, + stream_id: stream.stream_id, + error: MetaString::from(format!("failed to transmit Foldspace batch: {error_message}")), + failed_payload: None, + }); + } + } + } +} + +fn spawn_response_reader( + sender_id: SenderId, stream_id: StreamId, mut responses: tonic::Streaming, + events_tx: mpsc::UnboundedSender, +) { + tokio::spawn(async move { + loop { + match responses.message().await { + Ok(Some(status)) if status.status == i32::from(batch_status::Status::Ok) => { + if events_tx + .send(StatefulEvent::Ack { + sender_id, + stream_id, + batch_id: u64::from(status.batch_id), + }) + .is_err() + { + return; + } + } + Ok(Some(status)) => { + let _ = events_tx.send(StatefulEvent::Failed { + sender_id, + stream_id, + error: MetaString::from(format!( + "Foldspace intake rejected batch {} with status {}", + status.batch_id, status.status + )), + failed_payload: None, + }); + return; + } + Ok(None) => { + let _ = events_tx.send(StatefulEvent::Failed { + sender_id, + stream_id, + error: MetaString::from_static("Foldspace response stream closed"), + failed_payload: None, + }); + return; + } + Err(error) => { + let _ = events_tx.send(StatefulEvent::Failed { + sender_id, + stream_id, + error: MetaString::from(format!("Foldspace acknowledgement failed: {error}")), + failed_payload: None, + }); + return; + } + } + } + }); +} + +#[cfg(test)] +mod tests { + use std::{path::PathBuf, sync::Arc}; + + use bytes::Bytes; + use http::header::{CONTENT_ENCODING, CONTENT_TYPE}; + use saluki_io::net::util::retry::{DiskUsageRetrieverImpl, PersistedQueueArgs, RetryQueue}; + use saluki_metrics::test::TestRecorder; + + use super::*; + use crate::common::datadog::{ + io::PendingTransaction, + telemetry::{SharedTransactionQueueTelemetry, TransactionQueueTelemetry, TransactionRetryTelemetry}, + }; + + const TEST_DOMAIN: &str = "http://127.0.0.1"; + const TEST_RETRY_BYTES: u64 = 64 * 1024; + + fn transaction(body: Vec, encoding: Option<&'static str>) -> Transaction { + let mut builder = Request::builder() + .method("POST") + .uri(LOGS_INTAKE_PATH) + .header(CONTENT_TYPE, "application/json") + .header("x-test-routing", "route-a"); + if let Some(encoding) = encoding { + builder = builder.header(CONTENT_ENCODING, encoding); + } + Transaction::from_original( + Metadata::from_event_and_data_point_count(1, 0), + builder.body(Bytes::from(body)).unwrap(), + ) + } + + fn parsed_body(transaction: Transaction) -> Vec { + let (_, request) = transaction.into_parts(); + let encoding = content_encoding(request.headers()).unwrap(); + let body = decompress_body(©_body(request.body()), encoding).unwrap(); + serde_json::from_slice(&body).unwrap() + } + + fn retain(transaction: Transaction) -> (RetainedRequest, Vec) { + let (metadata, request) = transaction.into_parts(); + let (templates, records) = parse_request_logs(&request).unwrap(); + let retained = RetainedRequest { + metadata, + request: request.map(|_| TransactionBody::from(Vec::new())), + templates, + }; + (retained, records) + } + + fn pending(max_high_priority: usize, max_retry_bytes: u64) -> PendingTransactions> { + let metrics_builder = MetricsBuilder::default(); + let shared = SharedTransactionQueueTelemetry::from_builder(&metrics_builder); + let telemetry = TransactionQueueTelemetry::from_builder(&metrics_builder, TEST_DOMAIN, shared); + PendingTransactions::new( + max_high_priority, + RetryQueue::new("stateful-logs-test".to_owned(), max_retry_bytes), + telemetry, + MetaString::from_static(TEST_DOMAIN), + 900, + ) + } + + fn sender(endpoint: &str, api_key: &str) -> StatefulLogsSender { + let endpoint = ResolvedEndpoint::from_raw_endpoint(endpoint, api_key).unwrap(); + StatefulLogsSender::new( + endpoint, + &MetricsBuilder::default(), + TEST_DOMAIN, + Duration::from_secs(1), + ) + .unwrap() + } + + fn retry_counters(metrics_builder: &MetricsBuilder) -> TransactionRetryCounters { + TransactionRetryTelemetry::from_builder(metrics_builder, TEST_DOMAIN).counters_for(LOGS_INTAKE_PATH) + } + + fn open_core_streams(sender: &mut StatefulLogsSender) -> Vec> { + let mut receivers = Vec::new(); + for effect in sender.core.start() { + let MplexEffect::OpenStream { sender_id, stream_id } = effect else { + panic!("starting Foldspace should only open streams"); + }; + let (outbound, receiver) = mpsc::channel(STATEFUL_CHANNEL_CAPACITY); + sender.streams.insert( + sender_id.get(), + SenderStream { + stream_id, + outbound, + inflight: None, + pending: VecDeque::new(), + }, + ); + assert!(sender.core.handle_stream_opened(sender_id, stream_id).is_empty()); + receivers.push(receiver); + } + receivers + } + + fn simple_log(message: &str, uuid: &str) -> Transaction { + transaction( + serde_json::to_vec(&serde_json::json!([{ + "message": message, + "@timestamp": "2026-08-03T12:00:00.000Z", + "dual-send-uuid": uuid, + "custom": "preserved" + }])) + .unwrap(), + None, + ) + } + + fn multiple_logs(messages: &[&str]) -> Transaction { + let logs = messages + .iter() + .enumerate() + .map(|(index, message)| { + serde_json::json!({ + "message": message, + "@timestamp": "2026-08-03T12:00:00.000Z", + "dual-send-uuid": format!("multi-{index}"), + "custom": "preserved" + }) + }) + .collect::>(); + let request = Request::builder() + .method("POST") + .uri(LOGS_INTAKE_PATH) + .header(CONTENT_TYPE, "application/json") + .header("x-test-routing", "route-a") + .body(Bytes::from(serde_json::to_vec(&logs).unwrap())) + .unwrap(); + Transaction::from_original(Metadata::from_event_and_data_point_count(logs.len(), 0), request) + } + + #[test] + fn recovered_transaction_preserves_fields_uuid_and_compression() { + let input = serde_json::json!([{ + "message": "user alice logged in from 10.0.0.2", + "status": "info", + "hostname": "host-a", + "service": "auth", + "ddsource": "rust", + "ddtags": "env:test,team:logs", + "timestamp": 1700000000000_i64, + "custom": {"nested": true}, + "dual-send-uuid": "uuid-p2" + }]); + let compressed = compress_body(&serde_json::to_vec(&input).unwrap(), BodyEncoding::Zstd).unwrap(); + let (retained, records) = retain(transaction(compressed, Some("zstd"))); + assert_eq!(records[0].uuid.as_deref(), Some("uuid-p2")); + + let decoded = foldspace_server::DecodedLog { + message: "user alice logged in from 10.0.0.2".to_owned(), + status: Some("info".to_owned()), + hostname: Some("host-a".to_owned()), + service: Some("auth".to_owned()), + ddsource: Some("rust".to_owned()), + tags: Some("env:test,team:logs".to_owned()), + timestamp_millis: 1_700_000_000_000, + uuid: Some("uuid-p2".to_owned()), + original_json: None, + }; + let output = retained.into_transaction(vec![decoded]).unwrap(); + assert_eq!(output.metadata().event_count, 1); + assert_eq!(output.request_uri().path(), LOGS_INTAKE_PATH); + let (_, output_request) = output.clone().into_parts(); + assert_eq!(output_request.method(), "POST"); + assert_eq!(output_request.headers()["x-test-routing"], "route-a"); + assert_eq!(output_request.headers()[CONTENT_ENCODING], "zstd"); + let values = parsed_body(output); + assert_eq!(values, input.as_array().unwrap().clone()); + } + + #[test] + fn missing_uuid_is_added_to_both_stateful_and_recovered_forms() { + let input = serde_json::to_vec(&serde_json::json!([{ + "message": "hello 42", + "@timestamp": "2026-08-03T12:00:00.000Z" + }])) + .unwrap(); + let (retained, records) = retain(transaction(input, None)); + let uuid = records[0].uuid.clone().unwrap(); + let output = retained + .into_transaction(vec![foldspace_server::DecodedLog { + message: "hello 42".to_owned(), + timestamp_millis: 1_775_390_400_000, + uuid: Some(uuid.clone()), + ..foldspace_server::DecodedLog::default() + }]) + .unwrap(); + assert_eq!(parsed_body(output)[0][DUAL_SEND_UUID_FIELD], uuid); + } + + #[tokio::test] + async fn high_priority_overflow_converts_before_retry_and_keeps_healthy_streams() { + let mut sender = sender("http://127.0.0.1:4317", "key-a"); + let _receivers = open_core_streams(&mut sender); + let mut pending = pending(0, TEST_RETRY_BYTES); + let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default()); + + assert!(sender + .try_send_transaction( + simple_log("hello 42", "overflow-uuid"), + &mut pending, + &telemetry, + TEST_DOMAIN + ) + .await + .is_ok()); + + assert!(sender.retained.is_empty()); + assert_eq!(sender.streams.len(), STATEFUL_SENDERS); + assert!(sender.streams.values().all(|stream| stream.inflight.is_none())); + let Some(PendingTransaction::LowPriority(retry)) = pending.pop().await else { + panic!("overflow should enqueue one stateless retry"); + }; + assert_eq!(parsed_body(retry)[0][DUAL_SEND_UUID_FIELD], "overflow-uuid"); + + let stream_count = sender.streams.len(); + let push_result = pending + .push_low_priority(simple_log("ordinary retry", "ordinary-uuid")) + .await + .unwrap(); + assert!(!push_result.had_drops()); + assert_eq!(sender.streams.len(), stream_count); + } + + #[tokio::test] + async fn converted_retry_uses_existing_stream_as_a_self_contained_batch() { + let recorder = TestRecorder::default(); + let _recorder_guard = metrics::set_default_local_recorder(&recorder); + let metrics_builder = MetricsBuilder::default(); + let mut sender = sender("http://127.0.0.1:4317", "key-a"); + let mut receivers = open_core_streams(&mut sender); + let original_streams = sender + .streams + .iter() + .map(|(sender_id, stream)| (*sender_id, stream.stream_id)) + .collect::>(); + let mut pending = pending(0, TEST_RETRY_BYTES); + let telemetry = ComponentTelemetry::from_builder(&metrics_builder); + + assert!(sender + .try_send_transaction( + simple_log("retry on grpc 42", "grpc-retry-uuid"), + &mut pending, + &telemetry, + TEST_DOMAIN, + ) + .await + .is_ok()); + assert!(sender.retained.is_empty()); + let retry = pending.pop_low_priority().await.unwrap(); + + assert!(sender + .try_send_retry_transaction(retry, retry_counters(&metrics_builder), &telemetry, TEST_DOMAIN) + .await + .is_ok()); + + assert!(!sender.retained.contains_key(&1)); + assert!(sender.retained.contains_key(&2)); + assert_eq!(sender.streams.len(), original_streams.len()); + assert!(sender + .streams + .iter() + .all(|(sender_id, stream)| original_streams.get(sender_id) == Some(&stream.stream_id))); + assert!(sender + .streams + .values() + .any(|stream| stream.inflight.as_ref().and_then(|batch| batch.payload_id) == Some(PayloadId(2)))); + let proto = receivers + .iter_mut() + .find_map(|receiver| receiver.try_recv().ok()) + .expect("the stateless retry should be written to an existing stream"); + let mut decoder = StatefulLogsDecoder::new(); + let logs = decoder.decode_batch(&proto, FoldspaceContentEncoding::Zstd).unwrap(); + assert_eq!(logs.len(), 1); + assert_eq!(logs[0].message, "retry on grpc 42"); + assert_eq!(logs[0].uuid.as_deref(), Some("grpc-retry-uuid")); + let original: JsonValue = serde_json::from_slice(logs[0].original_json.as_deref().unwrap()).unwrap(); + assert_eq!(original["custom"], "preserved"); + assert_eq!(decoder.state_bytes(), 0); + assert_eq!(decoder.stats().state_changes, 0); + let tags = &[("domain", TEST_DOMAIN), ("endpoint", LOGS_INTAKE_PATH)]; + assert_eq!(recorder.counter(("network_http_requests_retries_total", tags)), Some(1)); + assert_eq!( + recorder.counter(("network_http_requests_requeued_total", tags)), + Some(0) + ); + } + + #[tokio::test] + async fn stateless_retry_batches_are_sent_sequentially_on_one_existing_stream() { + let recorder = TestRecorder::default(); + let _recorder_guard = metrics::set_default_local_recorder(&recorder); + let metrics_builder = MetricsBuilder::default(); + let telemetry = ComponentTelemetry::from_builder(&metrics_builder); + let mut sender = sender("http://127.0.0.1:4317", "key-a"); + sender.max_stateless_batch_bytes = 700; + let mut receivers = open_core_streams(&mut sender); + let original_streams = sender + .streams + .iter() + .map(|(sender_id, stream)| (*sender_id, stream.stream_id)) + .collect::>(); + let first_message = "a".repeat(200); + let second_message = "b".repeat(200); + + assert!(sender + .try_send_retry_transaction( + multiple_logs(&[&first_message, &second_message]), + retry_counters(&metrics_builder), + &telemetry, + TEST_DOMAIN, + ) + .await + .is_ok()); + + let (sender_id, stream_id) = sender + .streams + .iter() + .find_map(|(sender_id, stream)| { + stream + .inflight + .is_some() + .then_some((SenderId(*sender_id), stream.stream_id)) + }) + .expect("the first retry batch should be inflight"); + let receiver = &mut receivers[sender_id.get() as usize]; + let first_batch = receiver.recv().await.unwrap(); + let first_bytes = zstd::stream::decode_all(first_batch.data.as_slice()).unwrap(); + assert!(first_bytes.len() <= sender.max_stateless_batch_bytes); + let mut decoder = StatefulLogsDecoder::new(); + let first_logs = decoder + .decode_batch(&first_batch, FoldspaceContentEncoding::Zstd) + .unwrap(); + assert_eq!(first_logs.len(), 1); + assert_eq!(first_logs[0].message, first_message); + + let mut pending = pending(16, TEST_RETRY_BYTES); + sender + .handle_event( + StatefulEvent::Ack { + sender_id, + stream_id, + batch_id: u64::from(first_batch.batch_id), + }, + &mut pending, + &telemetry, + TEST_DOMAIN, + ) + .await; + assert_eq!(sender.streams.len(), original_streams.len()); + assert!(sender + .streams + .iter() + .all(|(id, stream)| original_streams.get(id) == Some(&stream.stream_id))); + assert_eq!(sender.retained.len(), 1); + + let second_batch = receiver.recv().await.unwrap(); + let second_bytes = zstd::stream::decode_all(second_batch.data.as_slice()).unwrap(); + assert!(second_bytes.len() <= sender.max_stateless_batch_bytes); + let second_logs = decoder + .decode_batch(&second_batch, FoldspaceContentEncoding::Zstd) + .unwrap(); + assert_eq!(second_logs.len(), 1); + assert_eq!(second_logs[0].message, second_message); + sender + .handle_event( + StatefulEvent::Ack { + sender_id, + stream_id, + batch_id: u64::from(second_batch.batch_id), + }, + &mut pending, + &telemetry, + TEST_DOMAIN, + ) + .await; + + assert!(sender.retained.is_empty()); + assert!(pending.is_empty()); + let tags = &[("domain", TEST_DOMAIN), ("endpoint", LOGS_INTAKE_PATH)]; + assert_eq!(recorder.counter(("network_http_requests_retries_total", tags)), Some(2)); + } + + #[tokio::test] + async fn split_retry_failure_requeues_only_unacknowledged_logs() { + let metrics_builder = MetricsBuilder::default(); + let telemetry = ComponentTelemetry::from_builder(&metrics_builder); + let mut sender = sender("http://127.0.0.1:4317", "key-a"); + sender.max_stateless_batch_bytes = 700; + let mut receivers = open_core_streams(&mut sender); + let first_message = "a".repeat(200); + let second_message = "b".repeat(200); + + assert!(sender + .try_send_retry_transaction( + multiple_logs(&[&first_message, &second_message]), + retry_counters(&metrics_builder), + &telemetry, + TEST_DOMAIN, + ) + .await + .is_ok()); + let (sender_id, stream_id) = sender + .streams + .iter() + .find_map(|(sender_id, stream)| { + stream + .inflight + .is_some() + .then_some((SenderId(*sender_id), stream.stream_id)) + }) + .unwrap(); + let receiver = &mut receivers[sender_id.get() as usize]; + let first_batch = receiver.recv().await.unwrap(); + let mut pending = pending(16, TEST_RETRY_BYTES); + sender + .handle_event( + StatefulEvent::Ack { + sender_id, + stream_id, + batch_id: u64::from(first_batch.batch_id), + }, + &mut pending, + &telemetry, + TEST_DOMAIN, + ) + .await; + let _second_batch = receiver.recv().await.unwrap(); + + sender + .handle_event( + StatefulEvent::Failed { + sender_id, + stream_id, + error: MetaString::from_static("split retry transport failure"), + failed_payload: None, + }, + &mut pending, + &telemetry, + TEST_DOMAIN, + ) + .await; + + let retry = pending.pop_low_priority().await.unwrap(); + assert_eq!(retry.metadata().event_count, 1); + let logs = parsed_body(retry); + assert_eq!(logs.len(), 1); + assert_eq!(logs[0]["message"], second_message); + } + + #[tokio::test] + async fn single_log_larger_than_a_stateless_batch_is_dropped() { + let recorder = TestRecorder::default(); + let _recorder_guard = metrics::set_default_local_recorder(&recorder); + let metrics_builder = MetricsBuilder::default(); + let telemetry = ComponentTelemetry::from_builder(&metrics_builder); + let mut sender = sender("http://127.0.0.1:4317", "key-a"); + sender.max_stateless_batch_bytes = 1; + let receivers = open_core_streams(&mut sender); + + assert!(sender + .try_send_retry_transaction( + simple_log("too large", "too-large-uuid"), + retry_counters(&metrics_builder), + &telemetry, + TEST_DOMAIN, + ) + .await + .is_ok()); + + assert!(sender.retained.is_empty()); + assert!(receivers.iter().all(|receiver| receiver.is_empty())); + assert_eq!( + recorder.counter(("stateful_logs_stateless_oversized_total", &[("domain", TEST_DOMAIN)])), + Some(1) + ); + } + + #[tokio::test] + async fn oversized_log_is_dropped_while_fitting_logs_are_sent() { + let recorder = TestRecorder::default(); + let _recorder_guard = metrics::set_default_local_recorder(&recorder); + let metrics_builder = MetricsBuilder::default(); + let telemetry = ComponentTelemetry::from_builder(&metrics_builder); + let mut sender = sender("http://127.0.0.1:4317", "key-a"); + sender.max_stateless_batch_bytes = 700; + let mut receivers = open_core_streams(&mut sender); + let oversized = "x".repeat(2_000); + + assert!(sender + .try_send_retry_transaction( + multiple_logs(&["fits", &oversized]), + retry_counters(&metrics_builder), + &telemetry, + TEST_DOMAIN, + ) + .await + .is_ok()); + + let sender_id = sender + .streams + .iter() + .find_map(|(sender_id, stream)| stream.inflight.is_some().then_some(SenderId(*sender_id))) + .expect("the fitting log should be inflight"); + let batch = receivers[sender_id.get() as usize].recv().await.unwrap(); + let decompressed = zstd::stream::decode_all(batch.data.as_slice()).unwrap(); + assert!(decompressed.len() <= sender.max_stateless_batch_bytes); + let logs = StatefulLogsDecoder::new() + .decode_batch(&batch, FoldspaceContentEncoding::Zstd) + .unwrap(); + assert_eq!(logs.len(), 1); + assert_eq!(logs[0].message, "fits"); + assert_eq!(sender.retained.values().next().unwrap().metadata().event_count, 1); + assert_eq!( + recorder.counter(("stateful_logs_stateless_oversized_total", &[("domain", TEST_DOMAIN)])), + Some(1) + ); + } + + #[tokio::test] + async fn retry_without_an_open_stream_remains_stateless_without_mutation() { + let metrics_builder = MetricsBuilder::default(); + let mut sender = sender("http://127.0.0.1:4317", "key-a"); + let telemetry = ComponentTelemetry::from_builder(&metrics_builder); + let transaction = simple_log("queued retry 42", "queued-retry-uuid"); + + let transaction = sender + .try_send_retry_transaction(transaction, retry_counters(&metrics_builder), &telemetry, TEST_DOMAIN) + .await + .expect_err("a retry without an open stream should remain queued"); + + assert!(sender.streams.is_empty()); + assert!(sender.retained.is_empty()); + assert_eq!(parsed_body(transaction)[0][DUAL_SEND_UUID_FIELD], "queued-retry-uuid"); + } + + #[tokio::test] + async fn retry_without_immediate_stream_capacity_remains_queued_without_rotating_streams() { + let metrics_builder = MetricsBuilder::default(); + let mut sender = sender("http://127.0.0.1:4317", "key-a"); + let _receivers = open_core_streams(&mut sender); + let original_streams = sender + .streams + .iter() + .map(|(sender_id, stream)| (*sender_id, stream.stream_id)) + .collect::>(); + let mut pending = pending(16, TEST_RETRY_BYTES); + let telemetry = ComponentTelemetry::from_builder(&metrics_builder); + for index in 0..STATEFUL_SENDERS { + assert!(sender + .try_send_transaction( + simple_log(&format!("busy stream {index}"), &format!("busy-{index}")), + &mut pending, + &telemetry, + TEST_DOMAIN, + ) + .await + .is_ok()); + } + + let transaction = sender + .try_send_retry_transaction( + simple_log("busy retry 42", "busy-retry-uuid"), + retry_counters(&metrics_builder), + &telemetry, + TEST_DOMAIN, + ) + .await + .expect_err("a retry should remain queued when every stream is busy"); + + assert_eq!(sender.retained.len(), STATEFUL_SENDERS); + assert_eq!(sender.streams.len(), original_streams.len()); + assert!(sender + .streams + .iter() + .all(|(sender_id, stream)| original_streams.get(sender_id) == Some(&stream.stream_id))); + assert_eq!(parsed_body(transaction)[0][DUAL_SEND_UUID_FIELD], "busy-retry-uuid"); + } + + #[tokio::test] + async fn grpc_retry_failed_before_transmission_requeues_complete_stateless_logs() { + let recorder = TestRecorder::default(); + let _recorder_guard = metrics::set_default_local_recorder(&recorder); + let metrics_builder = MetricsBuilder::default(); + let mut sender = sender("http://127.0.0.1:4317", "key-a"); + let mut receivers = open_core_streams(&mut sender); + drop(receivers.remove(0)); + let mut pending = pending(16, TEST_RETRY_BYTES); + let telemetry = ComponentTelemetry::from_builder(&metrics_builder); + + assert!(sender + .try_send_retry_transaction( + simple_log("grpc failed 42", "failed-grpc-retry-uuid"), + retry_counters(&metrics_builder), + &telemetry, + TEST_DOMAIN, + ) + .await + .is_ok()); + let event = sender.next_event().await.unwrap(); + sender.handle_event(event, &mut pending, &telemetry, TEST_DOMAIN).await; + + let retry = pending.pop_low_priority().await.unwrap(); + let value = &parsed_body(retry)[0]; + assert_eq!(value["message"], "grpc failed 42"); + assert_eq!(value[DUAL_SEND_UUID_FIELD], "failed-grpc-retry-uuid"); + let tags = &[("domain", TEST_DOMAIN), ("endpoint", LOGS_INTAKE_PATH)]; + assert_eq!(recorder.counter(("network_http_requests_retries_total", tags)), Some(0)); + assert_eq!( + recorder.counter(("network_http_requests_requeued_total", tags)), + Some(1) + ); + } + + #[tokio::test] + async fn transmitted_grpc_retry_failure_requeues_complete_stateless_logs() { + let recorder = TestRecorder::default(); + let _recorder_guard = metrics::set_default_local_recorder(&recorder); + let metrics_builder = MetricsBuilder::default(); + let mut sender = sender("http://127.0.0.1:4317", "key-a"); + let _receivers = open_core_streams(&mut sender); + let mut pending = pending(16, TEST_RETRY_BYTES); + let telemetry = ComponentTelemetry::from_builder(&metrics_builder); + + assert!(sender + .try_send_retry_transaction( + simple_log("ambiguous retry 42", "ambiguous-retry-uuid"), + retry_counters(&metrics_builder), + &telemetry, + TEST_DOMAIN, + ) + .await + .is_ok()); + let (sender_id, stream_id) = sender + .streams + .iter() + .find_map(|(sender_id, stream)| { + stream + .inflight + .is_some() + .then_some((SenderId(*sender_id), stream.stream_id)) + }) + .unwrap(); + sender + .handle_event( + StatefulEvent::Failed { + sender_id, + stream_id, + error: MetaString::from_static("ambiguous retry transport failure"), + failed_payload: None, + }, + &mut pending, + &telemetry, + TEST_DOMAIN, + ) + .await; + + let retry = pending.pop_low_priority().await.unwrap(); + let value = &parsed_body(retry)[0]; + assert_eq!(value["message"], "ambiguous retry 42"); + assert_eq!(value[DUAL_SEND_UUID_FIELD], "ambiguous-retry-uuid"); + let tags = &[("domain", TEST_DOMAIN), ("endpoint", LOGS_INTAKE_PATH)]; + assert_eq!(recorder.counter(("network_http_requests_retries_total", tags)), Some(1)); + assert_eq!( + recorder.counter(("network_http_requests_requeued_total", tags)), + Some(0) + ); + } + + #[tokio::test] + async fn ambiguous_stream_failure_converts_exactly_one_payload() { + let mut sender = sender("http://127.0.0.1:4317", "key-a"); + let _receivers = open_core_streams(&mut sender); + let mut pending = pending(16, TEST_RETRY_BYTES); + let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default()); + + assert!(sender + .try_send_transaction( + simple_log("ambiguous 42", "ambiguous-uuid"), + &mut pending, + &telemetry, + TEST_DOMAIN + ) + .await + .is_ok()); + let (sender_id, stream_id) = sender + .streams + .iter() + .find_map(|(sender_id, stream)| { + stream + .inflight + .is_some() + .then_some((SenderId(*sender_id), stream.stream_id)) + }) + .expect("one sender should have an in-flight payload"); + + sender + .handle_event( + StatefulEvent::Failed { + sender_id, + stream_id, + error: MetaString::from_static("test stream failure"), + failed_payload: None, + }, + &mut pending, + &telemetry, + TEST_DOMAIN, + ) + .await; + + assert!(sender.retained.is_empty()); + assert_eq!(sender.streams.len(), STATEFUL_SENDERS - 1); + let Some(PendingTransaction::LowPriority(retry)) = pending.pop().await else { + panic!("ambiguous failure should enqueue one stateless retry"); + }; + assert_eq!(parsed_body(retry)[0][DUAL_SEND_UUID_FIELD], "ambiguous-uuid"); + assert!(pending.pop().await.is_none()); + } + + #[tokio::test] + async fn ack_and_failure_races_resolve_once() { + let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default()); + + let mut ack_first = sender("http://127.0.0.1:4317", "key-a"); + let _ack_first_receivers = open_core_streams(&mut ack_first); + let mut ack_first_pending = pending(16, TEST_RETRY_BYTES); + assert!(ack_first + .try_send_transaction( + simple_log("ack first", "ack-first-uuid"), + &mut ack_first_pending, + &telemetry, + TEST_DOMAIN, + ) + .await + .is_ok()); + let (sender_id, stream_id, batch_id) = ack_first + .streams + .iter() + .find_map(|(sender_id, stream)| { + stream + .inflight + .as_ref() + .map(|batch| (SenderId(*sender_id), stream.stream_id, batch.batch_id)) + }) + .unwrap(); + ack_first + .handle_event( + StatefulEvent::Ack { + sender_id, + stream_id, + batch_id, + }, + &mut ack_first_pending, + &telemetry, + TEST_DOMAIN, + ) + .await; + ack_first + .handle_event( + StatefulEvent::Failed { + sender_id, + stream_id, + error: MetaString::from_static("failure after acknowledgement"), + failed_payload: None, + }, + &mut ack_first_pending, + &telemetry, + TEST_DOMAIN, + ) + .await; + assert!(ack_first_pending.is_empty()); + + let mut failure_first = sender("http://127.0.0.1:4318", "key-b"); + let _failure_first_receivers = open_core_streams(&mut failure_first); + let mut failure_first_pending = pending(16, TEST_RETRY_BYTES); + assert!(failure_first + .try_send_transaction( + simple_log("failure first", "failure-first-uuid"), + &mut failure_first_pending, + &telemetry, + TEST_DOMAIN, + ) + .await + .is_ok()); + let (sender_id, stream_id, batch_id) = failure_first + .streams + .iter() + .find_map(|(sender_id, stream)| { + stream + .inflight + .as_ref() + .map(|batch| (SenderId(*sender_id), stream.stream_id, batch.batch_id)) + }) + .unwrap(); + failure_first + .handle_event( + StatefulEvent::Failed { + sender_id, + stream_id, + error: MetaString::from_static("failure before acknowledgement"), + failed_payload: None, + }, + &mut failure_first_pending, + &telemetry, + TEST_DOMAIN, + ) + .await; + failure_first + .handle_event( + StatefulEvent::Ack { + sender_id, + stream_id, + batch_id, + }, + &mut failure_first_pending, + &telemetry, + TEST_DOMAIN, + ) + .await; + assert!(failure_first_pending.pop().await.is_some()); + assert!(failure_first_pending.pop().await.is_none()); + } + + #[tokio::test] + async fn shutdown_converts_retained_payload_before_flushing() { + let mut sender = sender("http://127.0.0.1:4317", "key-a"); + let mut pending = pending(16, TEST_RETRY_BYTES); + let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default()); + assert!(sender + .try_send_transaction( + simple_log("shutdown 42", "shutdown-uuid"), + &mut pending, + &telemetry, + TEST_DOMAIN + ) + .await + .is_ok()); + assert!(pending.is_empty()); + + sender.shutdown(&mut pending, &telemetry, TEST_DOMAIN).await; + + assert!(sender.retained.is_empty()); + let Some(PendingTransaction::LowPriority(retry)) = pending.pop().await else { + panic!("shutdown should enqueue one stateless retry"); + }; + assert_eq!(parsed_body(retry)[0][DUAL_SEND_UUID_FIELD], "shutdown-uuid"); + } + + #[tokio::test] + async fn oversized_stateless_conversion_is_dropped_explicitly() { + let mut sender = sender("http://127.0.0.1:4317", "key-a"); + let mut pending = pending(0, 1); + let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default()); + + assert!(sender + .try_send_transaction( + simple_log("too large", "oversized-uuid"), + &mut pending, + &telemetry, + TEST_DOMAIN + ) + .await + .is_ok()); + + assert!(sender.retained.is_empty()); + assert!(sender.core.payload_sender(PayloadId(1)).is_none()); + assert!(pending.is_empty()); + } + + #[tokio::test] + async fn endpoint_and_api_key_state_are_independent() { + let mut first = sender("http://127.0.0.1:4317", "key-a"); + let mut second = sender("http://127.0.0.1:4318", "key-b"); + let _first_receivers = open_core_streams(&mut first); + let _second_receivers = open_core_streams(&mut second); + let mut first_pending = pending(16, TEST_RETRY_BYTES); + let mut second_pending = pending(16, TEST_RETRY_BYTES); + let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default()); + assert_ne!(first.endpoint.cached_api_key(), second.endpoint.cached_api_key()); + + assert!(first + .try_send_transaction( + simple_log("first 42", "first-uuid"), + &mut first_pending, + &telemetry, + TEST_DOMAIN + ) + .await + .is_ok()); + assert!(second + .try_send_transaction( + simple_log("second 42", "second-uuid"), + &mut second_pending, + &telemetry, + TEST_DOMAIN + ) + .await + .is_ok()); + assert!(first.retained.contains_key(&1)); + assert!(second.retained.contains_key(&1)); + + let (sender_id, stream_id, batch_id) = first + .streams + .iter() + .find_map(|(sender_id, stream)| { + stream + .inflight + .as_ref() + .map(|batch| (SenderId(*sender_id), stream.stream_id, batch.batch_id)) + }) + .unwrap(); + first + .handle_event( + StatefulEvent::Ack { + sender_id, + stream_id, + batch_id, + }, + &mut first_pending, + &telemetry, + TEST_DOMAIN, + ) + .await; + assert!(first.retained.is_empty()); + assert!(second.retained.contains_key(&1)); + assert!(second_pending.is_empty()); + } + + #[tokio::test] + async fn persisted_recovery_is_a_complete_stateless_transaction() { + let original = simple_log("persisted 42", "persisted-uuid"); + let (retained, records) = retain(original); + let recovered = retained + .into_transaction(vec![foldspace_server::DecodedLog { + message: String::from_utf8(records[0].body.clone()).unwrap(), + timestamp_millis: records[0].timestamp_millis, + uuid: records[0].uuid.clone(), + ..foldspace_server::DecodedLog::default() + }]) + .unwrap(); + let temp_dir = tempfile::tempdir().unwrap(); + let root_path = temp_dir.path().to_path_buf(); + let queue_name = "stateful-disk-recovery"; + let persisted_args = |root_path: PathBuf| PersistedQueueArgs { + root_path: root_path.clone(), + max_on_disk_bytes: TEST_RETRY_BYTES, + storage_max_disk_ratio: 0.8, + disk_usage_retriever: Arc::new(DiskUsageRetrieverImpl::new(root_path)), + max_age_days: 10, + }; + let queue = RetryQueue::new(queue_name.to_owned(), TEST_RETRY_BYTES) + .with_disk_persistence(persisted_args(root_path.clone())) + .await + .unwrap(); + let metrics_builder = MetricsBuilder::default(); + let shared = SharedTransactionQueueTelemetry::from_builder(&metrics_builder); + let telemetry = TransactionQueueTelemetry::from_builder(&metrics_builder, TEST_DOMAIN, shared); + let mut pending = PendingTransactions::new(0, queue, telemetry, MetaString::from_static(TEST_DOMAIN), 900); + let push_result = pending.push_low_priority(recovered).await.unwrap(); + assert!(!push_result.had_drops()); + let flush_result = pending.flush().await.unwrap(); + assert!(!flush_result.had_drops()); + + let mut restarted = RetryQueue::>::new(queue_name.to_owned(), TEST_RETRY_BYTES) + .with_disk_persistence(persisted_args(root_path)) + .await + .unwrap(); + let recovered = restarted.pop().await.unwrap().expect("persisted retry should recover"); + let value = &parsed_body(recovered.clone())[0]; + assert_eq!(value["message"], "persisted 42"); + assert_eq!(value[DUAL_SEND_UUID_FIELD], "persisted-uuid"); + assert_eq!(value["custom"], "preserved"); + + let metrics_builder = MetricsBuilder::default(); + let mut restarted_sender = sender("http://127.0.0.1:4317", "key-a"); + let _receivers = open_core_streams(&mut restarted_sender); + assert!(restarted_sender + .try_send_retry_transaction( + recovered, + retry_counters(&metrics_builder), + &ComponentTelemetry::from_builder(&metrics_builder), + TEST_DOMAIN, + ) + .await + .is_ok()); + assert!(restarted_sender.retained.contains_key(&1)); + } +} diff --git a/lib/saluki-components/src/common/datadog/transaction.rs b/lib/saluki-components/src/common/datadog/transaction.rs index 544d58fd825..2644b9c7d4b 100644 --- a/lib/saluki-components/src/common/datadog/transaction.rs +++ b/lib/saluki-components/src/common/datadog/transaction.rs @@ -248,6 +248,10 @@ where self.request.uri() } + pub(crate) const fn request(&self) -> &Request> { + &self.request + } + /// Consumes the `Transaction` and returns the transaction metadata and original request. pub fn into_parts(self) -> (Metadata, http::Request>) { (self.metadata, self.request) diff --git a/lib/saluki-components/src/forwarders/datadog/mod.rs b/lib/saluki-components/src/forwarders/datadog/mod.rs index fbd357804c6..16c600aaeaf 100644 --- a/lib/saluki-components/src/forwarders/datadog/mod.rs +++ b/lib/saluki-components/src/forwarders/datadog/mod.rs @@ -1,3 +1,4 @@ +use agent_data_plane_config::domains::logs::StatefulEncoding; use agent_data_plane_config::shared::{Endpoints, MetricsEncoding}; use async_trait::async_trait; use http::Uri; @@ -38,6 +39,7 @@ pub struct DatadogForwarderConfiguration { forwarder_config: ForwarderConfiguration, configuration: Option, + stateful_logs: StatefulEncoding, } impl DatadogForwarderConfiguration { @@ -47,9 +49,16 @@ impl DatadogForwarderConfiguration { Ok(Self { forwarder_config, configuration: Some(config.clone()), + stateful_logs: StatefulEncoding::default(), }) } + /// Configures stateful logs transport. + pub fn with_stateful_logs(mut self, stateful_logs: StatefulEncoding) -> Self { + self.stateful_logs = stateful_logs; + self + } + /// Creates a forwarder using authoritative typed metrics-routing configuration. pub fn from_configuration_with_metrics_routing( config: &GenericConfiguration, metrics: &MetricsEncoding, endpoints: &Endpoints, @@ -102,7 +111,8 @@ impl ForwarderBuilder for DatadogForwarderConfiguration { get_dd_endpoint_name, telemetry.clone(), metrics_builder, - )?; + )? + .with_stateful_logs(self.stateful_logs.clone()); Ok(Box::new(Datadog { forwarder })) }