From f2f523d6260591ef44c6c521dab843b87be155af Mon Sep 17 00:00:00 2001 From: Snow Pettersen Date: Fri, 31 Jul 2026 19:06:08 -0600 Subject: [PATCH] logger: improve workflow fidelity for next app launch crash logs --- bd-logger/src/async_log_buffer.rs | 323 ++++++++++++-- bd-logger/src/async_log_buffer_test.rs | 533 +++++++++++++++++++++++- bd-logger/src/builder.rs | 2 +- bd-logger/src/init_buffer.rs | 366 ++++++++++++++++ bd-logger/src/init_buffer_test.rs | 69 +++ bd-logger/src/lib.rs | 2 +- bd-logger/src/logger.rs | 12 + bd-logger/src/logging_state.rs | 49 ++- bd-logger/src/pre_config_buffer.rs | 143 ------- bd-logger/src/pre_config_buffer_test.rs | 57 --- bd-logger/src/test/setup.rs | 12 +- bd-runtime/src/runtime.rs | 17 + 12 files changed, 1308 insertions(+), 277 deletions(-) create mode 100644 bd-logger/src/init_buffer.rs create mode 100644 bd-logger/src/init_buffer_test.rs delete mode 100644 bd-logger/src/pre_config_buffer.rs delete mode 100644 bd-logger/src/pre_config_buffer_test.rs diff --git a/bd-logger/src/async_log_buffer.rs b/bd-logger/src/async_log_buffer.rs index c681a6de3..c086880c7 100644 --- a/bd-logger/src/async_log_buffer.rs +++ b/bd-logger/src/async_log_buffer.rs @@ -10,13 +10,13 @@ mod async_log_buffer_test; use crate::device_id::DeviceIdInterceptor; +use crate::init_buffer::{InitItem, PendingInitBuffer, PendingStateOperation, ReplayReason}; use crate::log_replay::{LogReplay, LogReplayResult}; use crate::logger::{ReportProcessingRequest, with_thread_local_logger_guard}; use crate::logging_state::{ConfigUpdate, LoggingState, UninitializedLoggingContext}; use crate::metadata::MetadataCollector; use crate::network::{NetworkQualityInterceptor, SystemTimeProvider}; use crate::ordered_receiver::{OrderedMessage, OrderedReceiver, SequencedMessage}; -use crate::pre_config_buffer::{PendingStateOperation, PreConfigBuffer, PreConfigItem}; use crate::{Block, battery, internal_report, network}; use anyhow::anyhow; use bd_api::DataUpload; @@ -48,7 +48,7 @@ use bd_proto::protos::client::api::debug_data_request::{ }; use bd_proto::protos::client::api::{DebugDataRequest, debug_data_request}; use bd_proto::protos::logging::payload::LogType; -use bd_runtime::runtime::ConfigLoader; +use bd_runtime::runtime::{ConfigLoader, DurationWatch, init_buffer}; use bd_session::{PreparedSessionCallback, PreparedSessionOperation}; use bd_session_replay::CaptureScreenshotHandler; use bd_shutdown::{ComponentShutdown, ComponentShutdownTrigger, ComponentShutdownTriggerHandle}; @@ -67,7 +67,7 @@ use std::collections::{HashMap, VecDeque}; use std::mem::size_of_val; use std::pin::Pin; use std::sync::Arc; -use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::time::Duration as StdDuration; use time::OffsetDateTime; use time::ext::NumericalDuration; @@ -99,8 +99,19 @@ impl ReportProcessor for () { } } +/// Test event used to help coordinate specific sequences of events in the async log buffer during +/// testing. +#[cfg(test)] +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum InitReplayTestEvent { + Configured, + CrashDelayApplied, + Replayed, +} + #[derive(Debug)] pub enum StateUpdateMessage { + CrashPending, AddLogField(String, DataValue), RemoveLogField(String), SetFeatureFlagExposure(String, Option), @@ -113,23 +124,27 @@ pub enum StateUpdateMessage { bd_completion::Sender, ), FlushState(Option>), + #[cfg(test)] + TestBarrier(tokio::sync::oneshot::Sender<()>), } impl MemorySized for StateUpdateMessage { fn size(&self) -> usize { size_of_val(self) + match self { + Self::CrashPending | Self::SetMemoryPressureLevel { .. } => 0, Self::AddLogField(key, value) => key.size() + value.size(), Self::RemoveLogField(field_name) => field_name.len(), Self::SetFeatureFlagExposure(flag, variant) => { flag.len() + variant.as_ref().map_or(0, String::len) }, - Self::SetMemoryPressureLevel { .. } => 0, Self::SetEntityId(entity_id) => entity_id.as_ref().map_or(0, String::len), Self::PersistPreparedSession(operation, response_tx) => { operation.estimated_size() + size_of_val(response_tx) }, Self::FlushState(sender) => size_of_val(sender), + #[cfg(test)] + Self::TestBarrier(sender) => size_of_val(sender), } } } @@ -244,6 +259,7 @@ pub struct Sender { log_buffer_tx: bd_bounded_buffer::Sender, state_buffer_tx: bd_bounded_buffer::Sender, sequence: Arc, + crash_pending: Arc, } impl Sender { @@ -256,6 +272,7 @@ impl Sender { log_buffer_tx, state_buffer_tx, sequence: Arc::new(AtomicU64::new(0)), + crash_pending: Arc::new(AtomicBool::new(false)), } } @@ -279,6 +296,26 @@ impl Sender { self.state_buffer_tx.try_send(sequenced) } + pub fn hint_crash_report_pending(&self) -> Result<(), TrySendError> { + self.crash_pending.store(true, Ordering::Release); + self.try_send_state_update(StateUpdateMessage::CrashPending) + } + + #[cfg(test)] + pub(crate) fn set_crash_pending_without_enqueuing(&self) { + self.crash_pending.store(true, Ordering::Release); + } + + #[cfg(test)] + pub(crate) async fn wait_until_processed(&self) -> Result<(), TrySendError> { + let (completion_tx, completion_rx) = tokio::sync::oneshot::channel(); + self.try_send_state_update(StateUpdateMessage::TestBarrier(completion_tx))?; + completion_rx + .await + .expect("async log buffer stopped before processing test barrier"); + Ok(()) + } + pub fn flush_state(&self, block: Block) -> Result<(), TrySendError> { let (completion_tx, completion_rx) = if matches!(block, Block::Yes { .. }) { let (tx, rx) = bd_completion::Sender::new(); @@ -361,13 +398,20 @@ pub struct AsyncLogBuffer { replayer: R, interceptors: Vec>, - logging_state: LoggingState, + logging_state: LoggingState, + pending_init_buffer: Option, + init_replay_delay: DurationWatch, + crash_pending_replay_delay: DurationWatch, + crash_pending: Arc, global_state_tracker: global_state::Tracker, global_state_reader: global_state::Reader, time_provider: Arc, lifecycle_state: InitLifecycleState, sdk_status_tracker: bd_client_common::sdk_status::SdkStatusTracker, + #[cfg(test)] + init_replay_test_events: Option>, + pending_workflow_debug_state: HashMap, send_workflow_debug_state_delay: Option>>, last_session_id: Option, @@ -375,7 +419,7 @@ pub struct AsyncLogBuffer { impl AsyncLogBuffer { pub(crate) fn new( - uninitialized_logging_context: UninitializedLoggingContext, + uninitialized_logging_context: UninitializedLoggingContext, replayer: R, session_strategy: Arc, metadata_provider: Arc, @@ -395,11 +439,7 @@ impl AsyncLogBuffer { sdk_status_tracker: bd_client_common::sdk_status::SdkStatusTracker, data_upload_tx: mpsc::Sender, ) -> (Self, Sender) { - let (log_tx, log_rx) = channel( - uninitialized_logging_context - .pre_config_log_buffer - .max_size(), - ); + let (log_tx, log_rx) = channel(uninitialized_logging_context.init_buffer.max_size()); // Larger channel for state updates as they are less frequent and we want // to avoid dropping any state updates if possible. @@ -433,6 +473,8 @@ impl AsyncLogBuffer { Arc::new(NetworkQualityInterceptor::new(network_quality_resolver)); let device_id_interceptor = Arc::new(DeviceIdInterceptor::new(device_id)); + let crash_pending = Arc::new(AtomicBool::new(false)); + ( Self { ordered_rx: OrderedReceiver::new(log_rx, state_rx), @@ -468,6 +510,10 @@ impl AsyncLogBuffer { // The size of the pre-config buffer matches the size of the enclosing // async log buffer. logging_state: LoggingState::Uninitialized(uninitialized_logging_context), + pending_init_buffer: None, + init_replay_delay: runtime_loader.register_duration_watch(), + crash_pending_replay_delay: runtime_loader.register_duration_watch(), + crash_pending: crash_pending.clone(), global_state_tracker: global_state::Tracker::new( store.clone(), runtime_loader.register_duration_watch(), @@ -477,6 +523,9 @@ impl AsyncLogBuffer { lifecycle_state, sdk_status_tracker, + #[cfg(test)] + init_replay_test_events: None, + pending_workflow_debug_state: HashMap::new(), send_workflow_debug_state_delay: None, last_session_id: None, @@ -485,10 +534,27 @@ impl AsyncLogBuffer { log_buffer_tx: log_tx, state_buffer_tx: state_tx, sequence: Arc::new(AtomicU64::new(0)), + crash_pending, }, ) } + #[cfg(test)] + pub(crate) fn with_init_replay_test_events( + mut self, + events: mpsc::UnboundedSender, + ) -> Self { + self.init_replay_test_events = Some(events); + self + } + + #[cfg(test)] + fn send_init_replay_test_event(&self, event: InitReplayTestEvent) { + if let Some(events) = &self.init_replay_test_events { + let _ = events.send(event); + } + } + pub fn enqueue_log( tx: &Sender, log_level: LogLevel, @@ -637,6 +703,11 @@ impl AsyncLogBuffer { block: bool, state_store: &bd_state::Store, ) -> anyhow::Result { + let is_previous_run_log = matches!( + &log.attributes_overrides, + Some(LogAttributesOverrides::PreviousRunSessionID(_)) + ); + // Prevent re-entrancy when we are evaluating the log metadata. let result = with_thread_local_logger_guard(|| { if let Some(LogAttributesOverrides::PreviousRunSessionID(_)) = &log.attributes_overrides { @@ -724,7 +795,9 @@ impl AsyncLogBuffer { capture_session: log.capture_session, }; - self.write_log(processed_log, block, state_store).await + self + .write_log(processed_log, block, state_store, is_previous_run_log) + .await }, Err(e) => { // TODO(Augustyniak): Consider logging as error so that SDK customers can see these @@ -740,17 +813,46 @@ impl AsyncLogBuffer { log: Log, block: bool, state_store: &bd_state::Store, + is_previous_run_log: bool, ) -> anyhow::Result { + if self.pending_init_buffer.is_some() { + let must_replay_early = self + .pending_init_buffer + .as_ref() + .is_some_and(|pending_init_buffer| !pending_init_buffer.buffer().can_push_size(log.size())); + if must_replay_early { + log::info!("replaying init buffer early to avoid dropping a log due to buffer pressure"); + self + .replay_init_buffer(state_store, ReplayReason::CapacityLog, true) + .await; + } else if let Some(pending_init_buffer) = &mut self.pending_init_buffer { + let item = if is_previous_run_log { + InitItem::PreviousRunLog(log) + } else { + InitItem::Log(log) + }; + let result = pending_init_buffer.push(item); + if let Err(e) = result { + anyhow::bail!("failed to push log to an init buffer: {e}"); + } + + return Ok(LogReplayResult::default()); + } + } + let log_replay_result = match &mut self.logging_state { LoggingState::Uninitialized(uninitialized_logging_context) => { - let result = uninitialized_logging_context - .pre_config_log_buffer - .push(PreConfigItem::Log(log)); + let item = if is_previous_run_log { + InitItem::PreviousRunLog(log) + } else { + InitItem::Log(log) + }; + let result = uninitialized_logging_context.init_buffer.push(item); uninitialized_logging_context .stats - .pre_config_log_buffer - .record(&result); + .init_buffer + .record_push(&result); if let Err(e) = result { anyhow::bail!("failed to push log to a pre-config buffer: {e}"); } @@ -792,19 +894,19 @@ impl AsyncLogBuffer { } } - async fn update( - mut self, - config: ConfigUpdate, - ) -> (Self, Option>) { - let (initialized_logging_context, maybe_pre_config_log_buffer) = match self.logging_state { + async fn update(mut self, config: ConfigUpdate) -> (Self, Option) { + let (initialized_logging_context, pending_init_buffer) = match self.logging_state { LoggingState::Uninitialized(uninitialized_logging_context) => { - let (initialized_logging_context, pre_config_log_buffer) = uninitialized_logging_context + let (initialized_logging_context, buffer, stats) = uninitialized_logging_context .updated( config, self.session_replay_capture_screenshot_handler.clone(), ) .await; - (initialized_logging_context, Some(pre_config_log_buffer)) + ( + initialized_logging_context, + Some(PendingInitBuffer::new(buffer, stats)), + ) }, LoggingState::Initialized(mut initialized_logging_context) => { initialized_logging_context.update(config); @@ -814,23 +916,46 @@ impl AsyncLogBuffer { self.logging_state = LoggingState::Initialized(initialized_logging_context); - (self, maybe_pre_config_log_buffer) + (self, pending_init_buffer) } - async fn maybe_replay_pre_config_buffer( + async fn maybe_replay_init_buffer( &mut self, - pre_config_buffer: PreConfigBuffer, state_store: &bd_state::Store, + reason: ReplayReason, ) { if !matches!(self.logging_state, LoggingState::Initialized(_)) { return; } + let Some(pending_init_buffer) = self.pending_init_buffer.take() else { + return; + }; + let now = self.time_provider.now(); - for item in pre_config_buffer.pop_all() { + for item in pending_init_buffer.drain(reason) { match item { - PreConfigItem::Log(log) => { + InitItem::PreviousRunLog(log) => { + let LoggingState::Initialized(initialized_logging_context) = &mut self.logging_state + else { + return; + }; + if let Err(e) = self + .replayer + .replay_log( + log, + false, + &mut initialized_logging_context.processing_pipeline, + state_store, + now, + ) + .await + { + log::debug!("failed to replay prior-run init log: {e}"); + } + }, + InitItem::Log(log) => { self .update_system_session_id(state_store, &log.session_id) .await; @@ -852,7 +977,7 @@ impl AsyncLogBuffer { log::debug!("failed to replay pre-config log: {e}"); } }, - PreConfigItem::StateOperation(operation) => match operation { + InitItem::StateOperation(operation) => match operation { PendingStateOperation::SetFeatureFlagExposure { name, variant, @@ -881,6 +1006,62 @@ impl AsyncLogBuffer { } } + async fn replay_init_buffer( + &mut self, + state_store: &bd_state::Store, + reason: ReplayReason, + force_replay: bool, + ) { + if self.pending_init_buffer.is_none() { + return; + } + + // A platform thread can set the atomic hint while the ordered control message is still + // waiting behind the replay timer. Honor it before draining so that race cannot bypass the + // configured crash-report extension. + if !force_replay + && self.crash_pending.load(Ordering::Acquire) + && self.apply_crash_pending_hint() + { + return; + } + + self.maybe_replay_init_buffer(state_store, reason).await; + #[cfg(test)] + self.send_init_replay_test_event(InitReplayTestEvent::Replayed); + self + .lifecycle_state + .set(InitLifecycle::LogProcessingStarted); + self.sdk_status_tracker.record_running(); + } + + fn schedule_init_replay(&mut self) -> bool { + self + .pending_init_buffer + .as_mut() + .is_some_and(|pending_init_buffer| { + pending_init_buffer.schedule( + *self.init_replay_delay.read(), + self.crash_pending.load(Ordering::Acquire), + *self.crash_pending_replay_delay.read(), + ) + }) + } + + fn apply_crash_pending_hint(&mut self) -> bool { + let applied = self + .pending_init_buffer + .as_mut() + .is_some_and(|pending_init_buffer| { + pending_init_buffer.apply_crash_pending_hint(*self.crash_pending_replay_delay.read()) + }); + #[cfg(test)] + if applied { + self.send_init_replay_test_event(InitReplayTestEvent::CrashDelayApplied); + } + applied + } + pub async fn run( self, state_store: bd_state::Store, @@ -927,21 +1108,25 @@ impl AsyncLogBuffer { tokio::select! { Some(config) = self.config_update_rx.recv() => { - let (updated_self, maybe_pre_config_buffer) + let (updated_self, pending_init_buffer) = self.update(config).await; self = updated_self; - if let Some(pre_config_buffer) = maybe_pre_config_buffer { - self.lifecycle_state.set(InitLifecycle::LogProcessingStarted); - self.sdk_status_tracker.record_running(); - self - .maybe_replay_pre_config_buffer(pre_config_buffer, &state_store) - .await; + if let Some(pending_init_buffer) = pending_init_buffer { + self.pending_init_buffer = Some(pending_init_buffer); + let replay_immediately = self.schedule_init_replay(); + #[cfg(test)] + self.send_init_replay_test_event(InitReplayTestEvent::Configured); + if replay_immediately { + self + .replay_init_buffer(&state_store, ReplayReason::Scheduled, false) + .await; + } } }, - Some(ReportProcessingRequest { - session - }) = self.report_processor_rx.recv() => { + Some(ReportProcessingRequest { + session, + }) = self.report_processor_rx.recv() => { // TODO(snowp): Once we move over to using the file watcher we can more accurately pick // current vs previous for all reports, but as we need to handle restarts etc we may // also want to embed the full information into the report. This should ensure that we @@ -1004,6 +1189,9 @@ impl AsyncLogBuffer { }, OrderedMessage::State(async_log_buffer_message) => { match async_log_buffer_message { + StateUpdateMessage::CrashPending => { + self.apply_crash_pending_hint(); + }, StateUpdateMessage::AddLogField(key, value) => { if let Err(e) = self .metadata_collector @@ -1026,7 +1214,41 @@ impl AsyncLogBuffer { continue; }, }; - if let LoggingState::Initialized(initialized_logging_context) = + let state_operation_size = + flag.len() + variant.as_ref().map_or(0, String::len) + session_id.len(); + let must_replay_early = self.pending_init_buffer.as_ref().is_some_and( + |pending_init_buffer| { + !pending_init_buffer.buffer().can_push_size(state_operation_size) + }, + ); + if must_replay_early { + log::info!( + "replaying init buffer early to avoid dropping a state operation due to \ + buffer pressure" + ); + self + .replay_init_buffer( + &state_store, + ReplayReason::CapacityStateOperation, + true, + ) + .await; + } + + if let Some(pending_init_buffer) = &mut self.pending_init_buffer { + let result = pending_init_buffer.push( + InitItem::StateOperation( + PendingStateOperation::SetFeatureFlagExposure { + name: flag, + variant, + session_id, + }, + ), + ); + if let Err(e) = result { + log::debug!("failed to enqueue state operation to init buffer: {e}"); + } + } else if let LoggingState::Initialized(initialized_logging_context) = &mut self.logging_state { // Initialized: update state store and replay through workflows @@ -1048,8 +1270,8 @@ impl AsyncLogBuffer { if let LoggingState::Uninitialized(uninitialized_logging_context) = &mut self.logging_state { - let result = uninitialized_logging_context.pre_config_log_buffer.push( - PreConfigItem::StateOperation( + let result = uninitialized_logging_context.init_buffer.push( + InitItem::StateOperation( PendingStateOperation::SetFeatureFlagExposure { name: flag, variant, @@ -1059,8 +1281,8 @@ impl AsyncLogBuffer { ); uninitialized_logging_context .stats - .pre_config_log_buffer - .record(&result); + .init_buffer + .record_push(&result); if let Err(e) = result { log::debug!("failed to enqueue state operation to pre-config buffer: {e}"); } @@ -1153,6 +1375,10 @@ impl AsyncLogBuffer { completion_tx.send(()); } }, + #[cfg(test)] + StateUpdateMessage::TestBarrier(completion_tx) => { + let _ = completion_tx.send(()); + }, } }, } @@ -1166,6 +1392,13 @@ impl AsyncLogBuffer { () = maybe_await(&mut self.send_workflow_debug_state_delay) => { self.send_debug_data().await; }, + () = maybe_await_map(self.pending_init_buffer.as_mut(), |pending_init_buffer| async { + maybe_await(pending_init_buffer.replay_sleep()).await; + }) => { + self + .replay_init_buffer(&state_store, ReplayReason::Scheduled, false) + .await; + }, () = self.resource_utilization_reporter.run() => {}, () = self.session_replay_recorder.run() => {}, () = self.events_listener.run() => {}, diff --git a/bd-logger/src/async_log_buffer_test.rs b/bd-logger/src/async_log_buffer_test.rs index 196f4f61e..0e3bb9928 100644 --- a/bd-logger/src/async_log_buffer_test.rs +++ b/bd-logger/src/async_log_buffer_test.rs @@ -5,19 +5,20 @@ // LICENSE file or at: // https://polyformproject.org/wp-content/uploads/2020/06/PolyForm-Shield-1.0.0.txt -use crate::Block; use crate::async_log_buffer::{ AsyncLogBuffer, + InitReplayTestEvent, LogLine, LogReplay, - PreConfigItem, Sender, StateUpdateMessage, }; use crate::buffer_selector::BufferSelector; use crate::client_config::TailConfigurations; +use crate::init_buffer::InitItem; use crate::log_replay::{LogReplayResult, LoggerReplay, ProcessingPipeline}; use crate::logging_state::{BufferProducers, ConfigUpdate, UninitializedLoggingContext}; +use crate::{Block, LogAttributesOverrides}; use bd_api::{DataUpload, SimpleNetworkQualityProvider}; use bd_client_common::init_lifecycle::InitLifecycleState; use bd_client_stats::{FlushTrigger, Stats}; @@ -80,6 +81,7 @@ struct Setup { shutdown: Option, store: Arc, session_strategy: Arc, + init_buffer_capacity: usize, } impl Setup { @@ -113,6 +115,7 @@ impl Setup { data_upload_tx, store: in_memory_store(), session_strategy, + init_buffer_capacity: 1_000_000, } } @@ -160,6 +163,23 @@ impl Setup { ) } + fn make_test_async_log_buffer_with_init_replay_events( + &mut self, + config_update_rx: tokio::sync::mpsc::Receiver, + ) -> ( + AsyncLogBuffer, + Sender, + mpsc::UnboundedReceiver, + ) { + let (events_tx, events_rx) = mpsc::unbounded_channel(); + let (buffer, sender) = self.make_test_async_log_buffer(config_update_rx); + ( + buffer.with_init_replay_test_events(events_tx), + sender, + events_rx, + ) + } + fn make_real_async_log_buffer( &self, config_update_rx: tokio::sync::mpsc::Receiver, @@ -189,7 +209,7 @@ impl Setup { ) } - fn make_logging_context(&self) -> UninitializedLoggingContext { + fn make_logging_context(&self) -> UninitializedLoggingContext { let (trigger_upload_tx, _) = tokio::sync::mpsc::channel(1); let (_remote_flush_streaming_tx, remote_flush_streaming_rx) = tokio::sync::mpsc::channel(1); let (data_upload_tx, _) = tokio::sync::mpsc::channel(1); @@ -206,7 +226,7 @@ impl Setup { data_upload_tx, flush_buffers_tx, flush_stats_trigger, - 1_000_000, + self.init_buffer_capacity, Arc::new(AtomicBool::new(false)), Arc::new(ProcessLocalPendingFlushState::default()), ) @@ -988,6 +1008,511 @@ async fn pre_config_logs_trigger_session_id_update() { task.join().unwrap(); } +async fn expect_init_replay_event( + events: &mut mpsc::UnboundedReceiver, + expected: InitReplayTestEvent, +) { + assert_eq!(Some(expected), events.recv().await); +} + +#[tokio::test(start_paused = true)] +async fn init_buffer_replays_after_the_configured_delay() { + let mut setup = Setup::new(); + setup + .runtime + .update_snapshot(bd_test_helpers::runtime::make_simple_update(vec![( + bd_runtime::runtime::init_buffer::ReplayDelay::path(), + ValueKind::Int(100), + )])) + .await + .unwrap(); + + let (config_update_tx, config_update_rx) = tokio::sync::mpsc::channel(1); + let (buffer, sender, mut events) = + setup.make_test_async_log_buffer_with_init_replay_events(config_update_rx); + let test_store = TestStore::new().await; + let shutdown_trigger = ComponentShutdownTrigger::default(); + let handle = tokio::task::spawn(buffer.run_with_shutdown( + (*test_store).clone(), + (), + shutdown_trigger.make_shutdown(), + )); + config_update_tx + .send(setup.make_config_update(WorkflowsConfiguration::default())) + .await + .unwrap(); + expect_init_replay_event(&mut events, InitReplayTestEvent::Configured).await; + + assert_ok!(AsyncLogBuffer::::enqueue_log( + &sender, + 0, + LogType::NORMAL, + "delayed".into(), + [].into(), + [].into(), + None, + Block::No, + None, + )); + sender.wait_until_processed().await.unwrap(); + assert_eq!(0, setup.replayer_log_count.load(Ordering::SeqCst)); + + tokio::time::advance(99.std_milliseconds()).await; + assert_eq!(0, setup.replayer_log_count.load(Ordering::SeqCst)); + + tokio::time::advance(1.std_milliseconds()).await; + expect_init_replay_event(&mut events, InitReplayTestEvent::Replayed).await; + assert_eq!(1, setup.replayer_log_count.load(Ordering::SeqCst)); + + shutdown_trigger.shutdown().await; + handle.await.unwrap(); +} + +#[tokio::test(start_paused = true)] +async fn init_buffer_replays_when_full_before_processing_another_log() { + let mut setup = Setup::new(); + setup.init_buffer_capacity = 800; + setup + .runtime + .update_snapshot(bd_test_helpers::runtime::make_simple_update(vec![( + bd_runtime::runtime::init_buffer::ReplayDelay::path(), + ValueKind::Int(100), + )])) + .await + .unwrap(); + + let (config_update_tx, config_update_rx) = tokio::sync::mpsc::channel(1); + let (buffer, sender, mut events) = + setup.make_test_async_log_buffer_with_init_replay_events(config_update_rx); + let test_store = TestStore::new().await; + let shutdown_trigger = ComponentShutdownTrigger::default(); + let handle = tokio::task::spawn(buffer.run_with_shutdown( + (*test_store).clone(), + (), + shutdown_trigger.make_shutdown(), + )); + + config_update_tx + .send(setup.make_config_update(WorkflowsConfiguration::default())) + .await + .unwrap(); + expect_init_replay_event(&mut events, InitReplayTestEvent::Configured).await; + + assert_ok!(AsyncLogBuffer::::enqueue_log( + &sender, + 0, + LogType::NORMAL, + "early".into(), + [].into(), + [].into(), + None, + Block::No, + None, + )); + sender.wait_until_processed().await.unwrap(); + assert_eq!(0, setup.replayer_log_count.load(Ordering::SeqCst)); + + assert_ok!(AsyncLogBuffer::::enqueue_log( + &sender, + 0, + LogType::NORMAL, + "within_allowance".into(), + [].into(), + [].into(), + None, + Block::No, + None, + )); + sender.wait_until_processed().await.unwrap(); + assert_eq!(0, setup.replayer_log_count.load(Ordering::SeqCst)); + + assert_ok!(AsyncLogBuffer::::enqueue_log( + &sender, + 0, + LogType::NORMAL, + "replay_early".into(), + [].into(), + [].into(), + None, + Block::No, + None, + )); + expect_init_replay_event(&mut events, InitReplayTestEvent::Replayed).await; + sender.wait_until_processed().await.unwrap(); + assert_eq!( + vec![ + "early".to_string(), + "within_allowance".to_string(), + "replay_early".to_string(), + ], + *setup.replayer_logs.lock() + ); + + shutdown_trigger.shutdown().await; + handle.await.unwrap(); +} + +#[tokio::test(start_paused = true)] +async fn crash_pending_retains_current_init_logs_with_a_soft_limit_allowance() { + let mut setup = Setup::new(); + setup.init_buffer_capacity = 2_000; + setup + .runtime + .update_snapshot(bd_test_helpers::runtime::make_simple_update(vec![ + ( + bd_runtime::runtime::init_buffer::ReplayDelay::path(), + ValueKind::Int(100), + ), + ( + bd_runtime::runtime::init_buffer::CrashPendingReplayDelay::path(), + ValueKind::Int(50), + ), + ])) + .await + .unwrap(); + + let (config_update_tx, config_update_rx) = tokio::sync::mpsc::channel(1); + let (buffer, sender, mut events) = + setup.make_test_async_log_buffer_with_init_replay_events(config_update_rx); + let test_store = TestStore::new().await; + let shutdown_trigger = ComponentShutdownTrigger::default(); + let handle = tokio::task::spawn(buffer.run_with_shutdown( + (*test_store).clone(), + (), + shutdown_trigger.make_shutdown(), + )); + + config_update_tx + .send(setup.make_config_update(WorkflowsConfiguration::default())) + .await + .unwrap(); + expect_init_replay_event(&mut events, InitReplayTestEvent::Configured).await; + + let current_message = "x".repeat(1_000); + assert_ok!(AsyncLogBuffer::::enqueue_log( + &sender, + 0, + LogType::NORMAL, + current_message.clone().into(), + [].into(), + [].into(), + None, + Block::No, + None, + )); + sender.wait_until_processed().await.unwrap(); + sender.hint_crash_report_pending().unwrap(); + sender.wait_until_processed().await.unwrap(); + expect_init_replay_event(&mut events, InitReplayTestEvent::CrashDelayApplied).await; + + sender + .try_send_log( + LogLine { + log_level: 0, + log_type: LogType::LIFECYCLE, + message: "previous".into(), + fields: [].into(), + matching_fields: [].into(), + attributes_overrides: Some(LogAttributesOverrides::PreviousRunSessionID( + OffsetDateTime::now_utc(), + )), + capture_session: None, + } + .into(), + ) + .unwrap(); + sender.wait_until_processed().await.unwrap(); + + tokio::time::advance(150.std_milliseconds()).await; + expect_init_replay_event(&mut events, InitReplayTestEvent::Replayed).await; + assert_eq!( + vec!["previous".to_string(), current_message], + *setup.replayer_logs.lock() + ); + + shutdown_trigger.shutdown().await; + handle.await.unwrap(); +} + +#[tokio::test(start_paused = true)] +async fn init_buffer_replays_when_full_before_processing_another_state_operation() { + let mut setup = Setup::new(); + setup.init_buffer_capacity = 800; + setup + .runtime + .update_snapshot(bd_test_helpers::runtime::make_simple_update(vec![ + ( + bd_runtime::runtime::init_buffer::ReplayDelay::path(), + ValueKind::Int(100), + ), + ( + bd_runtime::runtime::init_buffer::CrashPendingReplayDelay::path(), + ValueKind::Int(50), + ), + ])) + .await + .unwrap(); + + let (config_update_tx, config_update_rx) = tokio::sync::mpsc::channel(1); + let (buffer, sender, mut events) = + setup.make_test_async_log_buffer_with_init_replay_events(config_update_rx); + let test_store = TestStore::new().await; + let shutdown_trigger = ComponentShutdownTrigger::default(); + let handle = tokio::task::spawn(buffer.run_with_shutdown( + (*test_store).clone(), + (), + shutdown_trigger.make_shutdown(), + )); + + config_update_tx + .send(setup.make_config_update(WorkflowsConfiguration::default())) + .await + .unwrap(); + expect_init_replay_event(&mut events, InitReplayTestEvent::Configured).await; + + assert_ok!(AsyncLogBuffer::::enqueue_log( + &sender, + 0, + LogType::NORMAL, + "buffered log".into(), + [].into(), + [].into(), + None, + Block::No, + None, + )); + sender.wait_until_processed().await.unwrap(); + assert_eq!(0, setup.replayer_log_count.load(Ordering::SeqCst)); + + sender.hint_crash_report_pending().unwrap(); + sender.wait_until_processed().await.unwrap(); + expect_init_replay_event(&mut events, InitReplayTestEvent::CrashDelayApplied).await; + let flag = "x".repeat(1_000); + sender + .try_send_state_update(StateUpdateMessage::SetFeatureFlagExposure( + flag.clone(), + Some("enabled".to_string()), + )) + .unwrap(); + sender.wait_until_processed().await.unwrap(); + assert_eq!(0, setup.replayer_log_count.load(Ordering::SeqCst)); + + let next_flag = "next flag".to_string(); + sender + .try_send_state_update(StateUpdateMessage::SetFeatureFlagExposure( + next_flag.clone(), + Some("enabled".to_string()), + )) + .unwrap(); + expect_init_replay_event(&mut events, InitReplayTestEvent::Replayed).await; + sender.wait_until_processed().await.unwrap(); + + assert_eq!( + vec!["buffered log".to_string()], + *setup.replayer_logs.lock() + ); + let reader = test_store.read().await; + let value = reader.get(Scope::FeatureFlagExposure, &flag); + assert!(value.is_some_and(|value| value.string_value() == "enabled")); + let value = reader.get(Scope::FeatureFlagExposure, &next_flag); + assert!(value.is_some_and(|value| value.string_value() == "enabled")); + + shutdown_trigger.shutdown().await; + handle.await.unwrap(); +} + +#[tokio::test(start_paused = true)] +async fn crash_pending_hint_extends_init_replay_once() { + let mut setup = Setup::new(); + setup + .runtime + .update_snapshot(bd_test_helpers::runtime::make_simple_update(vec![ + ( + bd_runtime::runtime::init_buffer::ReplayDelay::path(), + ValueKind::Int(100), + ), + ( + bd_runtime::runtime::init_buffer::CrashPendingReplayDelay::path(), + ValueKind::Int(50), + ), + ])) + .await + .unwrap(); + + let (config_update_tx, config_update_rx) = tokio::sync::mpsc::channel(1); + let (buffer, sender, mut events) = + setup.make_test_async_log_buffer_with_init_replay_events(config_update_rx); + let test_store = TestStore::new().await; + let shutdown_trigger = ComponentShutdownTrigger::default(); + let handle = tokio::task::spawn(buffer.run_with_shutdown( + (*test_store).clone(), + (), + shutdown_trigger.make_shutdown(), + )); + sender.hint_crash_report_pending().unwrap(); + sender.hint_crash_report_pending().unwrap(); + config_update_tx + .send(setup.make_config_update(WorkflowsConfiguration::default())) + .await + .unwrap(); + expect_init_replay_event(&mut events, InitReplayTestEvent::Configured).await; + + assert_ok!(AsyncLogBuffer::::enqueue_log( + &sender, + 0, + LogType::NORMAL, + "delayed".into(), + [].into(), + [].into(), + None, + Block::No, + None, + )); + sender.wait_until_processed().await.unwrap(); + + tokio::time::advance(149.std_milliseconds()).await; + assert_eq!(0, setup.replayer_log_count.load(Ordering::SeqCst)); + + tokio::time::advance(1.std_milliseconds()).await; + expect_init_replay_event(&mut events, InitReplayTestEvent::Replayed).await; + assert_eq!(1, setup.replayer_log_count.load(Ordering::SeqCst)); + + shutdown_trigger.shutdown().await; + handle.await.unwrap(); +} + +#[tokio::test(start_paused = true)] +async fn late_crash_pending_hint_extends_scheduled_init_replay() { + let mut setup = Setup::new(); + setup + .runtime + .update_snapshot(bd_test_helpers::runtime::make_simple_update(vec![ + ( + bd_runtime::runtime::init_buffer::ReplayDelay::path(), + ValueKind::Int(100), + ), + ( + bd_runtime::runtime::init_buffer::CrashPendingReplayDelay::path(), + ValueKind::Int(50), + ), + ])) + .await + .unwrap(); + + let (config_update_tx, config_update_rx) = tokio::sync::mpsc::channel(1); + let (buffer, sender, mut events) = + setup.make_test_async_log_buffer_with_init_replay_events(config_update_rx); + let test_store = TestStore::new().await; + let shutdown_trigger = ComponentShutdownTrigger::default(); + let handle = tokio::task::spawn(buffer.run_with_shutdown( + (*test_store).clone(), + (), + shutdown_trigger.make_shutdown(), + )); + config_update_tx + .send(setup.make_config_update(WorkflowsConfiguration::default())) + .await + .unwrap(); + expect_init_replay_event(&mut events, InitReplayTestEvent::Configured).await; + + assert_ok!(AsyncLogBuffer::::enqueue_log( + &sender, + 0, + LogType::NORMAL, + "delayed".into(), + [].into(), + [].into(), + None, + Block::No, + None, + )); + sender.wait_until_processed().await.unwrap(); + + tokio::time::advance(50.std_milliseconds()).await; + // Model the interval where the replay timer fires after the atomic hint but before its ordered + // control message is consumed. + sender.set_crash_pending_without_enqueuing(); + + tokio::time::advance(50.std_milliseconds()).await; + expect_init_replay_event(&mut events, InitReplayTestEvent::CrashDelayApplied).await; + assert_eq!(0, setup.replayer_log_count.load(Ordering::SeqCst)); + + tokio::time::advance(50.std_milliseconds()).await; + expect_init_replay_event(&mut events, InitReplayTestEvent::Replayed).await; + assert_eq!(1, setup.replayer_log_count.load(Ordering::SeqCst)); + + shutdown_trigger.shutdown().await; + handle.await.unwrap(); +} + +#[tokio::test(start_paused = true)] +async fn prior_run_logs_replay_before_current_session_init_logs() { + let mut setup = Setup::new(); + setup + .runtime + .update_snapshot(bd_test_helpers::runtime::make_simple_update(vec![( + bd_runtime::runtime::init_buffer::ReplayDelay::path(), + ValueKind::Int(100), + )])) + .await + .unwrap(); + + let (config_update_tx, config_update_rx) = tokio::sync::mpsc::channel(1); + let (buffer, sender, mut events) = + setup.make_test_async_log_buffer_with_init_replay_events(config_update_rx); + let test_store = TestStore::new().await; + let shutdown_trigger = ComponentShutdownTrigger::default(); + let handle = tokio::task::spawn(buffer.run_with_shutdown( + (*test_store).clone(), + (), + shutdown_trigger.make_shutdown(), + )); + config_update_tx + .send(setup.make_config_update(WorkflowsConfiguration::default())) + .await + .unwrap(); + expect_init_replay_event(&mut events, InitReplayTestEvent::Configured).await; + + assert_ok!(AsyncLogBuffer::::enqueue_log( + &sender, + 0, + LogType::NORMAL, + "current".into(), + [].into(), + [].into(), + None, + Block::No, + None, + )); + sender + .try_send_log( + LogLine { + log_level: 0, + log_type: LogType::LIFECYCLE, + message: "previous".into(), + fields: [].into(), + matching_fields: [].into(), + attributes_overrides: Some(LogAttributesOverrides::PreviousRunSessionID( + OffsetDateTime::now_utc(), + )), + capture_session: None, + } + .into(), + ) + .unwrap(); + sender.wait_until_processed().await.unwrap(); + + tokio::time::advance(100.std_milliseconds()).await; + expect_init_replay_event(&mut events, InitReplayTestEvent::Replayed).await; + assert_eq!( + vec!["previous".to_string(), "current".to_string()], + *setup.replayer_logs.lock() + ); + + shutdown_trigger.shutdown().await; + handle.await.unwrap(); +} + #[tokio::test] async fn processes_log_with_global_state_in_attributes_overrides() { let mut setup = Setup::new(); diff --git a/bd-logger/src/builder.rs b/bd-logger/src/builder.rs index 0f87f3204..42d51b184 100644 --- a/bd-logger/src/builder.rs +++ b/bd-logger/src/builder.rs @@ -607,7 +607,7 @@ impl LoggerBuilder { Ok(()) }, async move { - async_log_buffer.run(state_store, crash_monitor).await; + Box::pin(async_log_buffer.run(state_store, crash_monitor)).await; Ok(()) }, async move { diff --git a/bd-logger/src/init_buffer.rs b/bd-logger/src/init_buffer.rs new file mode 100644 index 000000000..ec8be2282 --- /dev/null +++ b/bd-logger/src/init_buffer.rs @@ -0,0 +1,366 @@ +// shared-core - bitdrift's common client/server libraries +// Copyright Bitdrift, Inc. All rights reserved. +// +// Use of this source code is governed by a source available license that can be found in the +// LICENSE file or at: +// https://polyformproject.org/wp-content/uploads/2020/06/PolyForm-Shield-1.0.0.txt + +#[cfg(test)] +#[path = "./init_buffer_test.rs"] +mod init_buffer_test; + +use bd_client_stats_store::{Counter, Histogram, Scope}; +use bd_log_primitives::size::MemorySized; +use bd_stats_common::{Counter as _, Histogram as _, labels}; +use std::pin::Pin; +use time::Duration; +use tokio::time::{Instant, Sleep}; + +#[derive(thiserror::Error, Debug, PartialEq, Eq)] +pub enum Error { + #[error("Full size overflow")] + FullSizeOverflow, +} + +// +// ReplayReason +// + +/// Identifies why buffered initialization work was replayed. +#[derive(Clone, Copy)] +pub enum ReplayReason { + Scheduled, + CapacityLog, + CapacityStateOperation, +} + +// +// PendingStateOperation +// + +/// Represents a state update that must be replayed in startup order with buffered logs. +#[derive(Debug, Clone)] +pub enum PendingStateOperation { + SetFeatureFlagExposure { + name: String, + variant: Option, + session_id: String, + }, +} + +// +// InitItem +// + +/// A log or state update held until initialization replay begins. +#[derive(Debug)] +pub enum InitItem { + PreviousRunLog(bd_log_primitives::Log), + Log(bd_log_primitives::Log), + StateOperation(PendingStateOperation), +} + +pub trait Prioritizable { + fn is_prioritized(&self) -> bool; +} + +impl MemorySized for InitItem { + fn size(&self) -> usize { + // Don't add extra discriminant overhead - the enum's memory layout already accounts for it. + // We only need to measure the actual data size of each variant. + match self { + Self::PreviousRunLog(log) | Self::Log(log) => log.size(), + Self::StateOperation(PendingStateOperation::SetFeatureFlagExposure { + name, + variant, + session_id, + }) => name.len() + variant.as_ref().map_or(0, String::len) + session_id.len(), + } + } +} + +impl Prioritizable for InitItem { + fn is_prioritized(&self) -> bool { + matches!(self, Self::PreviousRunLog(_)) + } +} + +// +// InitBuffer +// + +/// A FIFO of startup work with a soft memory limit. One item may exceed the limit, and prior-run +/// crash logs can be drained first while preserving their order and the order of all other items. +#[derive(Debug)] +pub struct InitBuffer { + max_size: usize, + current_size: usize, + over_limit_item_accepted: bool, + priority_items: Vec, + items: Vec, +} + +impl InitBuffer { + pub const fn new(max_size: usize) -> Self { + Self { + max_size, + current_size: 0, + over_limit_item_accepted: false, + priority_items: vec![], + items: vec![], + } + } + + pub fn can_push(&self, entry: &T) -> bool { + self.can_push_size(entry.size()) + } + + pub const fn can_push_size(&self, size: usize) -> bool { + self.current_size + size <= self.max_size || !self.over_limit_item_accepted + } + + pub fn push(&mut self, entry: T) -> Result<(), Error> { + let entry_size = entry.size(); + if !self.can_push(&entry) { + log::debug!( + "failed to enqueue init item due to size limit ({}), current size: {}, item size: {}", + self.max_size, + self.current_size, + entry_size, + ); + // Adding an item to the buffer would make it exceed the configured byte limit. + return Err(Error::FullSizeOverflow); + } + + if self.current_size + entry_size > self.max_size { + self.over_limit_item_accepted = true; + log::debug!( + "accepting init item beyond size limit ({}), current size: {}, item size: {}", + self.max_size, + self.current_size, + entry_size, + ); + } + self.current_size += entry_size; + if entry.is_prioritized() { + self.priority_items.push(entry); + } else { + self.items.push(entry); + } + Ok(()) + } + + pub fn drain(mut self) -> impl Iterator { + self.current_size = 0; + self.priority_items.into_iter().chain(self.items) + } + + pub const fn item_count(&self) -> usize { + self.priority_items.len() + self.items.len() + } + + pub const fn max_size(&self) -> usize { + self.max_size + } +} + +// +// PendingInitBuffer +// + +/// Holds startup work after configuration has created a processing pipeline but before that work +/// is replayed. The replay deadline is extended at most once for a pending crash report. +pub struct PendingInitBuffer { + buffer: InitBuffer, + stats: InitBufferStats, + replay_sleep: Option>>, + replay_deadline: Option, + crash_pending_delay_applied: bool, +} + +impl PendingInitBuffer { + pub fn new(buffer: InitBuffer, stats: InitBufferStats) -> Self { + Self { + buffer, + stats, + replay_sleep: None, + replay_deadline: None, + crash_pending_delay_applied: false, + } + } + + pub const fn buffer(&self) -> &InitBuffer { + &self.buffer + } + + pub fn push(&mut self, item: InitItem) -> Result<(), Error> { + let result = self.buffer.push(item); + self.stats.record_push(&result); + result + } + + pub fn schedule( + &mut self, + base_delay: Duration, + crash_pending: bool, + crash_delay: Duration, + ) -> bool { + let replay_delay = if crash_pending { + self.crash_pending_delay_applied = true; + self.stats.record_crash_hint(crash_delay); + base_delay + crash_delay + } else { + base_delay + }; + + if replay_delay.is_zero() { + return true; + } + + let deadline = Instant::now() + replay_delay.unsigned_abs(); + self.replay_deadline = Some(deadline); + self.replay_sleep = Some(Box::pin(tokio::time::sleep_until(deadline))); + false + } + + pub fn apply_crash_pending_hint(&mut self, crash_delay: Duration) -> bool { + if self.crash_pending_delay_applied { + self.stats.record_crash_hint_already_applied(); + return false; + } + self.crash_pending_delay_applied = true; + self.stats.record_crash_hint(crash_delay); + + if crash_delay.is_zero() { + return false; + } + + let deadline = self.replay_deadline.unwrap_or_else(Instant::now) + crash_delay.unsigned_abs(); + self.replay_deadline = Some(deadline); + if let Some(sleep) = &mut self.replay_sleep { + sleep.as_mut().reset(deadline); + } else { + self.replay_sleep = Some(Box::pin(tokio::time::sleep_until(deadline))); + } + true + } + + pub fn replay_sleep(&mut self) -> &mut Option>> { + &mut self.replay_sleep + } + + pub fn drain(mut self, reason: ReplayReason) -> impl Iterator { + self.replay_sleep = None; + self.replay_deadline = None; + self + .stats + .record_replay(reason, self.buffer.item_count(), self.buffer.current_size); + self.buffer.drain() + } +} + +// +// InitBufferStats +// + +/// Metrics for startup buffering. The legacy metric scope is retained for dashboard compatibility. +pub struct InitBufferStats { + pushes: PushCounters, + replay_scheduled: Counter, + replay_capacity_log: Counter, + replay_capacity_state_operation: Counter, + replay_item_count: Histogram, + replay_byte_count: Histogram, + crash_hint_extension_applied: Counter, + crash_hint_already_applied: Counter, + crash_hint_delay_disabled: Counter, +} + +impl InitBufferStats { + pub(crate) fn new(scope: &Scope) -> Self { + let scope = scope.scope("pre_config_log_buffer"); + Self { + pushes: PushCounters::new(&scope), + replay_scheduled: scope + .counter_with_labels("init_buffer_replay", labels!("reason" => "scheduled")), + replay_capacity_log: scope + .counter_with_labels("init_buffer_replay", labels!("reason" => "capacity_log")), + replay_capacity_state_operation: scope.counter_with_labels( + "init_buffer_replay", + labels!("reason" => "capacity_state_operation"), + ), + replay_item_count: scope.histogram("init_buffer_replay_item_count"), + replay_byte_count: scope.histogram("init_buffer_replay_byte_count"), + crash_hint_extension_applied: scope.counter_with_labels( + "init_buffer_crash_hint", + labels!("outcome" => "extension_applied"), + ), + crash_hint_already_applied: scope.counter_with_labels( + "init_buffer_crash_hint", + labels!("outcome" => "already_applied"), + ), + crash_hint_delay_disabled: scope.counter_with_labels( + "init_buffer_crash_hint", + labels!("outcome" => "delay_disabled"), + ), + } + } + + pub(crate) fn record_push(&self, result: &std::result::Result<(), Error>) { + self.pushes.record(result); + } + + fn record_replay(&self, reason: ReplayReason, item_count: usize, byte_count: usize) { + match reason { + ReplayReason::Scheduled => self.replay_scheduled.inc(), + ReplayReason::CapacityLog => self.replay_capacity_log.inc(), + ReplayReason::CapacityStateOperation => self.replay_capacity_state_operation.inc(), + } + self.replay_item_count.observe(histogram_value(item_count)); + self.replay_byte_count.observe(histogram_value(byte_count)); + } + + fn record_crash_hint(&self, crash_delay: Duration) { + if crash_delay.is_zero() { + self.crash_hint_delay_disabled.inc(); + } else { + self.crash_hint_extension_applied.inc(); + } + } + + fn record_crash_hint_already_applied(&self) { + self.crash_hint_already_applied.inc(); + } +} + +fn histogram_value(value: usize) -> f64 { + f64::from(u32::try_from(value).unwrap_or(u32::MAX)) +} + +// +// PushCounters +// + +struct PushCounters { + ok: Counter, + err_full_size_overflow: Counter, +} + +impl PushCounters { + fn new(scope: &Scope) -> Self { + Self { + ok: scope.counter_with_labels("log_enqueueing", labels!("result" => "success")), + err_full_size_overflow: scope.counter_with_labels( + "log_enqueueing", + labels!("result" => "failure_size_overflow"), + ), + } + } + + fn record(&self, result: &std::result::Result<(), Error>) { + match result { + Ok(()) => self.ok.inc(), + Err(Error::FullSizeOverflow) => self.err_full_size_overflow.inc(), + } + } +} diff --git a/bd-logger/src/init_buffer_test.rs b/bd-logger/src/init_buffer_test.rs new file mode 100644 index 000000000..91361025a --- /dev/null +++ b/bd-logger/src/init_buffer_test.rs @@ -0,0 +1,69 @@ +// shared-core - bitdrift's common client/server libraries +// Copyright Bitdrift, Inc. All rights reserved. +// +// Use of this source code is governed by a source available license that can be found in the +// LICENSE file or at: +// https://polyformproject.org/wp-content/uploads/2020/06/PolyForm-Shield-1.0.0.txt + +use crate::init_buffer::{self, InitBuffer, Prioritizable}; +use bd_log_primitives::size::MemorySized; + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +struct SimulatedSizeLog { + size: usize, +} + +impl MemorySized for SimulatedSizeLog { + fn size(&self) -> usize { + self.size + } +} + +impl Prioritizable for SimulatedSizeLog { + fn is_prioritized(&self) -> bool { + self.size.is_multiple_of(2) + } +} + +#[test] +fn buffer_accepts_one_item_beyond_its_soft_limit() { + let mut buffer = InitBuffer::new(1024); + + let logs = [ + SimulatedSizeLog { size: 1 }, + SimulatedSizeLog { size: 999 }, + SimulatedSizeLog { size: 101 }, + SimulatedSizeLog { size: 3 }, + ]; + + assert_eq!(Ok(()), buffer.push(logs[0])); + assert_eq!(Ok(()), buffer.push(logs[1])); + assert_eq!(Ok(()), buffer.push(logs[2])); + assert_eq!( + Err(init_buffer::Error::FullSizeOverflow), + buffer.push(logs[3]) + ); + + let items: Vec<_> = buffer.drain().collect(); + assert_eq!(vec![logs[0], logs[1], logs[2]], items); +} + +#[test] +fn buffer_can_prioritize_items_without_reordering_each_group() { + let mut buffer = InitBuffer::new(1024); + for size in [1, 2, 3, 4] { + buffer.push(SimulatedSizeLog { size }).unwrap(); + } + + let items: Vec<_> = buffer.drain().collect(); + + assert_eq!( + vec![ + SimulatedSizeLog { size: 2 }, + SimulatedSizeLog { size: 4 }, + SimulatedSizeLog { size: 1 }, + SimulatedSizeLog { size: 3 }, + ], + items + ); +} diff --git a/bd-logger/src/lib.rs b/bd-logger/src/lib.rs index 4c0880486..89a98b500 100644 --- a/bd-logger/src/lib.rs +++ b/bd-logger/src/lib.rs @@ -24,6 +24,7 @@ mod consumer; mod device_id; mod directory_lock; mod flush_registry; +mod init_buffer; pub mod internal; mod internal_report; mod log_replay; @@ -32,7 +33,6 @@ mod logging_state; mod metadata; mod network; mod ordered_receiver; -mod pre_config_buffer; mod service; mod state_upload; mod trigger_upload_artifact; diff --git a/bd-logger/src/logger.rs b/bd-logger/src/logger.rs index 42e382710..682cb6850 100644 --- a/bd-logger/src/logger.rs +++ b/bd-logger/src/logger.rs @@ -209,6 +209,12 @@ pub struct LoggerHandle { } impl LoggerHandle { + /// Informs the logger that the platform expects a prior-run crash report. The logger keeps + /// initialization work buffered for its configured crash-pending extension. + pub fn hint_crash_report_pending(&self) -> anyhow::Result<()> { + Ok(self.tx.hint_crash_report_pending()?) + } + /// Log a message with the given log level, log type, message, and fields. This will enqueue the /// log onto a bounded queue for further processing. pub fn log( @@ -718,6 +724,12 @@ impl Logger { ) } + /// Informs the logger that the platform expects a prior-run crash report. The logger keeps + /// initialization work buffered for its configured crash-pending extension. + pub fn hint_crash_report_pending(&self) -> anyhow::Result<()> { + Ok(self.async_log_buffer_tx.hint_crash_report_pending()?) + } + pub fn shutdown(&self, blocking: bool) { let shutdown_trigger = self.shutdown_state.lock().take(); diff --git a/bd-logger/src/logging_state.rs b/bd-logger/src/logging_state.rs index 9ce00fdd7..aa258d5a6 100644 --- a/bd-logger/src/logging_state.rs +++ b/bd-logger/src/logging_state.rs @@ -8,10 +8,10 @@ use crate::buffer_selector::BufferSelector; use crate::client_config::TailConfigurations; use crate::consumer::RemoteFlushStreamingRequest; +use crate::init_buffer::{InitBuffer, InitBufferStats, Prioritizable}; use crate::log_replay::{LogReplay, ProcessingPipeline}; use crate::logger::with_thread_local_logger_guard; use crate::metadata::MetadataCollector; -use crate::pre_config_buffer::{self, PreConfigBuffer}; use anyhow::anyhow; use bd_api::{DataUpload, TriggerUpload}; use bd_buffer::BuffersWithAck; @@ -47,32 +47,32 @@ use tokio::sync::mpsc::{Receiver, Sender}; /// that are needed to process incoming logs. #[derive(Debug)] #[allow(clippy::large_enum_variant)] -pub enum LoggingState { +pub enum LoggingState { /// The initial state that each `AsyncLogBuffer` starts in. While in this state /// the buffer takes incoming logs, populates them with extra information using /// its metadata provider and puts them on hold for further processing inside of - /// a `PreConfigBuffer`. The final processing of logs is postponed until after the + /// an `InitBuffer`. The final processing of logs is postponed until after the /// buffer moves to `Initialized` state. /// - /// The buffer stays in `Uninitialized` state until it gets a configuration update. - /// Configuration updates come from either a local cache (disk) or a Bitdrift control plane. - /// While loading from a local cache is extremely fast (measured in milliseconds), - /// the cached version of the configuration is not always available and in these - /// cases the `AsyncLogBuffer` waits for the configuration to be fetched - /// from the Bitdrift control plane (can potentially take seconds or even minutes). + /// The buffer stays in this state until it receives a configuration update from the local + /// cache or Bitdrift control plane. While it waits, the `InitBuffer` retains a bounded amount + /// of startup work. Cached configuration normally arrives within milliseconds, but when it is + /// unavailable the control-plane update can take seconds or minutes. Uninitialized(UninitializedLoggingContext), /// The state that `AsyncLogBuffer` moves to as soon as it receives any configuration /// update. /// While in this state the `AsyncLogBuffer` takes incoming logs, populates them with /// extra information its metadata provider and sends them for their final processing. - /// The first thing that the buffer does when it moves to this state is a replay of all - /// logs stored inside of its `PreConfigBuffer`. All replayed logs are sent for their final - /// processing to now initialized parts of the logs processing pipeline such as workflows engine - /// or various ring buffers. + /// Startup work buffered in the `Uninitialized` state is replayed through the initialized + /// pipeline after the runtime-configured `InitBuffer` replay delay. The crash handler can send + /// a crash-pending hint before or during that delay; the separately configured crash-pending + /// delay extends the replay deadline once. This gives prior-run crash logs time to arrive and + /// replay ahead of current-session startup activity, so persisted workflows process the crash + /// report first. Initialized(InitializedLoggingContext), } -impl LoggingState { +impl LoggingState { pub(crate) const fn flush_buffers_trigger(&self) -> &Sender { match self { Self::Uninitialized(context) => &context.flush_buffers_tx, @@ -101,8 +101,8 @@ impl LoggingState { // UninitializedLoggingContext // -pub struct UninitializedLoggingContext { - pub(crate) pre_config_log_buffer: PreConfigBuffer, +pub struct UninitializedLoggingContext { + pub(crate) init_buffer: InitBuffer, data_upload_tx: Sender, trigger_upload_tx: Sender, @@ -118,10 +118,10 @@ pub struct UninitializedLoggingContext { } // Skip `stats` and `runtime` fields that does not implement `std::fmt::Debug`. -impl Debug for UninitializedLoggingContext { +impl Debug for UninitializedLoggingContext { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { f.debug_struct("UninitializedLoggingContext") - .field("pre_config_log_buffer", &self.pre_config_log_buffer) + .field("init_buffer", &self.init_buffer) .field("trigger_upload_tx", &self.trigger_upload_tx) .field("flush_buffers_tx", &self.flush_buffers_tx) .field("flush_stats_trigger", &self.flush_stats_trigger) @@ -130,7 +130,7 @@ impl Debug for UninitializedLoggingContext { } } -impl UninitializedLoggingContext { +impl UninitializedLoggingContext { pub(crate) fn new( sdk_directory: &Path, runtime: &Arc, @@ -146,7 +146,7 @@ impl UninitializedLoggingContext { process_local_pending_flush_state: Arc, ) -> Self { Self { - pre_config_log_buffer: PreConfigBuffer::new(max_size), + init_buffer: InitBuffer::new(max_size), data_upload_tx, trigger_upload_tx, remote_flush_streaming_rx, @@ -164,7 +164,7 @@ impl UninitializedLoggingContext { self, config: ConfigUpdate, capture_screenshot_handler: CaptureScreenshotHandler, - ) -> (InitializedLoggingContext, PreConfigBuffer) { + ) -> (InitializedLoggingContext, InitBuffer, InitBufferStats) { let processing_pipeline = ProcessingPipeline::new( self.data_upload_tx, self.flush_buffers_tx, @@ -183,7 +183,7 @@ impl UninitializedLoggingContext { let context = InitializedLoggingContext::new(processing_pipeline); - (context, self.pre_config_log_buffer) + (context, self.init_buffer, self.stats.init_buffer) } } @@ -291,7 +291,7 @@ impl InitializedLoggingContext { // pub struct UninitializedLoggingContextStats { - pub(crate) pre_config_log_buffer: pre_config_buffer::PushCounters, + pub(crate) init_buffer: InitBufferStats, pub(crate) scope: Scope, root_scope: Scope, stats: Arc, @@ -300,10 +300,9 @@ pub struct UninitializedLoggingContextStats { impl UninitializedLoggingContextStats { fn new(root_scope: Scope, stats: Arc) -> Self { let stats_scope = root_scope.scope("logger"); - let pre_config_buffer_scope = stats_scope.scope("pre_config_log_buffer"); Self { - pre_config_log_buffer: pre_config_buffer::PushCounters::new(&pre_config_buffer_scope), + init_buffer: InitBufferStats::new(&stats_scope), scope: stats_scope, root_scope, stats, diff --git a/bd-logger/src/pre_config_buffer.rs b/bd-logger/src/pre_config_buffer.rs deleted file mode 100644 index 15a24d26b..000000000 --- a/bd-logger/src/pre_config_buffer.rs +++ /dev/null @@ -1,143 +0,0 @@ -// shared-core - bitdrift's common client/server libraries -// Copyright Bitdrift, Inc. All rights reserved. -// -// Use of this source code is governed by a source available license that can be found in the -// LICENSE file or at: -// https://polyformproject.org/wp-content/uploads/2020/06/PolyForm-Shield-1.0.0.txt - -#[cfg(test)] -#[path = "./pre_config_buffer_test.rs"] -mod pre_config_buffer_test; -use bd_client_stats_store::{Counter, Scope}; -use bd_log_primitives::size::MemorySized; -use bd_stats_common::{Counter as _, labels}; - -#[derive(thiserror::Error, Debug, PartialEq, Eq)] -pub enum Error { - #[error("Full size overflow")] - FullSizeOverflow, -} - -// -// PendingStateOperation -// - -/// Represents a state update operation that occurred before initialization. -/// These operations are queued and replayed after initialization to ensure -/// proper ordering with logs and proper workflow state transitions. -#[derive(Debug, Clone)] -pub enum PendingStateOperation { - SetFeatureFlagExposure { - name: String, - variant: Option, - session_id: String, - }, -} - -// -// PreConfigItem -// - -/// An item that can be stored in the pre-config buffer, representing either -/// a log or a state operation that occurred before initialization. This allows -/// both logs and state changes to be replayed in the exact order they arrived. -#[derive(Debug)] -pub enum PreConfigItem { - Log(bd_log_primitives::Log), - StateOperation(PendingStateOperation), -} - -impl MemorySized for PreConfigItem { - fn size(&self) -> usize { - // Don't add extra discriminant overhead - the enum's memory layout already accounts for it. - // We only need to measure the actual data size of each variant. - match self { - Self::Log(log) => log.size(), - Self::StateOperation(PendingStateOperation::SetFeatureFlagExposure { - name, - variant, - session_id, - }) => name.len() + variant.as_ref().map_or(0, String::len) + session_id.len(), - } - } -} - -/// An in-memory buffer that is used to buffer events that arrive before an initial logger -/// configuration has been received. This allows us to capture some number of events that can be -/// "replayed" once the configuration has been applied. -#[derive(Debug)] -pub struct PreConfigBuffer { - max_size: usize, - - current_size: usize, - items: Vec, -} - -impl PreConfigBuffer { - pub const fn new(max_size: usize) -> Self { - Self { - max_size, - current_size: 0, - items: vec![], - } - } - - pub fn push(&mut self, entry: T) -> Result<(), Error> { - let log_size = entry.size(); - if self.current_size + log_size > self.max_size { - log::debug!( - "failed to enqueue log due to items size limit ({}), current size: {}, log size: {}", - self.max_size, - self.current_size, - log_size, - ); - // Adding a log to the buffer would make it exceed - // configured bytes size limit. - return Err(Error::FullSizeOverflow); - } - - self.current_size += log_size; - self.items.push(entry); - - Ok(()) - } - - pub fn pop_all(mut self) -> impl Iterator { - self.current_size = 0; - self.items.into_iter() - } - - pub const fn max_size(&self) -> usize { - self.max_size - } -} - -// -// PushCounters -// - -pub struct PushCounters { - ok: Counter, - err_full_size_overflow: Counter, -} - -impl PushCounters { - pub(crate) fn new(scope: &Scope) -> Self { - Self { - ok: scope.counter_with_labels("log_enqueueing", labels!("result" => "success")), - err_full_size_overflow: scope.counter_with_labels( - "log_enqueueing", - labels!("result" => "failure_size_overflow"), - ), - } - } - - pub(crate) fn record(&self, result: &std::result::Result<(), Error>) { - match result { - Ok(()) => self.ok.inc(), - Err(Error::FullSizeOverflow) => { - self.err_full_size_overflow.inc(); - }, - } - } -} diff --git a/bd-logger/src/pre_config_buffer_test.rs b/bd-logger/src/pre_config_buffer_test.rs deleted file mode 100644 index 2f4e46d96..000000000 --- a/bd-logger/src/pre_config_buffer_test.rs +++ /dev/null @@ -1,57 +0,0 @@ -// shared-core - bitdrift's common client/server libraries -// Copyright Bitdrift, Inc. All rights reserved. -// -// Use of this source code is governed by a source available license that can be found in the -// LICENSE file or at: -// https://polyformproject.org/wp-content/uploads/2020/06/PolyForm-Shield-1.0.0.txt - -use crate::pre_config_buffer::{self, PreConfigBuffer}; -use bd_log_primitives::size::MemorySized; - -#[derive(Clone, Copy, Debug, PartialEq, Eq)] -struct SimulatedSizeLog { - size: usize, -} - -impl MemorySized for SimulatedSizeLog { - fn size(&self) -> usize { - self.size - } -} - -#[test] -fn buffer_acts_as_fifo_queue_with_limits() { - let mut buffer = PreConfigBuffer::new(1024); - - let logs = [ - SimulatedSizeLog { size: 1 }, - SimulatedSizeLog { size: 2000 }, - SimulatedSizeLog { size: 2 }, - SimulatedSizeLog { size: 3 }, - ]; - - assert_eq!(Ok(()), buffer.push(logs[0])); - // Too big to be accepted by the buffer. - assert_eq!( - Err(pre_config_buffer::Error::FullSizeOverflow), - buffer.push(logs[1]) - ); - assert_eq!(Ok(()), buffer.push(logs[2])); - assert_eq!(Ok(()), buffer.push(logs[3])); - - // Can add many small items as long as we're within memory limit - for _ in 0 .. 500 { - assert_eq!(Ok(()), buffer.push(SimulatedSizeLog { size: 2 })); - } - - // Now at 1 + 2 + 3 + (500 * 2) = 1006 bytes - // Adding another 20 bytes would exceed 1024 - assert_eq!( - Err(pre_config_buffer::Error::FullSizeOverflow), - buffer.push(SimulatedSizeLog { size: 20 }) - ); - - // First few items - let items: Vec<_> = buffer.pop_all().take(3).collect(); - assert_eq!(vec![logs[0], logs[2], logs[3]], items); -} diff --git a/bd-logger/src/test/setup.rs b/bd-logger/src/test/setup.rs index f3004aff2..bc114e3c2 100644 --- a/bd-logger/src/test/setup.rs +++ b/bd-logger/src/test/setup.rs @@ -401,13 +401,23 @@ impl Setup { } pub fn upload_individual_logs(&mut self) { + self.upload_individual_logs_with_runtime_values(vec![]); + } + + pub fn upload_individual_logs_with_runtime_values( + &mut self, + runtime_values: Vec<(&'static str, ValueKind)>, + ) { self .current_api_stream() .blocking_stream_action(StreamAction::SendRuntime(make_update( vec![( bd_runtime::runtime::log_upload::BatchSizeFlag::path(), bd_test_helpers::runtime::ValueKind::Int(1), - )], + )] + .into_iter() + .chain(runtime_values) + .collect(), "base".to_string(), ))); diff --git a/bd-runtime/src/runtime.rs b/bd-runtime/src/runtime.rs index fd19b374b..95c68f827 100644 --- a/bd-runtime/src/runtime.rs +++ b/bd-runtime/src/runtime.rs @@ -672,6 +672,23 @@ pub mod crash_reporting { ); } +pub mod init_buffer { + // Delays the first replay after logger configuration is available. This lets the platform add + // prior-run crash reports before current-session startup activity reaches workflows. + duration_feature_flag!( + ReplayDelay, + "logger.init_buffer.replay_delay_ms", + time::Duration::ZERO + ); + + // Extends the initial replay delay once when the platform expects a prior-run crash report. + duration_feature_flag!( + CrashPendingReplayDelay, + "logger.init_buffer.crash_pending_replay_delay_ms", + time::Duration::ZERO + ); +} + pub mod log_upload { use time::ext::NumericalDuration as _;